Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/ui/structure.py: 21%
66 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, status
20from sqlalchemy import select
21from sqlalchemy.orm import joinedload
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.parameters import QueryIncludeDownstream, QueryIncludeUpstream
26from airflow.api_fastapi.common.router import AirflowRouter
27from airflow.api_fastapi.core_api.datamodels.ui.structure import StructureDataResponse
28from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
29from airflow.api_fastapi.core_api.security import ReadableDagsFilterDep, requires_access_dag
30from airflow.api_fastapi.core_api.services.ui.structure import (
31 bind_output_assets_to_tasks,
32 get_upstream_assets,
33)
34from airflow.api_fastapi.core_api.services.ui.task_group import task_group_to_dict
35from airflow.models.dag import DagModel
36from airflow.models.dag_version import DagVersion
37from airflow.models.serialized_dag import SerializedDagModel
38from airflow.utils.dag_edges import dag_edges
40structure_router = AirflowRouter(tags=["Structure"], prefix="/structure")
43@structure_router.get(
44 "/structure_data",
45 responses=create_openapi_http_exception_doc(
46 [
47 status.HTTP_400_BAD_REQUEST,
48 status.HTTP_404_NOT_FOUND,
49 ]
50 ),
51 dependencies=[
52 Depends(requires_access_dag("GET")),
53 Depends(requires_access_dag("GET", DagAccessEntity.DEPENDENCIES)),
54 Depends(requires_access_dag("GET", DagAccessEntity.TASK_INSTANCE)),
55 ],
56)
57def structure_data(
58 session: SessionDep,
59 dag_id: str,
60 readable_dags_filter: ReadableDagsFilterDep,
61 include_upstream: QueryIncludeUpstream = False,
62 include_downstream: QueryIncludeDownstream = False,
63 depth: int | None = None,
64 root: str | None = None,
65 external_dependencies: bool = False,
66 version_number: int | None = None,
67) -> StructureDataResponse:
68 """Get Structure Data."""
69 if version_number is None:
70 dag_version_model = DagVersion.get_latest_version(dag_id)
71 if dag_version_model is None:
72 raise HTTPException(
73 status.HTTP_404_NOT_FOUND,
74 f"Dag with id {dag_id} was not found",
75 )
76 version_number = dag_version_model.version_number
78 serialized_dag: SerializedDagModel | None = session.scalar(
79 select(SerializedDagModel)
80 .join(DagVersion)
81 .where(SerializedDagModel.dag_id == dag_id, DagVersion.version_number == version_number)
82 .options(joinedload(SerializedDagModel.dag_model).joinedload(DagModel.task_outlet_asset_references)),
83 )
84 if serialized_dag is None:
85 raise HTTPException(
86 status.HTTP_404_NOT_FOUND,
87 f"Dag with id {dag_id} and version number {version_number} was not found",
88 )
89 dag = serialized_dag.dag
91 if root:
92 dag = dag.partial_subset(
93 task_ids=root,
94 include_upstream=include_upstream,
95 include_downstream=include_downstream,
96 depth=depth,
97 )
99 group_dict = dag.task_group.get_task_group_dict()
100 nodes = [
101 task_group_to_dict(child, group_dict=group_dict)
102 for child in dag.task_group.topological_sort(group_dict=group_dict)
103 ]
104 edges = dag_edges(dag)
106 data = {
107 "nodes": nodes,
108 "edges": edges,
109 }
111 if external_dependencies:
112 entry_node_ref = nodes[0] if nodes else None
113 exit_node_ref = nodes[-1] if nodes else None
115 start_edges: list[dict] = []
116 end_edges: list[dict] = []
118 readable_dag_ids = readable_dags_filter.value
119 for dependency_dag_id, dependencies in sorted(SerializedDagModel.get_dag_dependencies().items()):
120 if readable_dag_ids is not None and dependency_dag_id not in readable_dag_ids:
121 continue
122 for dependency in dependencies:
123 # Dependencies not related to `dag_id` are ignored
124 if dependency_dag_id != dag_id and dependency.target != dag_id:
125 continue
126 # When target is a real Dag ID (not a type label), hide it
127 # if the caller cannot read that Dag.
128 if (
129 readable_dag_ids is not None
130 and dependency.target != dependency.dependency_type
131 and dependency.target not in readable_dag_ids
132 ):
133 continue
135 # upstream assets are handled by the `get_upstream_assets` function.
136 if dependency.target != dependency.dependency_type and dependency.dependency_type in [
137 "asset-alias",
138 "asset",
139 ]:
140 continue
142 # Add edges
143 # start dependency
144 if (
145 dependency.source == dependency.dependency_type or dependency.target == dag_id
146 ) and entry_node_ref:
147 start_edges.append({"source_id": dependency.node_id, "target_id": entry_node_ref["id"]})
149 # end dependency
150 elif (
151 dependency.target == dependency.dependency_type or dependency.source == dag_id
152 ) and exit_node_ref:
153 end_edges.append(
154 {
155 "source_id": exit_node_ref["id"],
156 "target_id": dependency.node_id,
157 "resolved_from_alias": dependency.source.replace("asset-alias:", "", 1)
158 if dependency.source.startswith("asset-alias:")
159 else None,
160 }
161 )
163 # Add nodes
164 nodes.append(
165 {
166 "id": dependency.node_id,
167 "label": dependency.label,
168 "type": dependency.dependency_type,
169 }
170 )
172 if (asset_expression := serialized_dag.dag_model.asset_expression) and entry_node_ref:
173 try:
174 upstream_asset_nodes, upstream_asset_edges = get_upstream_assets(
175 asset_expression, entry_node_ref["id"]
176 )
177 except TypeError as e:
178 raise HTTPException(
179 status.HTTP_400_BAD_REQUEST,
180 f"Malformed asset_expression in Dag {dag_id!r} version {version_number}: {e}",
181 ) from e
182 data["nodes"] += upstream_asset_nodes
183 data["edges"] += upstream_asset_edges
185 data["edges"] += start_edges + end_edges
187 bind_output_assets_to_tasks(data["edges"], serialized_dag, version_number, session)
189 return StructureDataResponse(**data)