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 3e8318969e1 Add ADR for Lang SDK Dag and mixed-language Task
processing flow (#71929)
3e8318969e1 is described below
commit 3e8318969e1c82cee3e084d09efc8b7eff81863c
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Tue Sep 29 20:12:32 2026 +0800
Add ADR for Lang SDK Dag and mixed-language Task processing flow (#71929)
* Add ADR-0010 for mixed-language Dag processing flow
* Remove the usage of JavaDagImporter and the Dag-level is_mixed_language
flag
* Define mixed-language task handlers as a separate interface from Dags
A Lang-SDK artifact backing `@task.stub` tasks used to author a Dag under
the
same dag_id the Python file already owns, so two conflicting definitions
reached persistence and something downstream had to pick between them.
Removing
the Dag from the authoring interface answers that once, rather than leaving
every consumer of a serialized Dag to ask whether the Dag in front of it is
real.
Terms follow the Language SDK spec so the decision reads the same in Go,
Java,
and TypeScript, and the parse-side contract is grounded in the shipped
coordinator and supervisor-schema code rather than in a proposed shape.
* Separate the mixed-language decisions from their implementation mechanics
An ADR is read first to understand what was decided and why, and only later
to
build the thing. Exact message shapes, call sites, and the code that has to
change were interleaved with that reasoning, so neither audience could skim
its
half. The decisions now stand on their own and the mechanics wait in an
appendix.
* Split the Lang-SDK Dag processing decisions into one ADR each
The single document answered three questions with different audiences and
different review cycles: the wire protocol a runtime implements, how a
Lang-SDK
Dag source reaches the Dag processor, and how a Python stub is validated
against the handler behind it. Readers had to take all three to act on one,
and
a change to any one of them forced the other two back through review.
* Lead each Lang-SDK ADR with its diagrams and state the decisions plainly
The diagrams carry the decision faster than the prose around them did, so
they
belong where a reader lands rather than behind an appendix. The argument for
each choice is still worth keeping, but it is what someone reads second.
Coordinator lookup gains a named answer: for_bundle, beside the existing
for_queue, rather than an open question about how a registry finds the
coordinators serving a bundle.
* Route Lang-SDK parsing through one parse channel and one validation point
An earlier draft gave task-handler parsing its own message pair on the
theory that the runtime answered the coordinator. It does not: the
coordinator forwards bytes and decodes nothing, so the peer is whichever
process started the parse. Two reply unions would have meant two decoder
configurations and two relay paths for identical traffic.
Validation likewise had no home. It needs the parsed Dags and the
argument bindings that only exist after serialization, and the only place
holding both is the Dag-file parse itself — not an importer, which should
not have to know what a coordinator or a queue is.
Making a coordinator's importer optional restated in a return type what
the artifact source already decides, leaving two answers free to
disagree.
* Order the Lang-SDK Dag processing ADRs decisions-first
A reader meeting native Dag processing and mixed-language handlers for
the first time needs to know what each one decides before the message
shapes and subprocess classes that carry them mean anything. The parse
protocol is the mechanics both decisions share, so it reads better after
them than in front of them.
* Use dag and task as the TaskHandler argument names in ADR-0011
---
.../adr/lang-sdk/0010-native-dag-processing.md | 212 ++++++++++++++++++
.../lang-sdk/0011-mixed-language-dag-processing.md | 245 +++++++++++++++++++++
.../adr/lang-sdk/0012-lang-sdk-parse-protocol.md | 233 ++++++++++++++++++++
airflow-core/adr/lang-sdk/README.md | 3 +
4 files changed, 693 insertions(+)
diff --git a/airflow-core/adr/lang-sdk/0010-native-dag-processing.md
b/airflow-core/adr/lang-sdk/0010-native-dag-processing.md
new file mode 100644
index 00000000000..ff2acb00603
--- /dev/null
+++ b/airflow-core/adr/lang-sdk/0010-native-dag-processing.md
@@ -0,0 +1,212 @@
+<!--
+ Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+ Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+ -->
+
+# ADR-0010: Native Dag Processing — DagImporter Registration and Routing
+
+## Status
+
+Proposed
+
+## Context
+
+`AbstractDagImporter` (AIP-85, `task-sdk/src/airflow/sdk/importers/`) makes
the Dag processor aware of source formats beyond `.py`. It exposes
`can_handle`, `list_dag_definitions`,
+`import_definition` and `get_source_code`, and returns `DagImportResult.dags:
list[DAG]` — `airflow.sdk.DAG` objects.
+
+A Lang-SDK artifact can be one of those sources: a Dag authored entirely in Go
or Java, with `Dag(spec)` on the SDK side. Parsing it means launching a
runtime, which is what a
+coordinator does. This ADR settles where that importer comes from, which
coordinator instance backs it, and how the two avoid being configured twice.
+
+`JavaDagImporter` and `JavaCoordinator` are the examples throughout.
`ExecutableDagImporter` (Go) and `NodeDagImporter` (TypeScript) follow the same
shape.
+
+## Decision
+
+### The coordinator hands out its importer
+
+```
+BaseCoordinator.get_dag_importer() -> AbstractDagImporter
+ abstract — a coordinator names the importer for the artifacts it parses
+ JavaCoordinator.get_dag_importer() -> JavaDagImporter(coordinator=self)
+```
+
+The importer comes back already bound to the coordinator, so an operator never
configures which coordinator an importer uses. `[sdk] coordinators` stays the
one place a runtime is
+declared.
+
+The return is not optional. A coordinator that never contributes an importer
to any bundle already says so by the artifact source it was configured with —
`for_bundle` below
+answers that, and answering it a second time with a `None` return would let
the two disagree. What a coordinator can parse and where it gets registered are
separate questions.
+
+### Registration order
+
+```
+DagImporterRegistry.from_config(bundle_name)
+ ├── defaults PythonDagImporter, ZipImporter
+ ├── CoordinatorManager.for_bundle(bundle_name) (new)
+ │ └── register(coordinator.get_dag_importer())
+ ├── [dag_processor] dag_importer_configs (global, unchanged)
+ └── that bundle's own `importers` list (unchanged)
+```
+
+`dag_importer_configs` remains the door for importers with no runtime behind
them — a YAML importer, say. A Lang SDK never arrives that way.
+
+### `CoordinatorManager.for_bundle`
+
+```
+CoordinatorManager
+ ├── for_queue(queue) → one coordinator (shipped — task execution)
+ └── for_bundle(name) → the coordinators serving that bundle (new)
+```
+
+`for_queue` answers "who runs this task". `for_bundle` answers "who can parse
Dags in this bundle", which is what the registry tier above needs. It reads the
same `[sdk]
+coordinators` specs, selecting by the artifact source below.
+
+### One coordinator instance owns a DagBundle
+
+```
+EXPLICIT_ROOT jars_root / executables_root / bundles_root
+ → no DagBundle at all → not returned by for_bundle
+NAMED_BUNDLE dag_bundle_name
+ → that bundle, at the version current when work starts
+TASK_BUNDLE neither set
+ → the bundle the delegating Dag lives in, at the run's
version
+```
+
+Importers are keyed by file extension, one per extension, so two
`JavaCoordinator` instances on different JDKs would both claim `.jar`. Only
`NAMED_BUNDLE` names a bundle, which
+makes it the mode a deployment running two runtimes of the same language has
to use. `EXPLICIT_ROOT` has no DagBundle, so nothing scans its artifacts and it
cannot produce a Dag to
+persist; it serves mixed-language work only. Appendix A covers the three modes
in full.
+
+### The integration point is `import_definition`
+
+An importer runs inside the Dag-parsing child, so a native Lang-SDK Dag is
parsed by a process the importer itself starts. That process is
+`LangSDKDagFileProcessorProcess`, which differs from the one the manager
started only in the target callable it runs
+([ADR-0012](0012-lang-sdk-parse-protocol.md)).
+
+```
+DagFileProcessorProcess(analytics.jar) ← manager spawns,
as for any file
+ └── _parse_file_entrypoint → _parse_file
+ └── BundleDagBag → DagImporterRegistry.get_importer(".jar") →
JavaDagImporter
+ │
+ └── JavaDagImporter.import_definition(definition, bundle=...)
+ │
+ ├── LangSDKDagFileProcessorProcess.start(
+ │ target=_parse_lang_sdk_dag_entrypoint,
+ │ coordinator=JavaCoordinator(...),
path=analytics.jar)
+ │ │
+ │ ├── in the child: _build_parse_dag_command() →
(command, schema_version)
+ │ │ coordinator.parse_dag() — spawn
JVM, fd 0 ⇄ comm socket
+ │ │
+ │ │ ──DagFileParseRequest───▶ JVM
(ToDagProcessor)
+ │ │ ◀─DagFileParsingResult─── JVM
(ToManager)
+ │ │
┌────────────────────────────────────────────────────┐
+ │ │ │ serialized_dags: ["java_report"]
(@Builder.Dag) │
+ │ │
└────────────────────────────────────────────────────┘
+ │ │ TaskHandler registrations have no Dag to
serialize —
+ │ │ no "etl" entry exists to be discarded.
+ │ │
+ │ └── Get* from the JVM relayed up to the manager
unchanged
+ │
+ ├── LazyDeserializedDAG(data=...) → airflow.sdk.DAG
+ └── DagImportResult(dags=[DAG("java_report")])
+ ▼
+ _serialize_dags(bag) → DagFileParsingResult → manager
+ ▼
+ DagModelOperation → PERSIST "java_report" only
+```
+
+The coordinator is reached through the importer, never through the manager's
file-to-process routing. ADR-0004's `_resolve_processor_target` scan, which
picks a coordinator by
+asking each one `can_handle_dag_file`, is superseded here: extension-keyed
importer registration has already decided that `.jar` belongs to this
coordinator, and two mechanisms
+claiming the same file would have to agree.
+
+For comparison, a pure Python file with no stub tasks:
+
+```
+PythonDagImporter.import_definition(definition, bundle=...)
+ │
+ ├── Parse → DAG objects
+ ├── serialize_dag(dag) → no stub tasks, nothing to cross-validate
+ ├── Return DagImportResult(dags=[dag])
+ ▼
+DagModelOperation → PERSIST
+```
+
+`DagImportResult.dags` is `list[DAG]`, so the importer wraps each serialized
entry as a `LazyDeserializedDAG` and transforms it into an `airflow.sdk.DAG`.
+
+## Consequences
+
+- A Lang-SDK importer is never configured by hand. The runtime is declared
once, in `[sdk] coordinators`, and the importer follows from it.
+- Routing a Lang-SDK artifact to its runtime becomes the importer registry's
job. ADR-0004's `can_handle_dag_file` / `_resolve_processor_target` scan no
longer decides which
+ process parses a file, and the coordinator method it drove is replaced by
`parse_dag` ([ADR-0012](0012-lang-sdk-parse-protocol.md)).
+- One coordinator instance per DagBundle becomes a deployment constraint: two
JDKs mean two `dag_bundle_name` values and two bundle-scoped registries. This
is what keeps
+ extension-keyed registration unambiguous.
+- A coordinator in `EXPLICIT_ROOT` mode cannot back a Dag importer. Its
artifacts live outside any DagBundle, so nothing scans them. It still
implements `get_dag_importer`; the
+ importer is simply never asked for, because `for_bundle` does not return
that coordinator.
+- `CoordinatorManager` gains `for_bundle`, a second lookup axis beside
`for_queue`.
+- A packed Go bundle claims the empty extension, which three call sites
currently treat as absent rather than as a key. Appendix B lists them.
+- `get_source_code` is abstract, so every Lang-SDK importer must implement it,
and a native Lang-SDK Dag has no Python source to return. What it should
return, and how that squares
+ with [ADR-0006](0006-no-lang-sdk-source-display.md), is not settled here.
+- A native Dag crosses serialization twice over: the runtime serializes it,
the importer deserializes it into an `airflow.sdk.DAG`, and the enclosing parse
serializes it again.
+ That round trip is the price of `DagImportResult.dags` being `list[DAG]` —
the same cost AIP-85's `list[DAG]` / `list[LazyDeserializedDAG]` question is
about.
+- Parsing a native Dag costs two processes, not one: the Dag-processing child
the manager already starts, plus the runtime the importer starts inside it. The
inner one inherits its
+ request, result and logging from `DagFileProcessorProcess`, so the only
Lang-SDK-specific code is the target callable.
+
+## References
+
+- [ADR-0012](0012-lang-sdk-parse-protocol.md) — `parse_dag` and the
coordinator interface it belongs to
+- [ADR-0011](0011-mixed-language-dag-processing.md) — why `TaskHandler`
registrations never reach a `DagImporter`
+- [ADR-0003](0003-pure-java-dags.md) — `BundleScanner` / `BuilderProcessor`,
build-time artifact inventory
+- [ADR-0004](0004-dag-parsing.md) — `can_handle_dag_file`, the subprocess
bridge
+- [ADR-0006](0006-no-lang-sdk-source-display.md) — no Lang-SDK source display
+- `task-sdk/src/airflow/sdk/importers/` — `AbstractDagImporter`,
`DagImportResult`, `DagImporterRegistry`
+- `airflow-core/src/airflow/dag_processing/processor.py` —
`DagFileProcessorProcess`, the class `LangSDKDagFileProcessorProcess` extends
+- [AIP-85](https://cwiki.apache.org/confluence/x/_Q7OEg) — DagImporter
+- [AIP-108](https://cwiki.apache.org/confluence/x/pY4mGQ) — Language SDKs
+
+## Appendix
+
+### Appendix A — Artifact sources and what `for_bundle` returns
+
+`SubprocessCoordinator` classifies artifact ownership at construction. The
explicit root and `dag_bundle_name` are mutually exclusive, and both that
conflict and a
+`dag_bundle_name` naming an unconfigured bundle are rejected there.
+
+`NAMED_BUNDLE` is the unambiguous case. `for_bundle(name)` returns it when its
`dag_bundle_name` matches, so its importer lands in exactly one bundle-scoped
registry. Two
+`dag_bundle_name` values give two registries, and `.jar` is claimed once in
each.
+
+`TASK_BUNDLE` has no fixed bundle — its artifacts ride along with whichever
Dag delegates to them — so `for_bundle` returns it for every bundle and its
importer registers
+everywhere. That is sound only while it is the sole claimant of its extension,
which is the co-located single-runtime deployment.
+
+`EXPLICIT_ROOT` points at a filesystem path outside any DagBundle. The Dag
processor never scans it, so there is no file for an importer to claim and
`for_bundle` never returns it.
+Such a coordinator is reachable only by queue, for mixed-language work.
+
+`get_importer_registry(bundle_name)` is already cached per bundle, and
`CoordinatorManager` caches instances separately, so the two caches have to be
reset together.
+
+### Appendix B — Extensionless artifacts
+
+A packed Go bundle has no suffix. `ExecutableDagImporter` claims the empty
extension as a first-class key rather than depending on a `can_handle` scan,
whose winner varies with
+registration order because `_ordered_importers` is scanned in reverse.
+
+Three places assume a non-empty suffix today:
+
+- `_normalize_extensions` rewrites `""` to `"."`.
+- `get_importer` and `can_handle` guard on `if suffix:`, which skips the
extension map entirely for an extensionless file.
+- `find_file_dag_definitions` filters on `path.suffix.lower()`.
+
+Empty has to pass through all three, with the guards testing `suffix is not
None`.
+
+### Appendix C — Resolving artifact roots at parse time
+
+`_init_root_source` already resolves roots for all three modes, but publishes
them through `_get_scan_roots()`, which is scoped to an active task and raises
outside one. Both
+parse-side commands need the same roots with no `TaskInstance` in hand:
`EXPLICIT_ROOT` and `NAMED_BUNDLE` resolve from the coordinator's own
configuration, and `TASK_BUNDLE`
+resolves against the bundle the Dag processor is parsing. The scope that
publishes the roots has to open for a parse as well as for a task.
diff --git a/airflow-core/adr/lang-sdk/0011-mixed-language-dag-processing.md
b/airflow-core/adr/lang-sdk/0011-mixed-language-dag-processing.md
new file mode 100644
index 00000000000..c2acaadbedf
--- /dev/null
+++ b/airflow-core/adr/lang-sdk/0011-mixed-language-dag-processing.md
@@ -0,0 +1,245 @@
+<!--
+ Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+ Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+ -->
+
+# ADR-0011: Mixed-Language Dag Processing — Task Handlers Are Not Dags
+
+## Status
+
+Proposed
+
+## Context
+
+A Lang-SDK artifact can serve either of the two authoring features the
Language SDK spec fixes, or both at once:
+
+| Feature | Who owns the graph | Author writes
| Artifact contributes |
+|---------------------------------|----------------------------|------------------------------|----------------------|
+| **Mixed Language Task Handler** | Python, via `@task.stub` |
`TaskHandler(dag, task, fn)` | Only task bodies |
+| **Native Dag** | the Lang-SDK source itself | `Dag(spec)`
| The entire Dag |
+
+In the mixed-language role the artifact used to author a Dag too, under the
same `dag_id` the Python file already owns. Two conflicting definitions reached
persistence and
+something downstream had to choose between them. This ADR removes the conflict
at the authoring interface, and defines where a stub task is validated against
the handler that
+implements it.
+
+Terms follow the Language SDK spec (`task-sdk/docs/lang-sdk-spec.rst`, spec
version `1.0`). Note that the spec's `bundle` is the SDK-side registration
container, not Airflow's
+`DagBundle` — both appear below.
+
+## Decision
+
+### Two registration kinds, one bundle
+
+```
+┌──────────────────────────────────────────────────────────────────────────┐
+│ Per-task (Python): is_stub │
+│ AbstractOperator.is_stub = False (default) │
+│ _StubOperator.is_stub = True │
+│ │
+│ Per-registration (Lang-SDK): which interface the author wrote │
+│ Dag(spec) → DagRef │
+│ TaskHandler(dag, task, fn) → TaskHandlerRef │
+│ │
+│ Both kinds go into one bundle, through one verb: │
+│ bundle.register(dag, handler, ...) → bundle.serve() │
+└──────────────────────────────────────────────────────────────────────────┘
+```
+
+A `TaskHandlerRef` has no schedule, no task graph, and no `dag_id` of its own
to persist. A `DagRef` has all three. Dag parsing draws on the `Dag`
registrations
+([ADR-0010](0010-native-dag-processing.md)); validation draws on the
`TaskHandler` registrations. Because both kinds share one bundle and one
`register` verb, the bundle cannot be
+the discriminator — the registration kind is. Appendix A gives the rejected
alternative and why.
+
+### One artifact, both kinds
+
+```
+analytics.jar
+├── EtlTasks @Builder.TaskHandler(dagId = "etl", taskId = "extract")
+│ @Builder.TaskHandler(dagId = "etl", taskId =
"transform")
+│ backs etl.py's stub tasks — no Dag on the Java side
+└── JavaReportPipeline @Builder.Dag(id = "java_report")
+ native, persisted as "java_report"
+
+bundle.register(javaReport, etlExtract, etlTransform); // one verb, both
kinds
+bundle.serve(args);
+```
+
+The split is per registration, not per file and not per bundle.
+
+### Validation runs in the Dag-file parse, after the Dags exist
+
+Comparing a stub against its handler needs two things at once: the
`airflow.sdk.DAG` objects the Python file produced, and the `arg_bindings` that
only appear once those Dags are
+serialized. `_parse_file` holds both, between `_serialize_dags` and the
`DagFileParsingResult` it returns. That is where the handler query is issued.
+
+```
+DagFileProcessorProcess(etl.py) ← manager
spawns, as for any file
+ └── _parse_file_entrypoint → _parse_file
+ ├── BundleDagBag → PythonDagImporter → airflow.sdk.DAG objects ←
the Dags now exist
+ ├── _serialize_dags(bag) → is_stub tasks carry arg_bindings
(ADR-0007)
+ │
+ │
┌─────────────────────────────────────────────────────────────────────┐
+ │ │ Step 1: Collect the dag_ids to ask about — every Dag in this
│
+ │ │ file with at least one is_stub task → ["etl"]
│
+ │
└─────────────────────────────────────────────────────────────────────┘
+ │
+ │
┌─────────────────────────────────────────────────────────────────────┐
+ │ │ Step 2: Resolve stub tasks → Coordinator instances via
│
+ │ │ CoordinatorManager.for_queue(queue)
│
+ │ │ ([sdk] queue_to_coordinator)
│
+ │ │
│
+ │ │ stub task queue Coordinator instance
│
+ │ │ ──────────────────────────────────────────────────
│
+ │ │ extract → "jdk-11" → JavaCoordinator(name="jdk-11")
│
+ │ │ transform → "jdk-17" → JavaCoordinator(name="jdk-17")
│
+ │ │ load → "jdk-11" → JavaCoordinator(name="jdk-11")
│
+ │ │
│
+ │ │ Deduplicate → 2 distinct coordinator instances
│
+ │
└─────────────────────────────────────────────────────────────────────┘
+ │
+ │
┌─────────────────────────────────────────────────────────────────────┐
+ │ │ Step 3: Locate the artifact backing each dag_id, no JVM launch —
│
+ │ │ reuse BundleScanner (ADR-0003), whose build-time
inventory │
+ │ │ indexes registered (dagId, taskId) handler pairs
alongside │
+ │ │ the artifact's native Dag ids
│
+ │ │
│
+ │ │ JavaCoordinator(name="jdk-11")
│
+ │ │ ├── resolve artifact roots for its mode
│
+ │ │ │ EXPLICIT_ROOT → jars_root
│
+ │ │ │ NAMED_BUNDLE → dag_bundle_name, pinned
│
+ │ │ │ TASK_BUNDLE → the bundle being parsed
│
+ │ │ │
│
+ │ │ └── BundleScanner.scanBundles(roots)
│
+ │ │ → "etl" → ResolvedBundle(analytics.jar, mainClass, ...)
│
+ │ │
│
+ │ │ Group the dag_ids by (coordinator, artifact) — one group, one
│
+ │ │ process, one request
│
+ │
└─────────────────────────────────────────────────────────────────────┘
+ │
+ │
┌─────────────────────────────────────────────────────────────────────┐
+ │ │ Step 4: Query each group — one request, one response
│
+ │ │
│
+ │ │ SDKTaskHandlerProcessorProcess.start(
│
+ │ │ target=_parse_task_handler_entrypoint,
│
+ │ │ coordinator=JavaCoordinator("jdk-11"),
│
+ │ │ path=analytics.jar)
│
+ │ │ │
│
+ │ │ ├── in the child: _build_parse_task_handler_command()
│
+ │ │ │ coordinator.parse_task_handler() — spawn JVM
│
+ │ │ │
│
+ │ │ │ ──TaskHandlerParseRequest(file=analytics.jar,
│
+ │ │ │ dag_ids=["etl"])─────▶ JVM
│
+ │ │ │ (ToSDKTaskHandlerProcessor)
│
+ │ │ │
│
+ │ │ │ JVM answers from its own TaskHandler registrations
│
+ │ │ │ whose dagId is one of the requested ids
│
+ │ │ │
│
+ │ │ │ ◀─TaskHandlerParsingResult(task_handlers={
│
+ │ │ │ "etl": [extract, transform, load]})─────── JVM
│
+ │ │ │ (ToManager)
│
+ │ │ │
│
+ │ │ └── Get* from the JVM relayed up to the manager unchanged
│
+ │ │
│
+ │ │ JavaCoordinator(name="jdk-17")
│
+ │ │ └── its own process, its own resolved artifact and handler set
│
+ │
└─────────────────────────────────────────────────────────────────────┘
+ │
+ │
┌─────────────────────────────────────────────────────────────────────┐
+ │ │ Step 5: Compare per dag_id against the Dag just serialized —
│
+ │ │ union task_handlers[dag_id] across every process first
│
+ │ │
│
+ │ │ Python Dag "etl" (stub tasks) TaskHandlerDeclaration
│
+ │ │ ────────────────────────────── ──────────────────────────────
│
+ │ │ task_id ↔ task_id (sets must match)
│
+ │ │ arg_bindings[*].name ↔ params[*].name (in order)
│
+ │ │ arg_bindings[*].value_schema ↔ params[*].value_schema
│
+ │ │ compared only where neither side is null
│
+ │ │
│
+ │ │ On mismatch → import_errors["etl.py"]
│
+ │
└─────────────────────────────────────────────────────────────────────┘
+ ▼
+ DagFileParsingResult(serialized_dags=["etl"], import_errors={...})
+ ▼
+ manager → DagModelOperation → PERSIST (the Python Dag is the sole DB record)
+```
+
+The parse owns validation, not an importer. `PythonDagImporter` returns
`airflow.sdk.DAG` objects and knows nothing about coordinators or queues, so
`@task.stub` keeps working for
+any importer that can produce a Dag carrying stub tasks. `_parse_file` is also
the only place where the whole file's Dags are visible at once, which is what
lets one request cover
+every `dag_id` that resolved to the same artifact
([ADR-0012](0012-lang-sdk-parse-protocol.md)).
+
+Resolution goes through the coordinator registry, not the filesystem, so the
Python Dag and the Lang-SDK artifact **do not need to be in the same
DagBundle**. Nothing here needs an
+`airflow.sdk.DAG` round-trip either — validation compares against the Dag the
Python parser already built. Appendix B states exactly what is compared.
+
+### Decision matrix
+
+| Caller | Coordinator call
| What comes back | Action
|
+|--------------------------------------------|------------------------------------------------------|------------------------------------------|--------------------------------------------------|
+| `_parse_file` → `PythonDagImporter` | — (the Python file is parsed in
process) | its own parsed Dags | PERSIST
|
+| `_parse_file`, per (coordinator, artifact) | `parse_task_handler`, scoped to
that group's dag_ids | `TaskHandlerParsingResult` | VALIDATE only
— not a Dag, so nothing to persist |
+| `_parse_file` → `JavaDagImporter` | `parse_dag`
| `DagFileParsingResult`, native Dags only | PERSIST
|
+
+There is no fourth row. A `TaskHandlerRef` has no Dag, so no `DagImporter` —
and nothing reading a `DagImporter`'s results — ever sees one.
+
+## Consequences
+
+- Python leads, Lang-SDK follows. The Dag-file parse persists the Dag and
drives validation via `queue → Coordinator → parse_task_handler`.
+- No importer knows about coordinators. `PythonDagImporter` is unchanged by
this ADR; the stub-to-handler comparison sits in `_parse_file`, above every
importer.
+- A mixed-language `dag_id` never appears in Dag processing results. No `Dag`
registration exists for a `dag_id` a Python file already owns, so everything
downstream sees exactly
+ one record per `dag_id`, with no flag to interpret.
+- Stub/implementation mismatches — missing handler, extra handler, parameter
name or order, incompatible schema — surface as import errors against the
Python file at parse time,
+ alongside the errors the parse already reports. An unannotated stub argument
is checked by name and position only.
+- The Python Dag and Lang-SDK artifact can live in different DagBundles.
+- A single Dag can have stubs targeting different queues, some Java, some Go.
Each resolves to its own coordinator instance, and validation unions their
declarations per `dag_id`
+ before comparing task ids.
+- Validating a file costs one extra process per (coordinator, artifact) pair
its stubs resolve to — one for the common case of a file whose stubs all target
a single runtime, and
+ none at all for a file with no stub tasks.
+- Mixed-language is Python-primary only. Lang-SDK runtimes cannot define stub
operators; a native Dag cannot delegate tasks to Python.
+- No per-Dag flag, no schema migration, no new `DagModel` column, no REST/UI
change.
+- Terms track Language SDK spec `1.0`. A spec rename of `TaskHandler`, or of
the `register` / `serve` verbs, lands here too.
+
+## References
+
+- [ADR-0012](0012-lang-sdk-parse-protocol.md) — `parse_task_handler` and the
`TaskHandlerParsingResult` shape this ADR compares against
+- [ADR-0010](0010-native-dag-processing.md) — the `Dag`-registration half, and
the importer that persists it
+- [ADR-0003](0003-pure-java-dags.md) — `BundleScanner` / `BuilderProcessor`,
build-time artifact inventory
+- [ADR-0006](0006-no-lang-sdk-source-display.md) — no Lang-SDK source display
for mixed-language Dags
+- [ADR-0007](0007-taskflow-across-language-boundary.md) — `arg_bindings` /
`TaskArgBinding` / `ArgValueSchema`
+- Language SDK spec (`task-sdk/docs/lang-sdk-spec.rst`)
+- [AIP-108](https://cwiki.apache.org/confluence/x/pY4mGQ) — Language SDKs
+- [AIP-85](https://cwiki.apache.org/confluence/x/_Q7OEg) — DagImporter
+- `airflow-core/src/airflow/dag_processing/processor.py` — `_parse_file`,
`_serialize_dags`, `DagFileParsingResult`
+
+## Appendix
+
+### Appendix A — Why the split lives on the interface
+
+The alternative was to keep authoring the Dag on both sides and mark the
Lang-SDK copy with a per-Dag flag, `is_mixed_language_dag`, so importers knew
which copy to drop. That
+keeps producing the definition it then has to discard, and it pushes the
question "is this Dag real?" onto every consumer of a serialized Dag — the
importer, the persistence layer,
+anything later reading the record. Taking the Dag out of the authoring
interface answers the question once, at the only point where the answer is
known for free: the author already
+chose which interface to write against. The only marker left in serialization
is `is_stub`, and it is per-task.
+
+The bundle cannot carry the distinction either. The Language SDK spec puts
both kinds in one `bundle`, takes them through one `register` verb in any
mixture, and serves the process
+with one `bundle.serve()`. One artifact can hold both, so neither the file nor
the bundle tells a `DagImporter` what it is holding.
+
+### Appendix B — What is compared
+
+Handler declarations from every process spawned in Step 4 are unioned per
`dag_id` before comparison, since one Dag's stubs can target several queues.
+
+- `task_id` sets must match exactly. A missing or extra handler is an error.
+- `arg_bindings[*].name` against `params[*].name`, in order — both sides bind
positionally.
+- `arg_bindings[*].value_schema` against `params[*].value_schema`, compared
only where neither side is null. An unannotated `@task.stub` parameter produces
`null` today, so a
+ strict comparison would make every untyped stub argument a parse error.
+
+Any mismatch is reported against the Python file, which is the definition the
author can act on, and travels back on `DagFileParsingResult.import_errors`
with everything else the
+parse found.
diff --git a/airflow-core/adr/lang-sdk/0012-lang-sdk-parse-protocol.md
b/airflow-core/adr/lang-sdk/0012-lang-sdk-parse-protocol.md
new file mode 100644
index 00000000000..e144b87445b
--- /dev/null
+++ b/airflow-core/adr/lang-sdk/0012-lang-sdk-parse-protocol.md
@@ -0,0 +1,233 @@
+<!--
+ Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+ Unless required by applicable law or agreed to in writing,
+ software distributed under the License is distributed on an
+ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ KIND, either express or implied. See the License for the
+ specific language governing permissions and limitations
+ under the License.
+ -->
+
+# ADR-0012: Lang-SDK Parse Protocol — Handler Messages and Coordinator Verbs
+
+## Status
+
+Proposed
+
+## Context
+
+The Dag processor asks a Lang-SDK runtime two different questions. "Which Dags
does this artifact define?" is answered over the messages
[ADR-0004](0004-dag-parsing.md) already
+defines. "Which task handlers does this artifact register for a `dag_id`
Python already owns?" has no answer in those messages, because a `TaskHandler`
registration carries no Dag
+([ADR-0011](0011-mixed-language-dag-processing.md)).
+
+This ADR defines the request that carries the second question, the subprocess
classes that carry both, and the two parse-side entry points on the coordinator.
+
+Terms follow the Language SDK spec (`task-sdk/docs/lang-sdk-spec.rst`, spec
version `1.0`).
+
+## Decision
+
+### One channel shape, two request types
+
+```
+parse_dag parse_task_handler
+
+ parent process parent process
+ │ DagFileParseRequest │ TaskHandlerParseRequest
+ ▼ (ToDagProcessor) ▼
(ToSDKTaskHandlerProcessor)
+ coordinator (raw byte forward) coordinator (raw byte
forward)
+ ▼ ▼
+ runtime runtime
+ │ DagFileParsingResult │
TaskHandlerParsingResult
+ ▼ (ToManager) ▼ (ToManager)
+ parent process parent process
+```
+
+Both verbs are byte forwarders: the coordinator spawns the runtime, wires `fd
0` to the comm socket, and never decodes the payload. The process that spawned
the parse decodes the
+reply. They stay two methods, not one `parse(request)`, because a coordinator
can serve handlers without serving native Dag parsing. Appendix A has the
longer argument.
+
+### The reply travels on `ToManager`
+
+```
+_ParseSideResponses = shared tail — same members, same
`type` discriminator
+ ConnectionResult | VariableResult | VariableKeysResult | TaskStatesResult
+ | PreviousDagRunResult | PreviousTIResult | PrevSuccessfulDagRunResult
+ | ErrorResponse | OKResponse | XComCountResponse | XComResult
+ | XComSequenceIndexResult | XComSequenceSliceResult
+
+ToDagProcessor = DagFileParseRequest | _ParseSideResponses
parent → child
+ToSDKTaskHandlerProcessor = TaskHandlerParseRequest | _ParseSideResponses
parent → child (new)
+
+ToManager = DagFileParsingResult | TaskHandlerParsingResult
child → parent
+ | GetConnection | GetVariable | … | MaskSecret
+```
+
+`ToSDKTaskHandlerProcessor` is the only new union; `ToManager` gains one
member. The two parent → child unions differ in exactly one member, because the
child's questions about
+connections, variables and XComs do not depend on which parse it was asked for.
+
+`ToManager` is named for the process that usually holds the other end, but the
role it describes is "whoever spawned this parse". The Dag processor manager
fills it for
+`DagFileProcessorProcess`; a Dag-parsing child fills it for the two processes
below, relaying anything that is not a parsing result up its own `ToManager`
channel unchanged. That
+relay is only type-safe because both hops speak the same pair, which is the
reason not to mint a separate `ToCoordinator`.
+
+### Message shapes
+
+```
+class TaskHandlerParseRequest:
+ file: str # the artifact resolved for this
coordinator
+ dag_ids: list[str] # every Dag in the parsed file with
stub tasks that resolved here
+ bundle_path: Path
+ bundle_name: str
+ type: Literal["TaskHandlerParseRequest"]
+
+class TaskHandlerParsingResult:
+ fileloc: str
+ task_handlers: dict[str, list[TaskHandlerDeclaration]] # dag_id → its
declarations
+ import_errors: dict[str, str] | None = None
+ warnings: list | None = None
+ type: Literal["TaskHandlerParsingResult"]
+
+class TaskHandlerDeclaration:
+ task_id: str
+ params: list[TaskHandlerParam] # ordered — arg_bindings are positional
+
+class TaskHandlerParam:
+ name: str
+ value_schema: ArgValueSchema | None = None
+ required: bool # the handler declares no default
+```
+
+One request carries every `dag_id` that resolved to the same artifact under
the same coordinator, so a file whose stubs all target one runtime costs one
process. A `dag_id` the
+artifact registers nothing for is **omitted** from `task_handlers` rather than
returned empty: the key set is not required to match `dag_ids`, because it is
the union across
+coordinators that has to cover the stubs
([ADR-0011](0011-mixed-language-dag-processing.md)).
+
+`value_schema` reuses the `ArgValueSchema` definition `arg_bindings` already
carries ([ADR-0007](0007-taskflow-across-language-boundary.md)), so both sides
of a comparison are the
+same type. Two properties matter to validation: the field is nullable on both
sides, and `params` is ordered. Appendix B says what that forces.
+
+`task_handlers` is the counterpart to `DagFileParsingResult.serialized_dags`,
but fully typed. `serialized_dags` is `list[LazyDeserializedDAG]`, which is an
opaque object in the
+schema snapshot. A handler declaration carries no Dag, so it code-generates
and schema-validates in every SDK, and nothing on this path needs a
DagSerialization implementation.
+
+### Parse processes
+
+```
+WatchedSubprocess
+ └── BaseParsingProcess socket lifecycle · ToManager
decoding · Get* handling
+ │ · log forwarding under
dag_processor.*
+ ├── DagFileProcessorProcess
(shipped, now a subclass)
+ │ │ target = _parse_file_entrypoint
+ │ │ DagFileParseRequest → DagFileParsingResult
+ │ │
+ │ └── LangSDKDagFileProcessorProcess (new —
ADR-0010)
+ │ target = _parse_lang_sdk_dag_entrypoint
+ │ └── coordinator.parse_dag() — spawn runtime, forward fd
0 ⇄ comm socket
+ │ same request and result types as its base class
+ │
+ └── SDKTaskHandlerProcessorProcess (new —
ADR-0011)
+ target = _parse_task_handler_entrypoint
+ └── coordinator.parse_task_handler() — same forwarding
+ TaskHandlerParseRequest → TaskHandlerParsingResult
+```
+
+`BaseParsingProcess` is `DagFileProcessorProcess` minus the Dag-specific
request and result: the comm socket, the `ToManager` decode, the `Get*`
dispatch, and the
+`task.` → `dag_processor.` log-forwarder rename. The subclasses supply the
first message they send, the result they collect, and the target the child runs.
+
+`LangSDKDagFileProcessorProcess` differs from its base in the target callable
alone. Everything else — the request, the result, the socket, the logging — is
inherited, because a
+native Lang-SDK Dag answers the same question a Python file does.
+
+Answering `Get*` needs a `Client`, which only the manager holds.
`BaseParsingProcess` therefore resolves a request one of two ways: directly
against `self.client` when the manager
+is the parent, or by relaying it up `SUPERVISOR_COMMS` when a Dag-parsing
child is.
+
+### Coordinator interface
+
+```
+BaseCoordinator execution_time/coordinator.py
+ ├── execute_task (shipped)
+ ├── parse_dag (new — native Dags, ADR-0010)
+ └── parse_task_handler (new — handlers, ADR-0011)
+ │
+SubprocessCoordinator coordinators/_subprocess.py
+ implements all three; each resolves (command, subprocess_schema_version)
+ from a hook and owns the socket lifecycle:
+ ├── _build_execute_task_command (shipped)
+ ├── _build_parse_dag_command (new)
+ └── _build_parse_task_handler_command (new)
+ │
+JavaCoordinator · ExecutableCoordinator · NodeCoordinator
+ supply the three commands; no socket or protocol code
+```
+
+Names follow the shipped `execute_task` / `_build_execute_task_command` pair
and supersede ADR-0004's `run_dag_parsing` / `dag_parsing_cmd`. Each hook
returns its own
+`subprocess_schema_version`, so handler parsing negotiates the schema the same
way task execution does.
+
+## Consequences
+
+- The handler channel is typed; the Dag channel is opaque. Every SDK gets
generated models for the declaration shape.
+- Each runtime implements a second request type instead of a new field on the
existing one. A `DagRef` and a `TaskHandlerRef` have different shapes, so one
request/result pair
+ could not carry both.
+- `ToSDKTaskHandlerProcessor` extends the set of unions the supervisor schema
package introspects, from four to five. It adds two union members —
`TaskHandlerParseRequest` and
+ `TaskHandlerParsingResult`, the latter carrying `TaskHandlerDeclaration` and
`TaskHandlerParam` as nested definitions — because the shared responses are
classes the registry
+ already holds.
+- Nothing in the protocol distinguishes a coordinator-backed parse from a
Python one. A runtime's `Get*` request is answered by the same handlers that
answer a Python parser's,
+ through however many relay hops lie between it and the manager.
+- `DagFileProcessorProcess` becomes a subclass. Its public surface does not
move, but the shipped `_handle_request` and socket code shifts to
`BaseParsingProcess`.
+- Neither verb is reached through ADR-0004's `can_handle_dag_file` scan.
`parse_dag` is reached through the importer registered for the artifact's
extension
+ ([ADR-0010](0010-native-dag-processing.md)); `parse_task_handler` through
`queue → coordinator` ([ADR-0011](0011-mixed-language-dag-processing.md)).
+- Terms track Language SDK spec `1.0`. A spec rename of `TaskHandler` lands
here too.
+
+## References
+
+- [ADR-0004](0004-dag-parsing.md) — coordinator subprocess bridge,
`DagFileParseRequest` / `DagFileParsingResult`, `can_handle_dag_file`
+- [ADR-0006](0006-no-lang-sdk-source-display.md) — no Lang-SDK source display
+- [ADR-0007](0007-taskflow-across-language-boundary.md) — `arg_bindings` /
`TaskArgBinding` / `ArgValueSchema`
+- [ADR-0010](0010-native-dag-processing.md) — who calls `parse_dag`
+- [ADR-0011](0011-mixed-language-dag-processing.md) — who calls
`parse_task_handler`, and what it compares the reply against
+- Language SDK spec (`task-sdk/docs/lang-sdk-spec.rst`)
+- `airflow-core/src/airflow/dag_processing/processor.py` —
`DagFileProcessorProcess`, `ToManager` / `ToDagProcessor`
+- `task-sdk/src/airflow/sdk/execution_time/schema/` — supervisor schema
version bundle, union registry, generated snapshot
+- [AIP-108](https://cwiki.apache.org/confluence/x/pY4mGQ) — Language SDKs
+
+## Appendix
+
+### Appendix A — Why two verbs, and why one reply union
+
+The two verbs differ in what they ask for and in which command starts the
runtime, not in how they talk. Collapsing them into one `parse(request)` would
hide a real deployment
+distinction: a coordinator can serve mixed-language handlers with no interest
in native Dag parsing, and an absent method states that better than a runtime
rejection does.
+
+The reply direction does not need the same split. An earlier draft gave
handler parsing its own `ToRuntime` / `ToCoordinator` pair on the theory that
the runtime was answering the
+coordinator rather than the manager. It is not: the coordinator forwards bytes
in both directions and decodes nothing, so the peer at the far end of the
socket is whichever process
+spawned the parse. Giving that peer two unions to decode would mean two
`CommsDecoder` configurations, two relay paths for the identical `Get*`
traffic, and two registry entries for
+bodies that never differ. Reusing `ToManager` leaves one reply union with one
new member.
+
+A boolean on `DagFileParseRequest` was the other alternative. It cannot work:
a `DagRef` and a `TaskHandlerRef` are different payloads, not two subsets of
one, so the flag would
+select between shapes the result type cannot both hold.
+
+### Appendix B — What the nullable, ordered parameter list forces
+
+`value_schema` is nullable on both sides. An unannotated `@task.stub`
parameter produces `value_schema: null` today, so validation compares schemas
only where neither side is null,
+and falls back to name-and-arity otherwise. A strict comparison would turn
every untyped stub argument into a parse error.
+
+`params` is ordered because `LiteralArgBinding` and `XComArgBinding` are each
documented as "one positional stub-task argument". Position is part of the
contract, not incidental,
+and both sides bind positionally.
+
+A declaration carries no class, method, or source location.
[ADR-0006](0006-no-lang-sdk-source-display.md) rules out Lang-SDK source
display, and putting it on the wire would
+invite a consumer to render it. It carries no `dag_id` either — the
`task_handlers` key supplies it, so a declaration cannot disagree with the
bucket it arrived in.
+
+### Appendix C — Schema registration mechanics
+
+`registered_models_by_name()` in
`task-sdk/src/airflow/sdk/execution_time/schema/` introspects a fixed set of
unions — `ToTask` / `ToSupervisor` and `ToManager` / `ToDagProcessor`
+— so adding `ToSDKTaskHandlerProcessor` extends that set to five. The shared
responses appear in two unions now; the registry keys by class name and rejects
two *distinct* classes
+under one name, so a member reached twice is not a clash.
+
+Two prek hooks guard the generated snapshot. Their interaction on a
first-introduction body needs checking against the hooks rather than against
that package's `AGENTS.md`: the doc
+says no `VersionChange` is required for a new body, while
`check-supervisor-schemas-versions` fails when the snapshot moves and nothing
under `versions/` was touched.
+
+Artifact roots are resolved per mode by `_init_root_source`, but published
through `_get_scan_roots()`, which is scoped to an active task and raises
outside one. Both parse-side
+commands need those roots with no `TaskInstance` in hand, so the scope that
publishes them has to open for a parse as well as for a task.
[ADR-0010](0010-native-dag-processing.md)
+covers the modes themselves.
diff --git a/airflow-core/adr/lang-sdk/README.md
b/airflow-core/adr/lang-sdk/README.md
index 359275ac269..ec80aa58732 100644
--- a/airflow-core/adr/lang-sdk/README.md
+++ b/airflow-core/adr/lang-sdk/README.md
@@ -35,6 +35,9 @@ bind core interfaces and apply to every language SDK, not
just the Java SDK.
- [ADR-0007](0007-taskflow-across-language-boundary.md): TaskFlow across the
language boundary — argument binding for Lang-SDK tasks.
- [ADR-0008](0008-control-flow-constructs.md): control-flow constructs —
grouping, conditional skipping, branching, and triggering a Dag run.
- [ADR-0009](0009-provider-operators-as-generated-dsl.md): provider operators
as generated, serialization-only DSL in every Lang SDK.
+- [ADR-0010](0010-native-dag-processing.md): native Dag processing —
DagImporter registration and routing.
+- [ADR-0011](0011-mixed-language-dag-processing.md): mixed-language Dag
processing — task handlers are not Dags.
+- [ADR-0012](0012-lang-sdk-parse-protocol.md): Lang-SDK parse protocol — task
handler messages and coordinator verbs.
Decisions specific to a single SDK stay next to that SDK — for example, the Go
SDK's bundle-format
decisions live in [`go-sdk/adr/`](../../../go-sdk/adr). Java-SDK-only
interface-design decisions —