Skip to content

Commit f86a495

Browse files
committed
TS SDK: default a task id to the handler's name
A native Dag handler takes one object of named arguments, and a call names each input. withArgList(...) supplies the same inputs in the order the handler destructures them, zipped with the keys read off its argument pattern, so nothing labels an argument arg0. The task id still defaults to the handler's function name, which is why airflow-ts-pack passes esbuild's keepNames; argument names are property names, which minification leaves alone, so the pack step minifies identifiers again.
1 parent ea1f472 commit f86a495

10 files changed

Lines changed: 700 additions & 263 deletions

File tree

‎airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst‎

Lines changed: 35 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -247,8 +247,8 @@ A ``Dag`` is declared on this side rather than in Python: its schedule, its task
247247
the edges between them are all written in TypeScript. The surface is still growing, so a Dag declared
248248
this way is not served to Airflow yet.
249249

250-
``dag.task(taskId, handler)`` returns a *factory*. Calling it places the task in the Dag and supplies the
251-
handler's arguments, so the call graph is the task graph:
250+
``dag.task(taskId, handler)`` returns a *factory*. A handler takes one object of named arguments, and
251+
calling the factory names each input, so the call graph is the task graph:
252252

253253
.. code-block:: typescript
254254
@@ -257,26 +257,50 @@ handler's arguments, so the call graph is the task graph:
257257
const dag = new Dag("ts_etl");
258258
259259
const extract = dag.task("extract", async (): Promise<number> => 42);
260-
const transform = dag.task("transform", async (rows: number, region: string) => rows * 2);
261-
const load = dag.task("load", async (total: number) => {});
260+
const transform = dag.task(
261+
"transform",
262+
async ({ rows, region }: { rows: number; region: string }) => rows * 2,
263+
);
264+
const load = dag.task("load", async ({ total }: { total: number }) => {});
262265
263-
load(transform(extract(), "us"));
266+
const extracted = extract();
267+
const total = transform({ rows: extracted, region: "us" });
268+
load({ total });
264269
265-
Arguments are passed in the order the handler declares them. A handler that declares a single object of
266-
named arguments can also be called with that object, which names each input instead of ordering it:
270+
Naming the inputs is how a task is called. A handler that takes no arguments is called with none, and
271+
a single argument is named like any other, ``load({ total })``.
272+
273+
``withArgList`` supplies the same inputs in order, for a call that reads better that way:
267274

268275
.. code-block:: typescript
269276
270-
const store = dag.task("store", async ({ total }: { total: number }) => {});
277+
import { withArgList } from "apache-airflow-ts-sdk";
278+
279+
transform(withArgList(extracted, "us"));
271280
272-
store({ total: extract() });
281+
Each value binds to the argument in that position, and the order is the one the handler destructures
282+
its argument in, so the handler has to take a plain object pattern. A named call is the one the
283+
compiler checks in full: it reports an argument left out, a misspelled one, and a literal of the
284+
wrong type.
273285

274286
Each argument takes either an upstream reference or a literal JSON value. A reference has to be the
275287
argument itself: one buried inside an array or an object is a literal, and draws no edge.
276288

277289
Every task has to be called exactly once. An uncalled task fails when the Dag is read, so none can be
278290
left out of the graph by accident.
279291

292+
The task id may be omitted, in which case it is the handler's function name:
293+
294+
.. code-block:: typescript
295+
296+
const extract = dag.task(async function extract(): Promise<number> {
297+
return 42;
298+
});
299+
300+
``airflow-ts-pack`` keeps function names intact, so bundling cannot rename a task. A handler with no
301+
name of its own, such as an arrow function passed inline, has nothing to take an id from and needs
302+
one: either positionally or as ``taskId`` in its spec. Give it in one place only, not both.
303+
280304
``new Dag`` and ``dag.task`` both take a trailing spec of Airflow options:
281305
``{ schedule: "@daily", tags: ["etl"] }`` for the Dag, ``{ retries: 2, retryDelay: 30 }`` for a task.
282306

