Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
132 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
0669c86
Carry task handler bindings and probed answers on the parse result
jason810496 Sep 30, 2026
a7df3be
Regenerate TS SDK supervisor types for task handler bindings
jason810496 Sep 30, 2026
fb12a76
Sync the Java SDK supervisor schema for task handler bindings
jason810496 Sep 30, 2026
c6c7007
Regenerate Go SDK models for task handler bindings
jason810496 Oct 2, 2026
1665a67
Record probed artifacts and reconcile task handler bindings
jason810496 Sep 30, 2026
f7fd2ba
Index task handler artifacts by last_probed_at
jason810496 Sep 30, 2026
b1fb803
Delete task handler artifacts no stub task references
jason810496 Sep 30, 2026
f2c1b9e
Describe the binding reconcile and artifact sweep in ADR-0013
jason810496 Sep 30, 2026
36e5a7a
Give the task handler race tests their own sessions
jason810496 Sep 30, 2026
3cccb23
Warn only about bindings of Dags missing from the parse result
jason810496 Sep 30, 2026
459da04
Align the artifact sweep with the other cleanup checks
jason810496 Sep 30, 2026
aa3f771
Build the binding and probed-artifact test params without required
jason810496 Oct 1, 2026
e48f1fb
Let a coordinator list the artifacts it can probe
jason810496 Oct 1, 2026
ff88238
List Go task handler bundles by their trailer
jason810496 Oct 1, 2026
908d19c
List Java task handler JARs by their cache digest
jason810496 Oct 1, 2026
9ce2ed7
Read a TS bundle's cache digest from its layout line
jason810496 Oct 1, 2026
d851bdf
Decide per candidate whether a parse must probe it
jason810496 Oct 1, 2026
d238a18
Describe the per-candidate fast path in ADR-0013
jason810496 Oct 1, 2026
bd1191b
List only .zip archives as packaged Dags
jason810496 Oct 1, 2026
ec470f5
Add a significant newsfragment for listing only .zip packaged Dags
jason810496 Oct 2, 2026
e1cacea
Check stub task arguments against task handler declarations
jason810496 Oct 2, 2026
a857663
Describe how a stub task is checked against its task handler
jason810496 Oct 2, 2026
a71effa
Let the coordinator registry name the coordinator a queue routes to
jason810496 Oct 2, 2026
41383de
Read a Dag bundle's team from the bundle configuration
jason810496 Oct 2, 2026
3c89e6c
Stop a task handler probe at the Dag file's deadline
jason810496 Oct 2, 2026
9f872b8
Resolve a Dag file's stub tasks to task handler artifacts
jason810496 Oct 2, 2026
79eafdf
Check stub tasks against task handlers when a Dag file is parsed
jason810496 Oct 2, 2026
5274edc
Check the example Dags against the example bundles
jason810496 Oct 2, 2026
5f17455
Give the Dag processor the e2e task handler artifacts and runtimes
jason810496 Oct 2, 2026
ce1f57c
Document what the Dag processor checks in a stub task
jason810496 Oct 2, 2026
8efe1ef
Describe the stub-task check and its failures in the ADRs
jason810496 Oct 2, 2026
c0f84c4
Add the task handler artifact to StartupDetails
jason810496 Oct 1, 2026
8cb04dd
Regenerate TS SDK supervisor types for the task handler artifact
jason810496 Oct 1, 2026
ed6341a
Sync the Java SDK supervisor schema for the task handler artifact
jason810496 Oct 1, 2026
94fcf88
Regenerate Go SDK models for the task handler artifact
jason810496 Oct 2, 2026
71b451f
Resolve a stub task's artifact bundle from its workload
jason810496 Oct 1, 2026
9f61393
Add the task handler artifact to the ExecuteTask workload
jason810496 Oct 1, 2026
0e9b942
Name the task handler artifact reference in ADR-0013
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
67 changes: 39 additions & 28 deletions airflow-core/adr/lang-sdk/0011-mixed-language-dag-processing.md

