Skip to content

Add the task handler artifact to ExecuteTask and StartupDetails - #74136

Draft
jason810496 wants to merge 132 commits into
jason/core-taskhandler-refactor/14-task-handler-parse-validationfrom
jason/core-taskhandler-refactor/15-task-handler-artifact-ref
Draft

jason810496 wants to merge 132 commits into
jason/core-taskhandler-refactor/14-task-handler-parse-validationfrom
jason/core-taskhandler-refactor/15-task-handler-artifact-ref

Conversation

@jason810496

@jason810496 jason810496 commented Oct 3, 2026 •

Copy link
Copy Markdown
Member

Stack (bottom to top): #73970, #73971, #73972, #73973, #73974, #73975, #73976, #73977, #74030, #74031, #74032, #74033, #74067, #74135, #74136, #74137, #74138, #74139, #74140, #74141

Why

The worker reads no database, so a stub task can run the artifact the Dag processor bound it to only if its workload names that artifact (ADR-0013 "Objects on the wire"). This PR adds that reference to ExecuteTask and StartupDetails, and lets the subprocess coordinators resolve their bundle from it. Nothing fills it yet: a later PR has the scheduler read the binding, so behavior is unchanged today.

What changes

class TaskHandlerArtifactRef(BaseModel):
    bundle_info: BundleInfo | None = None  # None: the task's own Dag bundle
    rel_path: str  # POSIX path within that bundle


class ExecuteTask(BaseDagBundleWorkload):
    ...
    task_handler_artifact: TaskHandlerArtifactRef | None = None  # StartupDetails carries it too


# A JAR in a named bundle, read at the version current when the task starts
TaskHandlerArtifactRef(bundle_info=BundleInfo(name="java-task-handlers"), rel_path="libs/etl.jar")
# A Go bundle in the task's own Dag bundle, read at the version the run uses
TaskHandlerArtifactRef(rel_path="handlers/etl")
  • Two copies of the model, in airflow.executors.workloads and airflow.sdk.execution_time.comms, like BundleInfo, because the Task SDK does not import core workloads. bundle_info=None stands for the task's own Dag bundle, so a co-located artifact does not send the Dag bundle's version_data a second time.
  • Passing it on. BaseExecutor.run_workload passes the field to supervise_task, so no executor changes. supervise_task passes it to coordinator.execute_task(task_handler_artifact=...) only when it is set, so an out-of-tree coordinator whose execute_task takes no **kwargs still runs tasks without one. The Python coordinator ignores it.
  • Resolving the bundle. With a reference, SubprocessCoordinator reads the bundle it names instead of task_handler_bundle_name. The task's own bundle, also when a reference names it without a version, is read at the version the run uses: its pinned version, or the version current when the task starts if the run is not pinned. A named bundle is read at the version current when the task starts. BundleVersionLock holds the version for the whole task.
  • Checking the file. rel_path must be a file in that bundle. A missing file, a directory, an absolute path or a .. path raises FileNotFoundError naming the bundle and its version, before any runtime starts. A symlink in the bundle is followed, as the Dag processor follows it when it lists the bundle's artifacts.
  • Without a reference, as from an older scheduler, the configured bundle is used as before, and the existing directory scan still runs in whichever bundle was resolved. A later PR runs the referenced file directly and removes the scan.
  • Schema. AddTaskHandlerArtifactToStartupDetails in the in-progress 2026-10-30 version drops the field for a runtime pinned to 2026-06-16. The snapshot, ts-sdk/src/generated/supervisor.ts, java-sdk/sdk/schema/schema.json and the Go genmodels are regenerated through their hooks.
  • Edge. A worker on an older core ignores the new workload field. v2-edge-generated.yaml embeds ExecuteTask, so it is regenerated. The Edge UI openapi-gen types are left as they are: they already lag the spec (no version_data, no TestConnection), and the installed openapi-rq 2.2.0 rejects the codegen script's -c axios, so regenerating would rewrite the whole client.
  • Docs. lang-sdk-spec.rst says a runtime may ignore the field. ADR-0013 uses the new names, states the version rule once, and says the worker initializes the artifact's bundle in place of the Dag bundle on the coordinator path. The Go, Java and TypeScript pages use the same version wording.

How to test

uv run --project task-sdk --with-editable shared/secrets_masker pytest task-sdk/tests/task_sdk/coordinators/ task-sdk/tests/task_sdk/execution_time/schema/ task-sdk/tests/task_sdk/execution_time/test_supervisor.py task-sdk/tests/task_sdk/execution_time/test_coordinator.py -q
uv run --project airflow-core --with-editable shared/secrets_masker pytest airflow-core/tests/unit/executors/ -q
prek run check-ts-sdk-supervisor-schema sync-java-sdk-supervisor-schema check-go-sdk-generated-drift --all-files

Ran:

  • task-sdk coordinators/, schema, supervisor and coordinator tests: 707 passed, 1 skipped (non-Linux only).
  • airflow-core executors/: 258 passed, 1 skipped (non-fork spawning only).
  • go vet and go test in go-sdk, ./gradlew build in java-sdk, and pnpm build and pnpm test in ts-sdk (610 passed) passed against the regenerated schemas. The Edge spec generator reproduces v2-edge-generated.yaml unchanged.
  • prek run --from-ref <parent> --stage pre-commit passed, including check-supervisor-schemas-versions without SKIP. mypy-airflow-core and mypy-task-sdk passed with --all-files.

