Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/dag_stats.py: 95%

31 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 typing import TYPE_CHECKING, Annotated 

21 

22from fastapi import Depends, status 

23 

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

25from airflow.api_fastapi.common.db.common import ( 

26 SessionDep, 

27 paginated_select, 

28) 

29from airflow.api_fastapi.common.db.dag_runs import dagruns_select_with_state_count 

30from airflow.api_fastapi.common.parameters import ( 

31 FilterOptionEnum, 

32 FilterParam, 

33 filter_param_factory, 

34) 

35from airflow.api_fastapi.common.router import AirflowRouter 

36from airflow.api_fastapi.core_api.datamodels.dag_stats import ( 

37 DagStatsCollectionResponse, 

38 DagStatsResponse, 

39 DagStatsStateResponse, 

40) 

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

42from airflow.api_fastapi.core_api.security import ReadableDagRunsFilterDep, requires_access_dag 

43from airflow.models.dagrun import DagRun 

44from airflow.typing_compat import Unpack 

45from airflow.utils.state import DagRunState 

46 

47if TYPE_CHECKING: 47 ↛ 48line 47 didn't jump to line 48 because the condition on line 47 was never true

48 from sqlalchemy import Result 

49 

50dag_stats_router = AirflowRouter(tags=["DagStats"], prefix="/dagStats") 

51 

52 

53@dag_stats_router.get( 

54 "", 

55 responses=create_openapi_http_exception_doc( 

56 [ 

57 status.HTTP_400_BAD_REQUEST, 

58 status.HTTP_404_NOT_FOUND, 

59 ] 

60 ), 

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

62) 

63def get_dag_stats( 

64 readable_dag_runs_filter: ReadableDagRunsFilterDep, 

65 session: SessionDep, 

66 dag_ids: Annotated[ 

67 FilterParam[list[str]], 

68 Depends(filter_param_factory(DagRun.dag_id, list[str], FilterOptionEnum.IN, "dag_ids")), 

69 ], 

70) -> DagStatsCollectionResponse: 

71 """Get Dag statistics.""" 

72 dagruns_select, _ = paginated_select( 

73 statement=dagruns_select_with_state_count, 

74 filters=[dag_ids, readable_dag_runs_filter], 

75 session=session, 

76 return_total_entries=False, 

77 ) 

78 # The below type annotation is acceptable on SQLA2.1, but not on 2.0 

79 query_result: Result[Unpack[tuple[str, str, str, int]]] = session.execute(dagruns_select) # type: ignore[type-arg] 

80 

81 result_dag_ids = [] 

82 dag_display_names: dict[str, str] = {} 

83 dag_state_data = {} 

84 for dag_id, state, dag_display_name, count in query_result: 

85 dag_state_data[(dag_id, state)] = count 

86 if dag_id not in result_dag_ids: 

87 dag_display_names[dag_id] = dag_display_name 

88 result_dag_ids.append(dag_id) 

89 

90 dags = [ 

91 DagStatsResponse( 

92 dag_id=dag_id, 

93 dag_display_name=dag_display_names[dag_id], 

94 stats=[ 

95 DagStatsStateResponse( 

96 state=state, 

97 count=dag_state_data.get((dag_id, state), 0), 

98 ) 

99 for state in DagRunState 

100 ], 

101 ) 

102 for dag_id in result_dag_ids 

103 ] 

104 return DagStatsCollectionResponse(dags=dags, total_entries=len(dags))