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 5e476d7cf89 TS SDK: answer the Dag-parsing request from bundle.serve
(#73442)
5e476d7cf89 is described below
commit 5e476d7cf89ad6c6d435b0e461ab8d491159963a
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Fri Oct 2 15:31:51 2026 +0800
TS SDK: answer the Dag-parsing request from bundle.serve (#73442)
---
ts-sdk/src/coordinator/runtime.ts | 69 ++++++++++++++---
ts-sdk/tests/coordinator/integration.test.ts | 108 +++++++++++++++++++++++++--
2 files changed, 163 insertions(+), 14 deletions(-)
diff --git a/ts-sdk/src/coordinator/runtime.ts
b/ts-sdk/src/coordinator/runtime.ts
index bf9764f3e66..fe2b027105a 100644
--- a/ts-sdk/src/coordinator/runtime.ts
+++ b/ts-sdk/src/coordinator/runtime.ts
@@ -55,7 +55,10 @@ import {
type StartupDetails,
} from "./protocol.js";
import { getArgNames } from "../sdk/arg-names.js";
-import { bundleDagTaskIds, type Bundle } from "../sdk/bundle.js";
+import { bundleDags, bundleDagTaskIds, type Bundle } from "../sdk/bundle.js";
+import { finalizeDag } from "../sdk/dag.js";
+import { SERIALIZATION_VERSION } from "../generated/dag-schema-fields.js";
+import { computeRelativeFileloc, serializeDag } from "./serde.js";
import { runInTaskScope, type TaskContext } from "../sdk/task.js";
import type { JsonValue } from "../sdk/client-types.js";
@@ -262,22 +265,70 @@ export function createRuntimeAbort(
};
}
+/**
+ * Answer a parse request with the Dags this bundle declared in TypeScript.
+ *
+ * A Dag known only through task handlers is left out: its graph belongs to the
+ * Python Dag file that declares it, and serializing it here would register a
+ * second Dag with the same `dag_id` from a different `fileloc`.
+ *
+ * No handler body runs: a `TaskRef` is inert, so reading a Dag only walks what
+ * its module already built. Reading it is also what enforces that every task
+ * was called exactly once, which is why a Dag that is not fully laid out
+ * surfaces here.
+ *
+ * A Dag that cannot be finalized or serialized becomes an import error against
+ * this file, as a Python Dag file that raises does, rather than failing the
+ * whole parse: one broken Dag must not take out the others a bundle serves.
+ */
function handleParse(
request: { file: string; bundle_path: string },
bundle: Bundle,
logs: LogChannel,
): RuntimeDagFileParsingResult {
- // TypeScript-native Dag parsing is not yet supported.
- // Respond with an empty result so the Python-stub-Dag workflow works.
- logs.info("Parse-mode response (TS Dag parsing not yet supported)", {
- registered_tasks: Object.fromEntries(bundleDagTaskIds(bundle)),
+ const fileloc = request.file;
+ const relativeFileloc = computeRelativeFileloc(fileloc, request.bundle_path);
+ const serializedDags: { data: Record<string, unknown> }[] = [];
+ // Airflow keys an import error by the bundle-relative path and holds one row
+ // per file (`DagFileProcessorManager.update_import_errors`), so every
failure
+ // in this bundle is reported under that one key, naming its Dag in the
+ // message. An absolute path, or one with a Dag id appended, would give a row
+ // the UI cannot tie back to the file, and would leave the file itself
looking
+ // healthy while its Dags had vanished.
+ const failures: string[] = [];
+
+ const dags = [...bundleDags(bundle).values()];
+ for (const dag of dags) {
+ try {
+ finalizeDag(dag);
+ serializedDags.push({
+ data: {
+ __version: SERIALIZATION_VERSION,
+ dag: serializeDag(dag, fileloc, relativeFileloc),
+ },
+ });
+ } catch (err) {
+ const detail = err instanceof Error ? err.message : String(err);
+ logs.error("Dag could not be serialized", { dag_id: dag.dagId, detail });
+ failures.push(`Dag "${dag.dagId}": ${detail}`);
+ }
+ }
+
+ logs.info("Parse-mode response", {
+ fileloc,
+ dag_ids: dags.map((dag) => dag.dagId),
+ serialized: serializedDags.length,
+ import_errors: failures.length,
});
- const response: RuntimeDagFileParsingResult = {
+ const result: RuntimeDagFileParsingResult = {
type: "DagFileParsingResult",
- fileloc: request.file,
- serialized_dags: [],
+ fileloc,
+ serialized_dags: serializedDags,
};
- return response;
+ if (failures.length > 0) {
+ result.import_errors = { [relativeFileloc]: failures.join("\n") };
+ }
+ return result;
}
async function handleTask(
diff --git a/ts-sdk/tests/coordinator/integration.test.ts
b/ts-sdk/tests/coordinator/integration.test.ts
index d0b6d53f83f..31d48d01719 100644
--- a/ts-sdk/tests/coordinator/integration.test.ts
+++ b/ts-sdk/tests/coordinator/integration.test.ts
@@ -798,18 +798,116 @@ describe("coordinator runtime integration", () => {
expect(calledSecondDag).toBe(false);
});
- it("returns empty serialized_dags for DagFileParseRequest", async () => {
+ describe("DagFileParseRequest", () => {
const parseRequest = {
type: "DagFileParseRequest",
file: "/dags/test.mjs",
bundle_path: "/dags",
};
- const result = await driveSupervisor(parseRequest);
+ async function parse(): Promise<Record<string, unknown>> {
+ const result = await driveSupervisor(parseRequest);
+ return result.firstResponse!.body as Record<string, unknown>;
+ }
+
+ it("answers with the Dags the bundle declared in TypeScript", async () => {
+ testDag.task("extract", async () => undefined)();
+ otherDag.task("stage", async () => undefined)();
+
+ const body = await parse();
+
+ expect(body.type).toBe("DagFileParsingResult");
+ expect(body.fileloc).toBe("/dags/test.mjs");
+ const dags = body.serialized_dags as { data: Record<string, unknown> }[];
+ expect(dags.map((entry) => (entry.data.dag as Record<string,
unknown>).dag_id)).toEqual([
+ "test_dag",
+ "other_dag",
+ ]);
+ expect(dags[0]!.data.__version).toBe(3);
+ expect((dags[0]!.data.dag as Record<string,
unknown>).relative_fileloc).toBe("test.mjs");
+ expect(body.import_errors).toBeUndefined();
+ });
+
+ it("runs no handler body while parsing", async () => {
+ let ran = false;
+ testDag.task("extract", async () => {
+ ran = true;
+ })();
+ otherDag.task("stage", async () => undefined)();
+
+ await parse();
+
+ expect(ran).toBe(false);
+ });
+
+ it("answers with nothing when the bundle only binds handlers to Python
Dags", async () => {
+ // A Python Dag's graph belongs to the Python file that declares it, so a
+ // bundle of task handlers has no Dag of its own to serialize.
+ bundle = new Bundle(new TaskHandler("py_dag", "transform", async () =>
undefined));
+
+ const body = await parse();
+
+ expect(body.serialized_dags).toEqual([]);
+ expect(body.import_errors).toBeUndefined();
+ });
- const body = result.firstResponse!.body as Record<string, unknown>;
- expect(body.type).toBe("DagFileParsingResult");
- expect(body.serialized_dags).toEqual([]);
+ it("reports an uncalled task as an import error rather than failing the
parse", async () => {
+ testDag.task("extract", async () => undefined)();
+ testDag.task("orphan", async () => undefined);
+ otherDag.task("stage", async () => undefined)();
+
+ const body = await parse();
+
+ const dags = body.serialized_dags as { data: Record<string, unknown> }[];
+ expect(dags.map((entry) => (entry.data.dag as Record<string,
unknown>).dag_id)).toEqual([
+ "other_dag",
+ ]);
+ // Keyed by the bundle-relative path, which is how Airflow ties the row
+ // to the file it came from.
+ expect(body.import_errors).toEqual({
+ "test.mjs": expect.stringContaining('Task "orphan" of Dag "test_dag"
is never called'),
+ });
+ });
+
+ it("merges several failing Dags into the one row Airflow keeps per file",
async () => {
+ testDag.task("extract", async () => undefined)();
+ for (const dagId of ["broken_a", "broken_b"]) {
+ const broken = new Dag(dagId, { schedule: "" });
+ broken.task("t", async () => undefined)();
+ bundle.register(broken);
+ }
+
+ const body = await parse();
+
+ const errors = body.import_errors as Record<string, string>;
+ expect(Object.keys(errors)).toEqual(["test.mjs"]);
+ expect(errors["test.mjs"]).toContain('Dag "broken_a"');
+ expect(errors["test.mjs"]).toContain('Dag "broken_b"');
+ });
+
+ it("keeps the Dags it can serialize when one of them cannot be", async ()
=> {
+ testDag.task("extract", async () => undefined)();
+ // An empty schedule is rejected by the serializer, and only that Dag is
+ // lost: the rest of the bundle still parses.
+ const broken = new Dag("broken_dag", { schedule: "" });
+ broken.task("t", async () => undefined)();
+ bundle.register(broken);
+
+ const body = await parse();
+
+ const dags = body.serialized_dags as { data: Record<string, unknown> }[];
+ expect(dags.map((entry) => (entry.data.dag as Record<string,
unknown>).dag_id)).toEqual([
+ "test_dag",
+ "other_dag",
+ ]);
+ // One row per file, so the failing Dag is named in the message rather
+ // than appended to the key.
+ expect(body.import_errors).toEqual({
+ "test.mjs": expect.stringContaining(
+ 'Dag "broken_dag": schedule for Dag "broken_dag" is empty',
+ ),
+ });
+ });
});
it("auto-pushes return_value XCom when handler returns a value", async () =>
{