Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/ui/assets.py: 20%
68 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 TYPE_CHECKING, cast
21import structlog
22from fastapi import Depends, HTTPException, status
23from sqlalchemy import ColumnElement, and_, case, exists, func, select, true
25from airflow.api_fastapi.common.db.common import SessionDep
26from airflow.api_fastapi.common.partition_helpers import load_partitioned_timetable
27from airflow.api_fastapi.common.router import AirflowRouter
28from airflow.api_fastapi.core_api.datamodels.ui.assets import (
29 NextRunAssetEventResponse,
30 NextRunAssetsResponse,
31)
32from airflow.api_fastapi.core_api.security import requires_access_asset, requires_access_dag
33from airflow.models import DagModel
34from airflow.models.asset import (
35 AssetActive,
36 AssetDagRunQueue,
37 AssetEvent,
38 AssetModel,
39 AssetPartitionDagRun,
40 DagScheduleAssetReference,
41 PartitionedAssetKeyLog,
42)
44if TYPE_CHECKING: 44 ↛ 45line 44 didn't jump to line 45 because the condition on line 44 was never true
45 from airflow.partition_mappers.base import RollupMapper
47log = structlog.get_logger(logger_name=__name__)
49assets_router = AirflowRouter(tags=["Asset"])
52@assets_router.get(
53 "/next_run_assets/{dag_id}",
54 dependencies=[Depends(requires_access_asset(method="GET")), Depends(requires_access_dag(method="GET"))],
55)
56def next_run_assets(
57 dag_id: str,
58 session: SessionDep,
59) -> NextRunAssetsResponse:
60 dag_model = DagModel.get_dagmodel(dag_id, session=session)
61 if dag_model is None:
62 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Dag with id {dag_id} was not found")
64 latest_run = dag_model.get_last_dagrun(session=session)
65 event_filter = (
66 AssetEvent.timestamp >= latest_run.logical_date if latest_run and latest_run.logical_date else true()
67 )
69 pending_partition_count: int | None = None
71 queued_expr: ColumnElement[int]
72 if is_partitioned := dag_model.timetable_partitioned:
73 # Count pending APDRs directly without joining DagScheduleAssetReference /
74 # AssetModel — neither table appears in the WHERE clause, so the joins
75 # would collapse a cartesian product (N declared assets × M pending APDRs)
76 # before `count()` folded it back. `next_run_assets` runs on every home
77 # and Dags-list render, so this matters.
78 pending_partition_count = session.scalar(
79 select(func.count())
80 .select_from(AssetPartitionDagRun)
81 .where(
82 AssetPartitionDagRun.target_dag_id == dag_id,
83 AssetPartitionDagRun.created_dag_run_id.is_(None),
84 )
85 )
86 queued_expr = case(
87 (
88 exists(
89 select(PartitionedAssetKeyLog.id)
90 .join(
91 AssetPartitionDagRun,
92 PartitionedAssetKeyLog.asset_partition_dag_run_id == AssetPartitionDagRun.id,
93 )
94 .where(
95 PartitionedAssetKeyLog.asset_id == AssetModel.id,
96 PartitionedAssetKeyLog.target_dag_id == dag_id,
97 AssetPartitionDagRun.created_dag_run_id.is_(None),
98 )
99 ),
100 1,
101 ),
102 else_=0,
103 )
104 else:
105 queued_expr = func.max(case((AssetDagRunQueue.asset_id.is_not(None), 1), else_=0))
107 query = (
108 select(
109 AssetModel.id,
110 AssetModel.uri,
111 AssetModel.name,
112 func.max(AssetEvent.timestamp).label("last_update"),
113 queued_expr.label("queued"),
114 AssetActive.name.label("active_name"),
115 )
116 .join(DagScheduleAssetReference, DagScheduleAssetReference.asset_id == AssetModel.id)
117 .join(AssetEvent, and_(AssetEvent.asset_id == AssetModel.id, event_filter), isouter=True)
118 .outerjoin(
119 AssetActive,
120 and_(AssetActive.name == AssetModel.name, AssetActive.uri == AssetModel.uri),
121 )
122 .where(DagScheduleAssetReference.dag_id == dag_id)
123 .group_by(AssetModel.id, AssetModel.uri, AssetModel.name, AssetActive.name)
124 .order_by(AssetModel.uri)
125 )
127 if not is_partitioned:
128 query = query.join(
129 AssetDagRunQueue,
130 and_(
131 AssetDagRunQueue.asset_id == AssetModel.id,
132 AssetDagRunQueue.target_dag_id == DagScheduleAssetReference.dag_id,
133 ),
134 isouter=True,
135 )
137 raw_rows = list(session.execute(query))
139 if not is_partitioned:
140 events = [
141 NextRunAssetEventResponse(
142 id=row.id,
143 name=row.name,
144 uri=row.uri,
145 last_update=row.last_update if row.queued else None,
146 asset_inactive=(row.active_name is None),
147 )
148 for row in raw_rows
149 ]
150 return NextRunAssetsResponse(asset_expression=dag_model.asset_expression, events=events)
152 # Partitioned Dags: enrich with per-asset received/required counts and rollup flag.
153 # FIFO matches the scheduler's pending-APDR processing order
154 # (``_create_dagruns_for_partitioned_asset_dags``), so the "next run" the UI
155 # surfaces is the same one the scheduler will fire next.
156 pending_apdr = session.execute(
157 select(AssetPartitionDagRun.id, AssetPartitionDagRun.partition_key)
158 .where(
159 AssetPartitionDagRun.target_dag_id == dag_id,
160 AssetPartitionDagRun.created_dag_run_id.is_(None),
161 )
162 .order_by(AssetPartitionDagRun.created_at)
163 .limit(1)
164 ).one_or_none()
166 has_rollup_mappers = dag_model.has_rollup_mappers
168 if pending_apdr is None:
169 # No pending APDR yet — mark rollup assets so the UI can handle them
170 # correctly (e.g. skip "Asset Triggered" in favour of the asset name view).
171 # Reads from the cached partition_mapper_info so no timetable load is needed.
172 events = [
173 NextRunAssetEventResponse(
174 id=row.id,
175 name=row.name,
176 uri=row.uri,
177 last_update=row.last_update if row.queued else None,
178 is_rollup=has_rollup_mappers and dag_model.is_rollup_asset(name=row.name, uri=row.uri),
179 asset_inactive=(row.active_name is None),
180 )
181 for row in raw_rows
182 ]
183 return NextRunAssetsResponse(
184 asset_expression=dag_model.asset_expression,
185 events=events,
186 pending_partition_count=pending_partition_count,
187 )
189 # Collect received upstream partition keys per asset for this partition run.
190 # Use a set to deduplicate: multiple events for the same key count as one.
191 received_keys_by_asset: dict[int, set[str]] = {}
192 for log_row in session.execute(
193 select(
194 PartitionedAssetKeyLog.asset_id,
195 PartitionedAssetKeyLog.source_partition_key,
196 ).where(PartitionedAssetKeyLog.asset_partition_dag_run_id == pending_apdr.id)
197 ):
198 received_keys_by_asset.setdefault(log_row.asset_id, set()).add(log_row.source_partition_key)
200 # The timetable is only needed to call ``to_upstream`` for rollup mappers.
201 # When the cached info shows no rollup mappers, skip loading it entirely.
202 rollup_timetable = load_partitioned_timetable(dag_id, session) if has_rollup_mappers else None
204 events = []
205 for row in raw_rows:
206 received_keys = sorted(received_keys_by_asset.get(row.id, set()))
207 required_keys: list[str] = [pending_apdr.partition_key]
208 is_rollup = has_rollup_mappers and dag_model.is_rollup_asset(name=row.name, uri=row.uri)
209 mapper_failed = False
210 if is_rollup and rollup_timetable is not None:
211 try:
212 mapper = rollup_timetable.get_partition_mapper(name=row.name, uri=row.uri)
213 required_keys = sorted(cast("RollupMapper", mapper).to_upstream(pending_apdr.partition_key))
214 except Exception:
215 # Mirror the scheduler's ``_resolve_asset_partition_status``: a
216 # misconfigured rollup mapper marks the asset as not-yet-satisfied
217 # and the Dag run is held. Without this branch the UI would fall
218 # through to the non-rollup default of 1/1 and silently show
219 # "ready" for a run the scheduler will never fire.
220 log.warning(
221 "Failed to evaluate rollup mapper; treating asset as not-yet-satisfied",
222 dag_id=dag_id,
223 asset_name=row.name,
224 asset_uri=row.uri,
225 partition_key=pending_apdr.partition_key,
226 exc_info=True,
227 )
228 mapper_failed = True
229 if mapper_failed:
230 received_keys = []
231 required_keys = []
232 received_count = 0
233 required_count = 1
234 last_update = None
235 else:
236 received_count = len(received_keys)
237 required_count = len(required_keys)
238 # Only surface last_update once all required upstream keys have arrived.
239 last_update = row.last_update if row.queued and received_count >= required_count else None
240 events.append(
241 NextRunAssetEventResponse(
242 id=row.id,
243 name=row.name,
244 uri=row.uri,
245 last_update=last_update,
246 received_count=received_count,
247 required_count=required_count,
248 received_keys=received_keys,
249 required_keys=required_keys,
250 is_rollup=is_rollup,
251 mapper_error=mapper_failed,
252 asset_inactive=(row.active_name is None),
253 )
254 )
256 return NextRunAssetsResponse(
257 asset_expression=dag_model.asset_expression,
258 events=events,
259 pending_partition_count=pending_partition_count,
260 )