Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/ui/dag_runs.py: 71%

19 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 14:22 +0000

1# Licensed to the Apache Software Foundation (ASF) under one 

2# or more contributor license agreements. See the NOTICE file 

3# distributed with this work for additional information 

4# regarding copyright ownership. The ASF licenses this file 

5# to you under the Apache License, Version 2.0 (the 

6# "License"); you may not use this file except in compliance 

7# with the License. You may obtain a copy of the License at 

8# 

9# http://www.apache.org/licenses/LICENSE-2.0 

10# 

11# Unless required by applicable law or agreed to in writing, 

12# software distributed under the License is distributed on an 

13# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY 

14# KIND, either express or implied. See the License for the 

15# specific language governing permissions and limitations 

16# under the License. 

17 

18from __future__ import annotations 

19 

20from fastapi import Depends, HTTPException, status 

21from sqlalchemy import select 

22 

23from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity 

24from airflow.api_fastapi.common.db.common import SessionDep 

25from airflow.api_fastapi.common.router import AirflowRouter 

26from airflow.api_fastapi.core_api.datamodels.ui.dag_runs import DagRunStatsResponse 

27from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc 

28from airflow.api_fastapi.core_api.security import requires_access_dag 

29from airflow.api_fastapi.core_api.services.ui.dag_run import compute_duration_stats 

30from airflow.models.dagrun import DagRun 

31from airflow.utils.state import DagRunState 

32 

33dag_runs_router = AirflowRouter(prefix="/dags/{dag_id}/dagRuns", tags=["DagRun"]) 

34 

35 

36@dag_runs_router.get( 

37 "/{dag_run_id}/stats", 

38 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

39 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN))], 

40) 

41def get_dag_run_stats(dag_id: str, dag_run_id: str, session: SessionDep) -> DagRunStatsResponse: 

42 """Get duration statistics for a DAG based on its historical completed runs.""" 

43 if not session.scalar(select(DagRun.id).filter_by(dag_id=dag_id, run_id=dag_run_id)): 

44 raise HTTPException( 

45 status.HTTP_404_NOT_FOUND, 

46 f"The DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` was not found", 

47 ) 

48 

49 durations = [ 

50 d 

51 for d in session.scalars( 

52 select(DagRun.duration.expression) # type: ignore[attr-defined] 

53 .where( 

54 DagRun.dag_id == dag_id, 

55 DagRun.state.in_([DagRunState.SUCCESS, DagRunState.FAILED]), 

56 ) 

57 .order_by(DagRun.run_after.desc()) 

58 .limit(100) 

59 ) 

60 if d is not None 

61 ] 

62 

63 return DagRunStatsResponse(duration=compute_duration_stats(durations))