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 a3555c74cfc TS SDK: default a task id to the handler's name (#73437)
a3555c74cfc is described below
commit a3555c74cfcf90f40f43e399b558a5f5dcc2d069
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Tue Sep 29 23:48:29 2026 +0800
TS SDK: default a task id to the handler's name (#73437)
A native Dag handler takes one object of named arguments, and a call names
each input, so nothing labels an argument arg0. A call given more than one
argument is an error rather than a silently dropped value.
The task id defaults to the handler's function name, so airflow-ts-pack
passes esbuild's keepNames to keep that name through minification.
---
.../language-sdks/typescript.rst | 43 ++--
docs/spelling_wordlist.txt | 1 +
ts-sdk/adr/0002-native-dag-interface.md | 58 ++---
ts-sdk/src/cli/pack.ts | 4 +
ts-sdk/src/index.ts | 1 -
ts-sdk/src/sdk/dag.ts | 272 +++++++++++----------
ts-sdk/tests/cli/pack.test.ts | 68 ++++++
ts-sdk/tests/public-api.test.ts | 75 +++---
ts-sdk/tests/sdk/dag.test.ts | 197 ++++++++-------
9 files changed, 423 insertions(+), 296 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 d3ddca0c9b4..e40b1363c0a 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
@@ -247,8 +247,8 @@ A ``Dag`` is declared on this side rather than in Python:
its schedule, its task
the edges between them are all written in TypeScript. The surface is still
growing, so a Dag declared
this way is not served to Airflow yet.
-``dag.task(taskId, handler)`` returns a *factory*. Calling it places the task
in the Dag and supplies the
-handler's arguments, so the call graph is the task graph:
+``dag.task(taskId, handler)`` returns a *factory*. A handler takes one object
of named arguments, and
+calling the factory names each input, so the call graph is the task graph:
.. code-block:: typescript
@@ -257,19 +257,19 @@ handler's arguments, so the call graph is the task graph:
const dag = new Dag("ts_etl");
const extract = dag.task("extract", async (): Promise<number> => 42);
- const transform = dag.task("transform", async (rows: number, region:
string) => rows * 2);
- const load = dag.task("load", async (total: number) => {});
+ const transform = dag.task(
+ "transform",
+ async ({ rows, region }: { rows: number; region: string }) => rows * 2,
+ );
+ const load = dag.task("load", async ({ total }: { total: number }) => {});
- load(transform(extract(), "us"));
+ const extracted = extract();
+ const total = transform({ rows: extracted, region: "us" });
+ load({ total });
-Arguments are passed in the order the handler declares them. A handler that
declares a single object of
-named arguments can also be called with that object, which names each input
instead of ordering it:
-
-.. code-block:: typescript
-
- const store = dag.task("store", async ({ total }: { total: number }) =>
{});
-
- store({ total: extract() });
+Naming the inputs is how a task is called. A handler that takes no arguments
is called with none, and
+a single argument is named like any other, ``load({ total })``. The compiler
checks the call: it reports
+an argument left out, a misspelled one, and a literal of the wrong type.
Each argument takes either an upstream reference or a literal JSON value. A
reference has to be the
argument itself: one buried inside an array or an object is a literal, and
draws no edge.
@@ -277,6 +277,18 @@ argument itself: one buried inside an array or an object
is a literal, and draws
Every task has to be called exactly once. An uncalled task fails when the Dag
is read, so none can be
left out of the graph by accident.
+The task id may be omitted, in which case it is the handler's function name:
+
+.. code-block:: typescript
+
+ const extract = dag.task(async function extract(): Promise<number> {
+ return 42;
+ });
+
+``airflow-ts-pack`` keeps function names intact, so bundling cannot rename a
task. A handler with no
+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.
+
``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.
@@ -398,9 +410,8 @@ layout header. The layout records the byte ranges and
SHA-256 digests of the man
so there is one file to deploy, with no separate manifest or ``node_modules``.
The code is minified because an integrity digest is only worth taking over an
artifact nobody is expected to
-read or edit in place. The ``/*! */`` license banners of bundled dependencies
are kept. Nothing is identified by
-a function name, so minified names are safe: a Dag and a task are named by the
string ids their registration
-states, and a handler is dispatched by reference.
+read or edit in place. Function names are kept through minification, since a
task id defaults to its
+handler's name. The ``/*! */`` license banners of bundled dependencies are
kept.
Because the shipped code is not the code anyone wrote, the packer also embeds
the entry module verbatim in a
``/*# airflowSource ... #*/`` block comment, verified by its own digest, so
Airflow has something readable to
diff --git a/docs/spelling_wordlist.txt b/docs/spelling_wordlist.txt
index 8a1ed93032a..f4cc0e617dc 100644
--- a/docs/spelling_wordlist.txt
+++ b/docs/spelling_wordlist.txt
@@ -1102,6 +1102,7 @@ midnights
milli
millis
milton
+minification
minikube
misconfiguration
misconfigured
diff --git a/ts-sdk/adr/0002-native-dag-interface.md
b/ts-sdk/adr/0002-native-dag-interface.md
index eca1f5e2ea2..6614ee62091 100644
--- a/ts-sdk/adr/0002-native-dag-interface.md
+++ b/ts-sdk/adr/0002-native-dag-interface.md
@@ -27,16 +27,13 @@ Proposed. Revised after the review on #72047.
1. **`dag.task(handler)` returns a factory, and the task id is optional.**
With no id the task takes
the handler's function name (`dag.task(extract)` → task `"extract"`);
`dag.task(taskId, handler)`
- sets it explicitly, which an anonymous handler must do. Calling the factory
both places the task in
- the Dag and supplies its arguments, in the order the handler declares them
- (`load(transform(extract(), "us"))`). A handler that declares a single
object of named arguments can
- also be called with that object — the shape Python TaskFlow uses for
- `load(transformed=transform(...))`.
-2. **The call graph is the task graph.** `tsc` checks every wired key against
the handler's own
- parameter type, and a `TaskRef` exists only once its producing call has
returned, so a cycle
- through arguments is unrepresentable rather than rejected by a validator. A
reference passed by
- position is checked against the argument's own type, which is what tells
the two call shapes apart
- when a handler declares a single argument.
+ sets it explicitly, which an anonymous handler must do. A handler takes one
object of named
+ arguments, and calling the factory both places the task in the Dag and
names each of its inputs
+ (`load({ total: transform({ rows: extract(), region: "us" }) })`) — the
shape Python TaskFlow uses
+ for `load(transformed=transform(...))`.
+2. **The call graph is the task graph.** `tsc` checks every named input
against the handler's own
+ argument type, and a `TaskRef` exists only once its producing call has
returned, so a cycle
+ through arguments is unrepresentable rather than rejected by a validator.
3. **Every task is called exactly once.** An uncalled task fails when the Dag
is read, so none can be
silently left out of the graph.
4. **`before` and `after` draw order-only edges** — the TypeScript pair for
`>>` and `<<`, both
@@ -107,13 +104,12 @@ const extract = dag.task(async function extract():
Promise<number> {
// task id "extract"
```
-The id comes from the handler's *source* name, resolved when the bundle is
packed and written into
-the registration — not from `handler.name` at runtime, which minification
renames (see
-Implementation Notes). A handler with no source name — a bare anonymous arrow
passed inline,
-`dag.task(async () => 42)` — has nothing to resolve and is a compile error
until given an explicit
-id. This default is for native Dags, where both ends of every name are
TypeScript; a mixed-language
-handler names the Python-owned task explicitly and does not default from the
handler's function
-name ([ADR-0001](0001-mixed-lang-dag-interface.md), decision 3).
+The id is the handler's `name`, which `airflow-ts-pack` keeps through
minification (see
+Implementation Notes). A handler with no name, such as a bare anonymous arrow
passed inline,
+`dag.task(async () => 42)`, has nothing to take an id from and fails when the
Dag is declared until
+given an explicit id. This default is for native Dags, where both ends of
every name are
+TypeScript; a mixed-language handler names the Python-owned task explicitly
and does not default
+from the handler's function name
([ADR-0001](0001-mixed-lang-dag-interface.md), decision 3).
The `TaskSpec` also carries the task id, so it can be set alongside the other
task options:
@@ -149,9 +145,8 @@ convention.
by design. Native declaration is what fills them, generated from the
serialized-Dag JSON schema the
way `src/generated/supervisor.ts` is. This ADR does not choose those fields;
it fixes where an
author writes them.
-- `TaskOptions` carries the spec and the handler's positional argument names,
which the packer fills in
- from the parameter list so the Dag names each argument as its handler does.
With wiring moved to the
- factory call, `inputs` is no longer an option.
+- `TaskOptions` carries the task's spec and nothing else: the names on the
wire are the keys of the
+ call itself. With wiring moved to the factory call, `inputs` is no longer an
option.
- `TaskHandlerArgs` is removed from the public API, `DagRegistry` becomes
`Bundle`, and
`serveDags(registry)` becomes `bundle.serve()`, which breaks
0.1.0-beta1 authors; see [ADR-0001](0001-mixed-lang-dag-interface.md) for
the shipped call sites
@@ -159,10 +154,10 @@ convention.
## Alternatives
-- **Named-only wiring**, rejected in the review on #73435: naming every input
reads well at twenty
- tasks but forces an object around a single argument, and positional calls
are what TypeScript
- authors write. Both are offered, and the handler's own parameter list
decides which one a task can
- use.
+- **Positional handlers**, `async (rows: number, region: string) => ...`,
offered first and then
+ dropped: a positional parameter list has no names on the wire unless the SDK
reads them out of the
+ handler's source, and a single object of named arguments is what a
TypeScript library takes
+ anyway.
- **Injected `ctx`/`client` arguments**, mimicking the Python signature.
Rejected, per the above and
because feeling native to TypeScript matters more than matching Python's
parameter list.
@@ -181,17 +176,14 @@ convention.
- **The spec argument already has its slot.** `dag.task(taskId, handler,
options)` reads `{ spec = {} }`
and runs `validateEmptySpec` on it (`ts-sdk/src/sdk/dag.ts`), so task fields
land on a path that
exists rather than a new one.
-- **A positional argument binds by order, and its name is a label.** The
serialized Dag names each
- argument, so the packer reads the names from the handler's parameter list;
`arg0`, `arg1` and so on
- stand in for a name it cannot see, without changing which value reaches
which argument.
- **A `TaskRef` is inert** — a handle for wiring, not a promise. Nothing in a
Dag file executes a task
body.
-- **A defaulted task id is resolved at pack time, not read at runtime.**
esbuild renames function
- identifiers, so `handler.name` in a packed bundle is the minified name, not
the author's. The pack
- step (`ts-sdk/src/cli/pack.ts`) therefore reads an omitted id from the
handler's declared name in
- source and writes it into the registration, rather than depending on
`handler.name` or enabling
- esbuild's `keepNames` across the whole bundle. A handler with no source name
leaves nothing to
- read, which is why an anonymous handler must state its id.
+- **A defaulted task id is read off the handler itself.** The pack step
(`ts-sdk/src/cli/pack.ts`)
+ minifies and passes esbuild's `keepNames`, so `handler.name` is the author's
in a packed bundle as
+ much as in one run from source. Rewriting the call at pack time was tried
first and dropped: it
+ needed a TypeScript parser in the packer to tell a real `.task(` from one
inside a string or a
+ comment, and it could not see a handler declared in another module. A
handler with no name leaves
+ nothing to read, which is why an anonymous handler must state its id.
- **`withArgNames` and the name folding behind it**
([ADR-0001](0001-mixed-lang-dag-interface.md))
exist for the mixed-language case and are never needed here: both ends of
every name are
TypeScript, so `tsc` checks the wiring end to end and there is no foreign
name to reconcile.
diff --git a/ts-sdk/src/cli/pack.ts b/ts-sdk/src/cli/pack.ts
index 10e8c5fbbdc..3d33d3b0041 100644
--- a/ts-sdk/src/cli/pack.ts
+++ b/ts-sdk/src/cli/pack.ts
@@ -217,6 +217,10 @@ export async function runPack(argv: readonly string[]):
Promise<void> {
target: "node22",
// A digest is only worth taking over an artifact nobody reads or edits
in place.
minify: true,
+ // A task id comes from the handler's name, and bundling renames a symbol
+ // that two modules both declare, so the name the SDK reads is pinned to
+ // the one the author wrote.
+ keepNames: true,
// The manifest is read by running the staged bundle, so the metadata
describes what ships.
outfile: stagingPath,
});
diff --git a/ts-sdk/src/index.ts b/ts-sdk/src/index.ts
index d63bdf24186..ad45c3bad33 100644
--- a/ts-sdk/src/index.ts
+++ b/ts-sdk/src/index.ts
@@ -28,7 +28,6 @@ export type { ArgNameMap } from "./sdk/arg-names.js";
export type { Registerable } from "./sdk/bundle.js";
export type {
DagSpec,
- PositionalInputs,
TaskFactory,
TaskInput,
TaskInputs,
diff --git a/ts-sdk/src/sdk/dag.ts b/ts-sdk/src/sdk/dag.ts
index a264d306104..2903b2b6ac4 100644
--- a/ts-sdk/src/sdk/dag.ts
+++ b/ts-sdk/src/sdk/dag.ts
@@ -84,7 +84,9 @@ function kindOf(value: object): string {
}
const DAG_SPEC_KEYS: ReadonlySet<string> = new
Set(Object.keys(DAG_SCHEMA_FIELDS));
-const TASK_SPEC_KEYS: ReadonlySet<string> = new
Set(Object.keys(TASK_SCHEMA_FIELDS));
+// `taskId` is hand-written rather than generated: the schema's task_id is
+// serializer-owned, and this is the authoring surface's own way to set it.
+const TASK_SPEC_KEYS: ReadonlySet<string> = new
Set([...Object.keys(TASK_SCHEMA_FIELDS), "taskId"]);
/**
* Dag-level options: the schedule, the tags, how many runs may be active, and
@@ -103,9 +105,19 @@ export type DagSpec = GeneratedDagFields;
* Task-level options: the retries, the pool, the trigger rule, and the rest of
* what an operator takes in Python.
*
+ * The task id is here too, for a handler whose id is not given positionally.
* Optional and record-only on the same terms as {@link DagSpec}.
*/
-export type TaskSpec = GeneratedTaskFields;
+export interface TaskSpec extends GeneratedTaskFields {
+ /**
+ * Airflow task ID, when it should not be the handler's function name.
+ *
+ * `dag.task(handler)` takes the id from the handler's name. Set this for an
+ * id that has to outlive that name, or for an anonymous handler, which has
+ * no name to take one from.
+ */
+ readonly taskId?: string;
+}
// Carries a reference's return type without carrying a value. Not exported, so
// the property cannot be read or written from outside; it exists only so
@@ -197,75 +209,46 @@ type JsonCompatible<T> = T extends JsonValue
*/
export type TaskInput<TValue> = TaskRef<TValue> | JsonCompatible<TValue>;
-/** The inputs of a task that declares several arguments, in declaration
order. */
-export type PositionalInputs<TParams extends readonly unknown[]> = {
- [K in keyof TParams]: TaskInput<TParams[K]>;
-};
-
-/** The inputs of a task that declares one object of named arguments, by name.
*/
+/** The inputs of a task, keyed by the name of the argument each one supplies.
*/
export type TaskInputs<TArgs> = {
[K in keyof TArgs]: TaskRef | JsonCompatible<TArgs[K]>;
};
-// Offered only where it means something: `TaskInputs<number>` would map over
-// `number`'s own methods and accept `{ toFixed: ... }`.
-type NamedInputs<TOnly> = [TOnly] extends [object] ? TaskInputs<TOnly> : never;
-
/**
* What `dag.task(...)` returns: call it to declare where the task sits in the
Dag.
*
- * Pass one input per argument the handler declares, in order. An input is
- * either another task's reference, which makes this task wait for that task
and
- * receive its result, or a literal value:
+ * A handler takes one object of named arguments, and the call names each
+ * input. An input is either another task's reference, which makes this task
+ * wait for that task and receive its result, or a literal value:
*
* ```ts
* const extract = dag.task("extract", async (): Promise<number> => 42);
- * const transform = dag.task("transform", async (rows: number, region:
string) => rows);
- * const load = dag.task("load", async (total: number) => {});
+ * const transform = dag.task(
+ * "transform",
+ * async ({ rows, region }: { rows: number; region: string }) => rows,
+ * );
+ * const load = dag.task("load", async ({ total }: { total: number }) => {});
*
- * load(transform(extract(), "us"));
+ * const extracted = extract();
+ * load({ total: transform({ rows: extracted, region: "us" }) });
* ```
*
- * A handler that declares a single object of named arguments can also be
called
- * with that object, which names each input instead of ordering it:
- *
- * ```ts
- * const store = dag.task("store", async ({ total }: { total: number }) => {});
- *
- * store({ total: extract() });
- * ```
- *
- * The compiler checks that every argument is supplied and that each literal
- * matches its argument's type. A reference passed by position is checked
against
- * the argument's type as well, which is what tells the two call shapes apart
- * when a handler declares a single argument.
+ * The call is checked argument by argument: every one has to be supplied, and
+ * a literal has to match its argument's type.
*/
-export type TaskFactory<TParams extends readonly unknown[], TReturn = unknown>
= [TParams] extends [
- readonly [],
+export type TaskFactory<TArgs extends object | void = void, TReturn = unknown>
= [TArgs] extends [
+ void,
]
? () => TaskRef<TReturn>
- : TParams extends readonly [infer TOnly]
- ? (input: TaskInput<TOnly> | NamedInputs<TOnly>) => TaskRef<TReturn>
- : (...inputs: PositionalInputs<TParams>) => TaskRef<TReturn>;
+ : (inputs: TaskInputs<TArgs>) => TaskRef<TReturn>;
/**
- * The trailing argument of `dag.task()`: the task's own {@link TaskSpec}, plus
- * the names of the handler's positional arguments.
+ * The trailing argument of `dag.task()`: the task's own {@link TaskSpec}.
*
* Every field is optional, and an unknown key is rejected, so a misspelled
* field is an error rather than a task that quietly ignores it.
*/
-export type TaskOptions = TaskSpec & {
- /**
- * Names for the handler's positional arguments, in declaration order.
- *
- * `airflow-ts-pack` fills this in from the handler's parameter list, so the
- * Dag names each argument as its handler does. Positional inputs bind by
- * order, so a name left out only costs the label: `arg0`, `arg1` and so on
- * stand in for it.
- */
- readonly argBindings?: readonly string[];
-};
+export type TaskOptions = TaskSpec;
/** Per-task record a Dag retains: the reference, the handler, and its spec. */
export interface TaskRecord {
@@ -350,19 +333,86 @@ export class Dag {
/**
* Declare a task of this Dag, and return the factory that places it.
*
- * Every argument the handler declares becomes an input of the returned
- * {@link TaskFactory}; `getContext()` and `getClient()` reach the runtime
- * from inside the call, so neither is an argument. The trailing options
object
- * carries this task's own {@link TaskSpec}.
+ * A handler takes one object of named arguments, and every argument in it
+ * becomes an input of the returned {@link TaskFactory}; `getContext()` and
+ * `getClient()` reach the runtime from inside the call, so neither is an
+ * argument. The trailing options object carries this task's own
+ * {@link TaskSpec}.
+ *
+ * ```ts
+ * const extract = dag.task("extract", async () => 42);
+ * const transform = dag.task(async function transform() {}); // id
"transform"
+ * ```
*/
- task<TParams extends readonly unknown[] = [], TReturn = unknown>(
+ task<TArgs extends object | void = void, TReturn = unknown>(
taskId: string,
- handler: (...args: TParams) => TReturn | Promise<TReturn>,
- options: TaskOptions = {},
- ): TaskFactory<TParams, TReturn> {
+ handler: (args: TArgs) => TReturn | Promise<TReturn>,
+ options?: TaskOptions,
+ ): TaskFactory<TArgs, TReturn>;
+ /**
+ * Declare a task whose id is the handler's function name.
+ *
+ * `airflow-ts-pack` keeps handler names intact, so minification cannot
change
+ * a task id. An anonymous handler has no name to take one from, and needs
+ * {@link TaskSpec.taskId} to give it one.
+ */
+ task<TArgs extends object | void = void, TReturn = unknown>(
+ handler: (args: TArgs) => TReturn | Promise<TReturn>,
+ options?: TaskOptions,
+ ): TaskFactory<TArgs, TReturn>;
+ task<TArgs extends object | void = void, TReturn = unknown>(
+ 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 (
+ args: TArgs,
+ ) => TReturn | Promise<TReturn>;
+ const given = idGiven ? maybeOptions : (handlerOrOptions as TaskOptions |
undefined);
+ // Defaulted only when absent: an explicit `null` is a bad spec, not an
+ // omitted one, and #taskSpecOf is what reports it.
+ const options = given === undefined ? {} : given;
+ const specTaskId =
+ isPlainRecord(options) && typeof options.taskId === "string" ?
options.taskId : undefined;
+ // Two ids for one task disagree silently otherwise: the positional one
+ // wins and the spec's is dropped without a word.
+ if (idGiven && specTaskId !== undefined) {
+ throw new Error(
+ `Task "${taskIdOrHandler}" of Dag "${this.dagId}" also carries taskId
"${specTaskId}" in ` +
+ "its spec; give the id once, either positionally or in the spec",
+ );
+ }
+ const defaulted = idGiven ? undefined : (specTaskId ??
readFunctionName(handler));
+ const taskId = idGiven ? taskIdOrHandler : defaulted;
+ if (taskId === 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 ' +
+ "that drops function names also lands here; airflow-ts-pack keeps
them.",
+ );
+ }
+ // A name Airflow would reject is worth catching where it was taken, not in
+ // the server's answer: `fn.bind(...)` names itself "bound extract", and a
+ // method can be named anything at all.
+ if (defaulted !== undefined && !TASK_ID_CHARACTERS.test(defaulted)) {
+ throw new Error(
+ `A task of Dag "${this.dagId}" would take the id "${defaulted}" from
its handler's name, ` +
+ "which Airflow does not accept; give it an id of letters, digits,
dashes, dots and " +
+ 'underscores — dag.task("my_task", handler)',
+ );
+ }
if (typeof handler !== "function") {
throw new Error(`handler for Dag "${this.dagId}" task "${taskId}" must
be a function`);
}
+ // TypeScript already says so, but a plain-JavaScript author lands here
+ // with the argument list Python would take.
+ if (handler.length > 1) {
+ throw new Error(
+ `Handler for Dag "${this.dagId}" task "${taskId}" declares
${handler.length} parameters; ` +
+ "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}"`);
}
@@ -375,25 +425,30 @@ export class Dag {
);
}
const spec = this.#taskSpecOf(taskId, options);
- const argBindings = this.#validateArgBindings(taskId, options.argBindings);
const task = createTaskRef(this.dagId, taskId);
this.#tasks.set(taskId, {
task,
// The runtime dispatches every handler through one instantiation, as it
- // does a registered TaskHandler; a positional one is wrapped at wiring.
+ // does a registered TaskHandler.
fn: handler as unknown as TaskFunction,
spec: freezeSpec(spec, () => `The spec for Dag "${this.dagId}" task
"${taskId}"`),
});
return ((...inputs: unknown[]) => {
- this.#wire(taskId, inputs, argBindings);
+ // TypeScript already says so, but from plain JavaScript a second
+ // argument would be dropped without a word.
+ if (inputs.length > 1) {
+ throw new Error(
+ `Task "${taskId}" of Dag "${this.dagId}" was given ${inputs.length}
arguments; ` +
+ "it takes one object naming its inputs: myTask({ rows, region })",
+ );
+ }
+ this.#wire(taskId, inputs[0]);
return task;
- }) as TaskFactory<TParams, TReturn>;
+ }) as TaskFactory<TArgs, TReturn>;
}
// TypeScript is bypassable — from plain JavaScript, or an `as TaskSpec` cast
- // — so an unknown key is rejected rather than silently ignored.
`argBindings`
- // names the handler's arguments rather than configuring the task, so it is
- // taken out here instead of reaching the spec.
+ // — so an unknown key is rejected rather than silently ignored.
#taskSpecOf(taskId: string, options: TaskOptions): TaskSpec {
const value: unknown = options;
if (!isPlainRecord(value)) {
@@ -401,7 +456,6 @@ export class Dag {
}
const spec: Record<string, unknown> = {};
for (const key of Reflect.ownKeys(value)) {
- if (key === "argBindings") continue;
if (typeof key !== "string" || !TASK_SPEC_KEYS.has(key)) {
throw new Error(
`Unknown option "${String(key)}" in the spec for Dag "${this.dagId}"
task "${taskId}"`,
@@ -412,31 +466,7 @@ export class Dag {
return spec as TaskSpec;
}
- // TypeScript is bypassable, and these names become the keys the arguments
are
- // recorded under, so an integer-like one would reorder what it labels.
- #validateArgBindings(taskId: string, names: unknown): readonly string[] |
undefined {
- if (names === undefined) return undefined;
- const describe = `argBindings for Dag "${this.dagId}" task "${taskId}"`;
- if (!Array.isArray(names)) throw new Error(`${describe} must be an array
of names`);
- const seen = new Set<string>();
- for (const name of names as unknown[]) {
- if (typeof name !== "string" || name.length === 0 || /^\d+$/.test(name))
{
- throw new Error(
- `${describe} holds ${JSON.stringify(name)}; each name must be a
non-empty ` +
- "string that is not a number",
- );
- }
- if (seen.has(name)) throw new Error(`${describe} names "${name}" twice`);
- seen.add(name);
- }
- return Object.freeze([...(names as string[])]);
- }
-
- #wire(
- taskId: string,
- inputs: readonly unknown[],
- argBindings: readonly string[] | undefined,
- ): void {
+ #wire(taskId: string, inputs: unknown): void {
if (this.#finalized) {
throw new Error(
`Task "${taskId}" of Dag "${this.dagId}" was called after the Dag was
read; ` +
@@ -449,16 +479,20 @@ export class Dag {
"in a Dag, so call it once and reuse the reference",
);
}
- const positional = !isNamedCall(inputs);
- const recorded = this.#checkInputs(
- taskId,
- positional ? positionalInputs(inputs, argBindings) : (inputs[0] as
Record<string, unknown>),
- );
- this.#inputs.set(taskId, recorded);
- if (positional && inputs.length > 0) {
- const record = this.#tasks.get(taskId)!;
- this.#tasks.set(taskId, { ...record, fn: spreadArgs(record.fn) });
+ this.#inputs.set(taskId, this.#checkInputs(taskId,
this.#inputsByName(taskId, inputs)));
+ }
+
+ #inputsByName(taskId: string, inputs: unknown): Record<string, unknown> {
+ if (inputs === undefined) return {};
+ // A reference is a plain object too, so `load(extracted)` would otherwise
+ // read as a map of argument names.
+ if (!isPlainRecord(inputs) || isTaskRef(inputs)) {
+ throw new Error(
+ `Task "${taskId}" of Dag "${this.dagId}" takes one object naming its
inputs: ` +
+ "myTask({ rows, region })",
+ );
}
+ return inputs;
}
#checkInputs(taskId: string, inputs: Record<string, unknown>):
RecordedInputs {
@@ -525,39 +559,21 @@ export class Dag {
}
}
-/**
- * Whether a call named its inputs rather than ordering them.
- *
- * One plain object is the named form. A reference is a plain object too, so it
- * is ruled out first: `load(extract())` is one positional input, not a map of
- * argument names.
- */
-function isNamedCall(inputs: readonly unknown[]): boolean {
- return inputs.length === 1 && !isTaskRef(inputs[0]) &&
isPlainRecord(inputs[0]);
-}
-
-function positionalInputs(
- inputs: readonly unknown[],
- argBindings: readonly string[] | undefined,
-): Record<string, unknown> {
- const byName: Record<string, unknown> = {};
- inputs.forEach((value, index) => {
- byName[argBindings?.[index] ?? `arg${index}`] = value;
- });
- return byName;
-}
+// What Airflow accepts as a task id, and so what a name taken from a handler
+// has to look like.
+const TASK_ID_CHARACTERS = /^[\p{L}\p{N}_.-]+$/u;
/**
- * Dispatch a positional handler through the one call shape the runtime uses.
+ * A handler's own function name, or undefined when it has none.
*
- * A task is called with its bound arguments as a single object, in the order
- * they were recorded, so spreading its values back restores the argument list
- * the handler declared.
+ * Where a defaulted task id comes from. `airflow-ts-pack` bundles with
esbuild's
+ * `keepNames`, so the name survives minification; a bundler that drops names
+ * leaves the empty string, which is why the empty string is not an id.
*/
-function spreadArgs(fn: TaskFunction): TaskFunction {
- const handler = fn as unknown as (...args: unknown[]) => unknown;
- return ((args: Record<string, unknown>) =>
- handler(...Object.values(args ?? {}))) as unknown as TaskFunction;
+function readFunctionName(handler: unknown): string | undefined {
+ if (typeof handler !== "function") return undefined;
+ const { name } = handler as { name?: unknown };
+ return typeof name === "string" && name.length > 0 ? name : undefined;
}
function createTaskRef(dagId: string, taskId: string): TaskRef {
diff --git a/ts-sdk/tests/cli/pack.test.ts b/ts-sdk/tests/cli/pack.test.ts
index c9d8e6a9ea6..2ce54327cc8 100644
--- a/ts-sdk/tests/cli/pack.test.ts
+++ b/ts-sdk/tests/cli/pack.test.ts
@@ -557,6 +557,74 @@ describe("runPack", () => {
);
});
+ it("takes an omitted task id from the handler name, through minification",
async () => {
+ outdir = mkdtempSync(path.join(tmpdir(), "ts-pack-"));
+ const entry = path.join(outdir, "named-entry.ts");
+ writeFileSync(
+ entry,
+ [
+ `import { Bundle, Dag } from ${JSON.stringify(SDK_INDEX)};`,
+ 'const salesDag = new Dag("sales_dag");',
+ "salesDag.task(async function extractRows() {})();",
+ "async function loadRows() {}",
+ "salesDag.task(loadRows, { retries: 2 })();",
+ "await new Bundle(salesDag).serve();",
+ ].join("\n"),
+ );
+ const stderr = captureStderr();
+
+ await runPack([entry, "--outdir", outdir]);
+
+ // Read back off the packed artifact, so this asserts what minification
+ // left behind rather than what the source said.
+ expect(JSON.parse(readEmbeddedMetadata(path.join(outdir,
"bundle.min.mjs")))).toHaveProperty(
+ "task_handlers.sales_dag.tasks",
+ ["extractRows", "loadRows"],
+ );
+ expect(stderr()).toBe("");
+ });
+
+ it("fails the pack when a handler is anonymous and names no task", async ()
=> {
+ outdir = mkdtempSync(path.join(tmpdir(), "ts-pack-"));
+ const entry = path.join(outdir, "anonymous-entry.ts");
+ writeFileSync(
+ entry,
+ [
+ `import { Bundle, Dag } from ${JSON.stringify(SDK_INDEX)};`,
+ 'const salesDag = new Dag("sales_dag");',
+ "salesDag.task(async () => undefined)();",
+ "await new Bundle(salesDag).serve();",
+ ].join("\n"),
+ );
+
+ await expect(runPack([entry, "--outdir", outdir])).rejects.toThrow(
+ /has no id: its handler has no name/,
+ );
+ });
+
+ it("takes a computed task id as the id it evaluates to", async () => {
+ outdir = mkdtempSync(path.join(tmpdir(), "ts-pack-"));
+ const entry = path.join(outdir, "computed-entry.ts");
+ writeFileSync(
+ entry,
+ [
+ `import { Bundle, Dag } from ${JSON.stringify(SDK_INDEX)};`,
+ 'const salesDag = new Dag("sales_dag");',
+ 'const region = "north";',
+ "salesDag.task(`extract_${region}`, async () => undefined)();",
+ 'salesDag.task("load_".concat(region), async () => undefined)();',
+ "await new Bundle(salesDag).serve();",
+ ].join("\n"),
+ );
+
+ await runPack([entry, "--outdir", outdir]);
+
+ expect(JSON.parse(readEmbeddedMetadata(path.join(outdir,
"bundle.min.mjs")))).toHaveProperty(
+ "task_handlers.sales_dag.tasks",
+ ["extract_north", "load_north"],
+ );
+ });
+
it("packs only the Dags the served bundle holds", async () => {
outdir = mkdtempSync(path.join(tmpdir(), "ts-pack-"));
const entry = path.join(outdir, "forgotten-entry.ts");
diff --git a/ts-sdk/tests/public-api.test.ts b/ts-sdk/tests/public-api.test.ts
index b3a55e54f58..3ffba6332f9 100644
--- a/ts-sdk/tests/public-api.test.ts
+++ b/ts-sdk/tests/public-api.test.ts
@@ -27,7 +27,6 @@ import type {
SetXComOpts,
TaskClient,
Registerable,
- PositionalInputs,
TaskContext,
TaskFactory,
TaskFunction,
@@ -291,33 +290,33 @@ describe("public API", () => {
// a wider one is, and not the other way round.
expectTypeOf<TaskRef<boolean>>().toMatchTypeOf<TaskRef>();
expectTypeOf<TaskRef>().not.toMatchTypeOf<TaskRef<boolean>>();
- // Wiring moved to the factory call and the spec is the trailing argument
- // itself, alongside the argument names the packer fills in.
- expectTypeOf<TaskOptions>().toEqualTypeOf<
- TaskSpec & { readonly argBindings?: readonly string[] }
- >();
- // Each named argument takes any upstream reference or a literal of its
own type.
+ // Wiring moved to the factory call, and the spec is the trailing argument
+ // itself.
+ expectTypeOf<TaskOptions>().toEqualTypeOf<TaskSpec>();
+ // Each input takes any upstream reference or a literal of its own type.
expectTypeOf<TaskInputs<{ rows: number }>>().toEqualTypeOf<{ rows: TaskRef
| number }>();
- // A positional argument takes a literal or a reference of the argument's
own
- // type, which is what tells a one-argument positional call from a named
one.
expectTypeOf<TaskInput<number>>().toEqualTypeOf<TaskRef<number> |
number>();
- expectTypeOf<PositionalInputs<[number, string]>>().toEqualTypeOf<
- [TaskRef<number> | number, TaskRef<string> | string]
- >();
- // A handler with no arguments is called with none; one with several is
- // called with a value per argument, in order.
- expectTypeOf<TaskFactory<[]>>().toEqualTypeOf<() => TaskRef>();
- expectTypeOf<TaskFactory<[], boolean>>().toEqualTypeOf<() =>
TaskRef<boolean>>();
- expectTypeOf<TaskFactory<[number, string]>>().toEqualTypeOf<
- (...inputs: [TaskRef<number> | number, TaskRef<string> | string]) =>
TaskRef
+ // A handler with no arguments is called with none; one that takes an
object
+ // is called with its inputs named.
+ expectTypeOf<TaskFactory>().toEqualTypeOf<() => TaskRef>();
+ expectTypeOf<TaskFactory<void, boolean>>().toEqualTypeOf<() =>
TaskRef<boolean>>();
+ expectTypeOf<TaskFactory<{ rows: number }>>().toEqualTypeOf<
+ (inputs: { rows: TaskRef | number }) => TaskRef
>();
expectTypeOf<ConstructorParameters<typeof Dag>>().toEqualTypeOf<[string,
DagSpec?]>();
- expectTypeOf<Dag["task"]>().toEqualTypeOf<
- <TParams extends readonly unknown[] = [], TReturn = unknown>(
- taskId: string,
- handler: (...args: TParams) => TReturn | Promise<TReturn>,
- options?: TaskOptions,
- ) => TaskFactory<TParams, TReturn>
+ // Overloaded, so the two forms are checked by calling them rather than by
+ // matching one signature.
+ const overloadedDag = new Dag("overload_dag");
+ const withId = overloadedDag.task("extract", async (_: { rows: number })
=> undefined);
+ const withoutId = overloadedDag.task(async function transform(_: { rows:
number }) {});
+ expectTypeOf(withId).toEqualTypeOf<TaskFactory<{ rows: number },
undefined>>();
+ expectTypeOf(withoutId).toEqualTypeOf<TaskFactory<{ rows: number },
void>>();
+ // A handler's return type reaches the reference its factory hands back.
+ expectTypeOf(overloadedDag.task("decide", async () => true)).toEqualTypeOf<
+ TaskFactory<void, boolean>
+ >();
+ expectTypeOf(overloadedDag.task("plain", async () =>
undefined)).toEqualTypeOf<
+ TaskFactory<void, undefined>
>();
expectTypeOf<Dag["taskIds"]>().toEqualTypeOf<readonly string[]>();
// Both specs are all-optional, so `{}` stays assignable and a field the
@@ -333,9 +332,10 @@ describe("public API", () => {
// still a plain optional: `undefined` already means "unset".
expectTypeOf<TaskSpec["retryDelay"]>().toEqualTypeOf<number | undefined>();
expectTypeOf<TaskSpec["doXcomPush"]>().toEqualTypeOf<boolean |
undefined>();
- // Identity is positional, so it is not restated in either spec.
+ // Dag identity is positional, so the spec does not restate it. A task's is
+ // not: it is also what an anonymous handler sets instead of being named.
expectTypeOf<DagSpec>().not.toHaveProperty("dagId");
- expectTypeOf<TaskSpec>().not.toHaveProperty("taskId");
+ expectTypeOf<TaskSpec["taskId"]>().toEqualTypeOf<string | undefined>();
});
it("uses idiomatic TypeScript names for public client types", () => {
@@ -423,7 +423,7 @@ describe("public API", () => {
const rejectsPositionalMisuse = () => {
// @ts-expect-error dagId is positional, not an options object.
new Dag({ dagId: "example" });
- // @ts-expect-error a task handler is required.
+ // @ts-expect-error a task handler is required alongside an explicit id.
new Dag("example").task("extract");
const dag = new Dag("example");
const extract = dag.task("extract", async () => undefined);
@@ -449,17 +449,22 @@ describe("public API", () => {
// @ts-expect-error a literal has to match its argument's type.
transform({ rows: "many" });
const totals = dag.task("totals", async (): Promise<{ rows: number }> =>
({ rows: 1 }));
- // A single object of named arguments is given by name or by position.
+ // Inputs are named.
transform({ rows: 1 });
+ transform({ rows: totals() });
+ // @ts-expect-error a bare reference is not an object of named inputs.
transform(totals());
- // @ts-expect-error a positional reference has to return the argument's
type.
- transform(extract());
- const pair = dag.task("pair", async (rows: number, region: string) =>
`${region}${rows}`);
+ const pair = dag.task(
+ "pair",
+ async ({ rows, region }: { rows: number; region: string }) =>
`${region}${rows}`,
+ );
+ pair({ rows: 1, region: "us" });
+ // @ts-expect-error inputs are named, not given in order.
pair(1, "us");
- // @ts-expect-error a positional argument cannot be skipped.
- pair(1);
- // @ts-expect-error each positional literal has to match its own
argument.
- pair("many", "us");
+ // @ts-expect-error an argument cannot be skipped.
+ pair({ rows: 1 });
+ // @ts-expect-error each literal has to match its own argument.
+ pair({ rows: "many", region: "us" });
const stamp = dag.task("stamp", async (_: { at: Date }) => undefined);
// @ts-expect-error a Date cannot survive the serialized Dag, so only a
reference will do.
stamp({ at: new Date() });
diff --git a/ts-sdk/tests/sdk/dag.test.ts b/ts-sdk/tests/sdk/dag.test.ts
index 098a2064b21..1bd493b602d 100644
--- a/ts-sdk/tests/sdk/dag.test.ts
+++ b/ts-sdk/tests/sdk/dag.test.ts
@@ -193,103 +193,45 @@ describe("Dag", () => {
);
});
- it("records a positional call in argument order", () => {
- const dag = new Dag("example_dag");
- const extract = dag.task("extract", async (): Promise<number> => 1);
- const transform = dag.task("transform", async (rows: number, region:
string) => `${region}`);
-
- const extracted = extract();
- transform(extracted, "us");
-
- expect(getDagTaskInputs(dag).get("transform")).toEqual({ arg0: extracted,
arg1: "us" });
-
expect(Object.keys(getDagTaskInputs(dag).get("transform")!)).toEqual(["arg0",
"arg1"]);
- });
-
- it("names positional arguments from argBindings", () => {
- const dag = new Dag("example_dag");
- const extract = dag.task("extract", async (): Promise<number> => 1);
- const transform = dag.task("transform", async (rows: number, region:
string) => `${region}`, {
- argBindings: ["rows", "region"],
- });
-
- const extracted = extract();
- transform(extracted, "us");
-
- expect(getDagTaskInputs(dag).get("transform")).toEqual({ rows: extracted,
region: "us" });
- });
-
- it("labels the arguments argBindings does not reach", () => {
+ it.each([
+ ["a bare value", 7],
+ ["an array of values", [1, "us"]],
+ ])("rejects %s where a task's inputs belong", (_label, inputs) => {
const dag = new Dag("example_dag");
- const transform = dag.task("transform", async (rows: number, region:
string) => `${region}`, {
- argBindings: ["rows"],
- });
-
- transform(1, "us");
+ const load = dag.task("load", async ({ total }: { total: number }) =>
total);
- expect(getDagTaskInputs(dag).get("transform")).toEqual({ rows: 1, arg1:
"us" });
+ expect(() => load(inputs as never)).toThrowError(
+ /takes one object naming its inputs: myTask\({ rows, region }\)/,
+ );
});
- it.each([
- [
- "not an array",
- { argBindings: 1 },
- /argBindings for Dag "d" task "t" must be an array of names/,
- ],
- ["not a string", { argBindings: [1] }, /holds 1; each name must be a
non-empty string/],
- ["empty", { argBindings: [""] }, /holds ""; each name must be a non-empty
string/],
- ["a number", { argBindings: ["0"] }, /holds "0"; each name must be a
non-empty string/],
- [
- "a duplicate",
- { argBindings: ["a", "a"] },
- /argBindings for Dag "d" task "t" names "a" twice/,
- ],
- ])("rejects argBindings that are %s", (_label, options, expected) => {
- const dag = new Dag("d");
-
- expect(() => dag.task("t", async (a: number) => a, options as
never)).toThrowError(expected);
- });
-
- it("reads a single argument that is not a map of names as one positional
input", () => {
+ it("rejects inputs given as more than one argument", () => {
const dag = new Dag("example_dag");
- const when = new Date();
- const transform = dag.task("transform", async (at: Date) => at);
-
- transform(when as never);
+ const transform = dag.task(
+ "transform",
+ async ({ rows, region }: { rows: number; region: string }) =>
`${region}${rows}`,
+ );
- expect(getDagTaskInputs(dag).get("transform")).toEqual({ arg0: when });
+ // A positional call from plain JavaScript; TypeScript rejects it.
+ expect(() => (transform as (...inputs: unknown[]) => unknown)({ rows: 1 },
"us")).toThrowError(
+ /was given 2 arguments; it takes one object naming its inputs/,
+ );
});
- it("reads a single reference as one positional input rather than a map of
names", () => {
+ it("rejects a bare reference where a task's inputs belong", () => {
const dag = new Dag("example_dag");
const extract = dag.task("extract", async (): Promise<{ rows: number }> =>
({ rows: 1 }));
- const load = dag.task("load", async (totals: { rows: number }) =>
totals.rows);
+ const load = dag.task("load", async ({ totals }: { totals: { rows: number
} }) => totals.rows);
- const extracted = extract();
- load(extracted);
-
- expect(getDagTaskInputs(dag).get("load")).toEqual({ arg0: extracted });
+ expect(() => load(extract() as never)).toThrowError(/takes one object
naming its inputs/);
});
- it("spreads a positional task's bound arguments back into its argument
list", async () => {
+ it("rejects a handler that declares more than one parameter", () => {
const dag = new Dag("example_dag");
- const seen: unknown[] = [];
- const transform = dag.task(
- "transform",
- async (rows: number, region: string) => {
- seen.push(rows, region);
- },
- { argBindings: ["rows", "region"] },
- );
-
- transform(1, "us");
- // The order the runtime hands the bound arguments over in, which is the
- // order they were recorded.
- await new Bundle(dag).getTaskHandler("example_dag", "transform")!({
- rows: 1,
- region: "us",
- } as never);
- expect(seen).toEqual([1, "us"]);
+ expect(() =>
+ dag.task("transform", ((rows: number, region: string) =>
`${region}${rows}`) as never),
+ ).toThrowError(/declares 2 parameters; a handler takes one object of named
arguments/);
});
it("calls a task declaring one object of named arguments with that object",
async () => {
@@ -447,7 +389,6 @@ describe("Dag", () => {
it.each([
["a misspelling", "retry"],
["the raw schema key", "retry_delay"],
- ["the positional task_id", "taskId"],
])("rejects %s in the task spec", (_label, key) => {
const dag = new Dag("example_dag");
expect(() =>
@@ -480,6 +421,96 @@ describe("Dag", () => {
expect(dag.taskIds).toEqual([]);
});
+ describe("an omitted task id", () => {
+ it("takes the handler's function name", () => {
+ const dag = new Dag("named_dag");
+ dag.task(async function extract() {})();
+
+ expect(dag.taskIds).toEqual(["extract"]);
+ });
+
+ it("takes the name of a handler declared elsewhere", () => {
+ async function transform() {}
+ const dag = new Dag("named_dag");
+ dag.task(transform)();
+
+ expect(dag.taskIds).toEqual(["transform"]);
+ });
+
+ it("still accepts a spec as the second argument", () => {
+ const dag = new Dag("named_dag");
+ dag.task(async function extract() {}, { retries: 2 })();
+
+ expect(getDagTaskRecords(dag).get("extract")?.spec).toEqual({ retries: 2
});
+ });
+
+ it("is taken from the spec when one names the task", () => {
+ const dag = new Dag("specced_id_dag");
+ dag.task(async function extract() {}, { taskId: "extract_rows" })();
+
+ expect(dag.taskIds).toEqual(["extract_rows"]);
+ });
+
+ it("prefers the positional id over the handler name", () => {
+ const dag = new Dag("positional_dag");
+ dag.task("extract_rows", async function extract() {})();
+
+ expect(dag.taskIds).toEqual(["extract_rows"]);
+ });
+
+ it("rejects a positional id and a spec id together", () => {
+ // The positional one used to win and the spec's was dropped in silence.
+ const dag = new Dag("two_ids_dag");
+
+ expect(() =>
+ dag.task("extract_rows", async function extract() {}, { taskId:
"from_spec" }),
+ ).toThrowError(
+ /Task "extract_rows" of Dag "two_ids_dag" also carries taskId
"from_spec" in its spec/,
+ );
+ expect(dag.taskIds).toEqual([]);
+ });
+
+ it("fails for an anonymous handler, naming the two ways to give it an id",
() => {
+ const dag = new Dag("anonymous_dag");
+
+ expect(() => dag.task(async () => 42)).toThrowError(
+ /A task of Dag "anonymous_dag" has no id: its handler has no name/,
+ );
+ expect(dag.taskIds).toEqual([]);
+ });
+
+ it("points at the bundler when a handler's name was minified away", () => {
+ const dag = new Dag("minified_dag");
+ // What a bundle built by something other than airflow-ts-pack can hold:
+ // the function is real, but the bundler took its name.
+ const minified = Object.defineProperty(async () => 42, "name", { value:
"" });
+
+ expect(() => dag.task(minified)).toThrowError(
+ /A bundler that drops function names also lands here/,
+ );
+ });
+
+ it("rejects a handler name Airflow would not accept as an id", () => {
+ const dag = new Dag("bound_dag");
+ async function extract({ rows }: { rows: number }) {
+ return rows;
+ }
+
+ // `bind` names the result "bound extract", which has a space in it.
+ expect(() => dag.task(extract.bind(null))).toThrowError(
+ /would take the id "bound extract" from its handler's name, which
Airflow does not accept/,
+ );
+ });
+
+ it("rejects a non-function where a handler belongs", () => {
+ const dag = new Dag("bad_handler_dag");
+
+ expect(() => dag.task("x", 42 as unknown as () =>
Promise<void>)).toThrowError(
+ /handler for Dag "bad_handler_dag" task "x" must be a function/,
+ );
+ });
+ });
+
it("exposes its task IDs in attachment order", () => {
const dag = new Dag("ordered_dag");
expect(dag.taskIds).toEqual([]);