This is an automated email from the ASF dual-hosted git repository.
pierrejeambrun 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 53ca2a427ba TS SDK: add task groups to the native Dag surface (#73439)
53ca2a427ba is described below
commit 53ca2a427ba5684bae103ac407293a92686e30c1
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Wed Sep 30 16:54:10 2026 +0800
TS SDK: add task groups to the native Dag surface (#73439)
---
.../language-sdks/typescript.rst | 25 ++
ts-sdk/src/index.ts | 3 +
ts-sdk/src/sdk/dag.ts | 330 +++++++++++++++++--
ts-sdk/tests/public-api.test.ts | 27 +-
ts-sdk/tests/sdk/dag.test.ts | 357 ++++++++++++++++++++-
5 files changed, 703 insertions(+), 39 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 ed8532859ce..aea8077c1ee 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
@@ -311,6 +311,31 @@ already exists changes nothing. Each returns the reference
it was called on, so
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.
+Task groups
+~~~~~~~~~~~
+
+``dag.taskGroup(groupId)`` opens a scope with the same ``task`` and
``taskGroup`` methods as the Dag,
+prefixing the id of everything declared in it, as Python's ``prefix_group_id``
does:
+
+.. code-block:: typescript
+
+ const staging = dag.taskGroup("staging");
+ staging.task("stage_rows", stageRows)(); // task id
"staging.stage_rows"
+ staging.taskGroup("checks").task("nulls", checkNulls)(); //
"staging.checks.nulls"
+
+ staging.before(loaded); // staging >> loaded
+
+A group is an edge endpoint in its own right, so ``before`` and ``after``
order a whole group against
+a task or against another group.
+
+Tasks and groups share one id namespace, as they do in Python, so a Dag cannot
hold both a task and a
+group called ``staging``. A ``.`` is what separates a group from what it
holds, so it cannot appear in
+an id of either.
+
+Pass ``{ prefixGroupId: false }`` to keep the ids declared in a group as
written, as ``prefix_group_id=False``
+does in Python; they then have to be unique across the Dag. A group id is made
of letters, digits, dashes and
+underscores, and is at most 200 characters.
+
``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/index.ts b/ts-sdk/src/index.ts
index ad45c3bad33..4e7bee9a966 100644
--- a/ts-sdk/src/index.ts
+++ b/ts-sdk/src/index.ts
@@ -28,7 +28,10 @@ export type { ArgNameMap } from "./sdk/arg-names.js";
export type { Registerable } from "./sdk/bundle.js";
export type {
DagSpec,
+ Node,
TaskFactory,
+ TaskGroupOptions,
+ TaskGroupRef,
TaskInput,
TaskInputs,
TaskOptions,
diff --git a/ts-sdk/src/sdk/dag.ts b/ts-sdk/src/sdk/dag.ts
index aead708fe62..8f7a2042dea 100644
--- a/ts-sdk/src/sdk/dag.ts
+++ b/ts-sdk/src/sdk/dag.ts
@@ -136,7 +136,26 @@ declare const RETURN_TYPE: unique symbol;
* object is not one: an input also takes a plain JSON value, and without the
* brand such an object could not be told apart from an upstream reference.
*/
-export interface TaskRef<TReturn = unknown> {
+/**
+ * What an order-only edge can connect: a task, or a whole task group.
+ *
+ * The TypeScript counterpart of Python's `DAGNode` and the Go SDK's
+ * `airflow.Node`. A group carries edges as a task does, so `before` and
+ * `after` take either.
+ */
+export interface Node {
+ /** Identifier of the Dag this node belongs to. */
+ readonly dagId: string;
+ /**
+ * Run this node before each of `downstream`, carrying no value — the
+ * TypeScript spelling of Python's `>>`.
+ */
+ before(...downstream: readonly Node[]): Node;
+ /** Run this node after each of `upstream`, carrying no value — Python's
`<<`. */
+ after(...upstream: readonly Node[]): Node;
+}
+
+export interface TaskRef<TReturn = unknown> extends Node {
/** Identifier of the Dag this task belongs to. */
readonly dagId: string;
/** Airflow task ID, including any TaskGroup prefix. */
@@ -155,7 +174,7 @@ export interface TaskRef<TReturn = unknown> {
* 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>;
+ before(...downstream: readonly Node[]): TaskRef<TReturn>;
/**
* Run this task after each of `upstream`, carrying no value — Python's `<<`.
*
@@ -168,11 +187,59 @@ export interface TaskRef<TReturn = unknown> {
* direction has an answer: named keys when values flow, `after` when only
* order does.
*/
- after(...upstream: readonly TaskRef[]): TaskRef<TReturn>;
+ after(...upstream: readonly Node[]): TaskRef<TReturn>;
}
/**
- * An order-only edge of a Dag: upstream task ID, then downstream task ID.
+ * A scope that declares tasks under a shared id prefix, and carries edges as a
+ * whole — Python's `TaskGroup`, spelled to TypeScript convention.
+ *
+ * Offers the same `task` and `taskGroup` methods as the Dag, so nesting is the
+ * same call at every depth. Each id it declares is prefixed with the group id
+ * unless {@link TaskGroupOptions.prefixGroupId} turns that off, and the group
+ * itself stands at either end of an edge, so a whole group can be ordered
+ * against a task or against another group.
+ */
+export interface TaskGroupRef extends Node {
+ /** Identifier of the Dag this group belongs to. */
+ readonly dagId: string;
+ /** Group ID, including any enclosing group's prefix. */
+ readonly groupId: string;
+ /** Declare a task in this group; its id carries the group prefix unless it
is turned off. */
+ task<TArgs extends object | void = void, TReturn = unknown>(
+ taskId: string,
+ handler: (args: TArgs) => TReturn | Promise<TReturn>,
+ options?: TaskOptions,
+ ): TaskFactory<TArgs, TReturn>;
+ /** Declare a task whose id is the handler's function name, prefixed the
same way. */
+ task<TArgs extends object | void = void, TReturn = unknown>(
+ handler: (args: TArgs) => TReturn | Promise<TReturn>,
+ options?: TaskOptions,
+ ): TaskFactory<TArgs, TReturn>;
+ /** Nest a group inside this one. */
+ taskGroup(groupId: string, options?: TaskGroupOptions): TaskGroupRef;
+ before(...downstream: readonly Node[]): TaskGroupRef;
+ after(...upstream: readonly Node[]): TaskGroupRef;
+}
+
+/** Options of a task group, the trailing argument of `taskGroup()`. */
+export interface TaskGroupOptions {
+ /**
+ * Whether the ids declared in the group are prefixed with the group's, as
+ * in `staging.load`. On by default, as in Python; turning it off keeps them
+ * as written, as `prefix_group_id=False` does, so they have to be unique
+ * across the Dag.
+ */
+ readonly prefixGroupId?: boolean;
+}
+
+/**
+ * An order-only edge of a Dag, between two node IDs.
+ *
+ * An endpoint is a task ID or a group ID; the two share one namespace, so a
+ * bare ID names exactly one node. {@link TaskGroupRecord} is what tells a
+ * consumer which kind an endpoint is, and which tasks a group endpoint stands
+ * for.
*
* Kept apart from the wiring a factory call records, because an edge that
* carries no value has no argument name to be recorded under.
@@ -185,6 +252,41 @@ export interface OrderEdge {
// A task id cannot hold a NUL, so a joined pair cannot collide with one.
const EDGE_KEY_SEPARATOR = "\u0000";
+/** Internal: one group of a Dag, and the tree beneath it. */
+export interface TaskGroupRecord {
+ /** Group ID, including any enclosing group's prefix. */
+ readonly groupId: string;
+ /** Enclosing group's ID, absent for a group declared on the Dag itself. */
+ readonly parentGroupId?: string;
+ /** Whether the IDs declared in this group carry its ID as a prefix. */
+ readonly prefixGroupId: boolean;
+ /** Task IDs declared directly in this group, in declaration order. */
+ readonly taskIds: readonly string[];
+ /** Group IDs nested directly in this group, in declaration order. */
+ readonly childGroupIds: readonly string[];
+}
+
+/** Separates a group ID from what it contains, as Python's `prefix_group_id`
does. */
+const GROUP_SEPARATOR = ".";
+
+/** The Dag's own view of a group, which it appends to as an author declares.
*/
+interface MutableTaskGroupRecord extends TaskGroupRecord {
+ readonly taskIds: string[];
+ readonly childGroupIds: string[];
+}
+
+/** Whether `value` is a task group returned by any copy of this package. */
+function isTaskGroupRef(value: unknown): value is TaskGroupRef {
+ return hasBrand(value, "TaskGroupRef");
+}
+
+/** The ID an edge endpoint is recorded under: a task ID or a group ID. */
+function nodeId(node: Node): string | undefined {
+ if (isTaskRef(node)) return node.taskId;
+ if (isTaskGroupRef(node)) return node.groupId;
+ return undefined;
+}
+
/** Whether `value` is a TaskRef returned by any copy of this package. */
function isTaskRef(value: unknown): value is TaskRef {
return hasBrand(value, "TaskRef");
@@ -317,6 +419,7 @@ let taskRecordsOf: (dag: Dag) => ReadonlyMap<string,
TaskRecord>;
let inputsOf: (dag: Dag) => ReadonlyMap<string, RecordedInputs>;
let orderEdgesOf: (dag: Dag) => readonly OrderEdge[];
let definedInOf: (dag: Dag) => string | undefined;
+let groupsOf: (dag: Dag) => ReadonlyMap<string, TaskGroupRecord>;
let finalizeOf: (dag: Dag) => void;
/** Internal: whether `value` is a Dag built by any copy of this package. */
@@ -354,6 +457,12 @@ export class Dag {
// insertion-ordered so the serialized Dag reads as written.
readonly #orderEdges = new Map<string, OrderEdge>();
readonly #definedIn: string | undefined;
+ // Keyed by full group ID; a group's own record holds what it declares, so
+ // the tree is reconstructed by walking from the roots.
+ readonly #groups = new Map<string, MutableTaskGroupRecord>();
+ // The one reference each group was handed out as, so an edge endpoint can be
+ // checked by identity the way a task's is.
+ readonly #groupRefs = new Map<string, TaskGroupRef>();
#finalized = false;
static {
@@ -361,6 +470,7 @@ export class Dag {
inputsOf = (dag) => dag.#inputs;
orderEdgesOf = (dag) => [...dag.#orderEdges.values()];
definedInOf = (dag) => dag.#definedIn;
+ groupsOf = (dag) => dag.#groups;
finalizeOf = (dag) => dag.#finalize();
}
@@ -415,6 +525,31 @@ export class Dag {
taskIdOrHandler: string | ((args: TArgs) => TReturn | Promise<TReturn>),
handlerOrOptions?: ((args: TArgs) => TReturn | Promise<TReturn>) |
TaskOptions,
maybeOptions?: TaskOptions,
+ ): TaskFactory<TArgs, TReturn> {
+ return this.#addTask(undefined, taskIdOrHandler, handlerOrOptions,
maybeOptions);
+ }
+
+ /**
+ * Declare a task group of this Dag.
+ *
+ * The group prefixes the id of everything declared in it, unless
+ * {@link TaskGroupOptions.prefixGroupId} turns that off, and stands at
+ * either end of an order-only edge in its own right.
+ *
+ * ```ts
+ * const staging = dag.taskGroup("staging"); // tasks "staging.<id>"
+ * const checks = dag.taskGroup("checks", { prefixGroupId: false }); // ids
as written
+ * ```
+ */
+ taskGroup(groupId: string, options?: TaskGroupOptions): TaskGroupRef {
+ return this.#addGroup(undefined, groupId, options);
+ }
+
+ #addTask<TArgs extends object | void, TReturn>(
+ groupId: string | undefined,
+ taskIdOrHandler: string | ((args: TArgs) => TReturn | Promise<TReturn>),
+ handlerOrOptions?: ((args: TArgs) => TReturn | Promise<TReturn>) |
TaskOptions,
+ maybeOptions?: TaskOptions,
): TaskFactory<TArgs, TReturn> {
const idGiven = typeof taskIdOrHandler === "string";
const handler = (idGiven ? handlerOrOptions : taskIdOrHandler) as (
@@ -435,8 +570,8 @@ export class Dag {
);
}
const defaulted = idGiven ? undefined : (specTaskId ??
readFunctionName(handler));
- const taskId = idGiven ? taskIdOrHandler : defaulted;
- if (taskId === undefined) {
+ const declared = idGiven ? taskIdOrHandler : defaulted;
+ if (declared === undefined) {
throw new Error(
`A task of Dag "${this.dagId}" has no id: its handler has no name to
take one from. ` +
'Pass an id — dag.task("my_task", handler) — or give the handler a
name. A bundler ' +
@@ -453,6 +588,15 @@ export class Dag {
'underscores — dag.task("my_task", handler)',
);
}
+ if (declared.includes(GROUP_SEPARATOR)) {
+ // The separator is what joins a group to what it holds, so a task id
+ // carrying one would name a group that does not exist.
+ throw new Error(
+ `Task ID "${declared}" of Dag "${this.dagId}" cannot contain
"${GROUP_SEPARATOR}"; ` +
+ "declare a task group with taskGroup(...) and the prefix is added
for you",
+ );
+ }
+ const taskId = this.#childId(groupId, declared);
if (typeof handler !== "function") {
throw new Error(`handler for Dag "${this.dagId}" task "${taskId}" must
be a function`);
}
@@ -464,9 +608,6 @@ export class Dag {
"a handler takes one object of named arguments — async ({ rows,
region }) => ...",
);
}
- if (this.#tasks.has(taskId)) {
- throw new Error(`Task "${taskId}" is already registered for Dag
"${this.dagId}"`);
- }
// A task added after the Dag was read could no longer be wired into it,
and
// would sit in the Dag unplaced and unreported.
if (this.#finalized) {
@@ -476,7 +617,9 @@ export class Dag {
);
}
const spec = this.#taskSpecOf(taskId, options);
+ this.#reserveNodeId(taskId, "Task");
const task = this.#createTaskRef(taskId);
+ if (groupId !== undefined) this.#groups.get(groupId)!.taskIds.push(taskId);
this.#tasks.set(taskId, {
task,
// The runtime dispatches every handler through one instantiation, as it
@@ -517,6 +660,122 @@ export class Dag {
return spec as TaskSpec;
}
+ #addGroup(
+ parentGroupId: string | undefined,
+ groupId: string,
+ options: TaskGroupOptions = {},
+ ): TaskGroupRef {
+ if (typeof groupId !== "string" || groupId.length === 0) {
+ throw new Error(`A task group of Dag "${this.dagId}" must have a
non-empty ID`);
+ }
+ if (groupId.includes(GROUP_SEPARATOR)) {
+ // The separator is what joins a group to what it holds, so one inside an
+ // ID would make the resulting task id ambiguous.
+ throw new Error(
+ `Task group ID "${groupId}" of Dag "${this.dagId}" cannot contain ` +
+ `"${GROUP_SEPARATOR}"; nest groups with taskGroup(...) instead`,
+ );
+ }
+ // Python's validate_group_key. The length is in code points, as Python's
+ // len() counts, not in UTF-16 units.
+ const length = [...groupId].length;
+ if (length > GROUP_ID_MAX_LENGTH) {
+ throw new Error(
+ `Task group ID "${groupId}" of Dag "${this.dagId}" has ${length}
characters; ` +
+ `at most ${GROUP_ID_MAX_LENGTH} are allowed`,
+ );
+ }
+ if (!GROUP_ID_CHARACTERS.test(groupId)) {
+ throw new Error(
+ `Task group ID "${groupId}" of Dag "${this.dagId}" has to be made of
letters, digits, ` +
+ "dashes and underscores",
+ );
+ }
+ const prefixGroupId = this.#prefixGroupIdOf(groupId, options);
+ if (this.#finalized) {
+ throw new Error(
+ `Task group "${groupId}" cannot be added to Dag "${this.dagId}" after
the Dag was read; ` +
+ "declare every group while the module is loading",
+ );
+ }
+ const fullId = this.#childId(parentGroupId, groupId);
+ this.#reserveNodeId(fullId, "Task group");
+ this.#groups.set(fullId, {
+ groupId: fullId,
+ ...(parentGroupId !== undefined && { parentGroupId }),
+ prefixGroupId,
+ taskIds: [],
+ childGroupIds: [],
+ });
+ if (parentGroupId !== undefined)
this.#groups.get(parentGroupId)!.childGroupIds.push(fullId);
+ const group = this.#createTaskGroupRef(fullId);
+ this.#groupRefs.set(fullId, group);
+ return group;
+ }
+
+ // TypeScript is bypassable, as for #taskSpecOf, so a misspelled or mistyped
+ // option is rejected rather than read as the default.
+ #prefixGroupIdOf(groupId: string, options: TaskGroupOptions): boolean {
+ const value: unknown = options;
+ if (!isPlainRecord(value)) {
+ throw new Error(`options for Dag "${this.dagId}" task group "${groupId}"
must be an object`);
+ }
+ for (const key of Reflect.ownKeys(value)) {
+ if (key !== "prefixGroupId") {
+ throw new Error(
+ `Unknown option "${String(key)}" for Dag "${this.dagId}" task group
"${groupId}"`,
+ );
+ }
+ }
+ const { prefixGroupId = true } = value;
+ if (typeof prefixGroupId !== "boolean") {
+ throw new Error(
+ `prefixGroupId for Dag "${this.dagId}" task group "${groupId}" must be
a boolean`,
+ );
+ }
+ return prefixGroupId;
+ }
+
+ // Python's TaskGroup.child_id: a group adds its ID in front of what it holds
+ // only when prefixGroupId is on, and a nested group's ID is built the same
way.
+ #childId(groupId: string | undefined, id: string): string {
+ if (groupId === undefined || !this.#groups.get(groupId)!.prefixGroupId)
return id;
+ return `${groupId}${GROUP_SEPARATOR}${id}`;
+ }
+
+ // Tasks and groups share one namespace, as they do in Python: a serialized
+ // Dag addresses both by a bare ID, so `dag.task("x")` and
`dag.taskGroup("x")`
+ // cannot both exist.
+ #reserveNodeId(id: string, kind: "Task" | "Task group"): void {
+ if (this.#tasks.has(id) || this.#groups.has(id)) {
+ throw new Error(`${kind} "${id}" is already registered for Dag
"${this.dagId}"`);
+ }
+ }
+
+ #createTaskGroupRef(groupId: string): TaskGroupRef {
+ const group: TaskGroupRef = {
+ dagId: this.dagId,
+ groupId,
+ task: <TArgs extends object | void, TReturn>(
+ taskIdOrHandler: string | ((args: TArgs) => TReturn |
Promise<TReturn>),
+ handlerOrOptions?: ((args: TArgs) => TReturn | Promise<TReturn>) |
TaskOptions,
+ maybeOptions?: TaskOptions,
+ ) => this.#addTask<TArgs, TReturn>(groupId, taskIdOrHandler,
handlerOrOptions, maybeOptions),
+ taskGroup: (childId: string, options?: TaskGroupOptions) =>
+ this.#addGroup(groupId, childId, options),
+ before: (...downstream) => {
+ for (const other of downstream) this.#addOrderEdge(group, other,
"before");
+ return group;
+ },
+ after: (...upstream) => {
+ for (const other of upstream) this.#addOrderEdge(other, group,
"after");
+ return group;
+ },
+ } as TaskGroupRef;
+ brand(group, "TaskGroupRef");
+ return Object.freeze(group);
+ }
+
#createTaskRef(taskId: string): TaskRef {
const task: TaskRef = {
dagId: this.dagId,
@@ -534,51 +793,56 @@ export class Dag {
return Object.freeze(task);
}
- #addOrderEdge(upstream: TaskRef, downstream: TaskRef, verb: "before" |
"after"): void {
+ #addOrderEdge(upstream: Node, downstream: Node, 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.
+ // The argument is the one that can be foreign: the receiver is a node 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) {
+ this.#validateOwnNode(other, verb);
+ const upstreamId = nodeId(upstream);
+ const downstreamId = nodeId(downstream);
+ if (upstreamId === downstreamId) {
throw new Error(
- `${verb}() cannot draw an edge from task "${upstream.taskId}" of Dag
"${this.dagId}" to ` +
- "itself; an edge orders two different tasks",
+ `${verb}() cannot draw an edge from node "${upstreamId}" of Dag
"${this.dagId}" to ` +
+ "itself; an edge orders two different nodes",
);
}
- const key = `${upstream.taskId}${EDGE_KEY_SEPARATOR}${downstream.taskId}`;
+ const key = `${upstreamId}${EDGE_KEY_SEPARATOR}${downstreamId}`;
// 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 }),
+ Object.freeze({ upstream: upstreamId!, downstream: downstreamId! }),
);
}
}
- #validateOwnRef(ref: TaskRef, verb: string): void {
- if (!isTaskRef(ref)) {
+ #validateOwnNode(node: Node, verb: string): void {
+ const id = nodeId(node);
+ if (id === undefined) {
throw new Error(
- `${verb}() on Dag "${this.dagId}" takes task references returned by
calling a task, ` +
+ `${verb}() on Dag "${this.dagId}" takes tasks and task groups this Dag
handed out, ` +
"not arbitrary values",
);
}
- if (ref.dagId !== this.dagId) {
+ if (node.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`,
+ `${verb}() cannot draw an edge to Dag "${node.dagId}" node "${id}"
from Dag ` +
+ `"${this.dagId}"; an edge joins two nodes 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) {
+ // Identity, not the ID: two Dag objects can carry the same dagId, and a
+ // second resolved copy of this package brands its own nodes. A group is
+ // checked the same way — matching on the ID alone would silently retarget
+ // the edge at this Dag's own group of that name.
+ if (isTaskRef(node) ? this.#tasks.get(id)?.task !== node :
this.#groupRefs.get(id) !== node) {
throw new Error(
- `${verb}() was given a reference to "${ref.taskId}" that this Dag did
not hand out; ` +
+ `${verb}() was given a reference to "${id}" that this Dag did not hand
out; ` +
`it comes from another Dag object with the same ID, or
${DUPLICATE_COPY_HINT}`,
);
}
@@ -681,6 +945,11 @@ export class Dag {
// has to look like.
const TASK_ID_CHARACTERS = /^[\p{L}\p{N}_.-]+$/u;
+// Python's GROUP_KEY_REGEX and validate_group_key limit: a group ID has no
dot,
+// since the dot is what joins it to what it holds.
+const GROUP_ID_CHARACTERS = /^[\p{L}\p{N}_-]+$/u;
+const GROUP_ID_MAX_LENGTH = 200;
+
/**
* A handler's own function name, or undefined when it has none.
*
@@ -716,6 +985,11 @@ export function getDagTaskRecords(dag: Dag):
ReadonlyMap<string, TaskRecord> {
return taskRecordsOf(dag);
}
+/** Internal: every task group of a Dag, keyed by full group ID. */
+export function getDagTaskGroups(dag: Dag): ReadonlyMap<string,
TaskGroupRecord> {
+ return groupsOf(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);
diff --git a/ts-sdk/tests/public-api.test.ts b/ts-sdk/tests/public-api.test.ts
index be485998413..769a8342b0e 100644
--- a/ts-sdk/tests/public-api.test.ts
+++ b/ts-sdk/tests/public-api.test.ts
@@ -24,12 +24,15 @@ import type {
ConnectionResult,
DagSpec,
GetXComOpts,
+ Node,
SetXComOpts,
TaskClient,
Registerable,
TaskContext,
TaskFactory,
TaskFunction,
+ TaskGroupOptions,
+ TaskGroupRef,
TaskInput,
TaskInputs,
TaskOptions,
@@ -301,11 +304,31 @@ describe("public API", () => {
// Variadic, and each returns its own receiver rather than its arguments,
// return type included.
expectTypeOf<TaskRef<boolean>["before"]>().toEqualTypeOf<
- (...downstream: readonly TaskRef[]) => TaskRef<boolean>
+ (...downstream: readonly Node[]) => TaskRef<boolean>
>();
expectTypeOf<TaskRef<boolean>["after"]>().toEqualTypeOf<
- (...upstream: readonly TaskRef[]) => TaskRef<boolean>
+ (...upstream: readonly Node[]) => TaskRef<boolean>
>();
+ // A group is a scope and an edge endpoint, and nests the same way at
+ // every depth.
+ expectTypeOf<Extract<keyof TaskGroupRef, string>>().toEqualTypeOf<
+ "dagId" | "groupId" | "task" | "taskGroup" | "before" | "after"
+ >();
+ expectTypeOf<TaskGroupRef["groupId"]>().toEqualTypeOf<string>();
+ expectTypeOf<TaskGroupRef["taskGroup"]>().toEqualTypeOf<
+ (groupId: string, options?: TaskGroupOptions) => TaskGroupRef
+ >();
+ expectTypeOf<Dag["taskGroup"]>().toEqualTypeOf<
+ (groupId: string, options?: TaskGroupOptions) => TaskGroupRef
+ >();
+ // Python's prefix_group_id, and the only option a group takes so far.
+ expectTypeOf<TaskGroupOptions>().toEqualTypeOf<{ readonly prefixGroupId?:
boolean }>();
+ expectTypeOf<TaskGroupRef["before"]>().toEqualTypeOf<
+ (...downstream: readonly Node[]) => TaskGroupRef
+ >();
+ // Both satisfy Node, which is what lets an edge join either kind.
+ expectTypeOf<TaskRef>().toExtend<Node>();
+ expectTypeOf<TaskGroupRef>().toExtend<Node>();
// 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 612cdc9cdb0..72c064a3df4 100644
--- a/ts-sdk/tests/sdk/dag.test.ts
+++ b/ts-sdk/tests/sdk/dag.test.ts
@@ -22,6 +22,7 @@ import {
Dag,
finalizeDag,
getDagOrderEdges,
+ getDagTaskGroups,
getDagTaskInputs,
getDagTaskRecords,
type DagSpec,
@@ -585,10 +586,10 @@ describe("Dag", () => {
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/,
+ /before\(\) cannot draw an edge from node "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/,
+ /after\(\) cannot draw an edge from node "a" of Dag "d" to itself/,
);
});
@@ -630,7 +631,7 @@ describe("Dag", () => {
expect(() => draw(here.load!, there.cleanup!)).toThrowError(
new RegExp(
- `${verb}\\(\\) cannot draw an edge to Dag "there" task "cleanup"
from Dag "here"`,
+ `${verb}\\(\\) cannot draw an edge to Dag "there" node "cleanup"
from Dag "here"`,
),
);
});
@@ -652,7 +653,7 @@ describe("Dag", () => {
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/,
+ /before\(\) on Dag "d" takes tasks and task groups this Dag handed
out/,
);
});
@@ -678,6 +679,332 @@ describe("Dag", () => {
});
});
+ describe("task groups", () => {
+ it("prefixes the id of every task declared in it", () => {
+ const dag = new Dag("grouped");
+ const staging = dag.taskGroup("staging");
+ staging.task("stage_rows", async () => undefined)();
+
+ expect(dag.taskIds).toEqual(["staging.stage_rows"]);
+ });
+
+ it("prefixes a task id defaulted from the handler name", () => {
+ const dag = new Dag("grouped");
+ dag.taskGroup("staging").task(async function stageRows() {})();
+
+ expect(dag.taskIds).toEqual(["staging.stageRows"]);
+ });
+
+ it("nests, joining every enclosing group's id", () => {
+ const dag = new Dag("grouped");
+ const outer = dag.taskGroup("outer");
+ const inner = outer.taskGroup("inner");
+ inner.task("deep", async () => undefined)();
+
+ expect(dag.taskIds).toEqual(["outer.inner.deep"]);
+ expect(inner.groupId).toBe("outer.inner");
+ });
+
+ it("carries the Dag's own identity", () => {
+ const dag = new Dag("grouped");
+ expect(dag.taskGroup("staging").dagId).toBe("grouped");
+ });
+
+ it("records the tree, so nested groups serialize as one", () => {
+ const dag = new Dag("grouped");
+ const outer = dag.taskGroup("outer");
+ outer.task("first", async () => undefined)();
+ const inner = outer.taskGroup("inner");
+ inner.task("second", async () => undefined)();
+
+ expect([...getDagTaskGroups(dag).values()]).toEqual([
+ {
+ groupId: "outer",
+ prefixGroupId: true,
+ taskIds: ["outer.first"],
+ childGroupIds: ["outer.inner"],
+ },
+ {
+ groupId: "outer.inner",
+ parentGroupId: "outer",
+ prefixGroupId: true,
+ taskIds: ["outer.inner.second"],
+ childGroupIds: [],
+ },
+ ]);
+ });
+
+ it("lets the same handler name be reused across groups", () => {
+ const dag = new Dag("grouped");
+ dag.taskGroup("north").task(async function extract() {})();
+ dag.taskGroup("south").task(async function extract() {})();
+
+ expect(dag.taskIds).toEqual(["north.extract", "south.extract"]);
+ });
+
+ it("lets a task in a group be wired to one outside it", () => {
+ const dag = new Dag("grouped");
+ const extracted = dag.task("extract", async () => 1)();
+ const staged = dag
+ .taskGroup("staging")
+ .task(
+ "stage",
+ async (_: { extracted: number }) => undefined,
+ )({ extracted });
+
+ expect(staged.taskId).toBe("staging.stage");
+ expect(getDagTaskInputs(dag).get("staging.stage")).toEqual({ extracted
});
+ });
+
+ describe("as an edge endpoint", () => {
+ it("orders a whole group before a task", () => {
+ const dag = new Dag("grouped");
+ const staging = dag.taskGroup("staging");
+ staging.task("stage", async () => undefined)();
+ const loaded = dag.task("load", async () => undefined)();
+
+ staging.before(loaded);
+
+ expect(getDagOrderEdges(dag)).toEqual([{ upstream: "staging",
downstream: "load" }]);
+ });
+
+ it("orders a task before a whole group", () => {
+ const dag = new Dag("grouped");
+ const extracted = dag.task("extract", async () => undefined)();
+ const staging = dag.taskGroup("staging");
+ staging.task("stage", async () => undefined)();
+
+ extracted.before(staging);
+
+ expect(getDagOrderEdges(dag)).toEqual([{ upstream: "extract",
downstream: "staging" }]);
+ });
+
+ it("orders one group against another", () => {
+ const dag = new Dag("grouped");
+ const first = dag.taskGroup("first");
+ first.task("a", async () => undefined)();
+ const second = dag.taskGroup("second");
+ second.task("b", async () => undefined)();
+
+ second.after(first);
+
+ expect(getDagOrderEdges(dag)).toEqual([{ upstream: "first",
downstream: "second" }]);
+ });
+
+ it("returns its own receiver, as a task does", () => {
+ const dag = new Dag("grouped");
+ const staging = dag.taskGroup("staging");
+ staging.task("stage", async () => undefined)();
+ const loaded = dag.task("load", async () => undefined)();
+
+ expect(staging.before(loaded)).toBe(staging);
+ });
+
+ it("rejects a group from another Dag object with the same ID", () => {
+ const first = new Dag("same_id");
+ const second = new Dag("same_id");
+ const loaded = first.task("load", async () => undefined)();
+ const foreign = second.taskGroup("staging");
+
+ expect(() => loaded.before(foreign)).toThrowError(
+ /before\(\) was given a reference to "staging" that this Dag did not
hand out/,
+ );
+ });
+
+ it("rejects a foreign group even when this Dag has one of the same ID",
() => {
+ // Matching on the ID alone used to accept it, and the edge was then
+ // drawn at this Dag's own group of that name.
+ const first = new Dag("same_id");
+ const second = new Dag("same_id");
+ first.taskGroup("staging").task("stage", async () => undefined)();
+ const loaded = first.task("load", async () => undefined)();
+ const foreign = second.taskGroup("staging");
+
+ expect(() => loaded.before(foreign)).toThrowError(
+ /before\(\) was given a reference to "staging" that this Dag did not
hand out/,
+ );
+ expect(getDagOrderEdges(first)).toEqual([]);
+ });
+ });
+
+ it.each([
+ ["an empty ID", ""],
+ ["a non-string ID", 42],
+ ])("rejects %s", (_label, groupId) => {
+ const dag = new Dag("grouped");
+ expect(() => dag.taskGroup(groupId as string)).toThrowError(
+ /A task group of Dag "grouped" must have a non-empty ID/,
+ );
+ });
+
+ it("rejects a group ID holding the separator, which nesting is for", () =>
{
+ const dag = new Dag("grouped");
+ expect(() => dag.taskGroup("outer.inner")).toThrowError(
+ /Task group ID "outer.inner" of Dag "grouped" cannot contain "\.";
nest groups with taskGroup/,
+ );
+ });
+
+ it.each([
+ [
+ "a task and a group",
+ (dag: Dag) => [() => dag.task("x", async () => undefined), () =>
dag.taskGroup("x")],
+ ],
+ [
+ "a group and a task",
+ (dag: Dag) => [() => dag.taskGroup("x"), () => dag.task("x", async ()
=> undefined)],
+ ],
+ ["two groups", (dag: Dag) => [() => dag.taskGroup("x"), () =>
dag.taskGroup("x")]],
+ ])(
+ "rejects %s sharing one ID, since a serialized Dag addresses both by it",
+ (_label, build) => {
+ const dag = new Dag("grouped");
+ const [first, second] = build(dag);
+
+ first!();
+ expect(() => second!()).toThrowError(/"x" is already registered for
Dag "grouped"/);
+ },
+ );
+
+ it("allows the same group ID under different parents", () => {
+ const dag = new Dag("grouped");
+ dag
+ .taskGroup("north")
+ .taskGroup("shared")
+ .task("t", async () => undefined)();
+ dag
+ .taskGroup("south")
+ .taskGroup("shared")
+ .task("t", async () => undefined)();
+
+ expect(dag.taskIds).toEqual(["north.shared.t", "south.shared.t"]);
+ });
+
+ describe("with prefixGroupId off", () => {
+ it("keeps the ids declared in it as written", () => {
+ const dag = new Dag("grouped");
+ const checks = dag.taskGroup("checks", { prefixGroupId: false });
+ checks.task("nulls", async () => undefined)();
+ const inner = checks.taskGroup("inner");
+
+ expect(dag.taskIds).toEqual(["nulls"]);
+ expect(checks.groupId).toBe("checks");
+ expect(inner.groupId).toBe("inner");
+ expect(getDagTaskGroups(dag).get("checks")?.prefixGroupId).toBe(false);
+ });
+
+ it("lets a nested group prefix again, as Python does", () => {
+ const dag = new Dag("grouped");
+ const inner = dag.taskGroup("outer", { prefixGroupId: false
}).taskGroup("inner");
+ inner.task("t", async () => undefined)();
+
+ expect(inner.groupId).toBe("inner");
+ expect(dag.taskIds).toEqual(["inner.t"]);
+ });
+
+ it("still takes its own place under a prefixing parent", () => {
+ const dag = new Dag("grouped");
+ const inner = dag.taskGroup("outer").taskGroup("inner", {
prefixGroupId: false });
+ inner.task("t", async () => undefined)();
+
+ expect(inner.groupId).toBe("outer.inner");
+ expect(dag.taskIds).toEqual(["t"]);
+ });
+
+ it("rejects an id that clashes with one outside the group", () => {
+ const dag = new Dag("grouped");
+ dag.task("nulls", async () => undefined)();
+ const checks = dag.taskGroup("checks", { prefixGroupId: false });
+
+ expect(() => checks.task("nulls", async () => undefined)).toThrowError(
+ /Task "nulls" is already registered for Dag "grouped"/,
+ );
+ });
+ });
+
+ it.each([
+ [
+ "an unknown option",
+ { prefix: false },
+ /Unknown option "prefix" for Dag "grouped" task group "checks"/,
+ ],
+ [
+ "a prefixGroupId that is not a boolean",
+ { prefixGroupId: "no" },
+ /prefixGroupId for Dag "grouped" task group "checks" must be a
boolean/,
+ ],
+ [
+ "options that are not an object",
+ 7,
+ /options for Dag "grouped" task group "checks" must be an object/,
+ ],
+ ])("rejects %s", (_label, options, expected) => {
+ const dag = new Dag("grouped");
+
+ expect(() => dag.taskGroup("checks", options as
never)).toThrowError(expected);
+ expect(getDagTaskGroups(dag).size).toBe(0);
+ });
+
+ describe("its ID, which follows Python's validate_group_key", () => {
+ it("may hold letters of any script, digits, dashes and underscores", ()
=> {
+ const dag = new Dag("grouped");
+
+ expect(dag.taskGroup("región-1_組").groupId).toBe("región-1_組");
+ });
+
+ it.each([
+ ["a space", "two words"],
+ ["a slash", "a/b"],
+ ["a symbol", "cost$"],
+ ])("rejects %s", (_label, groupId) => {
+ const dag = new Dag("grouped");
+
+ expect(() => dag.taskGroup(groupId)).toThrowError(
+ `Task group ID "${groupId}" of Dag "grouped" has to be made of
letters, digits, dashes and underscores`,
+ );
+ });
+
+ it("allows 200 characters and rejects 201", () => {
+ const dag = new Dag("grouped");
+ dag.taskGroup("g".repeat(200));
+
+ expect(() => dag.taskGroup("h".repeat(201))).toThrowError(
+ /has 201 characters; at most 200 are allowed/,
+ );
+ });
+
+ it("counts characters as Python's len() does, not UTF-16 units", () => {
+ const dag = new Dag("grouped");
+ // Each of these letters is two UTF-16 units, so the id is 400 units
long.
+ const groupId = "\u{20000}".repeat(200);
+
+ expect(dag.taskGroup(groupId).groupId).toBe(groupId);
+ });
+ });
+
+ it("rejects a group declared after the Dag was read", () => {
+ const dag = new Dag("grouped");
+ dag.task("extract", async () => undefined)();
+ finalizeDag(dag);
+
+ expect(() => dag.taskGroup("late")).toThrowError(
+ /Task group "late" cannot be added to Dag "grouped" after the Dag was
read/,
+ );
+ });
+
+ it("holds its tasks to the same every-task-is-called rule", () => {
+ const dag = new Dag("grouped");
+ dag.taskGroup("staging").task("stage", async () => undefined);
+
+ expect(() => finalizeDag(dag)).toThrowError(
+ /Task "staging.stage" of Dag "grouped" is never called/,
+ );
+ });
+
+ it("has no groups before any are declared", () => {
+ expect(getDagTaskGroups(new Dag("plain")).size).toBe(0);
+ });
+ });
+
it("exposes its task IDs in attachment order", () => {
const dag = new Dag("ordered_dag");
expect(dag.taskIds).toEqual([]);
@@ -734,13 +1061,25 @@ describe("Dag", () => {
);
});
- it("treats a dotted TaskGroup taskId as a single taskId (group.task)", () =>
{
+ it("rejects a dotted task id, which names a group that does not exist", ()
=> {
+ // A dotted id used to be accepted verbatim, giving a task whose id carries
+ // a group prefix while the Dag holds no such group. `taskGroup(...)` is
+ // what puts a task under a prefix now.
const dag = new Dag("example_dag");
- dag.task("transforms.normalize", async () => "ok")();
- const bundle = new Bundle();
- bundle.register(dag);
+
+ expect(() => dag.task("transforms.normalize", async () =>
"ok")).toThrowError(
+ /Task ID "transforms.normalize" of Dag "example_dag" cannot contain
"\."/,
+ );
+ expect(dag.taskIds).toEqual([]);
+ });
+
+ it("names a task inside a group without the author writing the prefix", ()
=> {
+ const dag = new Dag("example_dag");
+ dag.taskGroup("transforms").task("normalize", async () => "ok")();
+ const bundle = new Bundle(dag);
+
expect(bundle.getTaskHandler("example_dag",
"transforms.normalize")).toBeDefined();
- // Should NOT accidentally match the prefix alone
+ // And not the prefix or the leaf on its own.
expect(bundle.getTaskHandler("example_dag", "transforms")).toBeUndefined();
expect(bundle.getTaskHandler("example_dag", "normalize")).toBeUndefined();
});