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 () => 
{

Reply via email to