Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/assets.py: 89%

174 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 datetime import datetime 

21from typing import TYPE_CHECKING, Annotated, cast 

22 

23from fastapi import Depends, HTTPException, status 

24from sqlalchemy import and_, delete, func, select 

25from sqlalchemy.engine import CursorResult 

26from sqlalchemy.orm import joinedload, subqueryload 

27 

28from airflow._shared.timezones import timezone 

29from airflow.api_fastapi.app import get_auth_manager 

30from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity, DagDetails 

31from airflow.api_fastapi.common.dagbag import DagBagDep, get_latest_version_of_dag 

32from airflow.api_fastapi.common.db.common import SessionDep, paginated_select 

33from airflow.api_fastapi.common.parameters import ( 

34 BaseParam, 

35 FilterParam, 

36 OptionalDateTimeQuery, 

37 QueryAssetAliasNamePatternSearch, 

38 QueryAssetAliasNamePrefixPatternSearch, 

39 QueryAssetDagIdPatternSearch, 

40 QueryAssetNamePatternSearch, 

41 QueryAssetNamePrefixPatternSearch, 

42 QueryLimit, 

43 QueryOffset, 

44 QueryUriPatternSearch, 

45 QueryUriPrefixPatternSearch, 

46 RangeFilter, 

47 SortParam, 

48 datetime_range_filter_factory, 

49 filter_param_factory, 

50) 

51from airflow.api_fastapi.common.router import AirflowRouter 

52from airflow.api_fastapi.core_api.datamodels.assets import ( 

53 AssetAliasCollectionResponse, 

54 AssetAliasResponse, 

55 AssetCollectionResponse, 

56 AssetEventCollectionResponse, 

57 AssetEventResponse, 

58 AssetResponse, 

59 CreateAssetEventsBody, 

60 MaterializeAssetBody, 

61 QueuedEventCollectionResponse, 

62 QueuedEventResponse, 

63) 

64from airflow.api_fastapi.core_api.datamodels.dag_run import DAGRunResponse 

65from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc 

66from airflow.api_fastapi.core_api.security import ( 

67 GetUserDep, 

68 ReadableAssetEventsFilterDep, 

69 ReadableDagsFilterDep, 

70 requires_access_asset, 

71 requires_access_asset_alias, 

72 requires_access_dag, 

73) 

74from airflow.api_fastapi.core_api.services.public.assets import serialize_asset_events 

75from airflow.api_fastapi.logging.decorators import action_logging 

76from airflow.assets.manager import asset_manager 

77from airflow.configuration import conf 

78from airflow.exceptions import ParamValidationError 

79from airflow.models.asset import ( 

80 AssetAliasModel, 

81 AssetDagRunQueue, 

82 AssetEvent, 

83 AssetModel, 

84 AssetWatcherModel, 

85 TaskOutletAssetReference, 

86) 

87from airflow.models.dag import DagModel 

88from airflow.typing_compat import Unpack 

89from airflow.utils.state import DagRunState 

90from airflow.utils.types import DagRunTriggeredByType, DagRunType 

91 

92if TYPE_CHECKING: 92 ↛ 93line 92 didn't jump to line 93 because the condition on line 92 was never true

93 from sqlalchemy.engine import Result 

94 from sqlalchemy.sql import Select 

95 

96assets_router = AirflowRouter(tags=["Asset"]) 

97 

98 

99def _generate_queued_event_where_clause( 

100 *, 

101 asset_id: int | None = None, 

102 dag_id: str | None = None, 

103 before: datetime | str | None = None, 

104 permitted_dag_ids: set[str] | None = None, 

105) -> list: 

106 """Get AssetDagRunQueue where clause.""" 

107 where_clause = [] 

108 if dag_id is not None: 

109 where_clause.append(AssetDagRunQueue.target_dag_id == dag_id) 

110 if asset_id is not None: 

111 where_clause.append(AssetDagRunQueue.asset_id == asset_id) 

112 if before is not None: 112 ↛ 113line 112 didn't jump to line 113 because the condition on line 112 was never true

113 where_clause.append(AssetDagRunQueue.created_at < before) 

114 if permitted_dag_ids is not None: 114 ↛ 116line 114 didn't jump to line 116 because the condition on line 114 was always true

115 where_clause.append(AssetDagRunQueue.target_dag_id.in_(permitted_dag_ids)) 

116 return where_clause 

117 

118 

