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");