Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,34 @@ one that pushed ``null`` binds ``null``.
A Python ``int`` beyond the ±9007199254740991 a JavaScript number holds exactly is refused rather than
bound, so carry such a value across the language boundary as a string.

Explicit renames
~~~~~~~~~~~~~~~~

``withArgNames`` states a binding when folding cannot reach it, for a name the Python side never used:
a clearer word than the Dag chose, or a TypeScript reserved word like ``enum``.
The mapping comes first, the handler second:

.. code-block:: typescript

interface ReportArgs {
summary: Summary;
label: string; // Python calls this `run_label`
}

const report = withArgNames({ label: "run_label" }, async ({ summary, label }: ReportArgs) => {
// `label` is the call's `run_label`; `summary` folded as usual.
});

bundle.register(new TaskHandler("etl", "report", report));

An entry beats folding, and everything the map does not mention still folds,
so ``withArgNames`` should be rare in a real Dag.
A mapped name the call did not pass misses rather than falling back to folding.

The map's keys are checked against the handler's own parameter type, so ``{ labl: "run_label" }`` is a
compile error naming the right key. Its values are Python names, which ``tsc`` cannot see and does not
check.

.. note::

Being upstream is not the same as being passed. As with the other language SDKs, an XCom *dependency*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,9 @@
round-trips, and task logs reaching the log store.

``typescript_taskflow_example`` covers TaskFlow arguments, including an upstream output pulled before
the handler runs, and shares a ``build_message`` task ID with ``typescript_example`` so that dispatch
keying on the task ID alone would run the wrong handler.
the handler runs and a ``withArgNames`` rename on its ``report`` task, and shares a ``build_message``
task ID with ``typescript_example`` so that dispatch keying on the task ID alone would run the wrong
handler.
"""

from __future__ import annotations
Expand Down Expand Up @@ -153,7 +154,12 @@ def test_second_dag_from_the_same_bundle_succeeded(completed_taskflow_run: _Comp
f"expected the run to succeed; got {completed_taskflow_run.state!r}. "
f"task states: {completed_taskflow_run.ti_states}"
)
expected = {"make_totals": "success", "summarize": "success", "build_message": "success"}
expected = {
"make_totals": "success",
"summarize": "success",
"report": "success",
"build_message": "success",
}
for task_id, want in expected.items():
assert completed_taskflow_run.ti_states.get(task_id) == want, (
f"{task_id!r} expected {want!r}. all task states: {completed_taskflow_run.ti_states}"
Expand Down Expand Up @@ -190,6 +196,24 @@ def test_summarize_binds_its_call_arguments(completed_taskflow_run: _CompletedRu
assert completed_taskflow_run.xcom("summarize", key="summary_line") == "uk: 12 orders"


def test_report_binds_an_explicitly_renamed_argument(completed_taskflow_run: _CompletedRun):
"""``report(summary, "nightly")`` reaches a handler that renamed one argument.

Python names it ``run_label``; the handler destructures ``label``, a word
the ``@task.stub`` signature never uses, so folding could not connect the
two and the binding is stated with ``withArgNames``. The handler throws
unless ``label`` is exactly ``"nightly"``, so a rename that did not take
effect fails this task rather than returning a null.

