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

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 

18 

19from typing import Annotated, NoReturn 

20 

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 

27 

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 

70 

71backfills_router = AirflowRouter(tags=["Backfill"], prefix="/backfills") 

72 

73 

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 ) 

85 

86 

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 ) 

114 

115 

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) 

133 

134 

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 

159 

160 

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 

185 

186 

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.") 

208 

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 

213 

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 

231 

232 # this is in separate transaction just to avoid potential conflicts 

233 session.refresh(b) 

234 b.completed_at = timezone.utcnow() 

235 return b 

236 

237 

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") 

281 

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 ) 

287 

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)) 

308 

309 

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) 

330 

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 ) 

362 

363 except ( 

364 InvalidReprocessBehavior, 

365 InvalidBackfillDirection, 

366 DagNonPeriodicScheduleException, 

367 InvalidBackfillDate, 

368 InvalidBackfillDateRange, 

369 InvalidBackfillConf, 

370 ) as e: 

371 raise RequestValidationError(str(e))