Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
94 commits
Select commit Hold shift + click to select a range
bd9be97
Replace coordinator artifact roots with task_handler_bundle_name
jason810496 Sep 30, 2026
ace4b2e
Serve lang-SDK harness artifacts from Dag bundles
jason810496 Sep 30, 2026
4bfc181
Document artifact Dag bundles for language SDK coordinators
jason810496 Sep 30, 2026
b9cde7f
Tidy names and mocks in the artifact bundle tests
jason810496 Sep 30, 2026
23898f0
Say artifact Dag bundles are registered on every component
jason810496 Sep 30, 2026
8997a35
Drop em-dashes from the rewritten doc sentences
jason810496 Sep 30, 2026
8983f6d
Say the Dag processor needs the language SDK setup too
jason810496 Oct 1, 2026
c705540
Validate task_handler_bundle_name only for routed coordinators
jason810496 Oct 2, 2026
9d00ead
Add newsfragment for the removed coordinator artifact roots
jason810496 Oct 2, 2026
68f03d6
Mark ADR-0013 as accepted
jason810496 Oct 2, 2026
c174aa4
Add lang-SDK task handler binding tables
jason810496 Sep 30, 2026
6762729
Match ADR-0013 task handler table SQL to the migration
jason810496 Sep 30, 2026
dc15082
Keep the artifact path hash in sync on ORM updates
jason810496 Sep 30, 2026
ab8f025
Test that a referenced task handler artifact cannot be deleted
jason810496 Sep 30, 2026
5da25f7
Describe LangSDKTaskHandler as a stub task binding
jason810496 Sep 30, 2026
396cac2
Regenerate the task handler migration revision id
jason810496 Sep 30, 2026
61b2d76
Store the task handler binding mode beside its params
jason810496 Sep 30, 2026
fce1160
Add handler_binding to ADR-0013's task handler table SQL
jason810496 Sep 30, 2026
d18d0af
Make session keyword-only in the task handler model test helpers
jason810496 Oct 1, 2026
7b572ab
Cache an artifact's task handlers on its row, not on each binding
jason810496 Oct 1, 2026
0361303
Move the handler declarations to the artifact in ADR-0013's SQL
jason810496 Oct 1, 2026
2c96b48
Let a task handler artifact store no cache digest
jason810496 Oct 1, 2026
da3372e
Allow a NULL cache digest in ADR-0013's artifact table SQL
jason810496 Oct 1, 2026
db4f847
Move the Dag file processor's shared plumbing into a base class
jason810496 Sep 29, 2026
7e26f0d
Test closing a Dag file processor that has no parse log file
jason810496 Oct 1, 2026
4f4b3bd
Add task-handler parse messages and their parse-channel union
jason810496 Sep 30, 2026
7e0327c
Regenerate TS SDK supervisor types for task-handler messages
jason810496 Sep 30, 2026
2826ec4
Document the version rules for a new supervisor schema body
jason810496 Sep 30, 2026
a2b0c29
Add a binding mode to task-handler declarations
jason810496 Sep 30, 2026
f3601d6
Regenerate TS SDK supervisor types for the binding mode
jason810496 Sep 30, 2026
7160a10
Describe the task-handler binding mode in ADR-0012
jason810496 Sep 30, 2026
de96d42
Describe the binding-mode parse check in ADR-0011 and ADR-0013
jason810496 Sep 30, 2026
06386fb
Add a named_open binding for non-exhaustive task-handler params
jason810496 Sep 30, 2026
12b142d
Regenerate TS SDK supervisor types for the named_open binding
jason810496 Sep 30, 2026
2700115
Add the named_open binding to the Lang-SDK ADRs
jason810496 Sep 30, 2026
d33f904
Ask a task handler runtime for every handler it registers
jason810496 Oct 1, 2026
d86cead
Regenerate TS SDK supervisor types for the probe without dag_ids
jason810496 Oct 1, 2026
c9eaf68
Describe the all-handlers probe in ADR-0011 and ADR-0012
jason810496 Oct 1, 2026
1bd5650
Bind task handlers only by position or by name
jason810496 Oct 1, 2026
a4598c4
Regenerate TS SDK supervisor types for the two binding modes
jason810496 Oct 1, 2026
1d793f9
Regenerate Go SDK models for task-handler messages
jason810496 Oct 2, 2026
07b5199
Describe the two task handler binding modes in the ADRs
jason810496 Oct 1, 2026
32db985
Keep a nested supervised child's standard streams open
jason810496 Sep 30, 2026
1ac623c
Let a subprocess coordinator exec a task-handler parse runtime
jason810496 Oct 2, 2026
4cb786b
Probe a Lang-SDK artifact for its task handlers
jason810496 Oct 2, 2026
bf18577
Read a Lang-SDK runtime's frames without blocking the parse loop
jason810496 Oct 2, 2026
896cc6f
Kill what a Lang-SDK runtime leaves behind when it exits
jason810496 Oct 2, 2026
1b93217
Do not exec a Lang-SDK runtime whose parent has already exited
jason810496 Oct 2, 2026
b45c6ec
Clean up a Lang-SDK runtime parse whose start fails after the fork
jason810496 Oct 2, 2026
aeff5da
Go SDK: describe task handler params with JSON Schema
jason810496 Sep 30, 2026
4619601
Go SDK: answer TaskHandlerParseRequest with registered handlers
jason810496 Sep 30, 2026
fc06e07
Go SDK: clarify parse timeout and nested-null docs, test anyOf
jason810496 Sep 30, 2026
f3d8034
Go SDK: record integrity and cache digests in bundle metadata
jason810496 Sep 30, 2026
c440bf7
Let the Executable coordinator probe a bundle's task handlers
jason810496 Sep 30, 2026
eb19610
Read a Go bundle's stored cache digest without hashing it
jason810496 Sep 30, 2026
6cf833f
Document the task handler parse for Language SDK authors
jason810496 Sep 30, 2026
9ecf47d
Go SDK: take the bundle digests from the staged executable
jason810496 Sep 30, 2026
262ca91
Resolve the probed bundle path and tighten the digest docs
jason810496 Sep 30, 2026
f731a54
Move Go SDK comments back to the code they describe
jason810496 Oct 1, 2026
a9d24ba
Tidy names, mocks and cases in the Go probe and digest tests
jason810496 Oct 1, 2026
f479c8a
Go SDK: answer a task handler parse with every handler
jason810496 Oct 1, 2026
5155067
State that a task handler parse answer depends only on the artifact
jason810496 Oct 1, 2026
3d9eaf2
Go SDK: declare every lone struct handler as named
jason810496 Oct 1, 2026
ccecc81
Describe the two task handler bindings for Language SDK authors
jason810496 Oct 1, 2026
d210b17
TS SDK: answer TaskHandlerParseRequest from bundle.serve
jason810496 Sep 30, 2026
8ea0048
Let the Node coordinator probe a bundle's task handlers
jason810496 Sep 30, 2026
e12ad5c
TS SDK: name the task handler lookup like its siblings
jason810496 Sep 30, 2026
57d85ff
Skip the TS real probe when node cannot report a version
jason810496 Sep 30, 2026
6d00642
Patch the TS real probe test with decorators
jason810496 Oct 1, 2026
f8acbd2
TS SDK: answer a task handler parse with every handler
jason810496 Oct 1, 2026
6c3e3e3
TS SDK: declare task handlers without listing their params
jason810496 Oct 1, 2026
0d37b86
Java SDK: sync the vendored supervisor schema from the monorepo
jason810496 Sep 30, 2026
d253959
Java SDK: record handler parameters in generated task classes
jason810496 Sep 30, 2026
078264d
Java SDK: answer TaskHandlerParseRequest with registered handlers
jason810496 Sep 30, 2026
8d360b8
Java SDK: stamp a cache digest on the bundle JAR manifest
jason810496 Sep 30, 2026
e4adc58
Java SDK: register stub Dags' Java tasks as task handlers
jason810496 Sep 30, 2026
f72d37f
Let the Java coordinator probe a JAR's task handlers
jason810496 Sep 30, 2026
11fb996
Java SDK: list a value schema's enum constants in a stable order
jason810496 Sep 30, 2026
e4c081c
Java SDK: check the task handler reply against the supervisor schema
jason810496 Sep 30, 2026
238799e
Say why the Java coordinator rejects a JAR it cannot probe
jason810496 Sep 30, 2026
5826da2
Test the Java SDK schema sync hook's snapshot copy
jason810496 Oct 1, 2026
8f942ec
Tidy names and mocks in the Java probe test
jason810496 Oct 1, 2026
18958ab
Test the cache digest of classpath directories and non-zip files
jason810496 Oct 1, 2026
628e134
Java SDK: answer a task handler parse with every handler
jason810496 Oct 1, 2026
f599ee6
Java SDK: declare task handler params without required
jason810496 Oct 1, 2026
580f68a
Let the coordinator registry list its task-handler bundles
jason810496 Sep 30, 2026
657bdc5
Add known task-handler artifacts to the Dag parse request
jason810496 Sep 30, 2026
90cfe6a
Regenerate TS SDK supervisor types for known artifacts
jason810496 Sep 30, 2026
f858597
Sync the Java SDK supervisor schema for known artifacts
jason810496 Sep 30, 2026
703d64f
Regenerate Go SDK models for known artifacts
jason810496 Oct 2, 2026
27370d0
Send the Dag-parsing child its known task-handler artifacts
jason810496 Sep 30, 2026
f78db1c
Describe how known artifacts are scoped per Dag-parsing child
jason810496 Sep 30, 2026
411c209
Test that a Dag file gets only named bundles without a fallback
jason810496 Sep 30, 2026
00586db
Build the known-artifact test params without required
jason810496 Oct 1, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 6 additions & 5 deletions .agents/skills/airflow-java-sdk/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,9 @@ subclasses must only import from `org.apache.airflow.sdk`; any import of

