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
« 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."""
20from __future__ import annotations
22from functools import cache
23from typing import TYPE_CHECKING
25from airflow.configuration import conf
26from airflow.serialization.definitions.baseoperator import SerializedBaseOperator
27from airflow.serialization.definitions.mappedoperator import SerializedMappedOperator, is_mapped
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
33 from airflow.serialization.definitions.taskgroup import SerializedTaskGroup
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()
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
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 ]
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"})
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"})
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
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.
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]
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
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 ]
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