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([]);