@@ -398,9 +422,8 @@ layout header. The layout records the byte ranges and SHA-256 digests of the man
398422
so there is one file to deploy, with no separate manifest or ``node_modules``.
399423

400424
The code is minified because an integrity digest is only worth taking over an artifact nobody is expected to
401-
read or edit in place. The ``/*! */`` license banners of bundled dependencies are kept. Nothing is identified by
402-
a function name, so minified names are safe: a Dag and a task are named by the string ids their registration
403-
states, and a handler is dispatched by reference.
425+
read or edit in place. Function names are kept through minification, since a task id defaults to its
426+
handler's name. The ``/*! */`` license banners of bundled dependencies are kept.
404427

405428
Because the shipped code is not the code anyone wrote, the packer also embeds the entry module verbatim in a
406429
``/*# airflowSource ... #*/`` block comment, verified by its own digest, so Airflow has something readable to

‎ts-sdk/adr/0002-native-dag-interface.md‎

Lines changed: 27 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -27,16 +27,15 @@ Proposed. Revised after the review on #72047.
2727

2828
1. **`dag.task(handler)` returns a factory, and the task id is optional.** With no id the task takes
2929
the handler's function name (`dag.task(extract)` → task `"extract"`); `dag.task(taskId, handler)`
30-
sets it explicitly, which an anonymous handler must do. Calling the factory both places the task in
31-
the Dag and supplies its arguments, in the order the handler declares them
32-
(`load(transform(extract(), "us"))`). A handler that declares a single object of named arguments can
33-
also be called with that object — the shape Python TaskFlow uses for
34-
`load(transformed=transform(...))`.
35-
2. **The call graph is the task graph.** `tsc` checks every wired key against the handler's own
36-
parameter type, and a `TaskRef` exists only once its producing call has returned, so a cycle
37-
through arguments is unrepresentable rather than rejected by a validator. A reference passed by
38-
position is checked against the argument's own type, which is what tells the two call shapes apart
39-
when a handler declares a single argument.
30+
sets it explicitly, which an anonymous handler must do. A handler takes one object of named
31+
arguments, and calling the factory both places the task in the Dag and names each of its inputs
32+
(`load({ total: transform({ rows: extract(), region: "us" }) })`) — the shape Python TaskFlow uses
33+
for `load(transformed=transform(...))`. `withArgList(...)` supplies the same inputs in the order
34+
the handler destructures them, for a call that reads better that way.
35+
2. **The call graph is the task graph.** `tsc` checks every named input against the handler's own
36+
argument type, and a `TaskRef` exists only once its producing call has returned, so a cycle
37+
through arguments is unrepresentable rather than rejected by a validator. A listed call is checked
38+
by value type, and which argument each value supplies is the position it was given in.
4039
3. **Every task is called exactly once.** An uncalled task fails when the Dag is read, so none can be
4140
silently left out of the graph.
4241
4. **`before` and `after` draw order-only edges** — the TypeScript pair for `>>` and `<<`, both
@@ -149,20 +148,20 @@ convention.
149148
by design. Native declaration is what fills them, generated from the serialized-Dag JSON schema the
150149
way `src/generated/supervisor.ts` is. This ADR does not choose those fields; it fixes where an
151150
author writes them.
152-
- `TaskOptions` carries the spec and the handler's positional argument names, which the packer fills in
153-
from the parameter list so the Dag names each argument as its handler does. With wiring moved to the
154-
factory call, `inputs` is no longer an option.
151+
- `TaskOptions` carries the task's spec and nothing else: the names on the wire are the keys of the
152+
call itself. With wiring moved to the factory call, `inputs` is no longer an option.
155153
- `TaskHandlerArgs` is removed from the public API, `DagRegistry` becomes `Bundle`, and
156154
`serveDags(registry)` becomes `bundle.serve()`, which breaks
157155
0.1.0-beta1 authors; see [ADR-0001](0001-mixed-lang-dag-interface.md) for the shipped call sites
158156
that change.
159157

160158
## Alternatives
161159

