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

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. 

17 

18from __future__ import annotations 

19 

20from typing import Annotated 

21 

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 

26 

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 

81 

82dags_router = AirflowRouter(tags=["DAG"], prefix="/dags") 

83 

84 

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 ) 

151 

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 ) 

177 

178 dags = session.scalars(dags_select) 

179 

180 return DAGCollectionResponse( 

181 dags=dags, 

182 total_entries=total_entries, 

183 ) 

184 

185 

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

207 

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) 

211 

212 return dag_model 

213 

214 

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) 

230 

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

234 

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) 

238 

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 ) 

247 

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 ) 

257 

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) 

261 

262 return DAGDetailsResponse.model_validate(dag_model) 

263 

264 

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) 

283 

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

286 

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

303 

304 data = patch_body.model_dump(include=fields_to_update, by_alias=True) 

305 

306 for key, val in data.items(): 

307 setattr(dag, key, val) 

308 

309 return dag 

310 

311 

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. 

338 

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

353 

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

371 

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

384 

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 ) 

391 

392 return DAGCollectionResponse( 

393 dags=dags, 

394 total_entries=total_entries, 

395 ) 

396 

397 

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

409 

410 user_id = str(user.get_id()) 

411 session.execute(insert(DagFavorite).values(dag_id=dag_id, user_id=user_id)) 

412 

413 

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

425 

426 user_id = str(user.get_id()) 

427 

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

436 

437 if not favorite_exists: 

438 raise HTTPException(status.HTTP_409_CONFLICT, detail="Dag is not marked as favorite") 

439 

440 session.execute( 

441 delete(DagFavorite).where( 

442 DagFavorite.dag_id == dag_id, 

443 DagFavorite.user_id == user_id, 

444 ) 

445 ) 

446 

447 

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)