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

64 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 collections.abc import Iterable, Sequence 

20from itertools import groupby 

21from operator import itemgetter 

22from typing import Annotated 

23 

24from fastapi import Depends, HTTPException, status 

25from sqlalchemy import and_, func, or_, select 

26 

27from airflow.api_fastapi.app import get_auth_manager 

28from airflow.api_fastapi.auth.managers.models.batch_apis import IsAuthorizedDagRequest 

29from airflow.api_fastapi.auth.managers.models.resource_details import ( 

30 DagDetails, 

31) 

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

33 SessionDep, 

34 apply_filters_to_select, 

35 paginated_select, 

36) 

37from airflow.api_fastapi.common.parameters import ( 

38 QueryLimit, 

39 QueryOffset, 

40 QueryParseImportErrorBundleNameFilter, 

41 QueryParseImportErrorFilenameFilter, 

42 QueryParseImportErrorFilenamePatternSearch, 

43 QueryParseImportErrorFilenamePrefixPatternSearch, 

44 SortParam, 

45) 

46from airflow.api_fastapi.common.router import AirflowRouter 

47from airflow.api_fastapi.core_api.datamodels.import_error import ( 

48 ImportErrorCollectionResponse, 

49 ImportErrorResponse, 

50) 

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

52from airflow.api_fastapi.core_api.security import ( 

53 AccessView, 

54 GetUserDep, 

55 requires_access_view, 

56) 

57from airflow.models import DagModel 

58from airflow.models.errors import ParseImportError 

59 

60REDACTED_STACKTRACE = "REDACTED - you do not have read permission on all Dags in the file" 

61import_error_router = AirflowRouter(tags=["Import Error"], prefix="/importErrors") 

62 

63 

64@import_error_router.get( 

65 "/{import_error_id}", 

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

67 dependencies=[ 

68 Depends(requires_access_view(AccessView.IMPORT_ERRORS)), 

69 ], 

70) 

71def get_import_error( 

72 import_error_id: int, 

73 session: SessionDep, 

74 user: GetUserDep, 

75) -> ImportErrorResponse: 

76 """Get an import error.""" 

77 error = session.scalar(select(ParseImportError).where(ParseImportError.id == import_error_id)) 

78 if error is None: 78 ↛ 83line 78 didn't jump to line 83 because the condition on line 78 was always true

79 raise HTTPException( 

80 status.HTTP_404_NOT_FOUND, 

81 f"The ImportError with import_error_id: `{import_error_id}` was not found", 

82 ) 

83 session.expunge(error) 

84 

85 auth_manager = get_auth_manager() 

86 readable_dag_ids = auth_manager.get_authorized_dag_ids(user=user) 

87 # ``ParseImportError.filename`` is a repository-relative path and 

88 # ``DagModel.fileloc`` is typically the absolute path those files were 

89 # loaded from, so matching on ``fileloc == filename`` would come back 

90 # empty in most real deployments. Match on ``relative_fileloc`` (and the 

91 # bundle name that scopes it) instead, which is the same key the list 

92 # endpoint already uses for the join below. 

93 file_dag_ids = set( 

94 session.scalars( 

95 select(DagModel.dag_id).where( 

96 DagModel.relative_fileloc == error.filename, 

97 DagModel.bundle_name == error.bundle_name, 

98 ) 

99 ).all() 

100 ) 

101 

102 # No Dags matched for this file -- either the file genuinely contains 

103 # no Dags (parse failed before any Dag was defined), or the name keys 

104 # did not resolve. Return the raw error in this case; a proper 

105 # admin-only path for unregistered files is tracked in a follow-up 

106 # issue (see https://github.com/apache/airflow/issues/67461). 

107 if not file_dag_ids: 

108 return error 

109 

110 # Can the user read any Dags in the file? 

111 if not readable_dag_ids.intersection(file_dag_ids): 

112 raise HTTPException( 

113 status.HTTP_403_FORBIDDEN, 

114 "You do not have read permission on any of the Dags in the file", 

115 ) 

116 

117 # Check if user has read access to all the Dags defined in the file 

118 if not file_dag_ids.issubset(readable_dag_ids): 

119 error.stacktrace = REDACTED_STACKTRACE 

120 return error 

121 

122 

123@import_error_router.get( 

124 "", 

125 dependencies=[ 

126 Depends(requires_access_view(AccessView.IMPORT_ERRORS)), 

127 ], 

128) 

129def get_import_errors( 

130 limit: QueryLimit, 

131 offset: QueryOffset, 

132 order_by: Annotated[ 

133 SortParam, 

134 Depends( 

135 SortParam( 

136 [ 

137 "id", 

138 "timestamp", 

139 "filename", 

140 "bundle_name", 

141 "stacktrace", 

142 ], 

143 ParseImportError, 

144 {"import_error_id": "id"}, 

145 ).dynamic_depends() 

146 ), 

147 ], 

148 filename_pattern: QueryParseImportErrorFilenamePatternSearch, 

149 filename_prefix_pattern: QueryParseImportErrorFilenamePrefixPatternSearch, 

150 filename: QueryParseImportErrorFilenameFilter, 

151 bundle_name: QueryParseImportErrorBundleNameFilter, 

152 session: SessionDep, 

153 user: GetUserDep, 

154) -> ImportErrorCollectionResponse: 

155 """Get all import errors.""" 

156 auth_manager = get_auth_manager() 

157 readable_dag_ids = auth_manager.get_authorized_dag_ids(method="GET", user=user) 

158 

159 # Subquery for files that have any Dags 

160 files_with_any_dags = select(DagModel.relative_fileloc).distinct().subquery() 

161 

