Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/routes/dag_runs.py: 54%
93 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 logging
21from typing import Annotated
23from cadwyn import VersionedAPIRouter
24from fastapi import HTTPException, Query, status
25from sqlalchemy import func, select
26from sqlalchemy.exc import NoResultFound
28from airflow.api.common.trigger_dag import trigger_dag
29from airflow.api_fastapi.common.dagbag import DagBagDep, get_dag_for_run, resolve_run_on_latest_version
30from airflow.api_fastapi.common.db.common import SessionDep
31from airflow.api_fastapi.common.types import UtcDateTime
32from airflow.api_fastapi.compat import HTTP_422_UNPROCESSABLE_CONTENT
33from airflow.api_fastapi.execution_api.datamodels.dagrun import DagRunStateResponse, TriggerDAGRunPayload
34from airflow.api_fastapi.execution_api.datamodels.taskinstance import DagRun
35from airflow.api_fastapi.execution_api.datamodels.token import TIToken
36from airflow.api_fastapi.execution_api.security import CurrentTIToken
37from airflow.exceptions import DagNotPartitionedError, DagRunAlreadyExists, InvalidPartitionKeyError
38from airflow.models.dag import DagModel
39from airflow.models.dagrun import DagRun as DagRunModel
40from airflow.models.taskinstance import TaskInstance
41from airflow.utils.state import DagRunState
42from airflow.utils.types import DagRunTriggeredByType, DagRunType
44router = VersionedAPIRouter()
46log = logging.getLogger(__name__)
49@router.only_exists_in_older_versions
50@router.get("/{dag_id}/previous")
51def get_previous_dagrun_compat(
52 dag_id: str,
53 logical_date: UtcDateTime,
54 session: SessionDep,
55 state: DagRunState | None = None,
56):
57 """
58 Redirect old previous dag run request to the new endpoint.
60 This endpoint must be put before ``get_dag_run`` so not to be shadowed.
61 Newer client versions would not see this endpoint, and be routed to
62 ``get_dag_run`` below instead.
63 """
64 return get_previous_dagrun(dag_id, logical_date, session, state)
67@router.get(
68 "/{dag_id}/{run_id}",
69 responses={status.HTTP_404_NOT_FOUND: {"description": "Dag run not found"}},
70)
71def get_dag_run(dag_id: str, run_id: str, session: SessionDep) -> DagRun:
72 """Get detail of a Dag run."""
73 dr = session.scalar(select(DagRunModel).where(DagRunModel.dag_id == dag_id, DagRunModel.run_id == run_id))
74 if dr is None:
75 raise HTTPException(
76 status.HTTP_404_NOT_FOUND,
77 detail={
78 "reason": "not_found",
79 "message": f"Dag run with dag_id '{dag_id}' and run_id '{run_id}' was not found",
80 },
81 )
82 return DagRun.model_validate(dr)
85@router.post(
86 "/{dag_id}/{run_id}",
87 status_code=status.HTTP_204_NO_CONTENT,
88 responses={
89 status.HTTP_400_BAD_REQUEST: {"description": "Dag has import errors or the run type is not allowed"},
90 status.HTTP_404_NOT_FOUND: {"description": "Dag not found for the given dag_id"},
91 status.HTTP_409_CONFLICT: {"description": "Dag run already exists for the given dag_id"},
92 HTTP_422_UNPROCESSABLE_CONTENT: {"description": "Invalid payload"},
93 },
94)
95def trigger_dag_run(
96 dag_id: str,
97 run_id: str,
98 payload: TriggerDAGRunPayload,
99 session: SessionDep,
100 token: TIToken = CurrentTIToken,
101) -> None:
102 """Trigger a Dag run."""
103 dm = session.scalar(select(DagModel).where(~DagModel.is_stale, DagModel.dag_id == dag_id).limit(1))
104 if not dm: 104 ↛ 105line 104 didn't jump to line 105 because the condition on line 104 was never true
105 raise HTTPException(
106 status.HTTP_404_NOT_FOUND,
107 detail={"reason": "not_found", "message": f"Dag with dag_id: '{dag_id}' not found"},
108 )
110 if dm.has_import_errors: 110 ↛ 111line 110 didn't jump to line 111 because the condition on line 110 was never true
111 raise HTTPException(
112 status.HTTP_400_BAD_REQUEST,
113 detail={
114 "reason": "import_errors",
115 "message": f"Dag with dag_id '{dag_id}' has import errors and cannot be triggered",
116 },
117 )
119 if dm.allowed_run_types is not None and DagRunType.OPERATOR_TRIGGERED not in dm.allowed_run_types: 119 ↛ 120line 119 didn't jump to line 120 because the condition on line 119 was never true
120 raise HTTPException(
121 status.HTTP_400_BAD_REQUEST,
122 detail={
123 "reason": "denied_run_type",
124 "message": f"Dag with dag_id '{dag_id}' does not allow operator-triggered runs",
125 },
126 )
128 # Inherit triggering_user_name from the calling task's DagRun so chains of
129 # TriggerDagRunOperator preserve the original human user across child runs.
130 parent_ti = session.get(TaskInstance, token.id)
131 triggering_user_name = parent_ti.dag_run.triggering_user_name if parent_ti else None
133 try:
134 trigger_dag(
135 dag_id=dag_id,
136 run_id=run_id,
137 run_type=DagRunType.OPERATOR_TRIGGERED,
138 conf=payload.conf,
139 logical_date=payload.logical_date,
140 run_after=payload.run_after,
141 triggered_by=DagRunTriggeredByType.OPERATOR,
142 triggering_user_name=triggering_user_name,
143 replace_microseconds=False,
144 partition_key=payload.partition_key,
145 note=payload.note,
146 session=session,
147 )
148 except DagRunAlreadyExists:
149 raise HTTPException(
150 status.HTTP_409_CONFLICT,
151 detail={
152 "reason": "already_exists",
153 "message": f"A run already exists for Dag '{dag_id}' with run_id '{run_id}'",
154 },
155 )
156 except DagNotPartitionedError as e:
157 raise HTTPException(
158 status.HTTP_400_BAD_REQUEST,
159 detail={"reason": "not_partitioned", "message": str(e)},
160 )
161 except InvalidPartitionKeyError as e:
162 raise HTTPException(
163 status.HTTP_400_BAD_REQUEST,
164 detail={"reason": "invalid_partition_key", "message": str(e)},
165 ) from e
166 except ValueError as e:
167 raise HTTPException(
168 status.HTTP_400_BAD_REQUEST,
169 detail={
170 "reason": "value_error",
171 "message": str(e),
172 },
173 ) from e
176@router.post(
177 "/{dag_id}/{run_id}/clear",
178 status_code=status.HTTP_204_NO_CONTENT,
179 responses={
180 status.HTTP_400_BAD_REQUEST: {"description": "Dag has import errors and cannot be triggered"},
181 status.HTTP_404_NOT_FOUND: {"description": "Dag not found for the given dag_id"},
182 HTTP_422_UNPROCESSABLE_CONTENT: {"description": "Invalid payload"},
183 },
184)
185def clear_dag_run(
186 dag_id: str,
187 run_id: str,
188 session: SessionDep,
189 dag_bag: DagBagDep,
190) -> None:
191 """Clear a Dag run."""
192 dm = session.scalar(select(DagModel).where(~DagModel.is_stale, DagModel.dag_id == dag_id).limit(1))
193 if not dm: 193 ↛ 194line 193 didn't jump to line 194 because the condition on line 193 was never true
194 raise HTTPException(
195 status.HTTP_404_NOT_FOUND,
196 detail={"reason": "not_found", "message": f"Dag with dag_id: '{dag_id}' not found"},
197 )
199 if dm.has_import_errors: 199 ↛ 200line 199 didn't jump to line 200 because the condition on line 199 was never true
200 raise HTTPException(
201 status.HTTP_400_BAD_REQUEST,
202 detail={
203 "reason": "import_errors",
204 "message": f"Dag with dag_id '{dag_id}' has import errors and cannot be triggered",
205 },
206 )
208 dag_run = session.scalar(
209 select(DagRunModel).where(DagRunModel.dag_id == dag_id, DagRunModel.run_id == run_id)
210 )
211 if dag_run is None: 211 ↛ 216line 211 didn't jump to line 216 because the condition on line 211 was always true
212 raise HTTPException(
213 status.HTTP_404_NOT_FOUND,
214 detail={"reason": "not_found", "message": f"Dag run with run_id: '{run_id}' not found"},
215 )
216 dag = get_dag_for_run(dag_bag, dag_run=dag_run, session=session)
218 resolved_run_on_latest = resolve_run_on_latest_version(None, dag_id, session)
219 dag.clear(run_id=run_id, run_on_latest_version=resolved_run_on_latest)
222@router.get(
223 "/{dag_id}/{run_id}/state",
224 responses={status.HTTP_404_NOT_FOUND: {"description": "Dag run not found"}},
225)
226def get_dagrun_state(
227 dag_id: str,
228 run_id: str,
229 session: SessionDep,
230) -> DagRunStateResponse:
231 """Get a Dag run State."""
232 try:
233 state: DagRunState = session.scalars(
234 select(DagRunModel.state).where(DagRunModel.dag_id == dag_id, DagRunModel.run_id == run_id)
235 ).one()
236 except NoResultFound:
237 raise HTTPException(
238 status.HTTP_404_NOT_FOUND,
239 detail={
240 "reason": "not_found",
241 "message": f"Dag run with dag_id '{dag_id}' and run_id '{run_id}' was not found",
242 },
243 )
244 return DagRunStateResponse(state=state)
247@router.get("/count", status_code=status.HTTP_200_OK)
248def get_dr_count(
249 dag_id: str,
250 session: SessionDep,
251 logical_dates: Annotated[list[UtcDateTime] | None, Query()] = None,
252 run_ids: Annotated[list[str] | None, Query()] = None,
253 states: Annotated[list[str] | None, Query()] = None,
254) -> int:
255 """Get the count of Dag runs matching the given criteria."""
256 stmt = select(func.count()).select_from(DagRunModel).where(DagRunModel.dag_id == dag_id)
257 if logical_dates:
258 stmt = stmt.where(DagRunModel.logical_date.in_(logical_dates))
259 if run_ids:
260 stmt = stmt.where(DagRunModel.run_id.in_(run_ids))
261 if states:
262 stmt = stmt.where(DagRunModel.state.in_(states))
263 return session.scalar(stmt) or 0
266@router.get("/previous", status_code=status.HTTP_200_OK)
267def get_previous_dagrun(
268 dag_id: str,
269 logical_date: UtcDateTime,
270 session: SessionDep,
271 state: Annotated[DagRunState | None, Query()] = None,
272) -> DagRun | None:
273 """Get the previous Dag run before the given logical date, optionally filtered by state."""
274 stmt = (
275 select(DagRunModel)
276 .where(DagRunModel.dag_id == dag_id, DagRunModel.logical_date < logical_date)
277 .order_by(DagRunModel.logical_date.desc())
278 .limit(1)
279 )
280 if state:
281 stmt = stmt.where(DagRunModel.state == state)
282 if not (dag_run := session.scalar(stmt)):
283 return None
284 return DagRun.model_validate(dag_run)