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
« 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.
18from __future__ import annotations
20from typing import TYPE_CHECKING, Annotated
22from fastapi import Depends, status
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
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
50dag_stats_router = AirflowRouter(tags=["DagStats"], prefix="/dagStats")
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]
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)
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))