## Bundle composition and coordinator discovery

A **bundle** is a directory of JAR files (typically `build/bundle/`) placed on the coordinator's
`jars_root`. The coordinator scans the directory at task-dispatch time to find:
A **bundle** is a directory of JAR files (typically `build/bundle/`) placed in the Dag bundle named
by the coordinator's `task_handler_bundle_name` (the task's own Dag bundle when unset). The
coordinator scans that Dag bundle at task-dispatch time to find:

1. **`Main-Class`** (standard JAR manifest attribute) — the fully-qualified class name of the
entry point that the coordinator invokes with `java -classpath … <Main-Class> --comm … --logs …`.
Expand All @@ -65,7 +66,7 @@ A **bundle** is a directory of JAR files (typically `build/bundle/`) placed on t
`runtimeClasspath` and copies it into the shadow JAR manifest. In thin-JAR mode (`fatJar =
false`), the value stays in the `airflow-sdk` JAR deployed alongside the bundle JAR.

The Python coordinator (`JavaCoordinator`) scans every JAR under `jars_root` with
The Python coordinator (`JavaCoordinator`) scans every JAR in that Dag bundle with
`_JarInfo.find()`, reads `META-INF/MANIFEST.MF` out of each ZIP, and collects `Main-Class` and
`Airflow-Supervisor-Schema-Version` from whichever JARs carry them. The resolved schema version
is then passed as the `schema_version` return value from `_build_execute_task_command`, which
Expand All @@ -74,7 +75,7 @@ the base `SubprocessCoordinator` uses to negotiate the supervisor wire protocol.
If `main_class` is set explicitly on the `JavaCoordinator` instance (via `[sdk] coordinators`
kwargs), the scan uses it as a filter; otherwise the first JAR with a `Main-Class` attribute
wins. Either way, `Airflow-Supervisor-Schema-Version` must be present in at least one JAR in
`jars_root` or startup fails.
the Dag bundle or startup fails. Every JAR in the Dag bundle goes on one classpath.

