Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/dags.py: 88%
116 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, Query, Response, status
23from fastapi.exceptions import RequestValidationError
24from pydantic import ValidationError
25from sqlalchemy import delete, func, insert, select, update
27from airflow.api.common import delete_dag as delete_dag_module
28from airflow.api_fastapi.common.dagbag import DagBagDep, get_latest_version_of_dag
29from airflow.api_fastapi.common.db.common import SessionDep, apply_filters_to_select, paginated_select
30from airflow.api_fastapi.common.db.dags import generate_dag_with_latest_run_query
31from airflow.api_fastapi.common.parameters import (
32 FilterOptionEnum,
33 FilterParam,
34 QueryAssetDependencyFilter,
35 QueryBundleNameFilter,
36 QueryBundleVersionFilter,
37 QueryDagDisplayNamePatternSearch,
38 QueryDagDisplayNamePrefixPatternSearch,
39 QueryDagIdPatternSearch,
40 QueryDagIdPatternSearchWithNone,
41 QueryDagIdPrefixPatternSearch,
42 QueryDagIdPrefixPatternSearchWithNone,
43 QueryExcludeStaleFilter,
44 QueryFavoriteFilter,
45 QueryHasAssetScheduleFilter,
46 QueryHasImportErrorsFilter,
47 QueryLastDagRunStateFilter,
48 QueryLimit,
49 QueryOffset,
50 QueryOwnersFilter,
51 QueryPausedFilter,
52 QueryTagsFilter,
53 RangeFilter,
54 SortParam,
55 _transform_dag_run_states,
56 datetime_range_filter_factory,
57 filter_param_factory,
58)
59from airflow.api_fastapi.common.router import AirflowRouter
60from airflow.api_fastapi.compat import HTTP_422_UNPROCESSABLE_CONTENT
61from airflow.api_fastapi.core_api.datamodels.dags import (
62 DAGCollectionResponse,
63 DAGDetailsResponse,
64 DAGPatchBody,
65 DAGPatchBodyPartial,
66 DAGResponse,
67)
68from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
69from airflow.api_fastapi.core_api.security import (
70 EditableDagsFilterDep,
71 GetUserDep,
72 ReadableDagsFilterDep,
73 requires_access_dag,
74)
75from airflow.api_fastapi.logging.decorators import action_logging
76from airflow.exceptions import AirflowException, DagNotFound
77from airflow.models import DagModel
78from airflow.models.dag_favorite import DagFavorite
79from airflow.models.dagrun import DagRun
80from airflow.utils.state import DagRunState
82dags_router = AirflowRouter(tags=["DAG"], prefix="/dags")
85@dags_router.get("", dependencies=[Depends(requires_access_dag(method="GET"))])
86def get_dags(
87 limit: QueryLimit,
88 offset: QueryOffset,
89 tags: QueryTagsFilter,
90 owners: QueryOwnersFilter,
91 dag_id_pattern: QueryDagIdPatternSearch,
92 dag_id_prefix_pattern: QueryDagIdPrefixPatternSearch,
93 dag_display_name_pattern: QueryDagDisplayNamePatternSearch,
94 dag_display_name_prefix_pattern: QueryDagDisplayNamePrefixPatternSearch,
95 exclude_stale: QueryExcludeStaleFilter,
96 paused: QueryPausedFilter,
97 has_import_errors: QueryHasImportErrorsFilter,
98 last_dag_run_state: QueryLastDagRunStateFilter,
99 bundle_name: QueryBundleNameFilter,
100 bundle_version: QueryBundleVersionFilter,
101 has_asset_schedule: QueryHasAssetScheduleFilter,
102 asset_dependency: QueryAssetDependencyFilter,
103 dag_run_start_date_range: Annotated[
104 RangeFilter, Depends(datetime_range_filter_factory("dag_run_start_date", DagRun, "start_date"))
105 ],
106 dag_run_end_date_range: Annotated[
107 RangeFilter, Depends(datetime_range_filter_factory("dag_run_end_date", DagRun, "end_date"))
108 ],
109 dag_run_state: Annotated[
110 FilterParam[list[str]],
111 Depends(
112 filter_param_factory(
113 DagRun.state,
114 list[str],
115 FilterOptionEnum.ANY_EQUAL,
116 "dag_run_state",
117 default_factory=list,
118 transform_callable=_transform_dag_run_states,
119 )
120 ),
121 ],
122 order_by: Annotated[
123 SortParam,
124 Depends(
125 SortParam(
126 ["dag_id", "dag_display_name", "next_dagrun", "state", "start_date"],
127 DagModel,
128 {"last_run_state": DagRun.state, "last_run_start_date": DagRun.start_date},
129 ).dynamic_depends()
130 ),
131 ],
132 readable_dags_filter: ReadableDagsFilterDep,
133 session: SessionDep,
134 is_favorite: QueryFavoriteFilter,
135 timetable_type: Annotated[
136 FilterParam[list[str] | None],
137 Depends(filter_param_factory(DagModel.timetable_type, list[str], FilterOptionEnum.IN)),
138 ],
139) -> DAGCollectionResponse:
140 """Get all Dags."""
141 query = generate_dag_with_latest_run_query(
142 max_run_filters=[
143 dag_run_start_date_range,
144 dag_run_end_date_range,
145 dag_run_state,
146 last_dag_run_state,
147 ],
148 order_by=order_by,
149 dag_ids=readable_dags_filter.value,
150 )
152 dags_select, total_entries = paginated_select(
153 statement=query,
154 filters=[
155 exclude_stale,
156 paused,
157 has_import_errors,
158 dag_id_pattern,
159 dag_id_prefix_pattern,
160 dag_display_name_pattern,
161 dag_display_name_prefix_pattern,
162 tags,
163 is_favorite,
164 owners,
165 readable_dags_filter,
166 bundle_name,
167 bundle_version,
168 timetable_type,
169 has_asset_schedule,
170 asset_dependency,
171 ],
172 order_by=order_by,
173 offset=offset,
174 limit=limit,
175 session=session,
176 )
178 dags = session.scalars(dags_select)
180 return DAGCollectionResponse(
181 dags=dags,
182 total_entries=total_entries,
183 )
186@dags_router.get(
187 "/{dag_id}",
188 responses=create_openapi_http_exception_doc(
189 [
190 status.HTTP_400_BAD_REQUEST,
191 status.HTTP_404_NOT_FOUND,
192 HTTP_422_UNPROCESSABLE_CONTENT,
193 ]
194 ),
195 dependencies=[Depends(requires_access_dag(method="GET"))],
196)
197def get_dag(
198 dag_id: str,
199 session: SessionDep,
200 dag_bag: DagBagDep,
201) -> DAGResponse:
202 """Get basic information about a Dag."""
203 dag = get_latest_version_of_dag(dag_bag, dag_id, session)
204 dag_model = session.get(DagModel, dag_id)
205 if not dag_model: 205 ↛ 206line 205 didn't jump to line 206 because the condition on line 205 was never true
206 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Unable to obtain dag with id {dag_id} from session")
208 for key, value in dag.__dict__.items():
209 if not key.startswith("_") and not hasattr(dag_model, key):
210 setattr(dag_model, key, value)
212 return dag_model
215@dags_router.get(
216 "/{dag_id}/details",
217 responses=create_openapi_http_exception_doc(
218 [
219 status.HTTP_400_BAD_REQUEST,
220 status.HTTP_404_NOT_FOUND,
221 ]
222 ),
223 dependencies=[Depends(requires_access_dag(method="GET"))],
224)
225def get_dag_details(
226 dag_id: str, session: SessionDep, dag_bag: DagBagDep, user: GetUserDep
227) -> DAGDetailsResponse:
228 """Get details of Dag."""
229 dag = get_latest_version_of_dag(dag_bag, dag_id, session)
231 dag_model = session.get(DagModel, dag_id)
232 if not dag_model: 232 ↛ 233line 232 didn't jump to line 233 because the condition on line 232 was never true
233 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Unable to obtain dag with id {dag_id} from session")
235 for key, value in dag.__dict__.items():
236 if not key.startswith("_") and not hasattr(dag_model, key):
237 setattr(dag_model, key, value)
239 # Check if this Dag is marked as favorite by the current user
240 user_id = str(user.get_id())
241 is_favorite = (
242 session.scalar(
243 select(DagFavorite.dag_id).where(DagFavorite.user_id == user_id, DagFavorite.dag_id == dag_id)
244 )
245 is not None
246 )
248 # Count only running Dag runs: this stat shows runs that are actually executing right now.
249 active_runs_count = (
250 session.scalar(
251 select(func.count())
252 .select_from(DagRun)
253 .where(DagRun.dag_id == dag_id, DagRun.state == DagRunState.RUNNING)
254 )
255 or 0
256 )
258 # Add is_favorite and active_runs_count fields to the Dag model
259 setattr(dag_model, "is_favorite", is_favorite)
260 setattr(dag_model, "active_runs_count", active_runs_count)
262 return DAGDetailsResponse.model_validate(dag_model)
265@dags_router.patch(
266 "/{dag_id}",
267 responses=create_openapi_http_exception_doc(
268 [
269 status.HTTP_400_BAD_REQUEST,
270 status.HTTP_404_NOT_FOUND,
271 ]
272 ),
273 dependencies=[Depends(requires_access_dag(method="PUT")), Depends(action_logging())],
274)
275def patch_dag(
276 dag_id: str,
277 patch_body: DAGPatchBody,
278 session: SessionDep,
279 update_mask: list[str] | None = Query(None),
280) -> DAGResponse:
281 """Patch the specific Dag."""
282 dag = session.get(DagModel, dag_id)
284 if dag is None: 284 ↛ 285line 284 didn't jump to line 285 because the condition on line 284 was never true
285 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Dag with id: {dag_id} was not found")
287 fields_to_update = patch_body.model_fields_set
288 if update_mask:
289 if update_mask != ["is_paused"]: 289 ↛ 293line 289 didn't jump to line 293 because the condition on line 289 was always true
290 raise HTTPException(
291 status.HTTP_400_BAD_REQUEST, "Only `is_paused` field can be updated through the REST API"
292 )
293 fields_to_update = fields_to_update.intersection(update_mask)
294 try:
295 DAGPatchBodyPartial(**patch_body.model_dump(include=fields_to_update))
296 except ValidationError as e:
297 raise RequestValidationError(errors=e.errors())
298 else:
299 try:
300 DAGPatchBody(**patch_body.model_dump())
301 except ValidationError as e:
302 raise RequestValidationError(errors=e.errors())
304 data = patch_body.model_dump(include=fields_to_update, by_alias=True)
306 for key, val in data.items():
307 setattr(dag, key, val)
309 return dag
312@dags_router.patch(
313 "",
314 responses=create_openapi_http_exception_doc(
315 [
316 status.HTTP_400_BAD_REQUEST,
317 status.HTTP_404_NOT_FOUND,
318 ]
319 ),
320 dependencies=[Depends(requires_access_dag(method="PUT")), Depends(action_logging())],
321)
322def patch_dags(
323 patch_body: DAGPatchBody,
324 limit: QueryLimit,
325 offset: QueryOffset,
326 tags: QueryTagsFilter,
327 owners: QueryOwnersFilter,
328 dag_id_pattern: QueryDagIdPatternSearchWithNone,
329 dag_id_prefix_pattern: QueryDagIdPrefixPatternSearchWithNone,
330 exclude_stale: QueryExcludeStaleFilter,
331 paused: QueryPausedFilter,
332 editable_dags_filter: EditableDagsFilterDep,
333 session: SessionDep,
334 update_mask: list[str] | None = Query(None),
335) -> DAGCollectionResponse:
336 """
337 Patch multiple Dags.
339 If neither `dag_id_pattern` nor `dag_id_prefix_pattern` is provided, no Dags will be
340 matched regardless of other filters. To match all Dags, pass a wildcard value such as
341 `~` or `%` for `dag_id_pattern`.
342 """
343 if update_mask:
344 if update_mask != ["is_paused"]: 344 ↛ 354line 344 didn't jump to line 354 because the condition on line 344 was always true
345 raise HTTPException(
346 status.HTTP_400_BAD_REQUEST, "Only `is_paused` field can be updated through the REST API"
347 )
348 else:
349 try:
350 DAGPatchBody.model_validate(patch_body)
351 except ValidationError as e:
352 raise RequestValidationError(errors=e.errors())
354 dags_select, total_entries = paginated_select(
355 statement=select(DagModel),
356 filters=[
357 exclude_stale,
358 paused,
359 dag_id_pattern,
360 dag_id_prefix_pattern,
361 tags,
362 owners,
363 editable_dags_filter,
364 ],
365 order_by=None,
366 offset=offset,
367 limit=limit,
368 session=session,
369 )
370 dags = session.scalars(dags_select).all()
372 filtered_dag_ids = apply_filters_to_select(
373 statement=select(DagModel.dag_id),
374 filters=[
375 exclude_stale,
376 paused,
377 dag_id_pattern,
378 dag_id_prefix_pattern,
379 tags,
380 owners,
381 editable_dags_filter,
382 ],
383 ).subquery()
385 session.execute(
386 update(DagModel)
387 .where(DagModel.dag_id.in_(select(filtered_dag_ids.c.dag_id)))
388 .values(is_paused=patch_body.is_paused)
389 .execution_options(synchronize_session="fetch")
390 )
392 return DAGCollectionResponse(
393 dags=dags,
394 total_entries=total_entries,
395 )
398@dags_router.post(
399 "/{dag_id}/favorite",
400 status_code=status.HTTP_204_NO_CONTENT,
401 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
402 dependencies=[Depends(requires_access_dag(method="GET")), Depends(action_logging())],
403)
404def favorite_dag(dag_id: str, session: SessionDep, user: GetUserDep):
405 """Mark the Dag as favorite."""
406 dag = session.get(DagModel, dag_id)
407 if not dag:
408 raise HTTPException(status.HTTP_404_NOT_FOUND, detail=f"Dag with id '{dag_id}' not found")
410 user_id = str(user.get_id())
411 session.execute(insert(DagFavorite).values(dag_id=dag_id, user_id=user_id))
414@dags_router.post(
415 "/{dag_id}/unfavorite",
416 status_code=status.HTTP_204_NO_CONTENT,
417 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND, status.HTTP_409_CONFLICT]),
418 dependencies=[Depends(requires_access_dag(method="GET")), Depends(action_logging())],
419)
420def unfavorite_dag(dag_id: str, session: SessionDep, user: GetUserDep):
421 """Unmark the Dag as favorite."""
422 dag = session.get(DagModel, dag_id)
423 if not dag:
424 raise HTTPException(status.HTTP_404_NOT_FOUND, detail=f"Dag with id '{dag_id}' not found")
426 user_id = str(user.get_id())
428 favorite_exists = session.execute(
429 select(DagFavorite)
430 .where(
431 DagFavorite.dag_id == dag_id,
432 DagFavorite.user_id == user_id,
433 )
434 .limit(1)
435 ).first()
437 if not favorite_exists:
438 raise HTTPException(status.HTTP_409_CONFLICT, detail="Dag is not marked as favorite")
440 session.execute(
441 delete(DagFavorite).where(
442 DagFavorite.dag_id == dag_id,
443 DagFavorite.user_id == user_id,
444 )
445 )
448@dags_router.delete(
449 "/{dag_id}",
450 responses=create_openapi_http_exception_doc(
451 [
452 status.HTTP_400_BAD_REQUEST,
453 status.HTTP_404_NOT_FOUND,
454 status.HTTP_409_CONFLICT,
455 HTTP_422_UNPROCESSABLE_CONTENT,
456 ]
457 ),
458 dependencies=[Depends(requires_access_dag(method="DELETE")), Depends(action_logging())],
459)
460def delete_dag(
461 dag_id: str,
462 session: SessionDep,
463) -> Response:
464 """Delete the specific Dag."""
465 try:
466 delete_dag_module.delete_dag(dag_id, session=session)
467 except DagNotFound:
468 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Dag with id: {dag_id} was not found")
469 except AirflowException:
470 raise HTTPException(
471 status.HTTP_409_CONFLICT, f"Task instances of dag with id: '{dag_id}' are still running"
472 )
473 return Response(status_code=status.HTTP_204_NO_CONTENT)