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();
   });

Reply via email to