Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/ui/dags.py: 37%

58 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, status 

23from sqlalchemy import select, union_all 

24from sqlalchemy.orm import defaultload 

25 

26from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity 

27from airflow.api_fastapi.common.db.common import ( 

28 SessionDep, 

29 paginated_select, 

30) 

31from airflow.api_fastapi.common.db.dags import generate_dag_with_latest_run_query 

32from airflow.api_fastapi.common.parameters import ( 

33 FilterOptionEnum, 

34 FilterParam, 

35 QueryAnyDagRunStateFilter, 

36 QueryAssetDependencyFilter, 

37 QueryBundleNameFilter, 

38 QueryBundleVersionFilter, 

39 QueryDagDisplayNamePatternSearch, 

40 QueryDagDisplayNamePrefixPatternSearch, 

41 QueryDagIdPatternSearch, 

42 QueryDagIdPrefixPatternSearch, 

43 QueryExcludeStaleFilter, 

44 QueryFavoriteFilter, 

45 QueryHasAssetScheduleFilter, 

46 QueryHasImportErrorsFilter, 

47 QueryLastDagRunStateFilter, 

48 QueryLimit, 

49 QueryOffset, 

50 QueryOwnersFilter, 

51 QueryPausedFilter, 

52 QueryPendingActionsFilter, 

53 QueryTagsFilter, 

54 SortParam, 

55 filter_param_factory, 

56) 

57from airflow.api_fastapi.common.router import AirflowRouter 

58from airflow.api_fastapi.core_api.datamodels.dags import DAG_ALIAS_MAPPING, DAGResponse 

59from airflow.api_fastapi.core_api.datamodels.ui.dag_runs import DAGRunLightResponse 

60from airflow.api_fastapi.core_api.datamodels.ui.dags import ( 

61 DAGWithLatestDagRunsCollectionResponse, 

62 DAGWithLatestDagRunsResponse, 

63) 

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

65from airflow.api_fastapi.core_api.security import ( 

66 GetUserDep, 

67 ReadableDagsFilterDep, 

68 requires_access_dag, 

69) 

70from airflow.models import DagModel, DagRun 

71from airflow.models.dag_favorite import DagFavorite 

72from airflow.models.hitl import HITLDetail 

73from airflow.models.taskinstance import TaskInstance 

74from airflow.utils.state import TaskInstanceState 

75 

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

77 

78# Per-dag run counts read at most this many rows per state; the UI shows "N+" at the cap. 

79STATE_COUNT_CAP = 1000 

80 

81 

82@dags_router.get( 

83 "", 

84 response_model_exclude_none=True, 

85 dependencies=[ 

86 Depends(requires_access_dag(method="GET")), 

87 Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN)), 

88 Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.HITL_DETAIL)), 

89 Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK_INSTANCE)), 

90 ], 

91 operation_id="get_dags_ui", 

92) 

93def get_dags( 

94 limit: QueryLimit, 

95 offset: QueryOffset, 

96 tags: QueryTagsFilter, 

97 owners: QueryOwnersFilter, 

98 dag_ids: Annotated[ 

99 FilterParam[list[str] | None], 

100 Depends(filter_param_factory(DagModel.dag_id, list[str] | None, FilterOptionEnum.IN, "dag_ids")), 

101 ], 

102 dag_id_pattern: QueryDagIdPatternSearch, 

103 dag_id_prefix_pattern: QueryDagIdPrefixPatternSearch, 

104 dag_display_name_pattern: QueryDagDisplayNamePatternSearch, 

105 dag_display_name_prefix_pattern: QueryDagDisplayNamePrefixPatternSearch, 

106 exclude_stale: QueryExcludeStaleFilter, 

107 paused: QueryPausedFilter, 

108 has_import_errors: QueryHasImportErrorsFilter, 

109 last_dag_run_state: QueryLastDagRunStateFilter, 

110 dag_run_state: QueryAnyDagRunStateFilter, 

111 bundle_name: QueryBundleNameFilter, 

112 bundle_version: QueryBundleVersionFilter, 

113 order_by: Annotated[ 

114 SortParam, 

115 Depends( 

116 SortParam( 

117 ["dag_id", "dag_display_name", "next_dagrun", "state", "start_date"], 

118 DagModel, 

119 {"last_run_state": DagRun.state, "last_run_start_date": DagRun.start_date}, 

120 ).dynamic_depends() 

121 ), 

122 ], 

123 is_favorite: QueryFavoriteFilter, 

124 has_asset_schedule: QueryHasAssetScheduleFilter, 

125 asset_dependency: QueryAssetDependencyFilter, 

126 has_pending_actions: QueryPendingActionsFilter, 

127 readable_dags_filter: ReadableDagsFilterDep, 

128 session: SessionDep, 

129 user: GetUserDep, 

130 dag_runs_limit: int = 10, 

131) -> DAGWithLatestDagRunsCollectionResponse: 

132 """Get Dags with recent DagRun.""" 

133 # Fetch Dags with their latest DagRun and apply filters 

134 query = generate_dag_with_latest_run_query( 

135 max_run_filters=[ 

136 last_dag_run_state, 

137 ], 

138 order_by=order_by, 

139 dag_ids=readable_dags_filter.value, 

140 ) 

141 

142 dags_select, total_entries = paginated_select( 

143 statement=query, 

144 filters=[ 

145 exclude_stale, 

146 paused, 

147 has_import_errors, 

148 dag_id_pattern, 

149 dag_id_prefix_pattern, 

150 dag_ids, 

151 dag_display_name_pattern, 

152 dag_display_name_prefix_pattern, 

153 tags, 

154 owners, 

155 last_dag_run_state, 

156 dag_run_state, 

157 is_favorite, 

158 has_asset_schedule, 

159 asset_dependency, 

160 has_pending_actions, 

161 readable_dags_filter, 

162 bundle_name, 

163 bundle_version, 

164 ], 

165 order_by=order_by, 

166 offset=offset, 

167 limit=limit, 

168 session=session, 

169 ) 

