Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/backfills.py: 68%
107 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.
17from __future__ import annotations
19from typing import Annotated, NoReturn
21from fastapi import Depends, HTTPException, status
22from fastapi.exceptions import RequestValidationError
23from pydantic import NonNegativeInt
24from sqlalchemy import select, update
25from sqlalchemy.exc import OperationalError
26from sqlalchemy.orm import joinedload
28from airflow._shared.timezones import timezone
29from airflow.api_fastapi.common.dagbag import resolve_run_on_latest_version
30from airflow.api_fastapi.common.db.common import (
31 SessionDep,
32 paginated_select,
33)
34from airflow.api_fastapi.common.parameters import QueryLimit, QueryOffset, SortParam
35from airflow.api_fastapi.common.router import AirflowRouter
36from airflow.api_fastapi.core_api.datamodels.backfills import (
37 BackfillCollectionResponse,
38 BackfillPostBody,
39 BackfillResponse,
40 DryRunBackfillCollectionResponse,
41 DryRunBackfillResponse,
42)
43from airflow.api_fastapi.core_api.openapi.exceptions import (
44 create_openapi_http_exception_doc,
45)
46from airflow.api_fastapi.core_api.security import (
47 BACKFILL_NOT_FOUND,
48 GetUserDep,
49 requires_access_backfill,
50)
51from airflow.api_fastapi.logging.decorators import action_logging
52from airflow.exceptions import DagNotFound, DagRunTypeNotAllowed
53from airflow.models import DagRun
54from airflow.models.backfill import (
55 AlreadyRunningBackfill,
56 Backfill,
57 BackfillDagRun,
58 DagNonPeriodicScheduleException,
59 InvalidBackfillConf,
60 InvalidBackfillDate,
61 InvalidBackfillDateRange,
62 InvalidBackfillDirection,
63 InvalidReprocessBehavior,
64 NoBackfillRunsToCreate,
65 _create_backfill,
66 _do_dry_run,
67)
68from airflow.utils.sqlalchemy import is_lock_not_available_error
69from airflow.utils.state import DagRunState
71backfills_router = AirflowRouter(tags=["Backfill"], prefix="/backfills")
74def _raise_locked_response_or_reraise(e: OperationalError, action: str) -> NoReturn:
75 """Map a database lock OperationalError to HTTP 503, or re-raise if not a lock error."""
76 if not is_lock_not_available_error(e):
77 raise
78 raise HTTPException(
79 status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
80 detail=(
81 f"Database is locked. Backfill {action} is not supported on SQLite "
82 "under concurrent access. Please use PostgreSQL or MySQL."
83 ),
84 )
87@backfills_router.get(
88 path="",
89 dependencies=[
90 Depends(requires_access_backfill(method="GET")),
91 ],
92)
93def list_backfills(
94 dag_id: str,
95 limit: QueryLimit,
96 offset: QueryOffset,
97 order_by: Annotated[
98 SortParam,
99 Depends(SortParam(["id"], Backfill).dynamic_depends()),
100 ],
101 session: SessionDep,
102) -> BackfillCollectionResponse:
103 select_stmt, total_entries = paginated_select(
104 statement=select(Backfill).where(Backfill.dag_id == dag_id).options(joinedload(Backfill.dag_model)),
105 order_by=order_by,
106 offset=offset,
107 limit=limit,
108 session=session,
109 )
110 return BackfillCollectionResponse(
111 backfills=session.scalars(select_stmt),
112 total_entries=total_entries,
113 )
116@backfills_router.get(
117 path="/{backfill_id}",
118 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
119 dependencies=[
120 Depends(requires_access_backfill(method="GET")),
121 ],
122)
123def get_backfill(
124 backfill_id: NonNegativeInt,
125 session: SessionDep,
126) -> BackfillResponse:
127 backfill = session.scalars(
128 select(Backfill).where(Backfill.id == backfill_id).options(joinedload(Backfill.dag_model))
129 ).one_or_none()
130 if backfill: 130 ↛ 132line 130 didn't jump to line 132 because the condition on line 130 was always true
131 return backfill
132 raise HTTPException(status.HTTP_404_NOT_FOUND, BACKFILL_NOT_FOUND)
135@backfills_router.put(
136 path="/{backfill_id}/pause",
137 responses=create_openapi_http_exception_doc(
138 [
139 status.HTTP_404_NOT_FOUND,
140 status.HTTP_409_CONFLICT,
141 ]
142 ),
143 dependencies=[
144 Depends(action_logging()),
145 Depends(requires_access_backfill(method="PUT")),
146 ],
147)
148def pause_backfill(backfill_id: NonNegativeInt, session: SessionDep) -> BackfillResponse:
149 b = session.scalars(
150 select(Backfill).where(Backfill.id == backfill_id).options(joinedload(Backfill.dag_model))
151 ).one_or_none()
152 if not b: 152 ↛ 153line 152 didn't jump to line 153 because the condition on line 152 was never true
153 raise HTTPException(status.HTTP_404_NOT_FOUND, BACKFILL_NOT_FOUND)
154 if b.completed_at: 154 ↛ 156line 154 didn't jump to line 156 because the condition on line 154 was always true
155 raise HTTPException(status.HTTP_409_CONFLICT, "Backfill is already completed.")
156 if b.is_paused is False:
157 b.is_paused = True
158 return b
161@backfills_router.put(
162 path="/{backfill_id}/unpause",
163 responses=create_openapi_http_exception_doc(
164 [
165 status.HTTP_404_NOT_FOUND,
166 status.HTTP_409_CONFLICT,
167 ]
168 ),
169 dependencies=[
170 Depends(action_logging()),
171 Depends(requires_access_backfill(method="PUT")),
172 ],
173)
174def unpause_backfill(backfill_id: NonNegativeInt, session: SessionDep) -> BackfillResponse:
175 b = session.scalars(
176 select(Backfill).where(Backfill.id == backfill_id).options(joinedload(Backfill.dag_model))
177 ).one_or_none()
178 if not b: 178 ↛ 179line 178 didn't jump to line 179 because the condition on line 178 was never true
179 raise HTTPException(status.HTTP_404_NOT_FOUND, BACKFILL_NOT_FOUND)
180 if b.completed_at: 180 ↛ 182line 180 didn't jump to line 182 because the condition on line 180 was always true
181 raise HTTPException(status.HTTP_409_CONFLICT, "Backfill is already completed.")
182 if b.is_paused:
183 b.is_paused = False
184 return b
187@backfills_router.put(
188 path="/{backfill_id}/cancel",
189 responses=create_openapi_http_exception_doc(
190 [
191 status.HTTP_404_NOT_FOUND,
192 status.HTTP_409_CONFLICT,
193 ]
194 ),
195 dependencies=[
196 Depends(action_logging()),
197 Depends(requires_access_backfill(method="PUT")),
198 ],
199)
200def cancel_backfill(backfill_id: NonNegativeInt, session: SessionDep) -> BackfillResponse:
201 b = session.scalars(
202 select(Backfill).where(Backfill.id == backfill_id).options(joinedload(Backfill.dag_model))
203 ).one_or_none()
204 if not b: 204 ↛ 205line 204 didn't jump to line 205 because the condition on line 204 was never true
205 raise HTTPException(status.HTTP_404_NOT_FOUND, BACKFILL_NOT_FOUND)
206 if b.completed_at is not None: 206 ↛ 210line 206 didn't jump to line 210 because the condition on line 206 was always true
207 raise HTTPException(status.HTTP_409_CONFLICT, "Backfill is already completed.")
209 # first, pause, and commit immediately to ensure no other dag runs are started
210 if not b.is_paused:
211 b.is_paused = True
212 session.commit() # ensure no new runs started
214 query = (
215 update(DagRun)
216 .where(
217 DagRun.id.in_(
218 select(
219 BackfillDagRun.dag_run_id,
220 ).where(
221 BackfillDagRun.backfill_id == b.id,
222 ),
223 ),
224 DagRun.state == DagRunState.QUEUED,
225 )
226 .values(state=DagRunState.FAILED)
227 .execution_options(synchronize_session=False)
228 )
229 session.execute(query)
230 session.commit() # this will fail all the queued dag runs in this backfill
232 # this is in separate transaction just to avoid potential conflicts
233 session.refresh(b)
234 b.completed_at = timezone.utcnow()
235 return b
238@backfills_router.post(
239 path="",
240 responses=create_openapi_http_exception_doc(
241 [
242 status.HTTP_400_BAD_REQUEST,
243 status.HTTP_404_NOT_FOUND,
244 status.HTTP_409_CONFLICT,
245 status.HTTP_503_SERVICE_UNAVAILABLE,
246 ]
247 ),
248 dependencies=[
249 Depends(action_logging()),
250 Depends(requires_access_backfill(method="POST")),
251 ],
252)
253def create_backfill(
254 backfill_request: BackfillPostBody,
255 user: GetUserDep,
256 session: SessionDep,
257) -> BackfillResponse:
258 from_date = timezone.coerce_datetime(backfill_request.from_date)
259 to_date = timezone.coerce_datetime(backfill_request.to_date)
260 resolved_run_on_latest = resolve_run_on_latest_version(
261 backfill_request.run_on_latest_version,
262 backfill_request.dag_id,
263 session,
264 fallback=True,
265 )
266 try:
267 backfill_obj = _create_backfill(
268 dag_id=backfill_request.dag_id,
269 from_date=from_date,
270 to_date=to_date,
271 max_active_runs=backfill_request.max_active_runs,
272 reverse=backfill_request.run_backwards,
273 dag_run_conf=backfill_request.dag_run_conf,
274 triggering_user_name=user.get_name(),
275 reprocess_behavior=backfill_request.reprocess_behavior,
276 run_on_latest_version=resolved_run_on_latest,
277 )
278 return BackfillResponse.model_validate(backfill_obj)
279 except OperationalError as e:
280 _raise_locked_response_or_reraise(e, "creation")
282 except AlreadyRunningBackfill:
283 raise HTTPException(
284 status_code=status.HTTP_409_CONFLICT,
285 detail=f"There is already a running backfill for dag {backfill_request.dag_id}",
286 )
288 except DagNotFound:
289 raise HTTPException(
290 status_code=status.HTTP_404_NOT_FOUND,
291 detail=f"Could not find dag {backfill_request.dag_id}",
292 )
293 except DagRunTypeNotAllowed as e:
294 raise HTTPException(
295 status_code=status.HTTP_400_BAD_REQUEST,
296 detail=str(e),
297 )
298 except (
299 InvalidReprocessBehavior,
300 InvalidBackfillDirection,
301 DagNonPeriodicScheduleException,
302 InvalidBackfillDate,
303 InvalidBackfillDateRange,
304 InvalidBackfillConf,
305 NoBackfillRunsToCreate,
306 ) as e:
307 raise RequestValidationError(str(e))
310@backfills_router.post(
311 path="/dry_run",
312 responses=create_openapi_http_exception_doc(
313 [
314 status.HTTP_400_BAD_REQUEST,
315 status.HTTP_404_NOT_FOUND,
316 status.HTTP_409_CONFLICT,
317 status.HTTP_503_SERVICE_UNAVAILABLE,
318 ]
319 ),
320 dependencies=[
321 Depends(requires_access_backfill(method="POST")),
322 ],
323)
324def create_backfill_dry_run(
325 body: BackfillPostBody,
326 session: SessionDep,
327) -> DryRunBackfillCollectionResponse:
328 from_date = timezone.coerce_datetime(body.from_date)
329 to_date = timezone.coerce_datetime(body.to_date)
331 try:
332 backfills_dry_run = _do_dry_run(
333 dag_id=body.dag_id,
334 from_date=from_date,
335 to_date=to_date,
336 reverse=body.run_backwards,
337 reprocess_behavior=body.reprocess_behavior,
338 dag_run_conf=body.dag_run_conf,
339 session=session,
340 )
341 backfills = [
342 DryRunBackfillResponse(
343 logical_date=d.logical_date,
344 partition_key=d.partition_key,
345 partition_date=d.partition_date,
346 )
347 for d in backfills_dry_run
348 ]
349 return DryRunBackfillCollectionResponse(backfills=backfills, total_entries=len(backfills))
350 except OperationalError as e:
351 _raise_locked_response_or_reraise(e, "dry-run")
352 except DagNotFound:
353 raise HTTPException(
354 status_code=status.HTTP_404_NOT_FOUND,
355 detail=f"Could not find dag {body.dag_id}",
356 )
357 except DagRunTypeNotAllowed as e:
358 raise HTTPException(
359 status_code=status.HTTP_400_BAD_REQUEST,
360 detail=str(e),
361 )
363 except (
364 InvalidReprocessBehavior,
365 InvalidBackfillDirection,
366 DagNonPeriodicScheduleException,
367 InvalidBackfillDate,
368 InvalidBackfillDateRange,
369 InvalidBackfillConf,
370 ) as e:
371 raise RequestValidationError(str(e))