162 # Files (identified by ``(relative_fileloc, bundle_name)``) where the 

163 # user can read at least one Dag. Used to decide which import errors 

164 # the user is allowed to see at all. 

165 readable_files_cte = ( 

166 select(DagModel.relative_fileloc, DagModel.bundle_name) 

167 .where(DagModel.dag_id.in_(readable_dag_ids)) 

168 .distinct() 

169 .cte() 

170 ) 

171 

172 # Full ``(relative_fileloc, dag_id, bundle_name)`` set for every file 

173 # the user can see at least one Dag of. Crucially this is **not** 

174 # filtered by ``readable_dag_ids`` -- the per-file authorization 

175 # check in the loop below needs the complete Dag set so it can 

176 # detect co-located Dags that the caller is not authorized to read 

177 # and redact the stacktrace accordingly. 

178 file_dags_cte = ( 

179 select(DagModel.relative_fileloc, DagModel.dag_id, DagModel.bundle_name) 

180 .join( 

181 readable_files_cte, 

182 and_( 

183 DagModel.relative_fileloc == readable_files_cte.c.relative_fileloc, 

184 DagModel.bundle_name == readable_files_cte.c.bundle_name, 

185 ), 

186 ) 

187 .cte() 

188 ) 

189 

190 # Prepare the import errors query by joining with the CTE above. 

191 # Each returned row will be a tuple: (ParseImportError, dag_id). 

192 # ``dag_id`` is NULL for import errors whose file has no Dags at all 

193 # in ``DagModel`` (parse failed before any Dag was defined). 

194 import_errors_stmt = ( 

195 select(ParseImportError, file_dags_cte.c.dag_id) 

196 .outerjoin( 

197 files_with_any_dags, 

198 ParseImportError.filename == files_with_any_dags.c.relative_fileloc, 

199 ) 

200 .outerjoin( 

201 file_dags_cte, 

202 and_( 

203 ParseImportError.filename == file_dags_cte.c.relative_fileloc, 

204 ParseImportError.bundle_name == file_dags_cte.c.bundle_name, 

205 ), 

206 ) 

207 .where( 

208 or_( 

209 files_with_any_dags.c.relative_fileloc.is_(None), 

210 file_dags_cte.c.dag_id.isnot(None), 

211 ) 

212 ) 

213 .order_by(ParseImportError.id) 

214 ) 

215 

216 filtered_import_errors_stmt = apply_filters_to_select( 

217 statement=import_errors_stmt, 

218 filters=[filename_pattern, filename_prefix_pattern, filename, bundle_name], 

219 ) 

220 import_error_ids_stmt = ( 

221 filtered_import_errors_stmt.with_only_columns( 

222 ParseImportError.id, 

223 ParseImportError.timestamp, 

224 ParseImportError.filename, 

225 ParseImportError.bundle_name, 

226 ParseImportError.stacktrace, 

227 ) 

228 .distinct() 

229 .order_by(None) 

230 ) 

231 total_entries = session.scalar(select(func.count()).select_from(import_error_ids_stmt.subquery())) or 0 

232 

233 # Paginate distinct import error IDs first so limit/offset apply to 

234 # import error objects, not to the joined Dag rows. 

235 paginated_import_error_ids_select, _ = paginated_select( 

236 statement=import_error_ids_stmt, 

237 filters=[], 

238 order_by=order_by, 

239 offset=offset, 

240 limit=limit, 

241 session=session, 

242 return_total_entries=False, 

243 ) 

244 paginated_import_error_ids = paginated_import_error_ids_select.subquery() 

245 

246 # Fetch all joined Dag rows for the paginated import error IDs before 

247 # grouping, so each returned import error still has the full Dag set. 

248 import_errors_select = apply_filters_to_select( 

249 statement=filtered_import_errors_stmt.where( 

250 ParseImportError.id.in_(select(paginated_import_error_ids.c.id)) 

251 ), 

252 filters=[order_by], 

253 ) 

254 import_errors_result: Iterable[tuple[ParseImportError, Iterable]] = groupby( 

255 session.execute(import_errors_select), itemgetter(0) 

256 ) 

257 

258 import_errors = [] 

259 for import_error, file_dag_ids_iter in import_errors_result: 259 ↛ 260line 259 didn't jump to line 260 because the loop on line 259 never started

260 dag_ids = [dag_id for _, dag_id in file_dag_ids_iter if dag_id is not None] 

261 

262 # No Dags matched for this file -- either the file genuinely has 

263 # no Dags yet (parse failed before any Dag was defined), or the 

264 # name keys did not resolve. Append the raw error in this case; 

265 # a proper admin-only path for unregistered files is tracked in 

266 # a follow-up issue 

267 # (see https://github.com/apache/airflow/issues/67461). 

268 if not dag_ids: 

269 import_errors.append(import_error) 

270 continue 

271 

272 dag_id_to_team = DagModel.get_dag_id_to_team_name_mapping(dag_ids, session=session) 

273 # Check if user has read access to all the Dags defined in the file 

274 requests: Sequence[IsAuthorizedDagRequest] = [ 

275 { 

276 "method": "GET", 

277 "details": DagDetails(id=dag_id, team_name=dag_id_to_team.get(dag_id)), 

278 } 

279 for dag_id in dag_ids 

280 ] 

281 if not auth_manager.batch_is_authorized_dag(requests, user=user): 

282 session.expunge(import_error) 

283 import_error.stacktrace = REDACTED_STACKTRACE 

284 import_errors.append(import_error) 

285 

286 return ImportErrorCollectionResponse( 

287 import_errors=import_errors, 

288 total_entries=total_entries, 

289 )