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
« 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 collections.abc import Iterable, Sequence
20from itertools import groupby
21from operator import itemgetter
22from typing import Annotated
24from fastapi import Depends, HTTPException, status
25from sqlalchemy import and_, func, or_, select
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
60REDACTED_STACKTRACE = "REDACTED - you do not have read permission on all Dags in the file"
61import_error_router = AirflowRouter(tags=["Import Error"], prefix="/importErrors")
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)
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 )
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
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 )
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
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)
159 # Subquery for files that have any Dags
160 files_with_any_dags = select(DagModel.relative_fileloc).distinct().subquery()
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 )
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 )
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 )
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
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()
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 )
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]
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
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)
286 return ImportErrorCollectionResponse(
287 import_errors=import_errors,
288 total_entries=total_entries,
289 )