119class OnlyActiveFilter(BaseParam[bool]): 

120 """Filter on asset activeness.""" 

121 

122 def to_orm(self, select: Select) -> Select: 

123 if self.value and self.skip_none: 

124 return select.where(AssetModel.active.has()) 

125 return select 

126 

127 @classmethod 

128 def depends(cls, only_active: bool = True) -> OnlyActiveFilter: 

129 return cls().set_value(only_active) 

130 

131 

132@assets_router.get( 

133 "/assets", 

134 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

135 dependencies=[ 

136 Depends(requires_access_asset(method="GET")), 

137 Depends(requires_access_asset_alias(method="GET")), 

138 ], 

139) 

140def get_assets( 

141 limit: QueryLimit, 

142 offset: QueryOffset, 

143 name_pattern: QueryAssetNamePatternSearch, 

144 name_prefix_pattern: QueryAssetNamePrefixPatternSearch, 

145 uri_pattern: QueryUriPatternSearch, 

146 uri_prefix_pattern: QueryUriPrefixPatternSearch, 

147 dag_ids: QueryAssetDagIdPatternSearch, 

148 only_active: Annotated[OnlyActiveFilter, Depends(OnlyActiveFilter.depends)], 

149 order_by: Annotated[ 

150 SortParam, 

151 Depends(SortParam(["id", "name", "uri", "created_at", "updated_at"], AssetModel).dynamic_depends()), 

152 ], 

153 session: SessionDep, 

154) -> AssetCollectionResponse: 

155 """Get assets.""" 

156 # Build a query that will be used to retrieve the ID and timestamp of the latest AssetEvent 

157 last_asset_events = ( 

158 select(AssetEvent.asset_id, func.max(AssetEvent.timestamp).label("last_timestamp")) 

159 .group_by(AssetEvent.asset_id) 

160 .subquery() 

161 ) 

162 

163 # First, we're pulling the Asset ID, AssetEvent ID, and AssetEvent timestamp for the latest (last) 

164 # AssetEvent. We'll eventually OUTER JOIN this to the AssetModel 

165 asset_event_query = ( 

166 select( 

167 AssetEvent.asset_id, # The ID of the Asset, which we'll need to JOIN to the AssetModel 

168 func.max(AssetEvent.id).label("last_asset_event_id"), # The ID of the last AssetEvent 

169 func.max(AssetEvent.timestamp).label("last_asset_event_timestamp"), 

170 ) 

171 .join( 

172 last_asset_events, 

173 and_( 

174 AssetEvent.asset_id == last_asset_events.c.asset_id, 

175 AssetEvent.timestamp == last_asset_events.c.last_timestamp, 

176 ), 

177 ) 

178 .group_by(AssetEvent.asset_id) 

179 .subquery() 

180 ) 

181 

182 assets_select_statement = select( 

183 AssetModel, 

184 asset_event_query.c.last_asset_event_id, # This should be the AssetEvent.id 

185 asset_event_query.c.last_asset_event_timestamp, 

186 ).outerjoin(asset_event_query, AssetModel.id == asset_event_query.c.asset_id) 

187 

188 assets_select, total_entries = paginated_select( 

189 statement=assets_select_statement, 

190 filters=[only_active, name_pattern, name_prefix_pattern, uri_pattern, uri_prefix_pattern, dag_ids], 

191 order_by=order_by, 

192 offset=offset, 

193 limit=limit, 

194 session=session, 

195 ) 

196 

197 # The below type annotation is acceptable on SQLA2.1, but not on 2.0 

198 assets_rows: Result[Unpack[tuple[AssetModel, int, datetime]]] = session.execute( # type: ignore[type-arg] 

199 assets_select.options( 

200 subqueryload(AssetModel.scheduled_dags), 

201 subqueryload(AssetModel.producing_tasks), 

202 subqueryload(AssetModel.consuming_tasks), 

203 subqueryload(AssetModel.aliases), 

204 subqueryload(AssetModel.watchers).joinedload(AssetWatcherModel.trigger), 

205 ) 

206 ) 

207 

208 assets = [] 

209 

210 for asset, last_asset_event_id, last_asset_event_timestamp in assets_rows: 

211 watchers_data = [ 

212 { 

213 "name": watcher.name, 

214 "trigger_id": watcher.trigger_id, 

215 "created_date": watcher.trigger.created_date, 

216 } 

217 for watcher in asset.watchers 

218 ] 