162-
- **Named-only wiring**, rejected in the review on #73435: naming every input reads well at twenty
163-
tasks but forces an object around a single argument, and positional calls are what TypeScript
164-
authors write. Both are offered, and the handler's own parameter list decides which one a task can
165-
use.
160+
- **Positional handlers**, `async (rows: number, region: string) => ...`, offered first and then
161+
dropped: a positional parameter list has no names on the wire unless the SDK reads them out of the
162+
handler's source, and a single object of named arguments is what a TypeScript library takes
163+
anyway. The handler is object-only, and `withArgList(...)` keeps the positional call site for the
164+
cases that read better in order.
166165
- **Injected `ctx`/`client` arguments**, mimicking the Python signature. Rejected, per the above and
167166
because feeling native to TypeScript matters more than matching Python's parameter list.
168167

@@ -181,17 +180,19 @@ convention.
181180
- **The spec argument already has its slot.** `dag.task(taskId, handler, options)` reads `{ spec = {} }`
182181
and runs `validateEmptySpec` on it (`ts-sdk/src/sdk/dag.ts`), so task fields land on a path that
183182
exists rather than a new one.
184-
- **A positional argument binds by order, and its name is a label.** The serialized Dag names each
185-
argument, so the packer reads the names from the handler's parameter list; `arg0`, `arg1` and so on
186-
stand in for a name it cannot see, without changing which value reaches which argument.
183+
- **A listed call binds by order, and the names still come from the handler.** The serialized Dag
184+
names every argument, so `withArgList(...)` is zipped with the keys the handler destructures, read
185+
off its argument pattern at the call. A pattern that cannot be read that way, such as one with a
186+
default or a rest element, is refused rather than guessed at, and its task is called by name.
187187
- **A `TaskRef` is inert** — a handle for wiring, not a promise. Nothing in a Dag file executes a task
188188
body.
189-
- **A defaulted task id is resolved at pack time, not read at runtime.** esbuild renames function
190-
identifiers, so `handler.name` in a packed bundle is the minified name, not the author's. The pack
191-
step (`ts-sdk/src/cli/pack.ts`) therefore reads an omitted id from the handler's declared name in
192-
source and writes it into the registration, rather than depending on `handler.name` or enabling
193-
esbuild's `keepNames` across the whole bundle. A handler with no source name leaves nothing to
194-
read, which is why an anonymous handler must state its id.
189+
- **A defaulted task id is read off the handler itself.** The pack step (`ts-sdk/src/cli/pack.ts`)
190+
minifies and passes esbuild's `keepNames`, so `handler.name` is the author's in a packed bundle as
191+
much as in one run from source; an argument name is a property name, which minification leaves
192+
alone. Rewriting the call at pack time was tried first and dropped: it needed a TypeScript parser
193+
in the packer to tell a real `.task(` from one inside a string or a comment, and it could not see a
194+
handler declared in another module. A handler with no name leaves nothing to read, which is why an
195+
anonymous handler must state its id.
195196
- **`withArgNames` and the name folding behind it** ([ADR-0001](0001-mixed-lang-dag-interface.md))
196197
exist for the mixed-language case and are never needed here: both ends of every name are
197198
TypeScript, so `tsc` checks the wiring end to end and there is no foreign name to reconcile.

‎ts-sdk/api-docs/dag-authoring-api.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,17 @@
1919

2020
/** @module Authoring */
2121

22-
export { Bundle, Dag, getClient, getContext, TaskHandler, withArgNames } from "../src/index.js";
22+
export {
23+
Bundle,
24+
Dag,
25+
getClient,
26+
getContext,
27+
TaskHandler,
28+
withArgList,
29+
withArgNames,
30+
} from "../src/index.js";
2331
export type {
32+
ArgList,
2433
ArgNameMap,
2534
DagSpec,
2635
Registerable,

‎ts-sdk/src/cli/pack.ts‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -215,8 +215,12 @@ export async function runPack(argv: readonly string[]): Promise<void> {
215215
platform: "node",
216216
format: "esm",
217217
target: "node22",
218-
// A digest is only worth taking over an artifact nobody reads or edits in place.
219218
minify: true,
219+
// A task id comes from the handler's name, and bundling renames a symbol
220+
// that two modules both declare, so the name the SDK reads is pinned to
221+
// the one the author wrote. Argument names are property names, which
222+
// minification leaves alone.
223+
keepNames: true,
220224
// The manifest is read by running the staged bundle, so the metadata describes what ships.
221225
outfile: stagingPath,
222226
});