Each new test fails without its change; the migrator test fails with AddTaskHandlerArtifactToStartupDetails removed from the bundle.


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code (Opus 5.5) following the guidelines

ADR-0013 locates Lang-SDK task handlers through a named Dag bundle, so the
unversioned jars_root, executables_root and bundles_root directories go away
instead of living beside it. dag_bundle_name becomes task_handler_bundle_name
because the native-Dag work will use the old name with a different meaning. It
has not shipped in a release, so there is no shim.

A coordinator reads handlers from the named bundle at its current version, or
from the task's own Dag bundle at the run's version, pinned for the task either
way. CoordinatorManager.from_config validates the name of every spec without
constructing coordinators, so a typo fails at config load instead of on the
first task.
The compose e2e and the k8s lang-SDK test configured coordinators with the
explicit artifact roots, which no longer exist. The Java e2e needs three
bundles because one bundle is one classpath, and its three coordinators each
need their own dependencies.

Artifacts stay mounted on the workers only, because only the workers run them;
the other components need just the bundle registration for the name to
resolve. The k8s init container keeps staging from S3 because the S3 download
drops the Go binary's execute bit, which the coordinator requires.
The Go, Java and TypeScript guides, the bundle specs, the READMEs and the
contributor guide still pointed users at the removed jars_root,
executables_root and bundles_root kwargs.

A Dag bundle now decides what a coordinator sees, so two of its limits reach
users. For Java, every JAR in the bundle goes on one classpath, so conflicting
dependency versions cannot share a bundle, and a second Main-Class in it makes
the entry point ambiguous. For Go, the coordinator skips files without the
execute bit, so an object-store Dag bundle that drops it leaves the
coordinator with no executable to find.
Review flagged a Java test fixture still named after the removed jars_root
kwarg, an e2e helper without a verb, and two execute_task patches that were
not autospecced like the ones beside them.
The guides said only the workers need the artifact Dag bundle registered, which
reads as permission to register it there alone. The [sdk] configuration is also
loaded by the scheduler under KubernetesExecutor, and it is rejected wherever the
name is missing, so the bundle belongs in the shared bundle list.
Two sentences rewritten for artifact Dag bundles kept their em-dashes, which the
docs style for this work avoids.
In Airflow 3.4 the Dag processor checks a Python Dag's stub tasks against
the task handlers the language SDK artifacts register, so a deployment that
follows the old worker-only advice fails that check. task_handler_bundle_name
only serves those stub tasks, which the docs did not say.
A coordinator no queue routes to is never used, so its kwargs are not checked.
Add from_config tests for an accepted name, an unrouted coordinator and a spec without the kwarg.
ADR-0013 records at parse time which Language SDK artifact runs each stub
task, so task execution can stop scanning artifact bundles. This adds the
two tables and their migration. Later layers write and read them.

The unique key and the owning-file index use an md5 hash of each path
instead of the path. A 2000-character path is over MySQL's 3072-byte key
limit and Postgres's 2704-byte btree entry limit. This follows
DagPriorityParsingRequest.
The tables as built need an md5 hash column beside each long path, a
wider cache_digest, and the repo's constraint names. Later layers read
this SQL as the spec, so it now matches the migration.
An ORM update of an artifact's relative_fileloc kept the old hash, so the
unique key no longer matched the path. The handler table already resyncs
its hash this way.
The handler-to-artifact foreign key has no ON DELETE on purpose, so a
cleanup that races with a new binding fails instead of dropping it. Only
the cascade from dag was tested.
The old docstring read as if the row were the artifact. A row is the
binding of one stub task to an artifact.
The id and create date of migration 0142 were not made by alembic. They
are now alembic's own rev_id() and a real timestamp, and the hooks updated
the revision map and migration reference to match.
A task handler declaration now says whether it binds stub arguments
positionally, by name, or by name or as a whole. The fast path checks a
changed Python file against the stored declaration without running the
runtime, so it needs the mode as well as the params.
The migration now stores the declaration's binding mode, and later
layers read this SQL as the spec.
airflow-core takes session keyword-only so that it can never be passed in the
wrong position, and review flagged the new test helpers for taking it
positionally.
The fast path trusts an artifact's fingerprint as the cache key, but a
binding stored only what its own file's probe asked about, so one file
probing a rebuilt or colliding artifact left another file's bindings
stale. Once the runtime is asked for every handler, its answer depends
on the artifact alone, so it belongs beside the fingerprint and is
written by the same probe. The binding's own copy then has no reader.
The migration now caches the probe answer on the artifact row, and
later layers read this SQL as the spec.
Some artifacts carry no digest: a Go bundle packed before digests.cache
existed, or a third-party coordinator that computes none. Such a
candidate is still probed and bound, so it needs an artifact row. A
NULL digest never matches, so the candidate is probed on every parse.
The migration now lets an artifact store no digest, and later layers
read this SQL as the spec.
The native Dag parse and the task-handler probe start parse children other
than the Python one, and they need the same parse log forwarding, request
handling, readiness and cleanup. No behavior change.
The base class makes logger_filehandle optional for parse children other than
the Python one, but DagFileProcessorProcess still requires it, so nothing in
this layer reached the close path without one.
ADR-0012 gives the Dag-file parse a second question for a Lang-SDK runtime: which task handlers an artifact registers for Dags that Python already owns. This adds its request and result, and ToSDKTaskHandlerProcessor, which shares every response with ToDagProcessor. ToManager gains the result. Nothing sends or receives these messages yet, so DagFileProcessorProcess still rejects the result as an unhandled request.

