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

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 fastapi import Depends, HTTPException, status 

20from sqlalchemy import select 

21from sqlalchemy.orm import joinedload 

22 

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 

39 

40structure_router = AirflowRouter(tags=["Structure"], prefix="/structure") 

41 

42 

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 

77 

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 

90 

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 ) 

98 

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) 

105 

106 data = { 

107 "nodes": nodes, 

108 "edges": edges, 

109 } 

110 

111 if external_dependencies: 

112 entry_node_ref = nodes[0] if nodes else None 

113 exit_node_ref = nodes[-1] if nodes else None 

114 

115 start_edges: list[dict] = [] 

116 end_edges: list[dict] = [] 

117 

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 

134 

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 

141 

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"]}) 

148 

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 ) 

162 

163 # Add nodes 

164 nodes.append( 

165 { 

166 "id": dependency.node_id, 

167 "label": dependency.label, 

168 "type": dependency.dependency_type, 

169 } 

170 ) 

171 

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 

184 

185 data["edges"] += start_edges + end_edges 

186 

187 bind_output_assets_to_tasks(data["edges"], serialized_dag, version_number, session) 

188 

189 return StructureDataResponse(**data)