Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/services/ui/task_group.py: 13%

60 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 14:22 +0000

1# 

2# Licensed to the Apache Software Foundation (ASF) under one 

3# or more contributor license agreements. See the NOTICE file 

4# distributed with this work for additional information 

5# regarding copyright ownership. The ASF licenses this file 

6# to you under the Apache License, Version 2.0 (the 

7# "License"); you may not use this file except in compliance 

8# with the License. You may obtain a copy of the License at 

9# 

10# http://www.apache.org/licenses/LICENSE-2.0 

11# 

12# Unless required by applicable law or agreed to in writing, 

13# software distributed under the License is distributed on an 

14# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY 

15# KIND, either express or implied. See the License for the 

16# specific language governing permissions and limitations 

17# under the License. 

18"""Task group utilities for UI API services.""" 

19 

20from __future__ import annotations 

21 

22from functools import cache 

23from typing import TYPE_CHECKING 

24 

25from airflow.configuration import conf 

26from airflow.serialization.definitions.baseoperator import SerializedBaseOperator 

27from airflow.serialization.definitions.mappedoperator import SerializedMappedOperator, is_mapped 

28 

29if TYPE_CHECKING: 29 ↛ 30line 29 didn't jump to line 30 because the condition on line 29 was never true

30 from collections.abc import Callable 

31 from typing import Any 

32 

33 from airflow.serialization.definitions.taskgroup import SerializedTaskGroup 

34 

35 

36@cache 

37def get_task_group_children_getter() -> Callable: 

38 """Get the Task Group Children Getter for the Dag.""" 

39 if conf.get("api", "grid_view_sorting_order") == "topological": 

40 return lambda task_group, group_dict=None: task_group.topological_sort(group_dict=group_dict) 

41 return lambda task_group, group_dict=None: task_group.hierarchical_alphabetical_sort() 

42 

43 

44def task_group_to_dict(task_item_or_group, *, group_dict=None, parent_group_is_mapped=False): 

45 """Create a nested dict representation of this TaskGroup and its children used to construct the Graph.""" 

46 if isinstance(task := task_item_or_group, (SerializedBaseOperator, SerializedMappedOperator)): 

47 # we explicitly want the short task ID here, not the full doted notation if in a group 

48 task_display_name = task.task_display_name if task.task_display_name != task.task_id else task.label 

49 node_operator = { 

50 "id": task.task_id, 

51 "label": task_display_name, 

52 "operator": task.operator_name, 

53 "type": "task", 

54 } 

55 if task.is_setup: 

56 node_operator["setup_teardown_type"] = "setup" 

57 elif task.is_teardown: 

58 node_operator["setup_teardown_type"] = "teardown" 

59 if is_mapped(task) or parent_group_is_mapped: 

60 node_operator["is_mapped"] = True 

61 return node_operator 

62 

63 task_group = task_item_or_group 

64 if group_dict is None: 

65 group_dict = task_group.dag.task_group.get_task_group_dict() 

66 mapped = is_mapped(task_group) 

67 children = [ 

68 task_group_to_dict( 

69 child, 

70 parent_group_is_mapped=parent_group_is_mapped or mapped, 

71 group_dict=group_dict, 

72 ) 

73 for child in get_task_group_children_getter()(task_group, group_dict) 

74 ] 

75 

76 if task_group.upstream_group_ids or task_group.upstream_task_ids: 

77 # This is the join node used to reduce the number of edges between two TaskGroup. 

78 children.append({"id": task_group.upstream_join_id, "label": "", "type": "join"}) 

79 

80 if task_group.downstream_group_ids or task_group.downstream_task_ids: 

81 # This is the join node used to reduce the number of edges between two TaskGroup. 

82 children.append({"id": task_group.downstream_join_id, "label": "", "type": "join"}) 

83 

84 node = { 

85 "id": task_group.group_id, 

86 "label": task_group.group_display_name or task_group.label, 

87 "tooltip": task_group.tooltip, 

88 "is_mapped": mapped, 

89 "children": children, 

90 "type": "task", 

91 } 

92 return node 

93 

94 

95def task_group_to_dict_grid( 

96 task_item_or_group, 

97 *, 

98 group_dict: dict[str | None, SerializedTaskGroup] | None = None, 

99 parent_group_is_mapped: bool = False, 

100) -> dict[str, Any]: 

101 """ 

102 Create a nested dict representation of this TaskGroup and its children used to construct the Grid. 

103 

104 :param group_dict: A ``{group_id: group}`` map used to resolve cross-group 

105 dependencies. Built once at the top of a render and threaded through the 

106 recursion so nested groups reuse it. 

107 :param parent_group_is_mapped: Whether an ancestor task group is mapped, propagated to children. 

108 """ 

109 node: dict[str, Any] 

110 

111 if isinstance(task := task_item_or_group, (SerializedMappedOperator, SerializedBaseOperator)): 

112 mapped = None 

113 if parent_group_is_mapped or is_mapped(task): 

114 mapped = True 

115 setup_teardown_type = None 

116 if task.is_setup is True: 

117 setup_teardown_type = "setup" 

118 elif task.is_teardown is True: 

119 setup_teardown_type = "teardown" 

120 # we explicitly want the short task ID here, not the full doted notation if in a group 

121 task_display_name = task.task_display_name if task.task_display_name != task.task_id else task.label 

122 node = { 

123 "id": task.task_id, 

124 "label": task_display_name, 

125 "is_mapped": mapped, 

126 "children": None, 

127 "setup_teardown_type": setup_teardown_type, 

128 } 

129 return node 

130 

131 task_group = task_item_or_group 

132 if group_dict is None: 

133 group_dict = task_group.dag.task_group.get_task_group_dict() 

134 task_group_sort = get_task_group_children_getter() 

135 mapped = is_mapped(task_group) 

136 children = [ 

137 task_group_to_dict_grid( 

138 child, 

139 group_dict=group_dict, 

140 parent_group_is_mapped=parent_group_is_mapped or mapped, 

141 ) 

142 for child in task_group_sort(task_group, group_dict) 

143 ] 

144 

145 node = { 

146 "id": task_group.group_id, 

147 "label": task_group.group_display_name or task_group.label, 

148 "is_mapped": mapped or None, 

149 "children": children or None, 

150 } 

151 if task_group.doc_md is not None: 

152 node["doc_md"] = task_group.doc_md 

153 return node