diff --git a/ts-sdk/src/coordinator/runtime.ts b/ts-sdk/src/coordinator/runtime.ts index bf9764f3e6618..fe2b027105a94 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 }[] = []; + // 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 d0b6d53f83fff..31d48d017194c 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> { + const result = await driveSupervisor(parseRequest); + return result.firstResponse!.body as Record; + } + + 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 }[]; + expect(dags.map((entry) => (entry.data.dag as Record).dag_id)).toEqual([ + "test_dag", + "other_dag", + ]); + expect(dags[0]!.data.__version).toBe(3); + expect((dags[0]!.data.dag as Record).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; - 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 }[]; + expect(dags.map((entry) => (entry.data.dag as Record).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; + 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 }[]; + expect(dags.map((entry) => (entry.data.dag as Record).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 () => {