Large diffs are not rendered by default.

43 changes: 28 additions & 15 deletions airflow-core/adr/lang-sdk/0012-lang-sdk-parse-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,8 @@ 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)).
defines. "Which task handlers does this artifact register, each 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.

Expand Down Expand Up @@ -82,7 +82,6 @@ relay is only type-safe because both hops speak the same pair, which is the reas
```
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"]
Expand All @@ -96,20 +95,24 @@ class TaskHandlerParsingResult:

class TaskHandlerDeclaration:
task_id: str
params: list[TaskHandlerParam] # ordered — arg_bindings are positional
binding: Literal["positional", "named"] # how stub-task arguments bind to params
params: list[TaskHandlerParam] | None # ordered; the order matters only for "positional". None: the runtime cannot list them

class TaskHandlerParam:
name: str
name: str | None # None: the runtime has no name for this positional parameter
value_schema: ArgValueSchema | None = None
required: bool # the handler declares no default
exact_name: bool = False # match as spelled, not case-insensitively with underscores ignored
```

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)).
The request names no Dags. The runtime answers with every task handler the artifact registers, keyed by `dag_id`, and with `{}` when it registers none. The answer must depend only
on the artifact, never on the request, so the Dag processor can cache it per artifact ([ADR-0013](0013-persisted-task-handler-bindings.md)). One request per candidate artifact
answers for every Dag it registers. The cost of one process per (coordinator, artifact) pair a file's stubs resolve to is superseded by
[ADR-0013](0013-persisted-task-handler-bindings.md) Flow 1: on a cold start the Dag processor probes each candidate without a recorded answer until a coordinator that lists it gets
an answer, and in the steady state it probes none. The key set is not required to match the file's Dags: it can hold other files' Dags, and each stub task needs its handler among
the answers of its own coordinator ([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.
same type. Two properties matter to validation: the field is nullable on both sides, and each declaration names its binding mode. 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.
Expand All @@ -129,8 +132,8 @@ WatchedSubprocess
│ └── 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
└── LangSDKTaskHandlerProcessorProcess (new — ADR-0011)
target = _start_task_handler_runtime_entrypoint
└── coordinator.parse_task_handler() — same forwarding
TaskHandlerParseRequest → TaskHandlerParsingResult
```
Expand Down Expand Up @@ -208,13 +211,23 @@ bodies that never differ. Reusing `ToManager` leaves one reply union with one ne
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
### Appendix B — What the nullable schema and the binding mode force

`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.
The stub side always has names and positions: `LiteralArgBinding` and `XComArgBinding` are each documented as "one positional stub-task argument". The handler side binds the way
its runtime does, so each declaration names its `binding` and the check follows it:

- `positional`: by position. Names are informative only, and absent where the runtime has none (Go flat params). Java's `TaskArgs` binds this way although it has names. An
argument count that matches `params` neither with every argument nor after dropping the defaulted ones, or a value type a parameter does not accept, is an import error.
- `named`: by name in any order, case-insensitively with underscores ignored unless `exact_name` is set (Go `arg:` tags, explicit Java names). A Go struct, tagged or untagged,
and Java's `TaskInput` bind this way. An argument no parameter takes, or a parameter no argument fills, is logged as a warning, and the task still runs: the runtimes allow
both, and an unfilled field keeps its default. When no parameter matches and exactly one argument was passed, it may be the whole value and is not warned about, unless no
parameter is declared, a parameter sets `exact_name` (a Go `arg:` tag), or the argument cannot be an object. A value type a parameter does not accept is an import error.

`params` is `None` when the runtime cannot list a handler's parameters, and then only the handler's presence is checked. TypeScript declares `named` with `params: None`: its
types are erased, so a handler cannot list what it takes.

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.
Expand Down
Loading
Loading