---

Expand Down Expand Up @@ -117,7 +118,7 @@ E2E_TEST_MODE=java_sdk uv run --project airflow-e2e-tests pytest \

`coordinator.py` extends `SubprocessCoordinator`. The only method subclasses must implement is
`_build_execute_task_command`, which returns `(argv, schema_version)`. Look at the existing
implementation for how `jars_root`, `java_executable`, `jvm_args`, and `main_class` are
implementation for how the scanned Dag bundle, `java_executable`, `jvm_args`, and `main_class` are
assembled into the command. Do not reach into the JVM process from Python beyond what this
method provides.

Expand Down
40 changes: 23 additions & 17 deletions airflow-core/adr/lang-sdk/0011-mixed-language-dag-processing.md
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ DagFileProcessorProcess(etl.py) ← manager spawn
├── _serialize_dags(bag) → is_stub tasks carry arg_bindings (ADR-0007)
│
│ ┌─────────────────────────────────────────────────────────────────────┐
│ │ Step 1: Collect the dag_ids to ask about — every Dag in this │
│ │ Step 1: Collect the Dags to validate — every Dag in this │
│ │ file with at least one is_stub task → ["etl"] │
│ └─────────────────────────────────────────────────────────────────────┘
│
Expand Down Expand Up @@ -123,27 +123,26 @@ DagFileProcessorProcess(etl.py) ← manager spawn
│ │ └── BundleScanner.scanBundles(roots) │
│ │ → "etl" → ResolvedBundle(analytics.jar, mainClass, ...) │
│ │ │
│ │ Group the dag_ids by (coordinator, artifact) — one group, one │
│ │ Group the stub tasks 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, │
│ │ LangSDKTaskHandlerProcessorProcess.start( │
│ │ target=_start_task_handler_runtime_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 │
│ │ │ ──TaskHandlerParseRequest(file=analytics.jar)─────▶ JVM │
│ │ │ (ToSDKTaskHandlerProcessor) │
│ │ │ │
│ │ │ JVM answers from its own TaskHandler registrations │
│ │ │ whose dagId is one of the requested ids │
│ │ │ JVM answers with every TaskHandler registration │
│ │ │ in the artifact, or {} when there is none │
│ │ │ │
│ │ │ ◀─TaskHandlerParsingResult(task_handlers={ │
│ │ │ "etl": [extract, transform, load]})─────── JVM │
Expand All @@ -162,7 +161,9 @@ DagFileProcessorProcess(etl.py) ← manager spawn
│ │ Python Dag "etl" (stub tasks) TaskHandlerDeclaration │
│ │ ────────────────────────────── ────────────────────────────── │
│ │ task_id ↔ task_id (sets must match) │
│ │ arg_bindings[*].name ↔ params[*].name (in order) │
│ │ arg_bindings[*] ↔ params[*] (per binding) │
│ │ by position, or by folded or exact name │
│ │ unmatched by name → a warning, not an error │
│ │ arg_bindings[*].value_schema ↔ params[*].value_schema │
│ │ compared only where neither side is null │
│ │ │
Expand All @@ -175,8 +176,8 @@ DagFileProcessorProcess(etl.py) ← manager spawn
```

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)).
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 per
artifact serve every Dag in the file whose stubs resolved to it ([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.
Expand All @@ -186,7 +187,7 @@ Resolution goes through the coordinator registry, not the filesystem, so the Pyt
| 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`, per (coordinator, artifact) | `parse_task_handler`, for every handler it registers | `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.
Expand All @@ -197,8 +198,9 @@ There is no fourth row. A `TaskHandlerRef` has no Dag, so no `DagImporter` — a
- 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.
- Stub/implementation mismatches — missing handler, extra handler, a `positional` argument count that does not match, incompatible schema — surface as import errors against
the Python file at parse time, alongside the errors the parse already reports. A `named` argument or parameter that matches nothing is only logged as a warning, because the
runtime runs the task anyway. An unannotated stub argument is checked only for how it binds, by position or by name.
- 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.
Expand Down Expand Up @@ -234,12 +236,16 @@ with one `bundle.serve()`. One artifact can hold both, so neither the file nor t

### 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.
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. Each reply covers
every Dag its artifact registers handlers for; only the parsed file's Dags are compared.

- `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` against `params` as each declaration's `binding` says: by position for `positional`, by name for `named`, compared case-insensitively with underscores ignored
unless `exact_name` is set. A `positional` count mismatch is an error. Under `named`, an argument or parameter that matches nothing is only logged as a warning, and a lone
unmatched argument, which may be the whole value, is not ([ADR-0012](0012-lang-sdk-parse-protocol.md) Appendix B). When `params` is `None`, only the handler's presence is
checked.
- `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
Every error 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.
Loading
Loading