Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/dag_run.py: 84%
231 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
20import datetime
21import textwrap
22from typing import Annotated, Literal, cast
24from fastapi import Depends, HTTPException, Query, Request, status
25from fastapi.exceptions import RequestValidationError
26from fastapi.responses import StreamingResponse
27from pydantic import ValidationError
28from sqlalchemy import select
29from sqlalchemy.orm import joinedload
31from airflow.api_fastapi.app import get_auth_manager
32from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity, DagDetails
33from airflow.api_fastapi.common.cursors import (
34 apply_cursor_filter,
35 encode_cursor,
36 make_backward_cursor,
37 parse_cursor,
38)
39from airflow.api_fastapi.common.dagbag import DagBagDep, get_dag_for_run, get_latest_version_of_dag
40from airflow.api_fastapi.common.db.common import (
41 SessionDep,
42 apply_filters_to_select,
43 bounded_total_entries,
44 paginated_select,
45)
46from airflow.api_fastapi.common.db.dag_runs import (
47 attach_dag_versions_to_runs,
48 eager_load_dag_run_for_list,
49)
50from airflow.api_fastapi.common.parameters import (
51 FilterOptionEnum,
52 FilterParam,
53 LimitFilter,
54 OffsetFilter,
55 QueryConsumingAssetPatternSearch,
56 QueryDagRunPartitionKeyPrefixSearch,
57 QueryDagRunPartitionKeySearch,
58 QueryDagRunRunTypesFilter,
59 QueryDagRunStateFilter,
60 QueryDagRunVersionFilter,
61 QueryLimit,
62 QueryOffset,
63 Range,
64 RangeFilter,
65 SortParam,
66 _PrefixSearchParam,
67 _SearchParam,
68 datetime_range_filter_factory,
69 filter_param_factory,
70 float_range_filter_factory,
71 prefix_search_param_factory,
72 search_param_factory,
73)
74from airflow.api_fastapi.common.router import AirflowRouter
75from airflow.api_fastapi.common.types import Mimetype
76from airflow.api_fastapi.core_api.base import OrmClause
77from airflow.api_fastapi.core_api.datamodels.assets import AssetEventCollectionResponse
78from airflow.api_fastapi.core_api.datamodels.common import BulkBody, BulkResponse
79from airflow.api_fastapi.core_api.datamodels.dag_run import (
80 BulkDAGRunBody,
81 BulkDAGRunClearBody,
82 ClearPartitionsBody,
83 ClearPartitionsResponse,
84 DAGRunClearBody,
85 DAGRunCollectionResponse,
86 DagRunMutableStates,
87 DAGRunPatchBody,
88 DAGRunResponse,
89 DAGRunsBatchBody,
90 TriggerDAGRunPostBody,
91)
92from airflow.api_fastapi.core_api.datamodels.task_instances import (
93 ClearTaskInstanceCollectionResponse,
94 NewTaskResponse,
95 TaskInstanceResponse,
96)
97from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
98from airflow.api_fastapi.core_api.security import (
99 GetUserDep,
100 ReadableDagRunsFilterDep,
101 requires_access_asset,
102 requires_access_dag,
103 requires_access_dag_run_bulk,
104 requires_access_dag_run_clear_bulk,
105)
106from airflow.api_fastapi.core_api.services.public.assets import serialize_asset_events
107from airflow.api_fastapi.core_api.services.public.dag_run import (
108 BulkDagRunService,
109 DagRunWaiter,
110 clear_partition_fields,
111 dry_run_clear_dag_run,
112 get_dag_run_and_dag_for_clear,
113 patch_dag_run_note,
114 patch_dag_run_state,
115 perform_clear_dag_run,
116)
117from airflow.api_fastapi.logging.decorators import action_logging
118from airflow.exceptions import ParamValidationError
119from airflow.models import DagModel, DagRun
120from airflow.models.asset import AssetEvent
121from airflow.models.dag_version import DagVersion
122from airflow.utils.state import DagRunState
123from airflow.utils.types import DagRunTriggeredByType, DagRunType
125dag_run_router = AirflowRouter(tags=["DagRun"], prefix="/dags/{dag_id}/dagRuns")
126dag_run_at_dag_router = AirflowRouter(tags=["DagRun"], prefix="/dags/{dag_id}")
129@dag_run_router.get(
130 "/{dag_run_id}",
131 responses=create_openapi_http_exception_doc(
132 [
133 status.HTTP_404_NOT_FOUND,
134 ]
135 ),
136 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN))],
137)
138def get_dag_run(dag_id: str, dag_run_id: str, session: SessionDep) -> DAGRunResponse:
139 dag_run = session.scalar(
140 select(DagRun).filter_by(dag_id=dag_id, run_id=dag_run_id).options(joinedload(DagRun.dag_model))
141 )
142 if dag_run is None:
143 raise HTTPException(
144 status.HTTP_404_NOT_FOUND,
145 f"The DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` was not found",
146 )
147 return dag_run
150@dag_run_router.delete(
151 "/{dag_run_id}",
152 status_code=status.HTTP_204_NO_CONTENT,
153 responses=create_openapi_http_exception_doc(
154 [
155 status.HTTP_400_BAD_REQUEST,
156 status.HTTP_404_NOT_FOUND,
157 status.HTTP_409_CONFLICT,
158 ],
159 ),
160 dependencies=[
161 Depends(requires_access_dag(method="DELETE", access_entity=DagAccessEntity.RUN)),
162 Depends(action_logging()),
163 ],
164)
165def delete_dag_run(dag_id: str, dag_run_id: str, session: SessionDep):
166 """Delete a Dag Run entry."""
167 dag_run = session.scalar(select(DagRun).filter_by(dag_id=dag_id, run_id=dag_run_id))
168 deletable_states = {s.value for s in DagRunMutableStates}
170 if dag_run is None:
171 raise HTTPException(
172 status.HTTP_404_NOT_FOUND,
173 f"The DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` was not found",
174 )
175 if dag_run.state not in deletable_states:
176 raise HTTPException(
177 status.HTTP_409_CONFLICT,
178 (
179 f"The DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` "
180 f"cannot be deleted in {dag_run.state} state"
181 ),
182 )
184 session.delete(dag_run)
187@dag_run_router.patch(
188 "/{dag_run_id}",
189 responses=create_openapi_http_exception_doc(
190 [
191 status.HTTP_400_BAD_REQUEST,
192 status.HTTP_404_NOT_FOUND,
193 ],
194 ),
195 dependencies=[
196 Depends(requires_access_dag(method="PUT", access_entity=DagAccessEntity.RUN)),
197 Depends(action_logging()),
198 ],
199)
200def patch_dag_run(
201 dag_id: str,
202 dag_run_id: str,
203 patch_body: DAGRunPatchBody,
204 session: SessionDep,
205 dag_bag: DagBagDep,
206 user: GetUserDep,
207 update_mask: list[str] | None = Query(None),
208) -> DAGRunResponse:
209 """Modify a Dag Run."""
210 dag_run = session.scalar(
211 select(DagRun).filter_by(dag_id=dag_id, run_id=dag_run_id).options(joinedload(DagRun.dag_model))
212 )
213 if dag_run is None:
214 raise HTTPException(
215 status.HTTP_404_NOT_FOUND,
216 f"The DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` was not found",
217 )
219 dag = get_dag_for_run(dag_bag, dag_run, session=session)
221 fields_to_update = patch_body.model_fields_set
223 if update_mask:
224 fields_to_update = fields_to_update.intersection(update_mask)
225 else:
226 try:
227 DAGRunPatchBody(**patch_body.model_dump())
228 except ValidationError as e:
229 raise RequestValidationError(errors=e.errors())
231 data = patch_body.model_dump(include=fields_to_update, by_alias=True)
233 # Apply "note" before "state" so listeners fired inside patch_dag_run_state() see the updated note.
234 if "note" in data:
235 updated_dag_run = session.get(DagRun, dag_run.id)
236 if updated_dag_run is not None: 236 ↛ 238line 236 didn't jump to line 238 because the condition on line 236 was always true
237 patch_dag_run_note(dag_run=updated_dag_run, note=data["note"], user=user)
238 if "state" in data and patch_body.state is not None:
239 patch_dag_run_state(dag=dag, dag_run=dag_run, state=patch_body.state, session=session)
241 final_dag_run = session.get(DagRun, dag_run.id)
242 if not final_dag_run: 242 ↛ 243line 242 didn't jump to line 243 because the condition on line 242 was never true
243 raise HTTPException(status.HTTP_404_NOT_FOUND, "Dag run not found after update")
245 return final_dag_run
248@dag_run_router.patch(
249 "",
250 dependencies=[Depends(requires_access_dag_run_bulk()), Depends(action_logging())],
251)
252def bulk_dag_runs(
253 request: BulkBody[BulkDAGRunBody],
254 session: SessionDep,
255 dag_id: str,
256 dag_bag: DagBagDep,
257 user: GetUserDep,
258) -> BulkResponse:
259 """Bulk update or delete Dag Runs."""
260 return BulkDagRunService(
261 session=session, request=request, dag_id=dag_id, dag_bag=dag_bag, user=user
262 ).handle_request()
265@dag_run_router.get(
266 "/{dag_run_id}/upstreamAssetEvents",
267 responses=create_openapi_http_exception_doc(
268 [
269 status.HTTP_404_NOT_FOUND,
270 ]
271 ),
272 dependencies=[
273 Depends(requires_access_asset(method="GET")),
274 Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN)),
275 ],
276)
277def get_upstream_asset_events(
278 dag_id: str, dag_run_id: str, session: SessionDep
279) -> AssetEventCollectionResponse:
280 """If dag run is asset-triggered, return the asset events that triggered it."""
281 dag_run: DagRun | None = session.scalar(
282 select(DagRun)
283 .where(
284 DagRun.dag_id == dag_id,
285 DagRun.run_id == dag_run_id,
286 )
287 .options(
288 joinedload(DagRun.consumed_asset_events).joinedload(AssetEvent.asset),
289 joinedload(DagRun.consumed_asset_events).subqueryload(AssetEvent.created_dagruns),
290 )
291 )
292 if dag_run is None:
293 raise HTTPException(
294 status.HTTP_404_NOT_FOUND,
295 f"The DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` was not found",
296 )
297 events = dag_run.consumed_asset_events
298 return AssetEventCollectionResponse(
299 asset_events=serialize_asset_events(events, session=session),
300 total_entries=len(events),
301 )
304@dag_run_router.post(
305 "/{dag_run_id}/clear",
306 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
307 dependencies=[
308 Depends(requires_access_dag(method="PUT", access_entity=DagAccessEntity.RUN)),
309 Depends(action_logging()),
310 ],
311)
312def clear_dag_run(
313 dag_id: str,
314 dag_run_id: str,
315 body: DAGRunClearBody,
316 dag_bag: DagBagDep,
317 session: SessionDep,
318 user: GetUserDep,
319) -> ClearTaskInstanceCollectionResponse | DAGRunResponse:
320 dag_run, dag = get_dag_run_and_dag_for_clear(
321 session=session, dag_bag=dag_bag, dag_id=dag_id, dag_run_id=dag_run_id
322 )
324 if body.dry_run:
325 task_instances = dry_run_clear_dag_run(
326 session=session,
327 dag_bag=dag_bag,
328 dag_id=dag_id,
329 dag_run_id=dag_run_id,
330 only_failed=body.only_failed,
331 only_new=body.only_new,
332 )
333 return ClearTaskInstanceCollectionResponse(
334 task_instances=task_instances,
335 total_entries=len(task_instances),
336 )
338 return perform_clear_dag_run(
339 session=session,
340 dag=dag,
341 dag_run=dag_run,
342 dag_id=dag_id,
343 only_failed=body.only_failed,
344 only_new=body.only_new,
345 run_on_latest_version=body.run_on_latest_version,
346 note=body.note,
347 user=user,
348 )
351@dag_run_at_dag_router.post(
352 "/clearDagRuns",
353 responses=create_openapi_http_exception_doc([status.HTTP_400_BAD_REQUEST, status.HTTP_404_NOT_FOUND]),
354 dependencies=[Depends(requires_access_dag_run_clear_bulk()), Depends(action_logging())],
355)
356def clear_dag_runs(
357 dag_id: str,
358 body: BulkDAGRunClearBody,
359 dag_bag: DagBagDep,
360 session: SessionDep,
361 user: GetUserDep,
362) -> ClearTaskInstanceCollectionResponse | DAGRunCollectionResponse:
363 """Clear multiple Dag Runs in a single request."""
364 url_dag_id_is_wildcard = dag_id == "~"
366 partition_mode = not body.dag_runs and body.has_partition_selectors
368 if partition_mode:
369 if url_dag_id_is_wildcard: 369 ↛ 370line 369 didn't jump to line 370 because the condition on line 369 was never true
370 raise HTTPException(
371 status.HTTP_400_BAD_REQUEST,
372 "Partition selectors require a concrete dag_id; '~' is not supported.",
373 )
374 dag = get_latest_version_of_dag(dag_bag, dag_id, session)
376 stmt = select(DagRun.run_id).where(DagRun.dag_id == dag_id)
377 if body.partition_key is not None:
378 stmt = stmt.where(DagRun.partition_key == body.partition_key)
379 else:
380 stmt = stmt.where(DagRun.partition_date.is_not(None))
381 stmt = DagRun.apply_partition_date_window(
382 stmt,
383 timetable=dag.timetable,
384 start=body.partition_date_start,
385 end=body.partition_date_end,
386 )
387 stmt = stmt.order_by(DagRun.partition_date, DagRun.run_id)
389 runs_to_clear: dict[tuple[str, str], None] = {
390 (dag_id, run_id): None for run_id in session.scalars(stmt)
391 }
392 else:
393 # No ordered set type in Python, using a dict with throwaway values as replacement.
394 runs_to_clear = {}
395 for run in body.dag_runs:
396 if url_dag_id_is_wildcard: 396 ↛ 397line 396 didn't jump to line 397 because the condition on line 396 was never true
397 if not run.dag_id or run.dag_id == "~":
398 raise HTTPException(
399 status.HTTP_400_BAD_REQUEST,
400 f"When the URL dag_id is '~', every entry must provide a concrete dag_id "
401 f"(missing on dag_run_id: {run.dag_run_id!r}).",
402 )
403 run_to_clear = (run.dag_id, run.dag_run_id)
404 else:
405 entity_dag_id = run.dag_id or dag_id
406 if entity_dag_id != dag_id:
407 raise HTTPException(
408 status.HTTP_400_BAD_REQUEST,
409 f"Entry dag_id {entity_dag_id!r} does not match the URL dag_id {dag_id!r}.",
410 )
411 run_to_clear = (dag_id, run.dag_run_id)
412 runs_to_clear[run_to_clear] = None
414 if body.dry_run:
415 affected: list[TaskInstanceResponse | NewTaskResponse] = []
416 for run_dag_id, run_id in runs_to_clear: 416 ↛ 417line 416 didn't jump to line 417 because the loop on line 416 never started
417 get_dag_run_and_dag_for_clear(
418 session=session, dag_bag=dag_bag, dag_id=run_dag_id, dag_run_id=run_id
419 )
420 affected.extend(
421 dry_run_clear_dag_run(
422 session=session,
423 dag_bag=dag_bag,
424 dag_id=run_dag_id,
425 dag_run_id=run_id,
426 only_failed=body.only_failed,
427 only_new=body.only_new,
428 )
429 )
430 return ClearTaskInstanceCollectionResponse(
431 task_instances=affected,
432 total_entries=len(affected),
433 )
435 cleared_runs: list[DagRun] = []
436 for run_dag_id, run_id in runs_to_clear:
437 dag_run, dag = get_dag_run_and_dag_for_clear(
438 session=session, dag_bag=dag_bag, dag_id=run_dag_id, dag_run_id=run_id
439 )
440 cleared_runs.append(
441 perform_clear_dag_run(
442 session=session,
443 dag=dag,
444 dag_run=dag_run,
445 dag_id=run_dag_id,
446 only_failed=body.only_failed,
447 only_new=body.only_new,
448 run_on_latest_version=body.run_on_latest_version,
449 note=body.note,
450 user=user,
451 )
452 )
453 return DAGRunCollectionResponse(
454 dag_runs=cleared_runs,
455 total_entries=len(cleared_runs),
456 )
459@dag_run_at_dag_router.post(
460 "/clearPartitions",
461 responses=create_openapi_http_exception_doc([status.HTTP_400_BAD_REQUEST, status.HTTP_404_NOT_FOUND]),
462 dependencies=[
463 Depends(requires_access_dag(method="PUT", access_entity=DagAccessEntity.RUN)),
464 Depends(action_logging()),
465 ],
466)
467def clear_dag_run_partitions(
468 dag_id: str,
469 body: ClearPartitionsBody,
470 dag_bag: DagBagDep,
471 session: SessionDep,
472) -> ClearPartitionsResponse:
473 """Reset partition_key and partition_date fields on matching Dag Runs."""
474 dag = get_latest_version_of_dag(dag_bag, dag_id, session)
475 dag_runs_cleared, task_instances_cleared = clear_partition_fields(
476 dag=dag,
477 body=body,
478 dag_id=dag_id,
479 session=session,
480 )
481 return ClearPartitionsResponse(
482 dag_runs_cleared=dag_runs_cleared,
483 task_instances_cleared=task_instances_cleared,
484 dry_run=body.dry_run,
485 )
488@dag_run_router.get(
489 "",
490 responses=create_openapi_http_exception_doc(
491 [
492 status.HTTP_400_BAD_REQUEST,
493 status.HTTP_404_NOT_FOUND,
494 ]
495 ),
496 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN))],
497)
498def get_dag_runs(
499 dag_id: str,
500 limit: QueryLimit,
501 offset: QueryOffset,
502 run_after: Annotated[RangeFilter, Depends(datetime_range_filter_factory("run_after", DagRun))],
503 logical_date: Annotated[RangeFilter, Depends(datetime_range_filter_factory("logical_date", DagRun))],
504 start_date_range: Annotated[RangeFilter, Depends(datetime_range_filter_factory("start_date", DagRun))],
505 end_date_range: Annotated[RangeFilter, Depends(datetime_range_filter_factory("end_date", DagRun))],
506 duration_range: Annotated[RangeFilter, Depends(float_range_filter_factory("duration", DagRun))],
507 update_at_range: Annotated[RangeFilter, Depends(datetime_range_filter_factory("updated_at", DagRun))],
508 conf_contains: Annotated[
509 FilterParam[str],
510 Depends(filter_param_factory(DagRun.conf, str, FilterOptionEnum.CONTAINS, "conf_contains")),
511 ],
512 run_type: QueryDagRunRunTypesFilter,
513 state: QueryDagRunStateFilter,
514 dag_version: QueryDagRunVersionFilter,
515 bundle_version: Annotated[
516 FilterParam[str | None], Depends(filter_param_factory(DagRun.bundle_version, str | None))
517 ],
518 order_by: Annotated[
519 SortParam,
520 Depends(
521 SortParam(
522 [
523 "id",
524 "state",
525 "dag_id",
526 "run_id",
527 "logical_date",
528 "partition_date",
529 "run_after",
530 "start_date",
531 "end_date",
532 "updated_at",
533 "conf",
534 "duration",
535 ],
536 DagRun,
537 {"dag_run_id": "run_id"},
538 ).dynamic_depends(default="id")
539 ),
540 ],
541 readable_dag_runs_filter: ReadableDagRunsFilterDep,
542 session: SessionDep,
543 dag_bag: DagBagDep,
544 run_id_pattern: Annotated[_SearchParam, Depends(search_param_factory(DagRun.run_id, "run_id_pattern"))],
545 run_id_prefix_pattern: Annotated[
546 _PrefixSearchParam,
547 Depends(prefix_search_param_factory(DagRun.run_id, "run_id_prefix_pattern")),
548 ],
549 triggering_user_name_pattern: Annotated[
550 _SearchParam,
551 Depends(search_param_factory(DagRun.triggering_user_name, "triggering_user_name_pattern")),
552 ],
553 triggering_user_name_prefix_pattern: Annotated[
554 _PrefixSearchParam,
555 Depends(
556 prefix_search_param_factory(DagRun.triggering_user_name, "triggering_user_name_prefix_pattern")
557 ),
558 ],
559 dag_id_pattern: Annotated[_SearchParam, Depends(search_param_factory(DagRun.dag_id, "dag_id_pattern"))],
560 dag_id_prefix_pattern: Annotated[
561 _PrefixSearchParam,
562 Depends(prefix_search_param_factory(DagRun.dag_id, "dag_id_prefix_pattern")),
563 ],
564 partition_key_pattern: QueryDagRunPartitionKeySearch,
565 partition_key_prefix_pattern: QueryDagRunPartitionKeyPrefixSearch,
566 consuming_asset_pattern: QueryConsumingAssetPatternSearch,
567 partition_date_gte: datetime.date | None = Query(
568 None,
569 description=(
570 "Inclusive lower bound of the partition_date window, interpreted as a local calendar "
571 "day in the Dag's timetable timezone. Runs from the start of this day onwards match."
572 ),
573 ),
574 partition_date_lte: datetime.date | None = Query(
575 None,
576 description=(
577 "Inclusive upper bound of the partition_date window, interpreted as a local calendar "
578 "day in the Dag's timetable timezone. The whole day is included: runs up to the end "
579 "of this day match."
580 ),
581 ),
582 cursor: str | None = Query(
583 None,
584 description="Cursor for keyset-based pagination. "
585 "Pass an empty string for the first page, then use ``next_cursor`` from the response. "
586 "When ``cursor`` is provided, ``offset`` is ignored.",
587 ),
588) -> DAGRunCollectionResponse:
589 """
590 Get all Dag Runs.
592 This endpoint allows specifying `~` as the dag_id to retrieve Dag Runs for all Dags.
594 Supports two pagination modes:
596 **Offset (default):** use `limit` and `offset` query parameters. Returns `total_entries`.
598 **Cursor:** pass `cursor` (empty string for the first page, then `next_cursor` from the response).
599 When `cursor` is provided, `offset` is ignored and `total_entries` is capped at
600 `total_entries_limit` (a value equal to that limit means at least that many runs match).
601 ``next_cursor`` is ``null`` when there are no more pages; ``previous_cursor`` is ``null``
602 on the first page.
603 """
604 use_cursor = cursor is not None
605 query = select(DagRun).options(*eager_load_dag_run_for_list())
607 has_partition_date_filter = partition_date_gte is not None or partition_date_lte is not None
609 if dag_id == "~": 609 ↛ 610line 609 didn't jump to line 610 because the condition on line 609 was never true
610 if has_partition_date_filter:
611 raise HTTPException(
612 status.HTTP_400_BAD_REQUEST,
613 "partition_date_gte and partition_date_lte require a specific dag_id.",
614 )
615 else:
616 dag = get_latest_version_of_dag(dag_bag, dag_id, session) # Check if the Dag exists.
617 query = query.filter(DagRun.dag_id == dag_id).options()
618 if has_partition_date_filter:
619 # Runs of a non-partitioned Dag never carry a partition_date (this includes
620 # partitioned-at-runtime Dags, whose runs keep it NULL), so the filter would
621 # silently match nothing; reject it instead.
622 if not dag.timetable.partitioned: 622 ↛ 630line 622 didn't jump to line 630 because the condition on line 622 was always true
623 raise HTTPException(
624 status.HTTP_400_BAD_REQUEST,
625 f"Dag with dag_id: '{dag_id}' is not partitioned; "
626 "partition_date_gte and partition_date_lte are not supported.",
627 )
628 # The bounds are calendar days, so the whole of partition_date_lte belongs to the
629 # window: widen it to the following local midnight and exclude that edge.
630 query = DagRun.apply_partition_date_window(
631 query,
632 timetable=dag.timetable,
633 start=(
634 datetime.datetime.combine(partition_date_gte, datetime.time.min)
635 if partition_date_gte is not None
636 else None
637 ),
638 end=(
639 datetime.datetime.combine(
640 partition_date_lte + datetime.timedelta(days=1), datetime.time.min
641 )
642 if partition_date_lte is not None
643 else None
644 ),
645 end_exclusive=True,
646 )
648 # Add join with DagVersion if dag_version filter is active
649 if dag_version.value:
650 query = query.join(DagVersion, DagRun.created_dag_version_id == DagVersion.id)
652 filters: list[OrmClause] = [
653 run_after,
654 logical_date,
655 start_date_range,
656 end_date_range,
657 update_at_range,
658 duration_range,
659 conf_contains,
660 state,
661 run_type,
662 dag_version,
663 bundle_version,
664 readable_dag_runs_filter,
665 run_id_pattern,
666 run_id_prefix_pattern,
667 triggering_user_name_pattern,
668 triggering_user_name_prefix_pattern,
669 dag_id_pattern,
670 dag_id_prefix_pattern,
671 partition_key_pattern,
672 partition_key_prefix_pattern,
673 consuming_asset_pattern,
674 ]
676 if use_cursor:
677 # Fetch one extra row so we can detect whether a next page exists.
678 page_limit = cast(
679 "int", limit.value
680 ) # LimitFilter value is guaranteed to be set to the default value of QueryLimit
681 cursor_limit = LimitFilter().set_value(page_limit + 1)
682 dag_run_select = apply_filters_to_select(statement=query, filters=[*filters, cursor_limit])
683 dag_run_select = order_by.to_orm(dag_run_select)
685 is_backward = False
686 if cursor:
687 token, is_backward = parse_cursor(cursor)
688 if is_backward: 688 ↛ 689line 688 didn't jump to line 689 because the condition on line 688 was never true
689 dag_run_select = order_by.to_orm(dag_run_select, reversed=True)
690 dag_run_select = apply_cursor_filter(
691 dag_run_select, token, order_by, session.get_bind().dialect.name, is_backward=is_backward
692 )
694 fetched = list(session.scalars(dag_run_select))
695 has_more = len(fetched) > page_limit
696 dag_runs = fetched[:page_limit]
698 if is_backward: 698 ↛ 699line 698 didn't jump to line 699 because the condition on line 698 was never true
699 dag_runs.reverse()
700 has_prev = has_more
701 has_next = True
702 else:
703 has_prev = bool(cursor)
704 has_next = has_more
706 attach_dag_versions_to_runs(dag_runs, session=session)
708 total_entries, total_entries_limit = bounded_total_entries(
709 statement=query, filters=filters, session=session
710 )
711 return DAGRunCollectionResponse(
712 dag_runs=dag_runs,
713 total_entries=total_entries,
714 total_entries_limit=total_entries_limit,
715 next_cursor=(encode_cursor(dag_runs[-1], order_by) if has_next and dag_runs else None),
716 previous_cursor=(
717 make_backward_cursor(encode_cursor(dag_runs[0], order_by)) if has_prev and dag_runs else None
718 ),
719 )
721 dag_run_select, total_entries = paginated_select(
722 statement=query,
723 filters=filters,
724 order_by=order_by,
725 offset=offset,
726 limit=limit,
727 session=session,
728 )
729 dag_runs = list(session.scalars(dag_run_select))
730 attach_dag_versions_to_runs(dag_runs, session=session)
732 return DAGRunCollectionResponse(
733 dag_runs=dag_runs,
734 total_entries=total_entries,
735 )
738@dag_run_router.post(
739 "",
740 responses=create_openapi_http_exception_doc(
741 [
742 status.HTTP_400_BAD_REQUEST,
743 status.HTTP_404_NOT_FOUND,
744 status.HTTP_409_CONFLICT,
745 ]
746 ),
747 dependencies=[
748 Depends(requires_access_dag(method="POST", access_entity=DagAccessEntity.RUN)),
749 Depends(action_logging()),
750 ],
751)
752def trigger_dag_run(
753 dag_id,
754 body: TriggerDAGRunPostBody,
755 dag_bag: DagBagDep,
756 user: GetUserDep,
757 session: SessionDep,
758 request: Request,
759) -> DAGRunResponse:
760 """Trigger a Dag."""
761 dm = session.scalar(select(DagModel).where(~DagModel.is_stale, DagModel.dag_id == dag_id).limit(1))
762 if not dm:
763 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Dag with dag_id: '{dag_id}' not found")
765 if dm.has_import_errors: 765 ↛ 766line 765 didn't jump to line 766 because the condition on line 765 was never true
766 raise HTTPException(
767 status.HTTP_400_BAD_REQUEST,
768 f"Dag with dag_id: '{dag_id}' has import errors and cannot be triggered",
769 )
771 if dm.allowed_run_types is not None and DagRunType.MANUAL not in dm.allowed_run_types: 771 ↛ 772line 771 didn't jump to line 772 because the condition on line 771 was never true
772 raise HTTPException(
773 status.HTTP_400_BAD_REQUEST,
774 f"Dag with dag_id: '{dag_id}' does not allow manual runs",
775 )
777 referer = request.headers.get("referer")
778 if referer: 778 ↛ 779line 778 didn't jump to line 779 because the condition on line 778 was never true
779 triggered_by = DagRunTriggeredByType.UI
780 else:
781 triggered_by = DagRunTriggeredByType.REST_API
783 dag = get_latest_version_of_dag(dag_bag, dag_id, session)
784 try:
785 params = body.validate_context(dag)
786 dag_run = dag.create_dagrun(
787 run_id=params["run_id"],
788 logical_date=params["logical_date"],
789 data_interval=params["data_interval"],
790 run_after=params["run_after"],
791 conf=params["conf"],
792 run_type=DagRunType.MANUAL,
793 triggered_by=triggered_by,
794 triggering_user_name=user.get_name(),
795 state=DagRunState.QUEUED,
796 partition_key=params["partition_key"],
797 partition_date=params["partition_date"],
798 session=session,
799 )
800 except (ParamValidationError, ValueError) as e:
801 raise HTTPException(status.HTTP_400_BAD_REQUEST, str(e)) from e
803 dag_run_note = body.note
804 if dag_run_note:
805 current_user_id = user.get_id()
806 dag_run.note = (dag_run_note, current_user_id)
807 return dag_run
810@dag_run_router.get(
811 "/{dag_run_id}/wait",
812 tags=["experimental"],
813 summary="Experimental: Wait for a dag run to complete, and return task results if requested.",
814 description="🚧 This is an experimental endpoint and may change or be removed without notice.Successful response are streamed as newline-delimited JSON (NDJSON). Each line is a JSON object representing the Dag run state.",
815 responses={
816 **create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
817 status.HTTP_200_OK: {
818 "description": "Successful Response",
819 "content": {
820 Mimetype.NDJSON: {
821 "schema": {
822 "type": "string",
823 "example": textwrap.dedent(
824 """\
825 {"state": "running"}
826 {"state": "success", "results": {"op": 42}}
827 """
828 ),
829 }
830 }
831 },
832 },
833 },
834 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN))],
835)
836def wait_dag_run_until_finished(
837 dag_id: str,
838 dag_run_id: str,
839 session: SessionDep,
840 user: GetUserDep,
841 interval: Annotated[float, Query(gt=0.0, description="Seconds to wait between dag run state checks")],
842 result_task_ids: Annotated[
843 list[str] | None,
844 Query(
845 alias="result",
846 description=(
847 "Collect result XCom from task. Can be set multiple times. "
848 "If unset, return value of the return task as specified in the "
849 "dag (in present) is returned by default."
850 ),
851 ),
852 ] = None,
853):
854 "Wait for a dag run until it finishes, and return its result(s)."
855 if not get_auth_manager().is_authorized_dag(
856 method="GET",
857 access_entity=DagAccessEntity.XCOM,
858 # The route dependency above already authorizes RUN access with the Dag's team resolved;
859 # this second, XCom-specific check has to resolve it the same way, or the two checks ask
860 # a team-aware auth manager about differently-scoped resources.
861 details=DagDetails(id=dag_id, team_name=DagModel.get_team_name(dag_id, session=session)),
862 user=user,
863 ):
864 if result_task_ids:
865 raise HTTPException(
866 status.HTTP_403_FORBIDDEN,
867 "User is not authorized to read XCom data for this Dag",
868 )
869 result_task_ids = [] # Explicitly not returning any XCom results.
870 if not session.scalar(select(1).where(DagRun.dag_id == dag_id, DagRun.run_id == dag_run_id)):
871 raise HTTPException(
872 status.HTTP_404_NOT_FOUND,
873 f"The DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` was not found",
874 )
875 waiter = DagRunWaiter(
876 dag_id=dag_id,
877 run_id=dag_run_id,
878 interval=interval,
879 result_task_ids=result_task_ids,
880 )
881 return StreamingResponse(waiter.wait())
884@dag_run_router.post(
885 "/list",
886 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
887 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN))],
888)
889def get_list_dag_runs_batch(
890 dag_id: Literal["~"],
891 body: DAGRunsBatchBody,
892 readable_dag_runs_filter: ReadableDagRunsFilterDep,
893 session: SessionDep,
894) -> DAGRunCollectionResponse:
895 """Get a list of Dag Runs."""
896 dag_ids = FilterParam(DagRun.dag_id, body.dag_ids, FilterOptionEnum.IN) # type: ignore[arg-type]
897 logical_date = RangeFilter(
898 Range(
899 lower_bound_gte=body.logical_date_gte,
900 lower_bound_gt=body.logical_date_gt,
901 upper_bound_lte=body.logical_date_lte,
902 upper_bound_lt=body.logical_date_lt,
903 ),
904 attribute=DagRun.logical_date, # type: ignore[arg-type]
905 )
906 run_after = RangeFilter(
907 Range(
908 lower_bound_gte=body.run_after_gte,
909 lower_bound_gt=body.run_after_gt,
910 upper_bound_lte=body.run_after_lte,
911 upper_bound_lt=body.run_after_lt,
912 ),
913 attribute=DagRun.run_after, # type: ignore[arg-type]
914 )
915 start_date = RangeFilter(
916 Range(
917 lower_bound_gte=body.start_date_gte,
918 lower_bound_gt=body.start_date_gt,
919 upper_bound_lte=body.start_date_lte,
920 upper_bound_lt=body.start_date_lt,
921 ),
922 attribute=DagRun.start_date, # type: ignore[arg-type]
923 )
924 end_date = RangeFilter(
925 Range(
926 lower_bound_gte=body.end_date_gte,
927 lower_bound_gt=body.end_date_gt,
928 upper_bound_lte=body.end_date_lte,
929 upper_bound_lt=body.end_date_lt,
930 ),
931 attribute=DagRun.end_date, # type: ignore[arg-type]
932 )
933 duration = RangeFilter(
934 Range(
935 lower_bound_gte=body.duration_gte,
936 lower_bound_gt=body.duration_gt,
937 upper_bound_lte=body.duration_lte,
938 upper_bound_lt=body.duration_lt,
939 ),
940 attribute=DagRun.duration, # type: ignore[arg-type]
941 )
942 conf_contains = FilterParam(DagRun.conf, body.conf_contains, FilterOptionEnum.CONTAINS) # type: ignore[arg-type]
943 state = FilterParam(DagRun.state, body.states, FilterOptionEnum.ANY_EQUAL) # type: ignore[arg-type]
945 offset = OffsetFilter(body.page_offset)
946 limit = LimitFilter(body.page_limit)
948 order_by = SortParam(
949 [
950 "id",
951 "state",
952 "dag_id",
953 "run_after",
954 "logical_date",
955 "run_id",
956 "start_date",
957 "end_date",
958 "updated_at",
959 "conf",
960 ],
961 DagRun,
962 {"dag_run_id": "run_id"},
963 ).set_value([body.order_by] if body.order_by else None)
965 base_query = select(DagRun).options(*eager_load_dag_run_for_list())
967 dag_runs_select, total_entries = paginated_select(
968 statement=base_query,
969 filters=[
970 dag_ids,
971 logical_date,
972 run_after,
973 start_date,
974 end_date,
975 duration,
976 conf_contains,
977 state,
978 readable_dag_runs_filter,
979 ],
980 order_by=order_by,
981 offset=offset,
982 limit=limit,
983 session=session,
984 )
986 dag_runs = list(session.scalars(dag_runs_select))
987 attach_dag_versions_to_runs(dag_runs, session=session)
989 return DAGRunCollectionResponse(
990 dag_runs=dag_runs,
991 total_entries=total_entries,
992 )