219 

220 asset_response = AssetResponse.model_validate( 

221 { 

222 **asset.__dict__, 

223 "aliases": asset.aliases, 

224 "watchers": watchers_data, 

225 "last_asset_event": { 

226 "id": last_asset_event_id, 

227 "timestamp": last_asset_event_timestamp, 

228 }, 

229 } 

230 ) 

231 assets.append(asset_response) 

232 

233 return AssetCollectionResponse( 

234 assets=assets, 

235 total_entries=total_entries, 

236 ) 

237 

238 

239@assets_router.get( 

240 "/assets/aliases", 

241 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

242 dependencies=[Depends(requires_access_asset_alias(method="GET"))], 

243) 

244def get_asset_aliases( 

245 limit: QueryLimit, 

246 offset: QueryOffset, 

247 name_pattern: QueryAssetAliasNamePatternSearch, 

248 name_prefix_pattern: QueryAssetAliasNamePrefixPatternSearch, 

249 order_by: Annotated[ 

250 SortParam, 

251 Depends(SortParam(["id", "name"], AssetAliasModel).dynamic_depends()), 

252 ], 

253 session: SessionDep, 

254) -> AssetAliasCollectionResponse: 

255 """Get asset aliases.""" 

256 asset_aliases_select, total_entries = paginated_select( 

257 statement=select(AssetAliasModel), 

258 filters=[name_pattern, name_prefix_pattern], 

259 order_by=order_by, 

260 offset=offset, 

261 limit=limit, 

262 session=session, 

263 ) 

264 

265 return AssetAliasCollectionResponse( 

266 asset_aliases=session.scalars(asset_aliases_select), 

267 total_entries=total_entries, 

268 ) 

269 

270 

271@assets_router.get( 

272 "/assets/aliases/{asset_alias_id}", 

273 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

274 dependencies=[Depends(requires_access_asset_alias(method="GET"))], 

275) 

276def get_asset_alias(asset_alias_id: int, session: SessionDep): 

277 """Get an asset alias.""" 

278 alias = session.scalar(select(AssetAliasModel).where(AssetAliasModel.id == asset_alias_id)) 

279 if alias is None: 

280 raise HTTPException( 

281 status.HTTP_404_NOT_FOUND, 

282 f"The Asset Alias with ID: `{asset_alias_id}` was not found", 

283 ) 

284 return AssetAliasResponse.model_validate(alias) 

285 

286 

287@assets_router.get( 

288 "/assets/events", 

289 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

290 dependencies=[Depends(requires_access_asset(method="GET"))], 

291) 

292def get_asset_events( 

293 limit: QueryLimit, 

294 offset: QueryOffset, 

295 order_by: Annotated[ 

296 SortParam, 

297 Depends( 

298 SortParam( 

299 [ 

300 "source_task_id", 

301 "source_dag_id", 

302 "source_run_id", 

303 "source_map_index", 

304 "timestamp", 

305 ], 

306 AssetEvent, 

307 ).dynamic_depends("timestamp") 

308 ), 

309 ], 

310 asset_id: Annotated[ 

311 FilterParam[int | None], Depends(filter_param_factory(AssetEvent.asset_id, int | None)) 

312 ], 

313 source_dag_id: Annotated[ 

314 FilterParam[str | None], Depends(filter_param_factory(AssetEvent.source_dag_id, str | None)) 

315 ], 

316 source_task_id: Annotated[ 

317 FilterParam[str | None], Depends(filter_param_factory(AssetEvent.source_task_id, str | None)) 

318 ], 

319 source_run_id: Annotated[ 

320 FilterParam[str | None], Depends(filter_param_factory(AssetEvent.source_run_id, str | None)) 

321 ], 

322 source_map_index: Annotated[ 

323 FilterParam[int | None], Depends(filter_param_factory(AssetEvent.source_map_index, int | None)) 

324 ], 

325 name_pattern: QueryAssetNamePatternSearch, 

326 name_prefix_pattern: QueryAssetNamePrefixPatternSearch, 

327 timestamp_range: Annotated[RangeFilter, Depends(datetime_range_filter_factory("timestamp", AssetEvent))], 

328 readable_asset_events_filter: ReadableAssetEventsFilterDep, 

329 session: SessionDep, 

330) -> AssetEventCollectionResponse: 

331 """Get asset events.""" 

332 base_statement = select(AssetEvent) 

