Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/ui/dags.py: 37%
58 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 Annotated
22from fastapi import Depends, HTTPException, status
23from sqlalchemy import select, union_all
24from sqlalchemy.orm import defaultload
26from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity
27from airflow.api_fastapi.common.db.common import (
28 SessionDep,
29 paginated_select,
30)
31from airflow.api_fastapi.common.db.dags import generate_dag_with_latest_run_query
32from airflow.api_fastapi.common.parameters import (
33 FilterOptionEnum,
34 FilterParam,
35 QueryAnyDagRunStateFilter,
36 QueryAssetDependencyFilter,
37 QueryBundleNameFilter,
38 QueryBundleVersionFilter,
39 QueryDagDisplayNamePatternSearch,
40 QueryDagDisplayNamePrefixPatternSearch,
41 QueryDagIdPatternSearch,
42 QueryDagIdPrefixPatternSearch,
43 QueryExcludeStaleFilter,
44 QueryFavoriteFilter,
45 QueryHasAssetScheduleFilter,
46 QueryHasImportErrorsFilter,
47 QueryLastDagRunStateFilter,
48 QueryLimit,
49 QueryOffset,
50 QueryOwnersFilter,
51 QueryPausedFilter,
52 QueryPendingActionsFilter,
53 QueryTagsFilter,
54 SortParam,
55 filter_param_factory,
56)
57from airflow.api_fastapi.common.router import AirflowRouter
58from airflow.api_fastapi.core_api.datamodels.dags import DAG_ALIAS_MAPPING, DAGResponse
59from airflow.api_fastapi.core_api.datamodels.ui.dag_runs import DAGRunLightResponse
60from airflow.api_fastapi.core_api.datamodels.ui.dags import (
61 DAGWithLatestDagRunsCollectionResponse,
62 DAGWithLatestDagRunsResponse,
63)
64from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
65from airflow.api_fastapi.core_api.security import (
66 GetUserDep,
67 ReadableDagsFilterDep,
68 requires_access_dag,
69)
70from airflow.models import DagModel, DagRun
71from airflow.models.dag_favorite import DagFavorite
72from airflow.models.hitl import HITLDetail
73from airflow.models.taskinstance import TaskInstance
74from airflow.utils.state import TaskInstanceState
76dags_router = AirflowRouter(prefix="/dags", tags=["DAG"])
78# Per-dag run counts read at most this many rows per state; the UI shows "N+" at the cap.
79STATE_COUNT_CAP = 1000
82@dags_router.get(
83 "",
84 response_model_exclude_none=True,
85 dependencies=[
86 Depends(requires_access_dag(method="GET")),
87 Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN)),
88 Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.HITL_DETAIL)),
89 Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK_INSTANCE)),
90 ],
91 operation_id="get_dags_ui",
92)
93def get_dags(
94 limit: QueryLimit,
95 offset: QueryOffset,
96 tags: QueryTagsFilter,
97 owners: QueryOwnersFilter,
98 dag_ids: Annotated[
99 FilterParam[list[str] | None],
100 Depends(filter_param_factory(DagModel.dag_id, list[str] | None, FilterOptionEnum.IN, "dag_ids")),
101 ],
102 dag_id_pattern: QueryDagIdPatternSearch,
103 dag_id_prefix_pattern: QueryDagIdPrefixPatternSearch,
104 dag_display_name_pattern: QueryDagDisplayNamePatternSearch,
105 dag_display_name_prefix_pattern: QueryDagDisplayNamePrefixPatternSearch,
106 exclude_stale: QueryExcludeStaleFilter,
107 paused: QueryPausedFilter,
108 has_import_errors: QueryHasImportErrorsFilter,
109 last_dag_run_state: QueryLastDagRunStateFilter,
110 dag_run_state: QueryAnyDagRunStateFilter,
111 bundle_name: QueryBundleNameFilter,
112 bundle_version: QueryBundleVersionFilter,
113 order_by: Annotated[
114 SortParam,
115 Depends(
116 SortParam(
117 ["dag_id", "dag_display_name", "next_dagrun", "state", "start_date"],
118 DagModel,
119 {"last_run_state": DagRun.state, "last_run_start_date": DagRun.start_date},
120 ).dynamic_depends()
121 ),
122 ],
123 is_favorite: QueryFavoriteFilter,
124 has_asset_schedule: QueryHasAssetScheduleFilter,
125 asset_dependency: QueryAssetDependencyFilter,
126 has_pending_actions: QueryPendingActionsFilter,
127 readable_dags_filter: ReadableDagsFilterDep,
128 session: SessionDep,
129 user: GetUserDep,
130 dag_runs_limit: int = 10,
131) -> DAGWithLatestDagRunsCollectionResponse:
132 """Get Dags with recent DagRun."""
133 # Fetch Dags with their latest DagRun and apply filters
134 query = generate_dag_with_latest_run_query(
135 max_run_filters=[
136 last_dag_run_state,
137 ],
138 order_by=order_by,
139 dag_ids=readable_dags_filter.value,
140 )
142 dags_select, total_entries = paginated_select(
143 statement=query,
144 filters=[
145 exclude_stale,
146 paused,
147 has_import_errors,
148 dag_id_pattern,
149 dag_id_prefix_pattern,
150 dag_ids,
151 dag_display_name_pattern,
152 dag_display_name_prefix_pattern,
153 tags,
154 owners,
155 last_dag_run_state,
156 dag_run_state,
157 is_favorite,
158 has_asset_schedule,
159 asset_dependency,
160 has_pending_actions,
161 readable_dags_filter,
162 bundle_name,
163 bundle_version,
164 ],
165 order_by=order_by,
166 offset=offset,
167 limit=limit,
168 session=session,
169 )
171 dags = [dag for dag in session.scalars(dags_select)]
173 # Fetch favorite status for each Dag for the current user
174 user_id = str(user.get_id())
175 favorites_select = select(DagFavorite.dag_id).where(
176 DagFavorite.user_id == user_id, DagFavorite.dag_id.in_([dag.dag_id for dag in dags])
177 )
178 favorite_dag_ids = set(session.scalars(favorites_select))
180 recent_dag_runs: list = []
181 if dags:
182 recent_runs_branches = [
183 select(
184 DagRun.id,
185 DagRun.dag_id,
186 DagRun.run_id,
187 DagRun.end_date,
188 DagRun.logical_date,
189 DagRun.run_after,
190 DagRun.start_date,
191 DagRun.state,
192 )
193 .where(DagRun.dag_id == dag.dag_id)
194 .order_by(DagRun.run_after.desc())
195 .limit(dag_runs_limit)
196 .subquery()
197 for dag in dags
198 ]
199 recent_runs_union = union_all(*(select(branch) for branch in recent_runs_branches)).subquery()
200 recent_dag_runs = list(
201 session.execute(select(recent_runs_union).order_by(recent_runs_union.c.run_after.desc()))
202 )
204 # Fetch pending HITL actions for each Dag if we are not certain whether some of the Dag might contain HITL actions
205 pending_actions_by_dag_id: dict[str, list[HITLDetail]] = {dag.dag_id: [] for dag in dags}
206 if has_pending_actions.value:
207 pending_actions_select = (
208 select(
209 TaskInstance.dag_id,
210 HITLDetail,
211 )
212 .join(TaskInstance, HITLDetail.ti_id == TaskInstance.id)
213 .options(
214 defaultload(HITLDetail.task_instance).joinedload(TaskInstance.rendered_task_instance_fields)
215 )
216 .where(
217 HITLDetail.responded_at.is_(None),
218 TaskInstance.state.in_((TaskInstanceState.DEFERRED, TaskInstanceState.AWAITING_INPUT)),
219 )
220 .where(TaskInstance.dag_id.in_([dag.dag_id for dag in dags]))
221 .order_by(TaskInstance.dag_id)
222 )
224 pending_actions = session.execute(pending_actions_select)
226 # Group pending actions by dag_id
227 for dag_id, hitl_detail in pending_actions:
228 pending_actions_by_dag_id[dag_id].append(hitl_detail)
230 # aggregate rows by dag_id
231 # Build the dict dynamically from DAGResponse.model_fields so that new fields
232 # added to DAGResponse are picked up automatically without code changes here.
233 dag_runs_by_dag_id: dict[str, DAGWithLatestDagRunsResponse] = {}
234 for dag in dags:
235 dag_data = {
236 DAG_ALIAS_MAPPING.get(field_name, field_name): getattr(
237 dag, DAG_ALIAS_MAPPING.get(field_name, field_name)
238 )
239 for field_name in DAGResponse.model_fields
240 }
241 dag_data.update(
242 {
243 "asset_expression": dag.asset_expression,
244 "latest_dag_runs": [],
245 "pending_actions": pending_actions_by_dag_id[dag.dag_id],
246 "is_favorite": dag.dag_id in favorite_dag_ids,
247 }
248 )
249 dag_runs_by_dag_id[dag.dag_id] = DAGWithLatestDagRunsResponse.model_validate(dag_data)
251 for row in recent_dag_runs:
252 dag_run_response = DAGRunLightResponse.model_validate(row)
253 dag_id = dag_run_response.dag_id
254 dag_runs_by_dag_id[dag_id].latest_dag_runs.append(dag_run_response)
256 return DAGWithLatestDagRunsCollectionResponse(
257 total_entries=total_entries,
258 dags=list(dag_runs_by_dag_id.values()),
259 )
262@dags_router.get(
263 "/{dag_id}/latest_run",
264 responses=create_openapi_http_exception_doc(
265 [
266 status.HTTP_400_BAD_REQUEST,
267 status.HTTP_404_NOT_FOUND,
268 ]
269 ),
270 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN))],
271)
272def get_latest_run_info(dag_id: str, session: SessionDep) -> DAGRunLightResponse | None:
273 """Get latest run."""
274 if dag_id == "~":
275 raise HTTPException(
276 status.HTTP_400_BAD_REQUEST,
277 "`~` was supplied as dag_id, but querying multiple dags is not supported.",
278 )
280 latest_run_info_select = (
281 select(
282 DagRun.id,
283 DagRun.dag_id,
284 DagRun.run_id,
285 DagRun.end_date,
286 DagRun.logical_date,
287 DagRun.run_after,
288 DagRun.start_date,
289 DagRun.state,
290 )
291 .where(DagRun.dag_id == dag_id)
292 .order_by(DagRun.run_after.desc())
293 .limit(1)
294 )
295 latest_run_info = session.execute(latest_run_info_select).one_or_none()
297 return DAGRunLightResponse(**latest_run_info._mapping) if latest_run_info else None