jason810496 commented on code in PR #73442:
URL: https://github.com/apache/airflow/pull/73442#discussion_r4152654494


##########
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:
   Fixed: each Dag is now finalized inside its own try/catch.



##########
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:
   Fixed: `other_dag` now has a task, and the test asserts it is the one Dag 
serialized while `test_dag` is reported as an import error.



##########
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);
+  } catch (err) {

Review Comment:
   Fixed: `finalizeDag(dag)` now runs in the per-Dag try/catch alongside 
`serializeDag`, and the bundle-wide try/catch is gone. The orphan-task test now 
gives `other_dag` a task and asserts it still serializes.



##########
ts-sdk/src/sdk/bundle.ts:
##########
@@ -239,3 +239,21 @@ export function bundleDagTaskIds(bundle: Bundle): 
Map<string, string[]> {
   }
   return byDag;
 }
+
+/**
+ * Internal: the Dags this bundle declared in TypeScript, in registration 
order.
+ *
+ * What a parse request answers with, and what closes each of them: a Dag known
+ * only through task handlers is not one of them, because 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`.
+ *
+ * Supersedes {@link finalizeBundleDags} for a caller that also wants the Dags.
+ */
+export function listBundleNativeDags(bundle: Bundle): Dag[] {
+  const dags = [...dagsOf(bundle).values()];
+  for (const dag of dags) {
+    finalizeDag(dag);
+  }
+  return dags;

Review Comment:
   `listBundleNativeDags` is removed: with finalization moved into the per-Dag 
loop, `handleParse` reads the Dags with the existing `bundleDags`.



##########
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);
+  } catch (err) {
+    // Reading the bundle closes every native Dag in it, so a fault in one is
+    // reported against the file rather than leaving the request unanswered.
+    const detail = err instanceof Error ? err.message : String(err);
+    logs.error("Bundle could not be read for parsing", { fileloc, detail });
+    return {
+      type: "DagFileParsingResult",
+      fileloc,
+      serialized_dags: [],
+      import_errors: { [relativeFileloc]: detail },
+    };
+  }
+
+  for (const dag of dags) {
+    try {
+      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 = {
+  return {
     type: "DagFileParsingResult",
-    fileloc: request.file,
-    serialized_dags: [],
-  };
-  return response;
+    fileloc,
+    serialized_dags: serializedDags,
+    ...(failures.length > 0 && { import_errors: { [relativeFileloc]: 
failures.join("\n") } }),
+  } as RuntimeDagFileParsingResult;
 }

Review Comment:
   The result is typed up front and `import_errors` is set only when there are 
failures, so the cast is gone.
   



##########
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);
+  } catch (err) {
+    // Reading the bundle closes every native Dag in it, so a fault in one is

Review Comment:
   Agreed. Both uses of "closes" were in code this change removes, and the new 
wording says "finalized".



-- 
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]

Reply via email to