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 5799660c4d4 TS SDK: draw order-only edges with before and after 
(#73438)
5799660c4d4 is described below

commit 5799660c4d41e95b8c3ce836bd9003c3e4f3a998
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Wed Sep 30 09:39:32 2026 +0800

    TS SDK: draw order-only edges with before and after (#73438)
---
 .../language-sdks/typescript.rst                   |  22 +++
 ts-sdk/src/sdk/dag.ts                              | 128 +++++++++++++--
 ts-sdk/tests/public-api.test.ts                    |  20 ++-
 ts-sdk/tests/sdk/dag.test.ts                       | 177 ++++++++++++++++++++-
 4 files changed, 330 insertions(+), 17 deletions(-)

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 e40b1363c0a..0900ba3a40a 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
@@ -289,6 +289,28 @@ The task id may be omitted, in which case it is the 
handler's function name:
 name of its own, such as an arrow function passed inline, has nothing to take 
an id from and needs
 one: either positionally or as ``taskId`` in its spec. Give it in one place 
only, not both.
 
+Order-only edges
+~~~~~~~~~~~~~~~~
+
+An edge that carries no value has no argument name to travel under, so it is 
drawn between the
+references themselves with ``before`` and ``after``, the TypeScript pair for 
Python's ``>>`` and
+``<<``:
+
+.. code-block:: typescript
+
+    const loaded = load({ transformed });
+    const cleaned = cleanup();
+
+    loaded.before(cleaned);                  // loaded >> cleaned
+    cleaned.after(loaded, transformed);      // [loaded, transformed] >> 
cleaned
+
+Both take any number of references, so one call draws several edges, and 
drawing an edge that
+already exists changes nothing. Each returns the reference it was called on, so
+``loaded.before(cleaned).before(notified)`` draws both edges from ``loaded``.
+
+Pass a value as an argument when the downstream task needs it, and use 
``before`` or ``after`` when
+it only needs to run in order.
+
 ``new Dag`` and ``dag.task`` both take a trailing spec of Airflow options:
 ``{ schedule: "@daily", tags: ["etl"] }`` for the Dag, ``{ retries: 2, 
retryDelay: 30 }`` for a task.
 
diff --git a/ts-sdk/src/sdk/dag.ts b/ts-sdk/src/sdk/dag.ts
index 2903b2b6ac4..6e5f274b0e1 100644
--- a/ts-sdk/src/sdk/dag.ts
+++ b/ts-sdk/src/sdk/dag.ts
@@ -127,9 +127,6 @@ declare const RETURN_TYPE: unique symbol;
 /**
  * A reference to the result of one task, returned by calling that task.
  *
- * Identity only: the handler and the value are deliberately not exposed. Pass 
a
- * reference as an input of a downstream task to make that task depend on it.
- *
  * `TReturn` is the handler's return type, so a construct that needs a
  * particular one can ask for it. A reference of a narrower type is usable
  * wherever a wider one is: a `TaskRef<number>` is a `TaskRef<unknown>`.
@@ -145,8 +142,48 @@ export interface TaskRef<TReturn = unknown> {
   readonly taskId: string;
   /** @internal Never set; see {@link RETURN_TYPE}. */
   readonly [RETURN_TYPE]?: TReturn;
+  /**
+   * Run this task before each of `downstream`, carrying no value — the
+   * TypeScript spelling of Python's `>>`.
+   *
+   * ```ts
+   * loaded.before(cleaned, notified); // loaded >> [cleanup, notify]
+   * ```
+   *
+   * Variadic, so one call fans out, and it returns its own receiver rather
+   * than its arguments: a fan-out has no single "next" reference to hand back.
+   * Declaring an edge that already exists changes nothing.
+   */
+  before(...downstream: readonly TaskRef[]): TaskRef<TReturn>;
+  /**
+   * Run this task after each of `upstream`, carrying no value — Python's `<<`.
+   *
+   * ```ts
+   * cleaned.after(loaded, transformed); // [load, transform] >> cleanup
+   * ```
+   *
+   * Fan-*in* that carries data is the wiring object instead
+   * (`summarize({ north: extractNorth(), south: extractSouth() })`), so each
+   * direction has an answer: named keys when values flow, `after` when only
+   * order does.
+   */
+  after(...upstream: readonly TaskRef[]): TaskRef<TReturn>;
+}
+
+/**
+ * An order-only edge of a Dag: upstream task ID, then downstream task ID.
+ *
+ * Kept apart from the wiring a factory call records, because an edge that
+ * carries no value has no argument name to be recorded under.
+ */
+export interface OrderEdge {
+  readonly upstream: string;
+  readonly downstream: string;
 }
 
+// A task id cannot hold a NUL, so a joined pair cannot collide with one.
+const EDGE_KEY_SEPARATOR = "\u0000";
+
 /** Whether `value` is a TaskRef returned by any copy of this package. */
 function isTaskRef(value: unknown): value is TaskRef {
   return hasBrand(value, "TaskRef");
@@ -277,6 +314,7 @@ export type RecordedInputs = Readonly<Record<string, 
TaskRef | JsonValue>>;
 // Dag's private state without public accessors on the class.
 let taskRecordsOf: (dag: Dag) => ReadonlyMap<string, TaskRecord>;
 let inputsOf: (dag: Dag) => ReadonlyMap<string, RecordedInputs>;
+let orderEdgesOf: (dag: Dag) => readonly OrderEdge[];
 let finalizeOf: (dag: Dag) => void;
 
 /** Internal: whether `value` is a Dag built by any copy of this package. */
@@ -310,11 +348,15 @@ export class Dag {
   readonly spec: DagSpec;
   readonly #tasks = new Map<string, TaskRecord>();
   readonly #inputs = new Map<string, RecordedInputs>();
+  // Keyed by the two task ids, so declaring an edge twice records it once, and
+  // insertion-ordered so the serialized Dag reads as written.
+  readonly #orderEdges = new Map<string, OrderEdge>();
   #finalized = false;
 
   static {
     taskRecordsOf = (dag) => dag.#tasks;
     inputsOf = (dag) => dag.#inputs;
+    orderEdgesOf = (dag) => [...dag.#orderEdges.values()];
     finalizeOf = (dag) => dag.#finalize();
   }
 
@@ -425,7 +467,7 @@ export class Dag {
       );
     }
     const spec = this.#taskSpecOf(taskId, options);
-    const task = createTaskRef(this.dagId, taskId);
+    const task = this.#createTaskRef(taskId);
     this.#tasks.set(taskId, {
       task,
       // The runtime dispatches every handler through one instantiation, as it
@@ -466,6 +508,73 @@ export class Dag {
     return spec as TaskSpec;
   }
 
+  #createTaskRef(taskId: string): TaskRef {
+    const task: TaskRef = {
+      dagId: this.dagId,
+      taskId,
+      before: (...downstream) => {
+        for (const other of downstream) this.#addOrderEdge(task, other, 
"before");
+        return task;
+      },
+      after: (...upstream) => {
+        for (const other of upstream) this.#addOrderEdge(other, task, "after");
+        return task;
+      },
+    };
+    brand(task, "TaskRef");
+    return Object.freeze(task);
+  }
+
+  #addOrderEdge(upstream: TaskRef, downstream: TaskRef, verb: "before" | 
"after"): void {
+    if (this.#finalized) {
+      throw new Error(
+        `An edge was drawn on Dag "${this.dagId}" after the Dag was read; ` +
+          "declare every edge while the module is loading",
+      );
+    }
+    // The argument is the one that can be foreign: the receiver is a reference
+    // this Dag handed out, since it is what carries the method.
+    const other = verb === "before" ? downstream : upstream;
+    this.#validateOwnRef(other, verb);
+    if (upstream.taskId === downstream.taskId) {
+      throw new Error(
+        `${verb}() cannot draw an edge from task "${upstream.taskId}" of Dag 
"${this.dagId}" to ` +
+          "itself; an edge orders two different tasks",
+      );
+    }
+    const key = `${upstream.taskId}${EDGE_KEY_SEPARATOR}${downstream.taskId}`;
+    // Idempotent, so an edge drawn from both ends is one edge.
+    if (!this.#orderEdges.has(key)) {
+      this.#orderEdges.set(
+        key,
+        Object.freeze({ upstream: upstream.taskId, downstream: 
downstream.taskId }),
+      );
+    }
+  }
+
+  #validateOwnRef(ref: TaskRef, verb: string): void {
+    if (!isTaskRef(ref)) {
+      throw new Error(
+        `${verb}() on Dag "${this.dagId}" takes task references returned by 
calling a task, ` +
+          "not arbitrary values",
+      );
+    }
+    if (ref.dagId !== this.dagId) {
+      throw new Error(
+        `${verb}() cannot draw an edge to Dag "${ref.dagId}" task 
"${ref.taskId}" from Dag ` +
+          `"${this.dagId}"; an edge joins two tasks of one Dag`,
+      );
+    }
+    // Identity, not the ID pair: two Dag objects can carry the same dagId, and
+    // a second resolved copy of this package brands its own references.
+    if (this.#tasks.get(ref.taskId)?.task !== ref) {
+      throw new Error(
+        `${verb}() was given a reference to "${ref.taskId}" that this Dag did 
not hand out; ` +
+          `it comes from another Dag object with the same ID, or 
${DUPLICATE_COPY_HINT}`,
+      );
+    }
+  }
+
   #wire(taskId: string, inputs: unknown): void {
     if (this.#finalized) {
       throw new Error(
@@ -576,12 +685,6 @@ function readFunctionName(handler: unknown): string | 
undefined {
   return typeof name === "string" && name.length > 0 ? name : undefined;
 }
 
-function createTaskRef(dagId: string, taskId: string): TaskRef {
-  const task: TaskRef = { dagId, taskId };
-  brand(task, "TaskRef");
-  return Object.freeze(task);
-}
-
 function validateDagSpec(dagId: string, spec: DagSpec): void {
   const value: unknown = spec;
   if (!isPlainRecord(value)) {
@@ -604,6 +707,11 @@ export function getDagTaskRecords(dag: Dag): 
ReadonlyMap<string, TaskRecord> {
   return taskRecordsOf(dag);
 }
 
+/** Internal: the order-only edges of a Dag, in the order they were drawn. */
+export function getDagOrderEdges(dag: Dag): readonly OrderEdge[] {
+  return orderEdgesOf(dag);
+}
+
 /** Internal: what each task of a Dag was called with, keyed by task ID.
  *  A task that has not been called is absent. */
 export function getDagTaskInputs(dag: Dag): ReadonlyMap<string, 
RecordedInputs> {
diff --git a/ts-sdk/tests/public-api.test.ts b/ts-sdk/tests/public-api.test.ts
index 3ffba6332f9..be485998413 100644
--- a/ts-sdk/tests/public-api.test.ts
+++ b/ts-sdk/tests/public-api.test.ts
@@ -59,8 +59,11 @@ describe("public API", () => {
     );
     const upstream = upstreamTask();
     const downstream = downstreamTask({ upstream });
-    expect(upstream).toEqual({ dagId: "public_api_dag", taskId: 
"public_api_task" });
-    expect(downstream).toEqual({ dagId: "public_api_dag", taskId: 
"public_api_downstream" });
+    expect(upstream).toMatchObject({ dagId: "public_api_dag", taskId: 
"public_api_task" });
+    expect(downstream).toMatchObject({
+      dagId: "public_api_dag",
+      taskId: "public_api_downstream",
+    });
     expect(dag.taskIds).toEqual(["public_api_task", "public_api_downstream"]);
     // serve() hands the bundle to the runtime, which needs the supervisor's
     // socket addresses that Airflow puts on argv.
@@ -283,6 +286,11 @@ describe("public API", () => {
   });
 
   it("keeps the Dag authoring signatures extensible via trailing specs", () => 
{
+    // Identity, plus the two order-only edge verbs; the handler and the value
+    // stay hidden.
+    expectTypeOf<Extract<keyof TaskRef, string>>().toEqualTypeOf<
+      "dagId" | "taskId" | "before" | "after"
+    >();
     expectTypeOf<TaskRef["dagId"]>().toEqualTypeOf<string>();
     expectTypeOf<TaskRef["taskId"]>().toEqualTypeOf<string>();
     // A reference carries its handler's return type, so a construct that needs
@@ -290,6 +298,14 @@ describe("public API", () => {
     // a wider one is, and not the other way round.
     expectTypeOf<TaskRef<boolean>>().toMatchTypeOf<TaskRef>();
     expectTypeOf<TaskRef>().not.toMatchTypeOf<TaskRef<boolean>>();
+    // Variadic, and each returns its own receiver rather than its arguments,
+    // return type included.
+    expectTypeOf<TaskRef<boolean>["before"]>().toEqualTypeOf<
+      (...downstream: readonly TaskRef[]) => TaskRef<boolean>
+    >();
+    expectTypeOf<TaskRef<boolean>["after"]>().toEqualTypeOf<
+      (...upstream: readonly TaskRef[]) => TaskRef<boolean>
+    >();
     // Wiring moved to the factory call, and the spec is the trailing argument
     // itself.
     expectTypeOf<TaskOptions>().toEqualTypeOf<TaskSpec>();
diff --git a/ts-sdk/tests/sdk/dag.test.ts b/ts-sdk/tests/sdk/dag.test.ts
index 1bd493b602d..612cdc9cdb0 100644
--- a/ts-sdk/tests/sdk/dag.test.ts
+++ b/ts-sdk/tests/sdk/dag.test.ts
@@ -21,6 +21,7 @@ import { describe, it, expect } from "vitest";
 import {
   Dag,
   finalizeDag,
+  getDagOrderEdges,
   getDagTaskInputs,
   getDagTaskRecords,
   type DagSpec,
@@ -36,7 +37,7 @@ describe("Dag", () => {
     expect(typeof myTask).toBe("function");
 
     const ref = myTask();
-    expect(ref).toEqual({ dagId: "example_dag", taskId: "my_task" });
+    expect(ref).toMatchObject({ dagId: "example_dag", taskId: "my_task" });
     expect(Object.isFrozen(ref)).toBe(true);
   });
 
@@ -53,9 +54,9 @@ describe("Dag", () => {
     const transformed = transform({ extracted });
     const loaded = load({ transformed });
 
-    expect(extracted).toEqual({ dagId: "chained_dag", taskId: "extract" });
-    expect(transformed).toEqual({ dagId: "chained_dag", taskId: "transform" });
-    expect(loaded).toEqual({ dagId: "chained_dag", taskId: "load" });
+    expect(extracted).toMatchObject({ dagId: "chained_dag", taskId: "extract" 
});
+    expect(transformed).toMatchObject({ dagId: "chained_dag", taskId: 
"transform" });
+    expect(loaded).toMatchObject({ dagId: "chained_dag", taskId: "load" });
 
     const inputs = getDagTaskInputs(dag);
     expect(inputs.get("extract")).toEqual({});
@@ -86,7 +87,7 @@ describe("Dag", () => {
     const transform = dag.task("transform", async (_: { upstream: unknown }) 
=> undefined);
     const lookalike = { dagId: "lookalike_dag", taskId: "ghost" };
 
-    transform({ upstream: lookalike });
+    transform({ upstream: lookalike } as unknown as { upstream: TaskRef });
 
     expect(getDagTaskInputs(dag).get("transform")).toEqual({ upstream: 
lookalike });
   });
@@ -511,6 +512,172 @@ describe("Dag", () => {
     });
   });
 
+  describe("order-only edges", () => {
+    /** A Dag whose tasks are all placed, ready for edges to be drawn on it. */
+    function placedDag(dagId: string, ...taskIds: string[]) {
+      const dag = new Dag(dagId);
+      const refs = Object.fromEntries(
+        taskIds.map((taskId) => [taskId, dag.task(taskId, async () => 
undefined)()]),
+      );
+      return { dag, refs };
+    }
+
+    it("draws an edge with before, from the receiver to the argument", () => {
+      const { dag, refs } = placedDag("d", "load", "cleanup");
+
+      refs.load!.before(refs.cleanup!);
+
+      expect(getDagOrderEdges(dag)).toEqual([{ upstream: "load", downstream: 
"cleanup" }]);
+    });
+
+    it("draws an edge with after, from the argument to the receiver", () => {
+      const { dag, refs } = placedDag("d", "load", "cleanup");
+
+      refs.cleanup!.after(refs.load!);
+
+      expect(getDagOrderEdges(dag)).toEqual([{ upstream: "load", downstream: 
"cleanup" }]);
+    });
+
+    it("fans out from one before call", () => {
+      const { dag, refs } = placedDag("d", "load", "cleanup", "notify");
+
+      refs.load!.before(refs.cleanup!, refs.notify!);
+
+      expect(getDagOrderEdges(dag)).toEqual([
+        { upstream: "load", downstream: "cleanup" },
+        { upstream: "load", downstream: "notify" },
+      ]);
+    });
+
+    it("fans in from one after call", () => {
+      const { dag, refs } = placedDag("d", "load", "transform", "cleanup");
+
+      refs.cleanup!.after(refs.load!, refs.transform!);
+
+      expect(getDagOrderEdges(dag)).toEqual([
+        { upstream: "load", downstream: "cleanup" },
+        { upstream: "transform", downstream: "cleanup" },
+      ]);
+    });
+
+    it.each([
+      ["before", (a: TaskRef, b: TaskRef) => a.before(b)],
+      ["after", (a: TaskRef, b: TaskRef) => b.after(a)],
+    ])("returns the receiver from %s, not the arguments", (_verb, draw) => {
+      const { refs } = placedDag("d", "load", "cleanup");
+
+      // A fan-out has no single "next" reference, so chaining continues from
+      // the same task rather than from what was just pointed at.
+      expect(draw(refs.load!, refs.cleanup!)).toBe(_verb === "before" ? 
refs.load : refs.cleanup);
+    });
+
+    it("records an edge once however many times it is drawn", () => {
+      const { dag, refs } = placedDag("d", "load", "cleanup");
+
+      refs.load!.before(refs.cleanup!);
+      refs.load!.before(refs.cleanup!);
+      refs.cleanup!.after(refs.load!);
+
+      expect(getDagOrderEdges(dag)).toEqual([{ upstream: "load", downstream: 
"cleanup" }]);
+    });
+
+    it("rejects an edge from a task to itself", () => {
+      const { refs } = placedDag("d", "a");
+
+      expect(() => refs.a!.before(refs.a!)).toThrowError(
+        /before\(\) cannot draw an edge from task "a" of Dag "d" to itself/,
+      );
+      expect(() => refs.a!.after(refs.a!)).toThrowError(
+        /after\(\) cannot draw an edge from task "a" of Dag "d" to itself/,
+      );
+    });
+
+    it("keeps the two directions apart", () => {
+      const { dag, refs } = placedDag("d", "a", "b");
+
+      refs.a!.before(refs.b!);
+      refs.b!.before(refs.a!);
+
+      // Two distinct edges, both recorded: rejecting the cycle they form is a
+      // Dag-level concern, not an edge-level one.
+      expect(getDagOrderEdges(dag)).toEqual([
+        { upstream: "a", downstream: "b" },
+        { upstream: "b", downstream: "a" },
+      ]);
+    });
+
+    it("leaves a frozen edge that a caller cannot rewrite", () => {
+      const { dag, refs } = placedDag("d", "load", "cleanup");
+      refs.load!.before(refs.cleanup!);
+
+      expect(Object.isFrozen(getDagOrderEdges(dag)[0])).toBe(true);
+    });
+
+    it("carries no value, so it records no input", () => {
+      const { dag, refs } = placedDag("d", "load", "cleanup");
+
+      refs.load!.before(refs.cleanup!);
+
+      expect(getDagTaskInputs(dag).get("cleanup")).toEqual({});
+    });
+
+    it.each([
+      ["before", (ref: TaskRef, other: TaskRef) => ref.before(other)],
+      ["after", (ref: TaskRef, other: TaskRef) => ref.after(other)],
+    ])("rejects a %s edge to a task of another Dag", (verb, draw) => {
+      const { refs: here } = placedDag("here", "load");
+      const { refs: there } = placedDag("there", "cleanup");
+
+      expect(() => draw(here.load!, there.cleanup!)).toThrowError(
+        new RegExp(
+          `${verb}\\(\\) cannot draw an edge to Dag "there" task "cleanup" 
from Dag "here"`,
+        ),
+      );
+    });
+
+    it("rejects a reference from another Dag object carrying the same Dag ID", 
() => {
+      const { refs: first } = placedDag("same_id", "load");
+      const { refs: second } = placedDag("same_id", "cleanup");
+
+      expect(() => first.load!.before(second.cleanup!)).toThrowError(
+        /before\(\) was given a reference to "cleanup" that this Dag did not 
hand out/,
+      );
+    });
+
+    it.each([
+      ["a plain object", { dagId: "d", taskId: "cleanup" }],
+      ["a string", "cleanup"],
+      ["null", null],
+    ])("rejects %s where a reference belongs", (_label, value) => {
+      const { refs } = placedDag("d", "load");
+
+      expect(() => refs.load!.before(value as unknown as 
TaskRef)).toThrowError(
+        /before\(\) on Dag "d" takes task references returned by calling a 
task/,
+      );
+    });
+
+    it("records nothing when an edge is rejected", () => {
+      const { dag, refs } = placedDag("d", "load");
+
+      expect(() => refs.load!.before("cleanup" as unknown as 
TaskRef)).toThrow();
+      expect(getDagOrderEdges(dag)).toEqual([]);
+    });
+
+    it("rejects an edge drawn after the Dag was read", () => {
+      const { dag, refs } = placedDag("d", "load", "cleanup");
+      finalizeDag(dag);
+
+      expect(() => refs.load!.before(refs.cleanup!)).toThrowError(
+        /An edge was drawn on Dag "d" after the Dag was read/,
+      );
+    });
+
+    it("has no edges before any are drawn", () => {
+      const { dag } = placedDag("d", "load");
+      expect(getDagOrderEdges(dag)).toEqual([]);
+    });
+  });
+
   it("exposes its task IDs in attachment order", () => {
     const dag = new Dag("ordered_dag");
     expect(dag.taskIds).toEqual([]);

Reply via email to