kaxil commented on code in PR #73724:
URL: https://github.com/apache/airflow/pull/73724#discussion_r4210407836


##########
airflow-core/tests/unit/utils/test_task_group.py:
##########
@@ -1351,6 +1351,88 @@ def spy(self, nodes, projected):
         assert position[f"r{i}"] < position[f"r{i + 1}"]
 
 
+def _make_sibling_groups_cycle():
+    with DAG("sibling_groups_cycle", schedule=None, start_date=DEFAULT_DATE) 
as dag:
+        start = EmptyOperator(task_id="start")
+        with TaskGroup("group1"):
+            a1 = EmptyOperator(task_id="a1")
+            a2 = EmptyOperator(task_id="a2")
+        with TaskGroup("group2"):
+            b1 = EmptyOperator(task_id="b1")
+            b2 = EmptyOperator(task_id="b2")
+        end = EmptyOperator(task_id="end")
+        start >> [a1, b2]
+        a1 >> b1
+        b2 >> a2
+        [a2, b1] >> end
+    return dag
+
+
+def _make_group_bridged_by_outside_task():
+    with DAG("group_bridged_by_outside_task", schedule=None, 
start_date=DEFAULT_DATE) as dag:
+        with TaskGroup("group"):
+            a = EmptyOperator(task_id="a")
+            b = EmptyOperator(task_id="b")
+        bridge = EmptyOperator(task_id="bridge")
+        a >> bridge >> b
+    return dag
+
+
+def _make_three_group_ring():
+    with DAG("three_group_ring", schedule=None, start_date=DEFAULT_DATE) as 
dag:
+        groups = {}
+        for group_id in ("g0", "g1", "g2"):
+            with TaskGroup(group_id):
+                groups[group_id] = (EmptyOperator(task_id="first"), 
EmptyOperator(task_id="second"))
+        groups["g1"][0] >> groups["g0"][1]
+        groups["g2"][0] >> groups["g1"][1]
+        groups["g0"][0] >> groups["g2"][1]
+    return dag
+
+
[email protected](
+    ("make_dag", "expected"),
+    [
+        pytest.param(
+            _make_sibling_groups_cycle,
+            {
+                None: ["start", "group1", "group2", "end"],
+                "group1": ["group1.a1", "group1.a2"],
+                "group2": ["group2.b1", "group2.b2"],
+            },
+            id="sibling-groups",
+        ),
+        pytest.param(
+            _make_group_bridged_by_outside_task,
+            {None: ["bridge", "group"], "group": ["group.a", "group.b"]},

Review Comment:
   The deserializer sorts children by label, so the expected order here, in 
`three-group-ring`, and in the two API tests is already the input order. If 
`_sort_cyclic_projection` just returned `list(nodes)`, all of those would still 
pass and only `sibling-groups` would fail. This is also the only case that 
reaches the fallback through `_sweep_projection`, so nothing checks that the 
sweep path places the cycle correctly among its other siblings.
   
   A shape where label order is wrong would pin it: `group(a, b)` plus 
top-level `bridge`, `after`, `x0`, `x1`, `x2`, wired `a >> bridge >> b >> 
after`. The input order is `[after, bridge, group, x0, x1, x2]`, this PR 
returns `[bridge, group, x0, x1, x2, after]`, and the identity version returns 
the input unchanged. The nested `outer.inner` / `outer.bridge` shape from the 
screenshots would be worth a param too, since every new case has the cycle at 
the root.



##########
airflow-core/src/airflow/serialization/definitions/taskgroup.py:
##########
@@ -236,10 +236,13 @@ def topological_sort(
         """
         Sort children topologically — a task always comes after its upstream 
dependencies.
 
-        See ``TaskGroup.topological_sort`` in task-sdk for the algorithm. 
Cycles are
-        treated as corrupt input: ``DAG.check_cycle`` rejects cyclic Dags 
before
-        serialization, so a cycle reaching this code indicates malformed 
serialized data,
-        and we raise ``ValueError`` rather than silently looping forever.
+        See ``TaskGroup.topological_sort`` in task-sdk for the algorithm. 
Unlike the task-sdk

Review Comment:
   The description says the Task SDK sort keeps raising so it can detect these 
Dags at parse time, but #73746 adds its own `_find_dependency_cycles` for that, 
and nothing at parse time calls the SDK `TaskGroup.topological_sort`. Its only 
in-tree caller is the deprecated `DAG.topological_sort()`, which keeps raising 
`AirflowDagCycleException` for Dags that pass `check_cycle()` and that #73746 
only warns about. Should the SDK side get the same fallback so the two sorts 
agree, or is the raise there intentional?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to