guan404ming commented on code in PR #73442:
URL: https://github.com/apache/airflow/pull/73442#discussion_r4144776497
##########
ts-sdk/src/coordinator/runtime.ts:
##########
@@ -262,22 +265,77 @@ export function createRuntimeAbort(
};
}
+/**
+ * Answer a parse request with the Dags this bundle declared in TypeScript.
+ *
+ * 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 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[] = [];
+
+ let dags: Dag[];
+ try {
+ dags = listBundleNativeDags(bundle);
Review Comment:
One Dag failing finalize drops every valid Dag here. Finalizing per Dag
inside the loop would isolate it.
##########
ts-sdk/tests/coordinator/integration.test.ts:
##########
@@ -798,18 +798,112 @@ 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();
+ });
+
+ 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);
- const body = result.firstResponse!.body as Record<string, unknown>;
- expect(body.type).toBe("DagFileParsingResult");
- expect(body.serialized_dags).toEqual([]);
+ const body = await parse();
+
+ expect(body.serialized_dags).toEqual([]);
Review Comment:
Maybe give `other_dag` a task and assert it still serializes.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]