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

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 TYPE_CHECKING, cast 

20 

21import structlog 

22from fastapi import Depends, HTTPException, status 

23from sqlalchemy import ColumnElement, and_, case, exists, func, select, true 

24 

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) 

43 

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 

46 

47log = structlog.get_logger(logger_name=__name__) 

48 

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

50 

51 

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

63 

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 ) 

68 

69 pending_partition_count: int | None = None 

70 

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

106 

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 ) 

126 

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 ) 

136 

137 raw_rows = list(session.execute(query)) 

138 

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) 

151 

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

165 

166 has_rollup_mappers = dag_model.has_rollup_mappers 

167 

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 ) 

188 

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) 

199 

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 

203 

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 ) 

255 

256 return NextRunAssetsResponse( 

257 asset_expression=dag_model.asset_expression, 

258 events=events, 

259 pending_partition_count=pending_partition_count, 

260 )