‎ts-sdk/src/index.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,15 +20,16 @@
2020
export { Dag } from "./sdk/dag.js";
2121
export { Bundle } from "./sdk/bundle.js";
2222
export { TaskHandler } from "./sdk/task-handler.js";
23+
export { withArgList } from "./sdk/arg-list.js";
2324
export { withArgNames } from "./sdk/arg-names.js";
2425
export { getClient, getContext } from "./sdk/task.js";
2526
export { ConnectionNotFoundError, VariableNotFoundError } from "./sdk/client.js";
2627
export { SUPERVISOR_API_VERSION } from "./coordinator/index.js";
28+
export type { ArgList } from "./sdk/arg-list.js";
2729
export type { ArgNameMap } from "./sdk/arg-names.js";
2830
export type { Registerable } from "./sdk/bundle.js";
2931
export type {
3032
DagSpec,
31-
PositionalInputs,
3233
TaskFactory,
3334
TaskInput,
3435
TaskInputs,

‎ts-sdk/src/sdk/arg-list.ts‎

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
1+
/*!
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
// Giving a task's inputs in order instead of by name.
21+
22+
import { brand, hasBrand } from "./brand.js";
23+
24+
// Not a declared field, so the values stay out of the public type: an arg list
25+
// is a carrier the Dag reads, never something an author takes apart. A global
26+
// symbol, as the brands are, so two resolved copies agree on the key.
27+
const VALUES = Symbol.for("airflow.ts-sdk.arg-list-values");
28+
29+
// Carries the value types without carrying a value, the way TaskRef carries
30+
// its return type.
31+
declare const LISTED_VALUES: unique symbol;
32+
33+
/**
34+
* A task's inputs in the order its handler destructures them, as
35+
* {@link withArgList} builds.
36+
*/
37+
export interface ArgList<TValues extends readonly unknown[] = readonly unknown[]> {
38+
/** @internal Never set; see {@link LISTED_VALUES}. */
39+
readonly [LISTED_VALUES]: TValues;
40+
}
41+
42+
/**
43+
* Give a task's inputs in order instead of naming them.
44+
*
45+
* Naming the inputs is the usual way to call a task, and reads best past a
46+
* couple of arguments:
47+
*
48+
* ```ts
49+
* transform({ rows: extracted, region: "us" });
50+
* ```
51+
*
52+
* `withArgList` supplies the same inputs in order, for a call that reads
53+
* better that way. Each value binds to the argument in that position:
54+
*
55+
* ```ts
56+
* transform(withArgList(extracted, "us"));
57+
* ```
58+
*
59+
* The order is the one the handler destructures its argument in, so the
60+
* handler has to take a plain object pattern — `async ({ rows, region }) =>
61+
* ...`. A handler written any other way has no order to read, and its task is
62+
* called by name.
63+
*/
64+
export function withArgList<const TValues extends readonly unknown[]>(
65+
...values: TValues
66+
): ArgList<TValues> {
67+
const list = {};
68+
Object.defineProperty(list, VALUES, { value: Object.freeze([...values]) });
69+
brand(list, "ArgList");
70+
return Object.freeze(list) as ArgList<TValues>;
71+
}
72+
73+
/** Internal: whether `value` is an arg list built by any copy of this package. */
74+
export function isArgList(value: unknown): value is ArgList {
75+
return hasBrand(value, "ArgList");
76+
}
77+
78+
/** Internal: the values an arg list carries, in the order they were given. */
79+
export function argListValues(list: ArgList): readonly unknown[] {
80+
const values: unknown = (list as unknown as Record<symbol, unknown>)[VALUES];
81+
return Array.isArray(values) ? (values as readonly unknown[]) : [];
82+
}

0 commit comments

Comments
 (0)