333 if name_pattern.value or name_prefix_pattern.value: 

334 base_statement = base_statement.join(AssetModel, AssetEvent.asset_id == AssetModel.id) 

335 

336 assets_event_select, total_entries = paginated_select( 

337 statement=base_statement, 

338 filters=[ 

339 asset_id, 

340 source_dag_id, 

341 source_task_id, 

342 source_run_id, 

343 source_map_index, 

344 name_pattern, 

345 name_prefix_pattern, 

346 timestamp_range, 

347 readable_asset_events_filter, 

348 ], 

349 order_by=order_by, 

350 offset=offset, 

351 limit=limit, 

352 session=session, 

353 ) 

354 

355 assets_event_select = assets_event_select.options( 

356 subqueryload(AssetEvent.created_dagruns), joinedload(AssetEvent.asset) 

357 ) 

358 assets_events = session.scalars(assets_event_select).all() 

359 

360 return AssetEventCollectionResponse( 

361 asset_events=serialize_asset_events(assets_events, session=session), 

362 total_entries=total_entries, 

363 ) 

364 

365 

366@assets_router.post( 

367 "/assets/events", 

368 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

369 dependencies=[Depends(requires_access_asset(method="POST")), Depends(action_logging())], 

370) 

371def create_asset_event( 

372 body: CreateAssetEventsBody, 

373 session: SessionDep, 

374 user: GetUserDep, 

375) -> AssetEventResponse: 

376 """Create asset events.""" 

377 asset_model = session.scalar(select(AssetModel).where(AssetModel.id == body.asset_id).limit(1)) 

378 if not asset_model: 

379 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Asset with ID: `{body.asset_id}` was not found") 

380 timestamp = timezone.utcnow() 

381 

382 api_user_teams: set[str] = set() 

383 api_allow_consumer_teams: list[str] | None = None 

384 api_allow_global_consumers: bool = True 

385 if conf.getboolean("core", "multi_team"): 385 ↛ 386line 385 didn't jump to line 386 because the condition on line 385 was never true

386 api_user_teams = get_auth_manager().get_authorized_teams(user=user) 

387 if body.access_control: 

388 api_allow_consumer_teams = body.access_control.consumer_teams 

389 api_allow_global_consumers = body.access_control.allow_global 

390 

391 assets_event = asset_manager.register_asset_change( 

392 asset=asset_model, 

393 timestamp=timestamp, 

394 extra=body.extra, 

395 partition_key=body.partition_key, 

396 source_is_api=True, 

397 api_user_teams=api_user_teams, 

398 api_allow_consumer_teams=api_allow_consumer_teams, 

399 api_allow_global_consumers=api_allow_global_consumers, 

400 session=session, 

401 ) 

402 

403 if not assets_event: 403 ↛ 404line 403 didn't jump to line 404 because the condition on line 403 was never true

404 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Asset with ID: `{body.asset_id}` was not found") 

405 return AssetEventResponse.model_validate(assets_event) 

406 

407 

408@assets_router.post( 

409 "/assets/{asset_id}/materialize", 

410 responses=create_openapi_http_exception_doc( 

411 [status.HTTP_400_BAD_REQUEST, status.HTTP_404_NOT_FOUND, status.HTTP_409_CONFLICT] 

412 ), 

413 dependencies=[Depends(requires_access_asset(method="POST")), Depends(action_logging())], 

414) 

415def materialize_asset( 

416 asset_id: int, 

417 dag_bag: DagBagDep, 

418 user: GetUserDep, 

419 session: SessionDep, 

420 body: MaterializeAssetBody | None = None, 

421) -> DAGRunResponse: 

422 """Materialize an asset by triggering a Dag run that produces it.""" 

423 dag_id_it = iter( 

424 session.scalars( 

425 select(TaskOutletAssetReference.dag_id) 

426 .where(TaskOutletAssetReference.asset_id == asset_id) 

427 .group_by(TaskOutletAssetReference.dag_id) 

428 .limit(2) 

429 ) 

430 ) 

431 

432 if (dag_id := next(dag_id_it, None)) is None: 

433 raise HTTPException(status.HTTP_404_NOT_FOUND, f"No Dag materializes asset with ID: {asset_id}") 

434 if next(dag_id_it, None) is not None: 

435 raise HTTPException( 

436 status.HTTP_409_CONFLICT, 

437 f"More than one Dag materializes asset with ID: {asset_id}", 

438 ) 