170 

171 dags = [dag for dag in session.scalars(dags_select)] 

172 

173 # Fetch favorite status for each Dag for the current user 

174 user_id = str(user.get_id()) 

175 favorites_select = select(DagFavorite.dag_id).where( 

176 DagFavorite.user_id == user_id, DagFavorite.dag_id.in_([dag.dag_id for dag in dags]) 

177 ) 

178 favorite_dag_ids = set(session.scalars(favorites_select)) 

179 

180 recent_dag_runs: list = [] 

181 if dags: 

182 recent_runs_branches = [ 

183 select( 

184 DagRun.id, 

185 DagRun.dag_id, 

186 DagRun.run_id, 

187 DagRun.end_date, 

188 DagRun.logical_date, 

189 DagRun.run_after, 

190 DagRun.start_date, 

191 DagRun.state, 

192 ) 

193 .where(DagRun.dag_id == dag.dag_id) 

194 .order_by(DagRun.run_after.desc()) 

195 .limit(dag_runs_limit) 

196 .subquery() 

197 for dag in dags 

198 ] 

199 recent_runs_union = union_all(*(select(branch) for branch in recent_runs_branches)).subquery() 

200 recent_dag_runs = list( 

201 session.execute(select(recent_runs_union).order_by(recent_runs_union.c.run_after.desc())) 

202 ) 

203 

204 # Fetch pending HITL actions for each Dag if we are not certain whether some of the Dag might contain HITL actions 

205 pending_actions_by_dag_id: dict[str, list[HITLDetail]] = {dag.dag_id: [] for dag in dags} 

206 if has_pending_actions.value: 

207 pending_actions_select = ( 

208 select( 

209 TaskInstance.dag_id, 

210 HITLDetail, 

211 ) 

212 .join(TaskInstance, HITLDetail.ti_id == TaskInstance.id) 

213 .options( 

214 defaultload(HITLDetail.task_instance).joinedload(TaskInstance.rendered_task_instance_fields) 

215 ) 

216 .where( 

217 HITLDetail.responded_at.is_(None), 

218 TaskInstance.state.in_((TaskInstanceState.DEFERRED, TaskInstanceState.AWAITING_INPUT)), 

219 ) 

220 .where(TaskInstance.dag_id.in_([dag.dag_id for dag in dags])) 

221 .order_by(TaskInstance.dag_id) 

222 ) 

223 

224 pending_actions = session.execute(pending_actions_select) 

225 

226 # Group pending actions by dag_id 

227 for dag_id, hitl_detail in pending_actions: 

228 pending_actions_by_dag_id[dag_id].append(hitl_detail) 

229 

230 # aggregate rows by dag_id 

231 # Build the dict dynamically from DAGResponse.model_fields so that new fields 

232 # added to DAGResponse are picked up automatically without code changes here. 

233 dag_runs_by_dag_id: dict[str, DAGWithLatestDagRunsResponse] = {} 

234 for dag in dags: 

235 dag_data = { 

236 DAG_ALIAS_MAPPING.get(field_name, field_name): getattr( 

237 dag, DAG_ALIAS_MAPPING.get(field_name, field_name) 

238 ) 

239 for field_name in DAGResponse.model_fields 

240 } 

241 dag_data.update( 

242 { 

243 "asset_expression": dag.asset_expression, 

244 "latest_dag_runs": [], 

245 "pending_actions": pending_actions_by_dag_id[dag.dag_id], 

246 "is_favorite": dag.dag_id in favorite_dag_ids, 

247 } 

248 ) 

249 dag_runs_by_dag_id[dag.dag_id] = DAGWithLatestDagRunsResponse.model_validate(dag_data) 

250 

251 for row in recent_dag_runs: 

252 dag_run_response = DAGRunLightResponse.model_validate(row) 

253 dag_id = dag_run_response.dag_id 

254 dag_runs_by_dag_id[dag_id].latest_dag_runs.append(dag_run_response) 

255 

256 return DAGWithLatestDagRunsCollectionResponse( 

257 total_entries=total_entries, 

258 dags=list(dag_runs_by_dag_id.values()), 

259 ) 

260 

261 

262@dags_router.get( 

263 "/{dag_id}/latest_run", 

264 responses=create_openapi_http_exception_doc( 

265 [ 

266 status.HTTP_400_BAD_REQUEST, 

267 status.HTTP_404_NOT_FOUND, 

268 ] 

269 ), 

270 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN))], 

271) 

272def get_latest_run_info(dag_id: str, session: SessionDep) -> DAGRunLightResponse | None: 

273 """Get latest run.""" 

274 if dag_id == "~": 

275 raise HTTPException( 

276 status.HTTP_400_BAD_REQUEST, 

277 "`~` was supplied as dag_id, but querying multiple dags is not supported.", 

278 ) 

279 

280 latest_run_info_select = ( 

281 select( 

282 DagRun.id, 

283 DagRun.dag_id, 

284 DagRun.run_id, 

285 DagRun.end_date, 

286 DagRun.logical_date, 

287 DagRun.run_after, 

288 DagRun.start_date, 

289 DagRun.state, 

290 ) 

291 .where(DagRun.dag_id == dag_id) 

292 .order_by(DagRun.run_after.desc()) 

293 .limit(1) 

294 ) 

295 latest_run_info = session.execute(latest_run_info_select).one_or_none() 

296 

297 return DAGRunLightResponse(**latest_run_info._mapping) if latest_run_info else None