This is an automated email from the ASF dual-hosted git repository.

jason810496 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 84f5e0990a0 TS SDK: Step over an empty nested task group to the group 
holding it (#74323)
84f5e0990a0 is described below

commit 84f5e0990a06cb61c08c665c225591dca475f8f0
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Tue Oct 6 16:17:03 2026 +0800

    TS SDK: Step over an empty nested task group to the group holding it 
(#74323)
    
    * TS SDK: Step over an empty nested group to the group holding it
    
    An ordering edge out of a group that holds no tasks expanded onto nothing
    when the group was nested inside another: the walk followed the group
    edges drawn around it and stopped there. Python's `find_leaves` takes one
    more step and stands the group in for its parent, so the edge leaves from
    whatever the parent finishes with.
    
    The serialized `task_group` already recorded the edge, so the UI drew a
    constraint the scheduler did not enforce.
    
    The shared conformance Dags gain the shape that catches it, along with a
    group edge whose ends sit inside another group, which is where resolving
    a group's endpoints eagerly and resolving them up front part company.
    
    * TS SDK: Say that a native Dag's cron schedule is read in UTC
    
    Python takes a Dag's timezone from `start_date`, so its cron runs in the
    zone that date was written in. `startDate` here is a `Date`, which is an
    instant and carries no zone, so there is nothing to take and the schedule
    is read in UTC. That was already the behaviour; nothing said so.
---
 .../language-sdks/typescript.rst                   |  4 +++
 scripts/ci/lang_sdk_serialization/test_dags.yaml   | 25 +++++++++++++
 ts-sdk/src/coordinator/serde.ts                    | 17 ++++++++-
 ts-sdk/tests/coordinator/serde.test.ts             | 42 ++++++++++++++++++++++
 4 files changed, 87 insertions(+), 1 deletion(-)

diff --git 
a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst 
b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
index 9364305ee2d..174ba3ebcff 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
@@ -360,6 +360,10 @@ without knowing which language declared it.
 expression. A cron preset such as ``@daily`` is recorded as the expression it 
stands for. Anything
 else names a Python object a TypeScript bundle cannot point at, and is 
rejected.
 
+A cron schedule is read in UTC. Python takes the Dag's timezone from 
``start_date``, but
+``startDate`` here is a ``Date``, which is an instant and carries no timezone, 
so there is none to
+take.
+
 Every task of a native Dag runs on the Node coordinator, so it needs the queue 
the deployment routes
 there. Set it once on the Dag and each task inherits it:
 
diff --git a/scripts/ci/lang_sdk_serialization/test_dags.yaml 
b/scripts/ci/lang_sdk_serialization/test_dags.yaml
index d80c90bbb40..bef2b95571e 100644
--- a/scripts/ci/lang_sdk_serialization/test_dags.yaml
+++ b/scripts/ci/lang_sdk_serialization/test_dags.yaml
@@ -164,3 +164,28 @@ dags:
       - [staging, publish]
       - [publish, load]
       - [load, notify]
+
+  # Group edges whose ends sit inside another group, which is where the two 
ways of resolving a
+  # group's endpoints part company. Each edge has to be expanded against the 
task edges the ones
+  # before it left behind, as Python does at every `>>`:
+  #
+  # - outer.inner before outer.t makes outer.inner.i the only root of outer, 
so the edge from
+  #   extract into outer reaches it alone rather than both tasks.
+  # - outer.spare holds nothing, so an edge out of it stands for the group 
holding it, and leaves
+  #   from outer's leaf, which is outer.t by then. Python's find_leaves walks 
up to the parent;
+  #   stopping at the group's own edges would draw no edge at all.
+  - dag_id: conformance_nested_group_edges
+    spec:
+      schedule: "@once"
+    groups: [outer, outer.inner, outer.spare]
+    tasks:
+      - task_id: extract
+      - task_id: t
+        group: outer
+      - task_id: i
+        group: outer.inner
+      - task_id: load
+    order_edges:
+      - [outer.inner, outer.t]
+      - [extract, outer]
+      - [outer.spare, load]
diff --git a/ts-sdk/src/coordinator/serde.ts b/ts-sdk/src/coordinator/serde.ts
index 900b41b6887..723a073b867 100644
--- a/ts-sdk/src/coordinator/serde.ts
+++ b/ts-sdk/src/coordinator/serde.ts
@@ -134,6 +134,10 @@ export function serializeDag(
     dag_id: dag.dagId,
     fileloc,
     relative_fileloc: relativeFileloc,
+    // Python derives this from `start_date.tzinfo`, so a cron schedule there
+    // runs in the zone the start date was written in. `startDate` is a Date,
+    // which is an instant and carries no zone, so there is nothing to derive:
+    // a native TypeScript Dag's cron is read in UTC.
     timezone: "UTC",
     timetable: serializeTimetable(dag.spec.schedule, dag.dagId),
     tasks: [...getDagTaskRecords(dag)].map(([taskId, record]) =>
@@ -394,6 +398,11 @@ function buildDagGraph(dag: Dag): DagGraph {
    * A group holding no tasks has neither, so the edge steps over it and
    * continues along the group edges beyond — `x >> empty >> y` still runs `y`
    * after `x`, as Python's `find_leaves` walk does.
+   *
+   * On the upstream side that walk has one more step: a group that still comes
+   * up empty stands for the group holding it, so an edge out of an empty
+   * nested group reaches the tasks around it. Python gives a group standing
+   * downstream no such fallback, and neither does this.
    */
   const tasksAt = (id: string, side: "upstream" | "downstream"): string[] => {
     if (!groups.has(id)) return [id];
@@ -408,11 +417,16 @@ function buildDagGraph(dag: Dag): DagGraph {
     if (seen.has(id)) return [];
     seen.add(id);
     const next = side === "upstream" ? upstreamsOf.get(id) : 
downstreamsOf.get(id);
-    return [...(next ?? [])].flatMap((other) => {
+    const beyond = [...(next ?? [])].flatMap((other) => {
       if (!groups.has(other)) return [other];
       const own = side === "upstream" ? ends.leaves(other) : ends.roots(other);
       return own.length > 0 ? own : tasksBeyond(other, side, seen);
     });
+    if (beyond.length > 0 || side === "downstream") return beyond;
+    const parent = groups.get(id)?.parentGroupId;
+    if (parent === undefined) return [];
+    const own = ends.leaves(parent);
+    return own.length > 0 ? own : tasksBeyond(parent, side, seen);
   };
   const setsFor = (groupId: string): GroupEdgeSets => {
     let sets = groupEdges.get(groupId);
@@ -553,6 +567,7 @@ function serializeTimetable(schedule: unknown, dagId: 
string): SerializedValue {
   // https://github.com/apache/airflow/issues/67938
   return {
     __type: CRON_TIMETABLE,
+    // UTC for the reason given where the Dag's own `timezone` is written.
     __var: { expression, timezone: "UTC", interval: 0, run_immediately: false 
},
   };
 }
diff --git a/ts-sdk/tests/coordinator/serde.test.ts 
b/ts-sdk/tests/coordinator/serde.test.ts
index f2ddc7ad22c..aa7f533c855 100644
--- a/ts-sdk/tests/coordinator/serde.test.ts
+++ b/ts-sdk/tests/coordinator/serde.test.ts
@@ -726,6 +726,48 @@ describe("serializeDag", () => {
       expect(tasks.get("before")?.["downstream_task_ids"]).toEqual(["after"]);
     });
 
+    it("steps over an empty nested group to the tasks of the group holding 
it", () => {
+      // Python's find_leaves falls back to the parent group, so an edge out of
+      // an empty nested group leaves from whatever the parent finishes with.
+      const dag = new Dag("d");
+      const outer = dag.taskGroup("outer");
+      place(outer, "t");
+      const inner = outer.taskGroup("inner");
+      const load = place(dag, "load");
+      inner.before(load);
+
+      const tasks = taskMap(serializeDag(dag, "", ".") as Json);
+      expect(tasks.get("outer.t")?.["downstream_task_ids"]).toEqual(["load"]);
+    });
+
+    it("steps over a chain of empty nested groups to the group holding them", 
() => {
+      // The fallback is a walk, not one step: mid holds only inner, so neither
+      // has leaves and the edge leaves from outer's.
+      const dag = new Dag("d");
+      const outer = dag.taskGroup("outer");
+      place(outer, "t");
+      const inner = outer.taskGroup("mid").taskGroup("inner");
+      const load = place(dag, "load");
+      inner.before(load);
+
+      const tasks = taskMap(serializeDag(dag, "", ".") as Json);
+      expect(tasks.get("outer.t")?.["downstream_task_ids"]).toEqual(["load"]);
+    });
+
+    it("draws nothing into an empty nested group, which has no roots to 
reach", () => {
+      // Only the upstream side falls back to the parent. Python resolves a
+      // group standing downstream through get_roots() alone, with no such
+      // walk, so an edge in reaches nothing rather than the parent's tasks.
+      const dag = new Dag("d");
+      const outer = dag.taskGroup("outer");
+      place(outer, "t");
+      const inner = outer.taskGroup("inner");
+      place(dag, "extract").before(inner);
+
+      const tasks = taskMap(serializeDag(dag, "", ".") as Json);
+      expect(tasks.get("extract")).not.toHaveProperty("downstream_task_ids");
+    });
+
     it("draws nothing when an empty group has no other side to reach", () => {
       const dag = new Dag("d");
       const before = place(dag, "before");

Reply via email to