439 

440 if not get_auth_manager().is_authorized_dag( 440 ↛ 449line 440 didn't jump to line 449 because the condition on line 440 was never true

441 method="POST", 

442 access_entity=DagAccessEntity.RUN, 

443 # The Dag is resolved from the asset here rather than named by the caller, so its team has 

444 # to be looked up too. A team-aware auth manager distinguishes a team-scoped Dag from a 

445 # global one by this field, so leaving it None asks the wrong question. 

446 details=DagDetails(id=dag_id, team_name=DagModel.get_team_name(dag_id, session=session)), 

447 user=user, 

448 ): 

449 raise HTTPException( 

450 status.HTTP_403_FORBIDDEN, 

451 f"User is not authorized to trigger a run for Dag: {dag_id} that materializes this asset", 

452 ) 

453 

454 dag = get_latest_version_of_dag(dag_bag, dag_id, session) 

455 

456 if dag.allowed_run_types is not None and DagRunType.ASSET_MATERIALIZATION not in dag.allowed_run_types: 456 ↛ 457line 456 didn't jump to line 457 because the condition on line 456 was never true

457 raise HTTPException( 

458 status.HTTP_400_BAD_REQUEST, 

459 f"Dag with dag_id: '{dag_id}' does not allow asset materialization runs", 

460 ) 

461 

462 try: 

463 params = (body or MaterializeAssetBody()).validate_context(dag) 

464 return dag.create_dagrun( 

465 run_id=params["run_id"], 

466 logical_date=params["logical_date"], 

467 data_interval=params["data_interval"], 

468 run_after=params["run_after"], 

469 conf=params["conf"], 

470 run_type=DagRunType.ASSET_MATERIALIZATION, 

471 triggered_by=DagRunTriggeredByType.REST_API, 

472 triggering_user_name=user.get_name(), 

473 state=DagRunState.QUEUED, 

474 partition_key=params["partition_key"], 

475 partition_date=params["partition_date"], 

476 note=params["note"], 

477 session=session, 

478 ) 

479 except (ParamValidationError, ValueError) as e: 

480 raise HTTPException(status.HTTP_400_BAD_REQUEST, str(e)) from e 

481 

482 

483@assets_router.get( 

484 "/assets/{asset_id}/queuedEvents", 

485 dependencies=[Depends(requires_access_asset(method="GET"))], 

486) 

487def get_asset_queued_events( 

488 asset_id: int, 

489 readable_dags_filter: ReadableDagsFilterDep, 

490 session: SessionDep, 

491 before: OptionalDateTimeQuery = None, 

492) -> QueuedEventCollectionResponse: 

493 """Get queued asset events for an asset.""" 

494 where_clause = _generate_queued_event_where_clause( 

495 asset_id=asset_id, before=before, permitted_dag_ids=readable_dags_filter.value 

496 ) 

497 query = select(AssetDagRunQueue).where(*where_clause).options(joinedload(AssetDagRunQueue.dag_model)) 

498 

499 dag_asset_queued_events_select, total_entries = paginated_select(statement=query) 

500 adrqs = session.scalars(dag_asset_queued_events_select).all() 

501 

502 queued_events = [ 

503 QueuedEventResponse( 

504 created_at=adrq.created_at, 

505 dag_id=adrq.target_dag_id, 

506 asset_id=adrq.asset_id, 

507 dag_display_name=adrq.dag_model.dag_display_name, 

508 ) 

509 for adrq in adrqs 

510 ] 

511 

512 return QueuedEventCollectionResponse( 

513 queued_events=queued_events, 

514 total_entries=total_entries, 

515 ) 

516 

517 

518@assets_router.get( 

519 "/assets/{asset_id}", 

520 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

521 dependencies=[ 

522 Depends(requires_access_asset(method="GET")), 

523 Depends(requires_access_asset_alias(method="GET")), 

524 ], 

525) 

526def get_asset( 

527 asset_id: int, 

528 session: SessionDep, 

529) -> AssetResponse: 

530 """Get an asset.""" 

531 # Build a subquery to be used to retrieve the latest AssetEvent by matching timestamp 

532 last_asset_event = ( 

533 select(func.max(AssetEvent.timestamp)).where(AssetEvent.asset_id == asset_id).scalar_subquery() 

534 ) 

535 

536 # Now, find the latest AssetEvent details using the subquery from above 

