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 e8d43932077 Add ADR for supporting Dynamic Dag Generation across
Language SDKs (#73518)
e8d43932077 is described below
commit e8d439320777c2dd2e708fef8ee73757fe4c98f0
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Tue Sep 29 21:00:07 2026 +0800
Add ADR for supporting Dynamic Dag Generation across Language SDKs (#73518)
* Add ADRs for persisted Lang-SDK task-handler bindings
Mixed-language Dags rediscover which artifact backs a stub task on every
task execution, by walking a filesystem root and matching against a Dag
inventory recorded when the artifact was packed. That inventory is only
correct when registration depends on nothing the build environment
lacks, so Dags whose ids are generated at runtime cannot be packed at
all -- both packers treat an empty inventory as fatal.
Record the decision to resolve the binding once during Dag processing
and persist it, so execution reads a path instead of searching for one,
and to replace the filesystem root with a named DagBundle that brings
download, refresh and versioning with it.
ADR-0011 is kept separate because it changes three SDK artifact formats
rather than core, and because it supersedes a requirement two accepted
go-sdk ADRs state outright.
* Strip Dag identifiers from Lang-SDK artifact metadata
An identifier recorded when an artifact is built is a claim about runtime
behaviour made by a process that cannot observe it. Dynamic Dag
rendering makes the claim plainly false, but a claim that happens to
hold today is still a second source of truth free to drift tomorrow.
Keeping the inventory as advisory metadata would preserve that drift
while removing the consumer that would have caught it, so record that
the artifact carries no dag_id or task_id in either role -- mixed
language task handler or native Dag -- rather than renaming the mapping.
Also record why a dedicated Execution API for the binding was rejected:
the API is a versioned contract every task SDK speaks, which makes a new
route there a far larger commitment than a field on a workload the
scheduler already builds.
Prefix the new wire objects with SDK so they read consistently beside
ToSDKTaskHandlerProcessor, and link every code reference to a pinned
commit so the citations stay meaningful as main moves.
* Sharpen the case against compile-time Dag identifiers
The two separate complaints about the build-time inventory -- that it is
wrong by construction, and that coordinators route on it -- are one
problem seen twice. Freezing dag_ids and task_ids when the artifact is
built is what makes runtime-generated Dag ids unroutable, and the
artifact unpackable at all. State it once.
Show that Dag processing records nothing about the artifact today, since
that absence is what forces the per-task scan, and show the artifact
metadata as it stands so the identifiers being removed are visible
rather than described.
Give the config example its DagBundle registration and the queue to
bundle to path resolution chain, so the relationship between a
coordinator and the bundle it draws artifacts from is on the page.
* Trim ADR-0010 to the decisions it is making
Several passages argued their case by walking through code that a
reader does not need in order to follow the decision: which coordinator
consults which mapping, which existing table the new one resembles,
which helper to copy a reconcile from, and why a neighbouring scheduler
branch handles its own failure badly. That detail belongs in review, or
in the code, not in the record of what was decided.
State the compile-time freeze as the flow it is -- packing records the
ids, so an id chosen at runtime can never match -- and let the rest go.
Move the note about deployments that mount artifacts themselves ahead of
the config example, where it answers the question the example raises
rather than trailing it.
* Simplify ADR-0010 wording and correct two rejection reasons
The Context paragraph argued from implementation where the flow was
enough, and several Consequences carried transport and per-SDK detail
that a reader does not need to weigh the decision.
Two rejected alternatives were also recorded for the wrong reason.
Letting the artifact write its own rows fails first because persisting
them is the manager's responsibility rather than the task subprocess's,
with the cold-start ordering a secondary cost; and the denormalised
table fails on a race between parse children writing disagreeing
digests, which is what makes the duplicate copies unsafe rather than
merely redundant.
Prefer parentheses to paired em dashes throughout, so asides read as
asides.
* Cut ADR-0011 back to its decision
The Decision section restated its own reasoning several times over: why
the identifiers are wrong, which coordinators read them today, and what
each accepted ADR said about discovery. The YAML example already shows
what is removed, so the prose around it can be brief.
Keep the metadata version at 1.0. Dropping a key that nothing reads is
not a breaking change, so it does not earn a version bump or the
compatibility window a bump would imply.
* Remove the useless content myself
* Remove the references and appendix
* Remove trailing blank lines from Lang-SDK ADRs
* Fix unbalanced parentheses in ADR-0010
* Renumber persisted-binding ADRs to 0013 and 0014
---
.../0013-persisted-task-handler-bindings.md | 593 +++++++++++++++++++++
.../0014-bundle-metadata-and-cache-digest.md | 122 +++++
airflow-core/adr/lang-sdk/README.md | 2 +
3 files changed, 717 insertions(+)
diff --git a/airflow-core/adr/lang-sdk/0013-persisted-task-handler-bindings.md
b/airflow-core/adr/lang-sdk/0013-persisted-task-handler-bindings.md
new file mode 100644
index 00000000000..cbcace7df52
--- /dev/null
+++ b/airflow-core/adr/lang-sdk/0013-persisted-task-handler-bindings.md
@@ -0,0 +1,593 @@
+<!--
+ 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-0013: Persisted Task-Handler Bindings (Resolving Lang-SDK Artifacts at
Parse Time)
+
+## Status
+
+Proposed
+
+## Context
+
+A mixed-language Dag is authored in Python with `@task.stub` tasks whose
bodies live in a Lang-SDK
+artifact (a packed Go binary, a JAR, a minified `.min.mjs`). Nothing in the
Dag says *which* artifact.
+Nothing is recorded about that artifact when the Dag is processed, so the link
has to be
+rediscovered on every single task execution by scanning a filesystem root:
+
+```
+DAG PROCESSING stores nothing about the
artifact
+ DagFileProcessorProcess(etl.py)
+ └── PythonDagImporter → Dags with @task.stub tasks
+ └── persist DagModel, SerializedDagModel, DagVersion, DagCode
+ ┌──────────────────────────────────────────────────────────┐
+ │ no artifact path recorded │
+ │ no artifact bundle recorded │
+ │ the parse never even looks at the Lang-SDK artifact │
+ └──────────────────────────────────────────────────────────┘
+
+TASK EXECUTION must therefore search, every
time
+ ExecutableCoordinator._build_execute_task_command(what=ti)
+ └── _Bundle.find(executables_root, what.dag_id)
+ └── walk every executable file under the root
+ read its trailer, verify SHA-256 over the binary region
+ parse its metadata, test `dag_id in metadata["dags"]`
+```
+
+Two problems compound here.
+
+**The scan is per task.** Because Dag processing records nothing, every task
execution re-walks the
+root and re-hashes candidates to answer a question whose answer changed only
when someone deployed.
+
+**The identifiers it scans are frozen when the artifact is built.** Packing an
artifact records the
+Dag ids and task ids it exposes into the artifact's own metadata at the
packing stage, and they are
+fixed from then on. Currently, a coordinator picks an artifact by looking a
`dag_id` up in that
+recorded list. A Dag whose id the artifact only decides on when it runs (for
example: generated from
+an external YAML) can never be matched to it.
+
+Separately, `[sdk] coordinators` locates artifacts through filesystem roots
(`jars_root`,
+`executables_root`, `bundles_root`
([ADR-0005](0005-coordinator-packaging.md))) which are
+unversioned mutable directories outside any `DagBundle`, with their own
delivery problem. Deployments
+already solve that delivery by staging a `DagBundle` *into* the root
+([`stage_artifacts.py`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/kubernetes-tests/lang_sdk/stage_artifacts.py)),
which makes the root a second addressing layer over
+a mechanism that already addresses and versions artifacts.
+
+This ADR replaces runtime discovery with a binding resolved once during Dag
processing and persisted,
+and replaces the filesystem root with a named `DagBundle`.
+
+Native Dags (a Dag authored entirely in a Lang SDK) are **out of scope** here;
they arrive through an
+ordinary `DagBundle` and a Dag importer, and are not yet recorded in an ADR.
+
+## Decision
+
+### Artifacts live in a named DagBundle, not a filesystem root
+
+`jars_root` / `executables_root` / `bundles_root` are replaced by a single
coordinator kwarg naming a
+`DagBundle`.
+
+Nothing is taken away from deployments that want to place artifacts
themselves. A `LocalDagBundle`
+pointed at the mount does exactly what an explicit root did (the directory is
still theirs to
+manage) but it arrives through the same mechanism as every other bundle rather
than beside it, so
+it inherits refresh and the rest without special-casing:
+
+```ini
+[dag_processor]
+# The artifact bundle is an ordinary DagBundle, registered like any other.
+dag_bundle_config_list = [
+ {
+ "name": "dags-folder",
+ "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle",
+ "kwargs": {}
+ },
+ {
+ "name": "java-task-handlers",
+ "classpath": "airflow.providers.amazon.aws.bundles.s3.S3DagBundle",
+ "kwargs": {"bucket_name": "artifacts", "prefix": "java",
"aws_conn_id": "aws_default"}
+ }
+]
+
+[sdk]
+coordinators = {
+ "jdk-17": {
+ "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
+ "kwargs": {
+ "java_executable": "/usr/lib/jvm/java-17/bin/java",
+ "task_handler_bundle_name": "java-task-handlers"
+ }
+ }
+}
+queue_to_coordinator = {"java": "jdk-17"}
+```
+
+```
[email protected](queue="java") the Dag author picks a queue
+ │
+ ▼ [sdk] queue_to_coordinator
+ "jdk-17" the coordinator instance
+ │
+ ▼ [sdk] coordinators → kwargs.task_handler_bundle_name
+ "java-task-handlers" the bundle name
+ │
+ ▼ [dag_processor] dag_bundle_config_list
+ S3DagBundle(bucket=artifacts, prefix=java)
+ │
+ ▼ DagBundlesManager().get_bundle(name).initialize()
+ bundle.path / <artifact_rel_path>
+```
+
+The Python Dag file and the artifact sit in different bundles (`dags-folder`
and
+`java-task-handlers` above) and that is the expected layout, not a workaround.
Binaries and JARs do
+not belong in the bundle holding `.py` files. Both are registered with the Dag
processor, because
+registration is what makes `get_bundle(name)` resolvable on the worker.
+
+The name is `task_handler_bundle_name`, not `..._bundle_path`: the point of
routing through a bundle
+is that `DagBundlesManager` owns download, refresh and versioning. A path
would keep the unversioned
+mutable directory and discard all of it. `sdk_` is omitted because the kwarg
is already scoped to a
+coordinator instance.
+
+The kwarg is mixed-language only. A native Lang-SDK Dag is delivered by
whichever `DagBundle` the Dag
+processor is scanning, exactly like a `.py` file, and never by coordinator
configuration. A single
+coordinator instance can serve both roles at once, so nothing may assume one
coordinator maps to one
+bundle.
+
+The name is validated eagerly when the coordinator registry is built,
alongside the existing
+validation of every `queue_to_coordinator` key
+([`from_config`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/task-sdk/src/airflow/sdk/execution_time/coordinator.py#L262-L264)),
so a typo surfaces at config load
+rather than as a lazy `InvalidCoordinatorError` on the first task.
+
+### Two tables
+
+The resolved binding is persisted. The artifact is normalised out, because one
artifact typically
+backs many handlers and its fingerprint must have exactly one value.
+
+```sql
+CREATE TABLE lang_sdk_task_handler_artifact (
+ id UUID NOT NULL,
+ bundle_name VARCHAR(250) NOT NULL, -- the
task_handler_bundle_name it was found in
+ relative_fileloc VARCHAR(2000) NOT NULL, -- path within that bundle
+ size_bytes BIGINT NOT NULL, -- cheap fingerprint tier
+ cache_digest VARCHAR(64) NOT NULL, -- content fingerprint tier;
see "The fast path"
+ last_probed_at TIMESTAMP NOT NULL,
+ PRIMARY KEY (id),
+ CONSTRAINT lstha_bundle_fileloc_uq UNIQUE (bundle_name, relative_fileloc)
+);
+
+CREATE TABLE lang_sdk_task_handler (
+ dag_id VARCHAR(250) NOT NULL,
+ task_id VARCHAR(250) NOT NULL,
+ artifact_id UUID NOT NULL,
+ dag_bundle_name VARCHAR(250) NOT NULL, -- the *Python* file that
owns this row
+ dag_relative_fileloc VARCHAR(2000) NOT NULL, -- ditto
+ handler_params JSON NOT NULL, -- list[TaskHandlerParam],
ordered
+ PRIMARY KEY (dag_id, task_id),
+ CONSTRAINT lsth_dag_fkey FOREIGN KEY (dag_id)
+ REFERENCES dag (dag_id) ON DELETE CASCADE,
+ CONSTRAINT lsth_artifact_fkey FOREIGN KEY (artifact_id)
+ REFERENCES lang_sdk_task_handler_artifact (id)
+);
+CREATE INDEX idx_lsth_dag_file ON lang_sdk_task_handler (dag_bundle_name,
dag_relative_fileloc);
+CREATE INDEX idx_lsth_artifact_id ON lang_sdk_task_handler (artifact_id);
+```
+
+`PRIMARY KEY (dag_id, task_id)` is the conflict guard: two artifacts claiming
the same task cannot
+both be recorded, and the collision is detected during the parse and reported
as an import error
+rather than resolved by scan order.
+
+`dag_bundle_name` / `dag_relative_fileloc` identify the Python file that owns
the row. They exist so
+the Dag processor manager can look up prior state by the file it is about to
dispatch, **without**
+joining through `DagModel`: `dag.relative_fileloc` is not indexed, and the
codebase already notes
+that querying it means "a sequential scan of dag"
+([`reassign_dags_with_unconfigured_bundles`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/bundles/manager.py#L478)).
+
+`handler_params` stores what the runtime declared, so a changed Python file
can be re-validated
+against a cached declaration with no subprocess. It is deliberately **not**
called `arg_bindings`:
+that name already denotes the Python side of the comparison (`XComArgBinding`
/ `LiteralArgBinding`,
+carrying wiring and values,
[`build_arg_bindings`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/serialization/stub_arg_bindings.py#L221-L286)),
+and reusing it would make the validation read as comparing a thing to itself.
+
+`cache_digest` is **opaque and coordinator-defined**, not "SHA-256 of the
file".
+
+### Objects on the wire
+
+**Manager → Dag-parsing child.**
[`DagFileParseRequest`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/processor.py#L113-L130)
gains the artifacts the manager
+already knows about. The parse child processor subprocesses run in the client
context without a
+database connection
([`_parse_file_entrypoint`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/processor.py#L208-L232)),
+so the prior cache state must be pushed down from the manager. The alternative
is a dedicated
+Execution API for the child processor process to retrieve
`KnownSDKTaskHandlerArtifact` itself, which
+was rejected on blast radius.
+
+```python
+class KnownSDKTaskHandlerArtifact(BaseModel):
+ bundle_name: str
+ relative_fileloc: str
+ size_bytes: int
+ cache_digest: str
+
+
+class DagFileParseRequest(BaseModel):
+ file: str
+ bundle_path: Path
+ bundle_name: str
+ callback_requests: list[CallbackRequest]
+ known_artifacts: list[KnownSDKTaskHandlerArtifact] = [] # new
+ type: Literal["DagFileParseRequest"]
+```
+
+`known_artifacts` is scoped by *artifact* bundle, not by Dag file, so the
manager reads it **once per
+parsing loop** for every configured `task_handler_bundle_name` and pushes the
same list to every
+child. Ten Dag files resolving against one twenty-jar bundle therefore probe
that bundle once in
+total, not once each.
+
+**Dag-parsing child → coordinator subprocess.** Introduced here. The child
spawns the runtime and
+forwards bytes in both directions, decoding nothing; the process that spawned
the parse decodes the
+reply. `ToSDKTaskHandlerProcessor` is a new parent-to-child union differing
from `ToDagProcessor` in
+one member, and `ToManager` gains `SDKTaskHandlerParsingResult`.
+
+```python
+class SDKTaskHandlerParseRequest(BaseModel): # parent -> runtime, on
ToSDKTaskHandlerProcessor
+ file: str # the candidate artifact being probed
+ dag_ids: list[str] # every Dag in this file with stub tasks routed here
+ bundle_path: Path
+ bundle_name: str
+ type: Literal["SDKTaskHandlerParseRequest"]
+
+
+class SDKTaskHandlerParsingResult(BaseModel): # runtime -> parent, on
ToManager
+ fileloc: str
+ task_handlers: dict[str, list[TaskHandlerDeclaration]] # dag_id ->
declarations
+ import_errors: dict[str, str] | None = None
+ warnings: list | None = None
+ type: Literal["SDKTaskHandlerParsingResult"]
+
+
+class TaskHandlerDeclaration(BaseModel):
+ task_id: str
+ params: list[TaskHandlerParam] # ordered; arg bindings are positional
+
+
+class TaskHandlerParam(BaseModel):
+ name: str
+ value_schema: JSONSchema | None = None
+ required: bool # the handler declares no default
+```
+
+A `dag_id` the artifact registers nothing for is **omitted** from
`task_handlers` rather than returned
+empty, so a probe that matches nothing is distinguishable from a probe that
matched a Dag with zero
+tasks.
+
+**Dag-parsing child → manager.**
[`DagFileParsingResult`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/processor.py#L133-L145)
gains the resolved
+bindings.
+
+```python
+class SDKTaskHandlerBinding(BaseModel):
+ dag_id: str
+ task_id: str
+ artifact_bundle_name: str
+ artifact_rel_path: str
+ artifact_size_bytes: int
+ artifact_cache_digest: str
+ handler_params: list[TaskHandlerParam]
+
+
+class DagFileParsingResult(BaseModel):
+ fileloc: str
+ serialized_dags: list[LazyDeserializedDAG]
+ warnings: list | None = None
+ import_errors: dict[str, str] | None = None
+ task_handler_bindings: list[SDKTaskHandlerBinding] | None = None # new
+```
+
+`None` and `[]` mean different things, and the difference is load-bearing:
+
+| value | meaning | manager does |
+|--------|------------------------------------------|-----------------------|
+| `None` | handlers were not evaluated in this parse | **nothing** — no
reconcile |
+| `[]` | evaluated, this file has no stub handlers | delete this file's rows
|
+| `[…]` | evaluated, these are the bindings | reconcile to this set
|
+
+`None` covers the stability-check early return
([`_parse_file`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/processor.py#L245-L251)),
callback-only runs
+([`_parse_file`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/processor.py#L260-L263)),
and validation failure (below). It mirrors the `files_parsed=None` semantics
+`persist_parsing_result` already uses
([`persist_parsing_result`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/airflow-core/src/airflow/dag_processing/manager.py#L1361-L1364)).
+Without it, a transient parse failure would silently wipe every binding the
file owns.
+
+On a fast-path skip the child **re-emits the bindings it was given**,
unchanged. It does not omit
+them. Request and result carrying the same information makes that a copy
rather than a special
+"keep these" signal the reconcile could get wrong.
+
+**Scheduler → worker.** `ExecuteTask` and `StartupDetails` each gain one
optional reference. The
+artifact bundle is a second, independent bundle, so it needs its own
`BundleInfo`.
+
+```python
+class SDKTaskHandlerRef(BaseModel):
+ bundle_info: BundleInfo # the artifact bundle: name, version, version_data
+ rel_path: str # path within it
+
+
+class ExecuteTask(BaseDagBundleWorkload):
+ ti: TaskInstanceDTO
+ dag_rel_path: os.PathLike[str] # the Python Dag file, unchanged
+ bundle_info: BundleInfo # the Dag bundle, unchanged
+ task_handler: SDKTaskHandlerRef | None = None # new
+ ...
+
+
+class StartupDetails(BaseModel):
+ ti: TaskInstance
+ dag_rel_path: str
+ bundle_info: BundleInfo
+ task_handler: SDKTaskHandlerRef | None = None # new
+ ...
+```
+
+`None` means "this task needs no Lang-SDK artifact" (an ordinary Python task).
It never means
+"unknown": a stub task that failed to resolve is not queued at all (see
"Failure handling").
+
+### Flow 1) Dag processing: the write path
+
+Resolution is driven by the Python file's parse, after `PythonDagImporter` has
produced the Dags.
+The artifact never parses itself into these tables. That ordering is
deliberate: a `dag_id` exists
+before any row referencing it, so the foreign key holds and a cold start has
no window in which a
+stub Dag is validated against bindings that have not been written yet.
+
+```
+DagProcessorManager [reads DB]
+ │
+ │ once per parsing loop, per configured task_handler_bundle_name:
+ │ SELECT bundle_name, relative_fileloc, size_bytes, cache_digest
+ │ FROM lang_sdk_task_handler_artifact
+ │ WHERE bundle_name IN (:configured bundles) ──▶ known_artifacts
+ │
+ │ per file about to be dispatched:
+ │ (the child re-reads nothing; it has no DB)
+ │
+ ├── DagFileParseRequest(file=etl.py, bundle_*, known_artifacts=[...])
+ ▼
+DagFileProcessorProcess(etl.py) [no DB — client
context]
+ └── _parse_file_entrypoint → _parse_file
+ │
+ ├─1─ BundleDagBag → PythonDagImporter → airflow.sdk.DAG objects
+ │
+ ├─2─ _serialize_dags(bag)
+ │ is_stub tasks now carry arg_bindings (ADR-0007)
+ │
+ ├─3─ collect stub tasks, group by coordinator
+ │ stub task queue coordinator task_handler_bundle_name
+ │
─────────────────────────────────────────────────────────────────
+ │ extract → "java" → jdk-17 → "java-task-handlers"
+ │ transform → "java" → jdk-17 → "java-task-handlers"
+ │ ingest → "go" → go-sdk → "go-task-handlers"
+ │
+ │ a queue with no coordinator entry is an import error here,
+ │ not a silent Python fallback at execution time
+ │
+ ├─4─ per coordinator: list candidates in its bundle
+ │ Go → files carrying the AFBNDL01 trailer magic
+ │ Java → *.jar with a Main-Class manifest attribute
+ │ TS → *.min.mjs with a valid //# airflowBundle= layout header
+ │ walk order is deterministic, so conflicts reproduce
+ │
+ ├─5─ FAST PATH, per candidate (see "The fast path")
+ │ size + cache_digest match known_artifacts, and the candidate
+ │ set is unchanged ──▶ skip the launch, echo the bindings
+ │ anything differs ──▶ probe
+ │
+ ├─6─ PROBE, per differing candidate — one subprocess
+ │ SDKTaskHandlerProcessorProcess.start(
+ │ target=_parse_task_handler_entrypoint,
+ │ coordinator=JavaCoordinator("jdk-17"),
+ │ path=<candidate>)
+ │ ──SDKTaskHandlerParseRequest(file=…, dag_ids=["etl"])──▶
runtime
+ │ ◀─SDKTaskHandlerParsingResult(task_handlers={"etl": […]})──
runtime
+ │ Get* from the runtime is relayed up ToManager unchanged
+ │
+ ├─7─ VALIDATE per dag_id, unioned across coordinators
+ │ task_id sets must match exactly
+ │ arg_bindings[*].name ↔ handler_params[*].name, in order
+ │ arg_bindings[*].schema ↔ handler_params[*].value_schema,
+ │ compared only where neither is null
+ │ two candidates claiming one (dag_id, task_id) → import error
+ │ naming both
paths
+ │
+ └─8─ on success → task_handler_bindings=[…]
+ on mismatch → import_errors[etl.py]=… AND bindings=None
+ │
+ ├── DagFileParsingResult(serialized_dags=[…], import_errors={…},
+ ▼ task_handler_bindings=[…] | [] | None)
+DagProcessorManager.persist_parsing_result [writes DB]
+ └── update_dag_parsing_results_in_db — one transaction, in this order:
+ 1. add_dags / update_dags → DagModel rows exist
+ 2. asset reference tables (existing)
+ 3. lang_sdk_task_handler_artifact UPSERT by (bundle_name,
relative_fileloc) ← new
+ 4. lang_sdk_task_handler reconcile by dag_id
← new
+ 5. SerializedDagModel / DagVersion / DagCode
+ 6. ParseImportError / DagWarning
+```
+
+Step 3 must upsert: two parse children can discover the same artifact in the
same loop and race on
+the unique key.
+
+Step 4 reconciles **by the `dag_id`s in the result**, not by file path: delete
rows for those
+`dag_id`s whose `task_id` is absent from the returned set, then insert or
update the rest. Path-keyed
+eviction breaks when a Dag moves between files — the old rows stay keyed to a
path nothing parses any
+more, and the primary key then blocks the new insert. With `dag_id` as the key
a move simply updates
+`dag_relative_fileloc`, and a Dag that disappears entirely is reclaimed by the
`ON DELETE CASCADE`.
+
+Artifact rows are **never** evicted from one file's result. One artifact backs
handlers owned by many
+Python files, so this file seeing fewer candidates says nothing about another
file's. They are a
+cache; they are reclaimed by orphan sweep or by `db clean`, never by a
per-file reconcile.
+
+### Flow 2) Scheduling: the read path
+
+```
+SchedulerJobRunner._executable_task_instances_to_queued [reads DB]
+ │
+ ├── SELECT TI JOIN dag_run JOIN dag … WHERE DR.state=RUNNING
+ │ AND TI.state=SCHEDULED …
+ │ .options(selectinload(TI.dag_model))
+ │ .options(joinedload(TI.dag_run).selectinload(DagRun.created_dag_version)
+ │ .load_only(DagVersion.version_data))
+ │ .options(<the task-handler binding>) ←
new
+ │
+ │ the binding MUST be loaded here, before make_transient below.
+ │ a lazy load on a transient object returns None silently instead of
+ │ raising DetachedInstanceError, so a late read yields a workload with
+ │ no artifact and no error.
+ │
+ ├── make_transient(ti) for every returned TI
+ ▼
+SchedulerJobRunner._enqueue_task_instances_with_queued_state [no further DB
reads on ti]
+ └── for ti in task_instances:
+ dag_run finished? → set_state(None); continue (existing)
+ no dag_version_id? → warn; continue (existing)
+ is_stub and no binding? → FAIL the TI with the reason ← new
+ never queue it
+ │
+ └── ExecuteTask.make(ti, task_handler=SDKTaskHandlerRef(...))
+ dag_rel_path ← ti.dag_model.relative_fileloc (existing)
+ bundle_info ← Dag bundle, pinned to the run (existing)
+ task_handler ← artifact bundle + rel_path ← new
+ │
+ └── executor.queue_workload(workload)
+```
+
+A stub task with no binding is **failed with its reason**, not skipped.
+
+### Flow 3) Task execution
+
+```
+executor worker process
+ └── BaseExecutor.run_workload(workload)
+ └── supervise_task(ti=…, bundle_info=…, dag_rel_path=…,
+ task_handler=workload.task_handler) ← new
+ │
+ ├── coordinator = get_coordinator_manager().for_queue(ti.queue)
+ │ unchanged: execution still routes on queue
+ │
+ └── coordinator.execute_task(what=ti, …, task_handler=…)
+ └── SubprocessCoordinator._build_execute_task_command(
+ what=ti, task_handler=task_handler) ←
signature change
+ │
+ ├── bundle = DagBundlesManager().get_bundle(
+ │ name=task_handler.bundle_info.name,
+ │ version=task_handler.bundle_info.version,
+ │
version_data=task_handler.bundle_info.version_data)
+ │ bundle.initialize() ← the SECOND bundle
+ │
+ ├── artifact = bundle.path / task_handler.rel_path
+ │ NO directory walk. NO dag_id match. NO
metadata["dags"].
+ │
+ ├── verify integrity, read supervisor_schema_version
+ │ Go → AFBNDL01 trailer: SHA-256 over the
binary
+ │ region, exactly as today
+ │ Java → Airflow-Supervisor-Schema-Version
manifest attr
+ │ TS → //# airflowBundle= layout header:
SHA-256 over
+ │ all three regions, exactly as today
+ │
+ └── command
+ Go → [artifact]
+ Java → [java, -classpath, <bundle root>/*, …,
Main-Class]
+ TS → [node, artifact]
+ │
+ └── _PopenActivitySubprocess.start(…)
+ ──StartupDetails(ti, dag_rel_path, bundle_info,
+ task_handler)──▶ runtime
+ runtime looks up its own registration by
+ (ti.dag_id, ti.task_id) — unchanged
+```
+
+No database is read in this flow. The worker never queries the binding tables;
everything it needs
+arrived on the workload. That is the property the whole design exists to buy.
+
+The integrity check is unchanged and stays on this path. It is a different
concern from the
+`cache_digest`: the digest answers "did this artifact change since I validated
it", the integrity
+hash answers "is this artifact intact". A truncated or half-downloaded file
still reports a plausible
+stored digest, so one cannot substitute for the other.
+
+### The fast path
+
+Validation is expensive (one subprocess per candidate) and a file is re-parsed
every
+`[dag_processor] min_file_process_interval` seconds, 30 by default. Re-probing
unchanged artifacts
+every 30 seconds forever is not acceptable, so the probe is skipped when
nothing relevant changed.
+
+Two tiers, cheapest first:
+
+```
+for each candidate in the coordinator's bundle:
+ stat(candidate).st_size ≠ known.size_bytes → PROBE
+ read stored cache_digest ≠ known.cache_digest → PROBE
+ otherwise → SKIP, echo the binding
+```
+
+`mtime` is deliberately absent. It is reset by an object-store download and by
container rebuilds, so
+it produces churn without adding certainty; size plus digest is sufficient,
with size acting only as
+a free pre-filter.
+
+Two conditions beyond the per-file comparison:
+
+**The candidate set must be unchanged.** A newly deployed artifact has no
known fingerprint, so a set
+that differs from `known_artifacts` forces a probe. Without this, an added
artifact that collides on
+`(dag_id, task_id)` would never be detected, and the conflict rule above would
be unenforceable.
+
+**The Python side is re-validated regardless.** The skip avoids the
*subprocess*, not the comparison.
+`handler_params` is stored precisely so a changed `.py` (a stub task that
gained an argument) is
+compared against the cached declaration in process. Skipping the comparison as
well would cache a
+verdict for a signature that no longer exists, and ship the mismatch to a
worker as a runtime
+argument error instead of catching it as an import error.
+
+### Failure handling
+
+**Validation fails.** Record the import error; leave every existing row
untouched. `bindings=None`
+already expresses "do not reconcile", so this needs no additional mechanism.
+
+**A stub task has no binding at all.** The scheduler fails it with the reason
rather than queueing a
+workload that would die on the worker at `ValueError("dag_path is required")`,
far from the cause.
+
+**Two artifacts claim one `(dag_id, task_id)`.** An import error against the
Python file (the
+definition whose author can act) naming both artifact paths, since the fix is
in the deployment.
+
+## Consequences
+
+- Task execution stops searching. A worker resolves its artifact by path, from
data that arrived on
+ the workload, with no directory walk and no database read.
+- Dag ids are never recorded at build time, so dynamic Dag rendering works:
whatever the artifact
+ registers when the Dag processor asks it is what gets recorded.
+- `jars_root` / `executables_root` / `bundles_root` are removed. Artifacts
inherit download, refresh
+ and versioning from `DagBundle`. Deployments that mount artifacts themselves
point a `LocalDagBundle` at the mount.
+- Two new tables and one migration. `lang_sdk_task_handler` is reconciled on
every parse of a file
+ that owns rows in it; `lang_sdk_task_handler_artifact` is a cache with no
per-file eviction.
+- The artifact bundle is a second bundle on the execution path. `ExecuteTask`
and `StartupDetails`
+ each grow one optional `SDKTaskHandlerRef`, and the worker performs a second
`initialize()`. Workload payloads grow by roughly one `BundleInfo`.
+- `DagFileParseRequest` and `DagFileParsingResult` each gain a field, and
`ToSDKTaskHandlerProcessor`
+ becomes a fifth union the supervisor-schema registry introspects. Both
messages already appear in
+ the generated schemas of all three SDKs, so the snapshot is regenerated and
the two prek hooks guarding it run.
+- Every Lang SDK runtime must answer `SDKTaskHandlerParseRequest`.
+- A misrouted queue becomes an import error at Dag-parsing stage instead of a
runtime failure. Today a stub task on a
+ queue absent from `queue_to_coordinator` silently falls back to the Python
coordinator
+
([`for_queue`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/task-sdk/src/airflow/sdk/execution_time/coordinator.py#L280-L284))
and dies in
+ `_StubOperator.execute()`.
+- One artifact bundle is one Java classpath. `_calculate_classpath` joins
every JAR under the root
+
([`_calculate_classpath`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/task-sdk/src/airflow/sdk/coordinators/java/coordinator.py#L85-L87)),
so all handlers in a bundle
+ share one dependency graph. Isolating conflicting dependency versions
requires a second bundle, a
+ second coordinator instance, and a second queue. This needs documenting.
+- Steady-state parsing costs no subprocesses. A cold start (empty tables, or a
bundle whose
+ candidate set changed) costs one subprocess per changed candidate per
coordinator, shared across
+ all Dag files in that parsing loop through `known_artifacts`.
+- No `DagVersion` coupling. An artifact rebuild does not bump a Dag's version,
and a Dag edit does
+ not invalidate an artifact fingerprint.
+- Mixed-language stays Python-primary. A Lang-SDK runtime cannot declare stub
tasks, and a native Dag
+ cannot delegate a task to Python.
diff --git a/airflow-core/adr/lang-sdk/0014-bundle-metadata-and-cache-digest.md
b/airflow-core/adr/lang-sdk/0014-bundle-metadata-and-cache-digest.md
new file mode 100644
index 00000000000..dfc63beec19
--- /dev/null
+++ b/airflow-core/adr/lang-sdk/0014-bundle-metadata-and-cache-digest.md
@@ -0,0 +1,122 @@
+<!--
+ 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-0014: Bundle Metadata — Retiring the Build-Time Inventory, Converging on
a Cache Digest
+
+## Status
+
+Proposed
+
+## Context
+
+[ADR-0013](0013-persisted-task-handler-bindings.md) resolves a stub task to
its artifact during Dag
+processing and persists the result, so nothing searches for an artifact at
execution time. Two
+consequences land on the artifact format: the build-time Dag inventory loses
its only purpose, and
+the Dag processor gains a new need — a stable value it can compare cheaply to
decide whether
+re-validation is required.
+
+Today that artifact carries a build-time inventory of the Dag and task ids it
exposes. This is what
+`airflow-go-pack` emits, and what the published schema requires
+([`$id`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/task-sdk/docs/airflow-metadata.schema.json#L3),
+[`required`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/task-sdk/docs/airflow-metadata.schema.json#L7)):
+
+```yaml
+airflow_bundle_metadata_version: "1.0"
+sdk:
+ language: "go"
+ version: "0.1.0"
+ supervisor_schema_version: "2026-06-16"
+source: "main.go"
+dags: # <-- frozen when the artifact was built
+ etl:
+ tasks:
+ - "extract"
+ - "transform"
+ reporting:
+ tasks:
+ - "publish"
+```
+
+## Decision
+
+### The artifact carries no Dag or task identifiers
+
+After this change an artifact contains exactly three things:
+
+```
+┌─────────────────────────────────────────────────────────────────────────┐
+│ compiled artifact the executable, JAR, or bundled code │
+│ entrypoint source the authored source, verbatim, for display │
+│ metadata only what is needed to launch and to trust │
+└─────────────────────────────────────────────────────────────────────────┘
+```
+
+and the metadata region is reduced to this:
+
+```yaml
+airflow_bundle_metadata_version: "1.0"
+sdk:
+ language: "go" # this is an Airflow Lang-SDK
artifact
+ version: "0.1.0"
+ supervisor_schema_version: "2026-06-16" # how to speak to it
+source: "main.go" # display name only
+digests:
+ integrity: "<sha256 of the executable region>"
+ cache: "<sha256 of all logical content>"
+# no dags:
+# no task_handlers:
+```
+
+No `dag_id` appears anywhere in it, and no `task_id`.
+
+Candidate *detection* is unaffected and still requires executing nothing — it
is exactly the "this is
+an Airflow Lang-SDK artifact" marker doing its job: the `AFBNDL01` trailer
magic for Go
+([`Magic`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/go-sdk/internal/bundlefooter/footer.go#L57)),
the `.min.mjs` suffix plus a valid layout header for
+TypeScript
([`BUNDLE_SUFFIX`](https://github.com/apache/airflow/blob/79991cd4db0c9346a28b23c453377f6df0c6b4ed/task-sdk/src/airflow/sdk/coordinators/node/coordinator.py#L43)),
a `Main-Class` attribute
+for Java.
+
+The entrypoint source region stays, for display.
+
+### A cache digest, distinct from the integrity hash
+
+The Dag processor must answer "has this artifact changed since I validated it"
tens of times a
+minute. That is a different question from "is this artifact intact", and the
two must not be
+conflated:
+
+| | integrity hash | cache digest |
+|---|---|---|
+| answers | is this artifact intact? | has this artifact changed? |
+| checked | at task execution, per launch | during Dag processing, per parse |
+| how | **computed** and compared | **read** |
+| must cover | the executable region, at minimum | all logical content |
+
+Reading a stored value cannot substitute for computing one — a truncated or
half-downloaded artifact
+still reports a plausible stored digest. Every existing integrity check stays
exactly where it is and
+keeps computing.
+
+**The digest is opaque and coordinator-defined.** It is not "SHA-256 of the
file". Each runtime
+supplies a value that is stable across rebuilds changing nothing and differs
across rebuilds changing
+something; consumers compare for equality and interpret nothing.
+
+
+## Consequences
+
+- Dynamic Dag rendering works.
+- The canonical schema drops the identifier mapping entirely, the coordinator
task execution side will rely on persisted rel_path instead of discovering the
artifact every time.
+- The packer no longer needs to execute the artifact at all.
`supervisor_schema_version` is a compile-time constant of the SDK.
diff --git a/airflow-core/adr/lang-sdk/README.md
b/airflow-core/adr/lang-sdk/README.md
index ec80aa58732..24afa22bd8c 100644
--- a/airflow-core/adr/lang-sdk/README.md
+++ b/airflow-core/adr/lang-sdk/README.md
@@ -38,6 +38,8 @@ bind core interfaces and apply to every language SDK, not
just the Java 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.
+- [ADR-0013](0013-persisted-task-handler-bindings.md): persisted task-handler
bindings — resolving Lang-SDK artifacts at parse time.
+- [ADR-0014](0014-bundle-metadata-and-cache-digest.md): bundle metadata —
dropping the Dag inventory, adding a cache digest.
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 —