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
« 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
20from datetime import datetime
21from typing import TYPE_CHECKING, Annotated, cast
23from fastapi import Depends, HTTPException, status
24from sqlalchemy import and_, delete, func, select
25from sqlalchemy.engine import CursorResult
26from sqlalchemy.orm import joinedload, subqueryload
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
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
96assets_router = AirflowRouter(tags=["Asset"])
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
119class OnlyActiveFilter(BaseParam[bool]):
120 """Filter on asset activeness."""
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
127 @classmethod
128 def depends(cls, only_active: bool = True) -> OnlyActiveFilter:
129 return cls().set_value(only_active)
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 )
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 )
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)
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 )
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 )
208 assets = []
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 ]
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)
233 return AssetCollectionResponse(
234 assets=assets,
235 total_entries=total_entries,
236 )
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 )
265 return AssetAliasCollectionResponse(
266 asset_aliases=session.scalars(asset_aliases_select),
267 total_entries=total_entries,
268 )
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)
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)
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 )
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()
360 return AssetEventCollectionResponse(
361 asset_events=serialize_asset_events(assets_events, session=session),
362 total_entries=total_entries,
363 )
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()
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
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 )
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)
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 )
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 )
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 )
454 dag = get_latest_version_of_dag(dag_bag, dag_id, session)
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 )
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
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))
499 dag_asset_queued_events_select, total_entries = paginated_select(statement=query)
500 adrqs = session.scalars(dag_asset_queued_events_select).all()
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 ]
512 return QueuedEventCollectionResponse(
513 queued_events=queued_events,
514 total_entries=total_entries,
515 )
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 )
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()
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 )
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
558 if asset is None:
559 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Asset with ID: `{asset_id}` was not found")
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 ]
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 )
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))
599 dag_asset_queued_events_select, total_entries = paginated_select(statement=query)
600 adrqs = session.scalars(dag_asset_queued_events_select).all()
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 ]
612 return QueuedEventCollectionResponse(
613 queued_events=queued_events,
614 total_entries=total_entries,
615 )
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 )
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 )
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 )
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 )
704 delete_statement = delete(AssetDagRunQueue).where(*where_clause)
705 result = cast("CursorResult", session.execute(delete_statement))
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")
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 )