The supervisor schema registry introspects the new union, so the snapshot gains four definitions. They are new bodies, so no VersionChange is needed.
check-ts-sdk-supervisor-schema regenerates supervisor.ts from the supervisor schema snapshot and fails on any difference. CI runs it with --all-files whenever a ts-sdk file changes, so leaving it stale after the snapshot gained the task-handler messages would fail the next unrelated ts-sdk change. The Type<N> renumbering is the generator's own output.
The versions module now says a brand-new body needs no VersionChange, but the schema guide did not say why, and check-supervisor-schemas-versions still fails a local commit that changes the snapshot without touching versions/. The next author of a new body could not tell that failure from a real one, so the guide gives the rule, the local hook behaviour, and the task-handler union.
The parse-time check has to bind stub-task arguments the way the runtime does, and names compared in order fit none of the SDKs. Go flat params bind by position and have no names. A Go struct binds by name in any order, ignoring case and underscores except for arg-tagged fields, and an untagged struct can take one unmatched argument as the whole value. Java's TaskArgs has names but binds by position; Java's TaskInput and TypeScript bind by name.

TaskHandlerDeclaration gains a required binding (positional, named or named_or_whole), TaskHandlerParam.name becomes nullable for nameless positional params, and exact_name marks a name that matches only as spelled. The bodies are still unreleased and first introduced in 2026-10-30, so AddTaskHandlerParseMessages still covers them, and check-supervisor-schemas-versions was skipped: it fails on any snapshot change that touches nothing under versions/.
check-ts-sdk-supervisor-schema regenerates supervisor.ts from the supervisor schema snapshot and fails on any difference, so the snapshot's new binding, nullable name and exact_name fields have to land in the generated types too.
The artifact a stub task runs must be the one the Dag processor validated it against, not whatever the worker's coordinator config points at when the task starts. A referenced file that is missing from the bundle at the resolved version, or whose path is absolute or has a `..` part, fails before any runtime starts, naming the bundle and its version. A reference that names the task's own bundle without a version gets the workload's bundle_info, so a pinned run keeps its version. A workload without a reference, such as one queued by an older scheduler, keeps using the configured bundle, and the existing directory scan still runs, now rooted in whichever bundle was resolved. The supervisor passes the reference only when there is one, so a coordinator that overrides execute_task without **kwargs keeps running those tasks.
The worker reads nothing from the database, so the artifact a stub task was bound to has to arrive on its workload. Every executor runs a task workload through BaseExecutor.run_workload, so passing the reference there reaches each of them without executor changes. A worker on an older core ignores the field and scans its configured bundle as before. Nothing fills the field yet. The Edge worker API spec embeds ExecuteTask, so it is regenerated.
The ADR still showed the SDKTaskHandlerRef sketch with a mandatory BundleInfo. The reference that landed uses None for the task's own Dag bundle, and the version rule it relies on was stated only in coordinator docstrings and the language pages. Those pages said the task's own bundle is pinned to the version the run was created with, which does not hold for a run that is not pinned, so they now say the version the run uses. The ADR also called the artifact's bundle a second bundle on the worker, but on the coordinator path it is initialized in place of the Dag bundle.
@boring-cyborg boring-cyborg Bot added area:coordinator Coordinator: The interface to spawn Lang-SDK subprocesses area:Executors-core LocalExecutor & SequentialExecutor area:go-sdk area:java-sdk area:providers area:task-sdk area:ts-sdk kind:documentation provider:edge Edge Executor / Worker (AIP-69) / edge3 labels Oct 3, 2026
@jason810496
jason810496 added this pull request to stack #73978 October 3, 2026 07:56
@jason810496 jason810496 changed the title jason/core taskhandler refactor/15 task handler artifact ref Add the task handler artifact to ExecuteTask and StartupDetails Oct 3, 2026
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/14-task-handler-parse-validation branch from 8efe1ef to a3646e8 Compare October 5, 2026 18:12
@jason810496
jason810496 removed this pull request from stack #73978 October 6, 2026 02:27
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/14-task-handler-parse-validation branch from a3646e8 to b94ed37 Compare October 6, 2026 02:44

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:coordinator Coordinator: The interface to spawn Lang-SDK subprocesses area:Executors-core LocalExecutor & SequentialExecutor area:go-sdk area:java-sdk area:providers area:task-sdk area:ts-sdk kind:documentation provider:edge Edge Executor / Worker (AIP-69) / edge3

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant