jason810496 commented on code in PR #73518:
URL: https://github.com/apache/airflow/pull/73518#discussion_r4130708768


##########
airflow-core/adr/lang-sdk/0011-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-0011: Bundle Metadata — Retiring the Build-Time Inventory, Converging on 
a Cache Digest
+
+## Status
+
+Proposed
+
+## Context
+
+[ADR-0010](0010-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"

Review Comment:
   No need to bump the schema version since the previous release is a beta 
version.



##########
airflow-core/adr/lang-sdk/0010-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-0010: 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.

Review Comment:
   Addressed in d7155562.



##########
airflow-core/adr/lang-sdk/0010-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-0010: Persisted Task-Handler Bindings (Resolving Lang-SDK Artifacts at 
Parse Time

Review Comment:
   Addressed in d7155562.



##########
airflow-core/adr/lang-sdk/0010-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-0010: 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

Review Comment:
   Addressed in d7155562.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to