``summary`` is not renamed: folding already covers it, which is the point
that keeps ``withArgNames`` rare.
"""
value = completed_taskflow_run.xcom("report")
assert value == {"label": "nightly", "regionCode": "uk", "healthy": True}, (
f"unexpected 'report' return_value: {value!r}"
)


def test_same_task_id_under_two_dags_runs_its_own_handler(
completed_run: _CompletedRun, completed_taskflow_run: _CompletedRun
):
Expand Down
24 changes: 24 additions & 0 deletions ts-sdk/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -241,6 +241,30 @@ const rows = await getClient().getXCom<number>({ key: "return_value", taskId: "e
A Python `int` beyond the ±9007199254740991 a JavaScript number holds exactly is refused rather than bound,
so carry such a value across the boundary as a string.

### Explicit renames

`withArgNames` states a binding folding cannot reach, for a name the Python side never used:
a clearer word than the Dag chose, or a TypeScript reserved word like `enum`. Mapping first, handler second:

```ts
interface ReportArgs {
summary: Summary;
label: string; // Python calls this `run_label`
}

const report = withArgNames({ label: "run_label" }, async ({ summary, label }: ReportArgs) => {
// `label` is the call's `run_label`; `summary` folded as usual.
});

bundle.register(new TaskHandler("etl", "report", report));
```

An entry beats folding, and everything the map does not mention still folds,
so `withArgNames` should be rare in a real Dag.
The map's keys are checked against the handler's own parameter type,
so `{ labl: "run_label" }` is a compile error naming the right key.
Its values are Python names, which `tsc` cannot see and does not check.

`Dag` is another interface, for a Dag declared natively in TypeScript, and is still a work in progress.

Airflow launches the bundled entrypoint with `--comm=host:port` and
Expand Down
3 changes: 2 additions & 1 deletion ts-sdk/api-docs/dag-authoring-api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,9 @@

/** @module Authoring */

export { Bundle, Dag, getClient, getContext, TaskHandler } from "../src/index.js";
export { Bundle, Dag, getClient, getContext, TaskHandler, withArgNames } from "../src/index.js";
export type {
ArgNameMap,
DagSpec,
Registerable,
TaskClient,
Expand Down
3 changes: 3 additions & 0 deletions ts-sdk/docs/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,9 @@ export async function transform({ regionCode, threshold }: TransformArgs) {
}
```

`withArgNames` states a binding folding cannot reach, for a name the Python side never used.
It should be rare, since folding covers ordinary spelling differences.

`Dag` is another interface, for a Dag declared in TypeScript rather than in Python, and is still a work in progress.

## Coordinators
Expand Down
15 changes: 13 additions & 2 deletions ts-sdk/example/dags/typescript_taskflow_example.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,9 @@

``summarize`` is called TaskFlow-style, and every argument its call passes reaches the TypeScript
handler by name, including ``make_totals``'s output, which the runtime pulls before the handler runs.
Its ``build_message`` stub shares a ``task_id`` with a task in ``typescript_example`` on purpose:
``report`` shows the one case folding cannot cover: its handler wants a name the Dag never used,
so it states that binding explicitly with ``withArgNames``.
The ``build_message`` stub shares a ``task_id`` with a task in ``typescript_example`` on purpose:
a handler binds the ``(dag_id, task_id)`` pair, so the two are different tasks.
See ``src/taskflow.ts``.
"""
Expand All @@ -46,6 +48,13 @@ def make_totals():
def summarize(totals: dict, region_code: str, currency: str, threshold: float, dry_run: bool = False): ...


# `run_label` is not a spelling difference. The handler wants to call it
# `label`, a word this signature never uses, which is what `withArgNames` is
# for; folding would never connect the two.
@task.stub(queue="typescript")
def report(summary: dict, run_label: str): ...


# Same task_id as `typescript_example.build_message`, on purpose.
@task.stub(queue="typescript")
def build_message(): ...
Expand All @@ -58,7 +67,9 @@ def build_message(): ...
tags=["typescript", "example", "taskflow"],
)
def typescript_taskflow_example():
summarize(make_totals(), "uk", "GBP", 280.0) >> build_message()
summary = summarize(make_totals(), "uk", "GBP", 280.0)
report(summary, "nightly")
summary >> build_message()


typescript_taskflow_example()
3 changes: 2 additions & 1 deletion ts-sdk/example/src/main.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@

import { Bundle, getClient, TaskHandler } from "apache-airflow-ts-sdk";

import { buildSummaryMessage, summarize } from "./taskflow.js";
import { buildSummaryMessage, report, summarize } from "./taskflow.js";

