Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/dag_sources.py: 87%
35 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 fastapi import Depends, HTTPException, Response, status
20from sqlalchemy import select
22from airflow.api_fastapi.app import get_auth_manager
23from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity
24from airflow.api_fastapi.common.db.common import SessionDep
25from airflow.api_fastapi.common.headers import HeaderAcceptJsonOrText
26from airflow.api_fastapi.common.router import AirflowRouter
27from airflow.api_fastapi.common.types import Mimetype
28from airflow.api_fastapi.core_api.datamodels.dag_sources import DAGSourceResponse
29from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
30from airflow.api_fastapi.core_api.security import GetUserDep, requires_access_dag
31from airflow.models import DagModel
32from airflow.models.dag_version import DagVersion
34REDACTED_SOURCE = "REDACTED - you do not have read permission on all Dags in the file"
35dag_sources_router = AirflowRouter(tags=["DagSource"], prefix="/dagSources")
38@dag_sources_router.get(
39 "/{dag_id}",
40 responses={
41 **create_openapi_http_exception_doc(
42 [
43 status.HTTP_400_BAD_REQUEST,
44 status.HTTP_404_NOT_FOUND,
45 status.HTTP_406_NOT_ACCEPTABLE,
46 ]
47 ),
48 "200": {
49 "description": "Successful Response",
50 "content": {
51 Mimetype.TEXT: {"schema": {"type": "string", "example": "dag code"}},
52 },
53 },
54 },
55 response_model=DAGSourceResponse,
56 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.CODE))],
57)
58def get_dag_source(
59 accept: HeaderAcceptJsonOrText,
60 dag_id: str,
61 session: SessionDep,
62 user: GetUserDep,
63 version_number: int | None = None,
64):
65 """Get source code using file token."""
66 dag_version = DagVersion.get_version(dag_id, version_number, session=session)
67 if not dag_version:
68 raise HTTPException(
69 status.HTTP_404_NOT_FOUND,
70 f"The source code of the Dag {dag_id}, version_number {version_number} was not found",
71 )
72 if not dag_version.dag_code: 72 ↛ 73line 72 didn't jump to line 73 because the condition on line 72 was never true
73 raise HTTPException(
74 status.HTTP_404_NOT_FOUND,
75 detail=f"Code not found. dag_id='{dag_id}' version_number='{version_number}'",
76 )
78 # Per-file authorization overlay on top of the ``DagAccessEntity.CODE``
79 # check above: a single source file may define multiple Dags, and the
80 # caller having CODE access to ``dag_id`` does not imply they may read
81 # every other Dag co-located in the file. Match the file by
82 # ``(relative_fileloc, bundle_name)`` -- the same keying
83 # ``import_error.py`` uses for its equivalent check -- and redact the
84 # response when any co-located Dag is not in the caller's readable set.
85 content = dag_version.dag_code.source_code
86 dag_model = dag_version.dag_model
87 if dag_model is not None and dag_model.relative_fileloc: 87 ↛ 101line 87 didn't jump to line 101 because the condition on line 87 was always true
88 file_dag_ids = set(
89 session.scalars(
90 select(DagModel.dag_id).where(
91 DagModel.relative_fileloc == dag_model.relative_fileloc,
92 DagModel.bundle_name == dag_model.bundle_name,
93 )
94 ).all()
95 )
96 if file_dag_ids: 96 ↛ 101line 96 didn't jump to line 101 because the condition on line 96 was always true
97 readable_dag_ids = get_auth_manager().get_authorized_dag_ids(user=user)
98 if not file_dag_ids.issubset(readable_dag_ids): 98 ↛ 99line 98 didn't jump to line 99 because the condition on line 98 was never true
99 content = REDACTED_SOURCE
101 dag_source_model = DAGSourceResponse(
102 dag_id=dag_id,
103 content=content,
104 version_number=dag_version.version_number,
105 dag_display_name=dag_version.dag_model.dag_display_name,
106 )
108 if accept == Mimetype.TEXT:
109 return Response(dag_source_model.content, media_type=Mimetype.TEXT)
110 return dag_source_model