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
69 changes: 60 additions & 9 deletions ts-sdk/src/coordinator/runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -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(
Expand Down
108 changes: 103 additions & 5 deletions ts-sdk/tests/coordinator/integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand Down
Loading