export async function buildMessage() {
const client = getClient();
Expand Down Expand Up @@ -63,6 +63,7 @@ bundle.register(
new TaskHandler("typescript_example", "build_message", buildMessage),
new TaskHandler("typescript_example", "read_connection", readConnection),
new TaskHandler("typescript_taskflow_example", "summarize", summarize),
new TaskHandler("typescript_taskflow_example", "report", report),
new TaskHandler("typescript_taskflow_example", "build_message", buildSummaryMessage),
);
await bundle.serve();
30 changes: 29 additions & 1 deletion ts-sdk/example/src/taskflow.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@
// `buildSummaryMessage` implements a task named `build_message`, exactly as the other Dag has, and
// the two share nothing else.

import { getClient, getContext } from "apache-airflow-ts-sdk";
import { getClient, getContext, withArgNames } from "apache-airflow-ts-sdk";

/** What `make_totals` returns on the Python side. */
export interface Totals {
Expand Down Expand Up @@ -88,6 +88,34 @@ export async function summarize({
};
}

/** Every argument the Dag's `report(...)` call binds, as the handler wants them. */
export interface ReportArgs {
summary: Summary;
/** The call's `run_label`, which folding cannot reach, so it is mapped below. */
label: string;
}

/**
* Renaming an argument the Python side named something else entirely.
*
* The mapping comes first, the handler second. `summary` is absent from the map
* because folding already reaches it.
*/
export const report = withArgNames(
{ label: "run_label" },
async ({ summary, label }: ReportArgs) => {
if (label !== "nightly") {
throw new Error(`expected run label "nightly" but got "${label}"`);
}

return {
label,
regionCode: summary.regionCode,
healthy: summary.passed,
};
},
);

export async function buildSummaryMessage() {
// Nothing was passed to this task, so its upstream's output is read explicitly.
const ctx = getContext();
Expand Down
18 changes: 14 additions & 4 deletions ts-sdk/src/coordinator/arg-binding.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,8 @@ export interface ArgBindingDeps {
/** The task's abort signal, so a terminated task stops mid-pull. */
readonly signal: AbortSignal;
readonly logs: LogChannel;
/** Renames the handler declared with `withArgNames`, which beat folding. */
readonly argNames: ReadonlyMap<string, string>;
}

/**
Expand Down Expand Up @@ -128,7 +130,7 @@ export async function resolveArgs(
// Literals need no request, so only a call that pulls races the abort signal.
const entries = pullsUpstream ? await abortable(resolveAll, deps.signal) : await resolveAll();

return { args: makeArgsProxy(names, byFold, new Map(entries), deps.logs), names };
return { args: makeArgsProxy(names, byFold, new Map(entries), deps), names };
}

/** Airflow omits `value` for a literal whose value is null. */
Expand Down Expand Up @@ -223,15 +225,21 @@ function abortError(signal: AbortSignal): Error {
* has no way to know which spelling a handler will use: it sees Python's names
* and nothing else. Folding on read means binding needs nothing declared on
* either side, and no guess about the TypeScript name is ever materialized.
* It is also what lets a `withArgNames` entry take precedence, decided per read.
*/
function makeArgsProxy(
names: readonly string[],
byFold: ReadonlyMap<string, string>,
values: ReadonlyMap<string, JsonValue>,
logs: LogChannel,
deps: ArgBindingDeps,
): object {
const resolve = (property: string): string | undefined =>
values.has(property) ? property : byFold.get(foldArgName(property));
const { argNames, logs } = deps;
const resolve = (property: string): string | undefined => {
// An explicit rename wins and never falls back to folding, so a wrong entry misses.
const renamed = argNames.get(property);
if (renamed !== undefined) return values.has(renamed) ? renamed : undefined;
return values.has(property) ? property : byFold.get(foldArgName(property));
};

// A null prototype so a read never reaches Object.prototype: a Python
// argument named `constructor` or `toString` must bind like any other, and a
Expand All @@ -249,6 +257,8 @@ function makeArgsProxy(
// one from a typo.
logs.warning("Task argument not bound by this task's call", {
requested: property,
// Tells a wrong `withArgNames` entry apart from an argument the call never passed.
renamed_to: argNames.get(property) ?? null,
bound: [...names],
});
return undefined;
Expand Down
2 changes: 2 additions & 0 deletions ts-sdk/src/coordinator/runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ import {
type RuntimeTaskState,
type StartupDetails,
} from "./protocol.js";
import { getArgNames } from "../sdk/arg-names.js";
import { listBundleTasks, type Bundle } from "../sdk/bundle.js";
import { runInTaskScope, type TaskContext } from "../sdk/task.js";
import type { JsonValue } from "../sdk/client-types.js";
Expand Down Expand Up @@ -315,6 +316,7 @@ async function handleTask(
client,
signal: ctx.signal,
logs,
argNames: getArgNames(handler),
});
} catch (err) {
// Before the handler ran, so nothing it might have written is at stake.
Expand Down
2 changes: 2 additions & 0 deletions ts-sdk/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,9 +20,11 @@
export { Dag } from "./sdk/dag.js";
export { Bundle } from "./sdk/bundle.js";
export { TaskHandler } from "./sdk/task-handler.js";
export { withArgNames } from "./sdk/arg-names.js";
export { getClient, getContext } from "./sdk/task.js";
export { ConnectionNotFoundError, VariableNotFoundError } from "./sdk/client.js";
export { SUPERVISOR_API_VERSION } from "./coordinator/index.js";
export type { ArgNameMap } from "./sdk/arg-names.js";
export type { Registerable } from "./sdk/bundle.js";
export type { DagSpec, TaskInputs, TaskOptions, TaskRef, TaskSpec } from "./sdk/dag.js";
export type { TaskClient } from "./sdk/client.js";
Expand Down
Loading