pierrejeambrun commented on code in PR #73442:
URL: https://github.com/apache/airflow/pull/73442#discussion_r4144832423
##########
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:
A few references to `closes`. I believe we should replace with `finalizes`
which is more explicite. "closing a dag" is ambiguous.
##########
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:
Nit: Per-Dag isolation is asymmetric: serialization failures are caught
per-Dag (line 320-324), but finalization failures are not —
listBundleNativeDags finalizes in a bare loop, throws on the first orphan-task
Dag, and the whole bundle's serialized_dags becomes []. Test case:
```ts
testDag.task("extract", async () => undefined)(); // OK
testDag.task("orphan", async () => undefined); // never called
otherDag.task("stage", async () => undefined)(); // OK on its own
```
Per the current code, otherDag never reaches serialization:
finalizeDag(testDag) throws first, listBundleNativeDags bails, the outer catch
returns empty. The comment at line 300-301 says "one broken Dag must not take
out the others a bundle serves" — but that's exactly what happens for
finalize-time failures. Fix is small: move finalizeDag(dag) into the per-Dag
try/catch alongside serializeDag,
##########
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 doc "Supersedes finalizeBundleDags" (bundle.ts:251) —
but finalizeBundleDags still has a live caller in coordinator/manifest.ts:41
(pack time). Not superseded, just a sibling for a different call site. Worth
softening the wording.
##########
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:
as RuntimeDagFileParsingResult cast at the return (line 336ish) — the
conditional-spread pattern makes TS lose the narrowing. Straightforward without
the cast:
```ts
const result: RuntimeDagFileParsingResult = {
type: "DagFileParsingResult",
fileloc,
serialized_dags: serializedDags,
};
if (failures.length > 0) {
result.import_errors = { [relativeFileloc]: failures.join("\n") };
}
return result;
```
--
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]