537 asset_event_rows = session.execute( 

538 select(AssetEvent.asset_id, AssetEvent.id, AssetEvent.timestamp).where( 

539 AssetEvent.asset_id == asset_id, AssetEvent.timestamp == last_asset_event 

540 ) 

541 ).one_or_none() 

542 

543 # Retrieve the Asset; there should only be one for that asset_id 

544 asset = session.scalar( 

545 select(AssetModel) 

546 .where(AssetModel.id == asset_id) 

547 .options( 

548 joinedload(AssetModel.scheduled_dags), 

549 joinedload(AssetModel.producing_tasks), 

550 joinedload(AssetModel.consuming_tasks), 

551 joinedload(AssetModel.watchers).joinedload(AssetWatcherModel.trigger), 

552 ) 

553 ) 

554 

555 last_asset_event_id = asset_event_rows[1] if asset_event_rows else None 

556 last_asset_event_timestamp = asset_event_rows[2] if asset_event_rows else None 

557 

558 if asset is None: 

559 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Asset with ID: `{asset_id}` was not found") 

560 

561 watchers_data = [ 

562 { 

563 "name": watcher.name, 

564 "trigger_id": watcher.trigger_id, 

565 "created_date": watcher.trigger.created_date, 

566 } 

567 for watcher in asset.watchers 

568 ] 

569 

570 return AssetResponse.model_validate( 

571 { 

572 **asset.__dict__, 

573 "aliases": asset.aliases, 

574 "watchers": watchers_data, 

575 "last_asset_event": { 

576 "id": last_asset_event_id, 

577 "timestamp": last_asset_event_timestamp, 

578 }, 

579 } 

580 ) 

581 

582 

583@assets_router.get( 

584 "/dags/{dag_id}/assets/queuedEvents", 

585 dependencies=[Depends(requires_access_asset(method="GET")), Depends(requires_access_dag(method="GET"))], 

586) 

587def get_dag_asset_queued_events( 

588 dag_id: str, 

589 readable_dags_filter: ReadableDagsFilterDep, 

590 session: SessionDep, 

591 before: OptionalDateTimeQuery = None, 

592) -> QueuedEventCollectionResponse: 

593 """Get queued asset events for a Dag.""" 

594 where_clause = _generate_queued_event_where_clause( 

595 dag_id=dag_id, before=before, permitted_dag_ids=readable_dags_filter.value 

596 ) 

597 query = select(AssetDagRunQueue).where(*where_clause).options(joinedload(AssetDagRunQueue.dag_model)) 

598 

599 dag_asset_queued_events_select, total_entries = paginated_select(statement=query) 

600 adrqs = session.scalars(dag_asset_queued_events_select).all() 

601 

602 queued_events = [ 

603 QueuedEventResponse( 

604 created_at=adrq.created_at, 

605 dag_id=adrq.target_dag_id, 

606 asset_id=adrq.asset_id, 

607 dag_display_name=adrq.dag_model.dag_display_name, 

608 ) 

609 for adrq in adrqs 

610 ] 

611 

612 return QueuedEventCollectionResponse( 

613 queued_events=queued_events, 

614 total_entries=total_entries, 

615 ) 

616 

617 

618@assets_router.get( 

619 "/dags/{dag_id}/assets/{asset_id}/queuedEvents", 

620 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

621 dependencies=[Depends(requires_access_asset(method="GET")), Depends(requires_access_dag(method="GET"))], 

622) 

623def get_dag_asset_queued_event( 

624 dag_id: str, 

625 asset_id: int, 

626 readable_dags_filter: ReadableDagsFilterDep, 

627 session: SessionDep, 

628 before: OptionalDateTimeQuery = None, 

629) -> QueuedEventResponse: 

630 """Get a queued asset event for a Dag.""" 

631 where_clause = _generate_queued_event_where_clause( 

632 dag_id=dag_id, asset_id=asset_id, before=before, permitted_dag_ids=readable_dags_filter.value 

633 ) 

634 query = select(AssetDagRunQueue).where(*where_clause) 

635 adrq = session.scalar(query) 

636 if not adrq: 636 ↛ 642line 636 didn't jump to line 642 because the condition on line 636 was always true

637 raise HTTPException( 

638 status.HTTP_404_NOT_FOUND, 

639 f"Queued event with dag_id: `{dag_id}` and asset_id: `{asset_id}` was not found", 

640 ) 

641 

642 return QueuedEventResponse( 

643 created_at=adrq.created_at, 

644 dag_id=adrq.target_dag_id, 

645 asset_id=asset_id, 

646 dag_display_name=adrq.dag_model.dag_display_name, 

647 ) 

648 

649 

650@assets_router.delete( 

651 "/assets/{asset_id}/queuedEvents", 

652 status_code=status.HTTP_204_NO_CONTENT, 

653 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

654 dependencies=[ 

655 Depends(requires_access_asset(method="DELETE")), 

656 Depends(requires_access_dag(method="PUT")), 

657 Depends(action_logging()), 

658 ], 

659) 

660def delete_asset_queued_events( 

661 asset_id: int, 

662 readable_dags_filter: ReadableDagsFilterDep, 

663 session: SessionDep, 

664 before: OptionalDateTimeQuery = None, 

665): 

666 """Delete queued asset events for an asset.""" 

667 where_clause = _generate_queued_event_where_clause( 

668 asset_id=asset_id, before=before, permitted_dag_ids=readable_dags_filter.value 

669 ) 

670 delete_stmt = delete(AssetDagRunQueue).where(*where_clause) 

671 result = cast("CursorResult", session.execute(delete_stmt)) 

672 if result.rowcount == 0: 

673 raise HTTPException( 

674 status.HTTP_404_NOT_FOUND, 

675 detail=f"Queue event with asset_id: `{asset_id}` was not found", 

676 ) 

677 

678 

679@assets_router.delete( 

680 "/dags/{dag_id}/assets/queuedEvents", 

681 status_code=status.HTTP_204_NO_CONTENT, 

682 responses=create_openapi_http_exception_doc( 

683 [ 

684 status.HTTP_400_BAD_REQUEST, 

685 status.HTTP_404_NOT_FOUND, 

686 ] 

687 ), 

688 dependencies=[ 

689 Depends(requires_access_asset(method="DELETE")), 

690 Depends(requires_access_dag(method="PUT")), 

691 Depends(action_logging()), 

692 ], 

693) 

694def delete_dag_asset_queued_events( 

695 dag_id: str, 

696 readable_dags_filter: ReadableDagsFilterDep, 

697 session: SessionDep, 

698 before: OptionalDateTimeQuery = None, 

699): 

700 where_clause = _generate_queued_event_where_clause( 

701 dag_id=dag_id, before=before, permitted_dag_ids=readable_dags_filter.value 

702 ) 

703 

704 delete_statement = delete(AssetDagRunQueue).where(*where_clause) 

705 result = cast("CursorResult", session.execute(delete_statement)) 

706 

707 if result.rowcount == 0: 707 ↛ exitline 707 didn't return from function 'delete_dag_asset_queued_events' because the condition on line 707 was always true

708 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Queue event with dag_id: `{dag_id}` was not found") 

709 

710 

711@assets_router.delete( 

712 "/dags/{dag_id}/assets/{asset_id}/queuedEvents", 

713 status_code=status.HTTP_204_NO_CONTENT, 

714 responses=create_openapi_http_exception_doc( 

715 [ 

716 status.HTTP_400_BAD_REQUEST, 

717 status.HTTP_404_NOT_FOUND, 

718 ] 

719 ), 

720 dependencies=[ 

721 Depends(requires_access_asset(method="DELETE")), 

722 Depends(requires_access_dag(method="PUT")), 

723 Depends(action_logging()), 

724 ], 

725) 

726def delete_dag_asset_queued_event( 

727 dag_id: str, 

728 asset_id: int, 

729 readable_dags_filter: ReadableDagsFilterDep, 

730 session: SessionDep, 

731 before: OptionalDateTimeQuery = None, 

732): 

733 """Delete a queued asset event for a Dag.""" 

734 where_clause = _generate_queued_event_where_clause( 

735 dag_id=dag_id, before=before, asset_id=asset_id, permitted_dag_ids=readable_dags_filter.value 

736 ) 

737 delete_statement = delete(AssetDagRunQueue).where(*where_clause) 

738 result = cast("CursorResult", session.execute(delete_statement)) 

739 if result.rowcount == 0: 739 ↛ exitline 739 didn't return from function 'delete_dag_asset_queued_event' because the condition on line 739 was always true

740 raise HTTPException( 

741 status.HTTP_404_NOT_FOUND, 

742 detail=f"Queued event with dag_id: `{dag_id}` and asset_id: `{asset_id}` was not found", 

743 )