Skip to content

Run exactly the artifact a stub task is bound to - #74138

Draft
jason810496 wants to merge 8 commits into
jason/core-taskhandler-refactor/16-send-task-handler-artifactfrom
jason/core-taskhandler-refactor/17-run-task-handler-artifact-by-path
Draft

jason810496 wants to merge 8 commits into
jason/core-taskhandler-refactor/16-send-task-handler-artifactfrom
jason/core-taskhandler-refactor/17-run-task-handler-artifact-by-path

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 scheduler now sends a stub task the artifact the Dag processor bound it to, but the worker still searches a bundle for an artifact that declares the task's Dag id, so it can run another artifact than the bound one. This PR has the worker run exactly the file it is pointed at (ADR-0013 "Flow 3"), deletes the search, and fails a task whose artifact cannot be resolved in the task runtime with the reason in its task log. It also removes the known limit of the previous PR: a binding that names a missing file used to fail the try with no reason.

What changes

A stub task that no Dag processor bound, queued for example before the first parse after an upgrade, now fails its try with the reason in its task log and as its state reason, and retries apply:

Task 'load' of Dag 'etl' has no task handler artifact, and its Dag file 'etl.py' is not an artifact that ExecutableCoordinator runs. Queue 'golang' routes it to a Lang-SDK coordinator, so it must be a @task.stub task the Dag processor bound to an artifact, or a task of a Dag defined in a Lang SDK. Check the import errors of 'etl.py', and that the scheduler has the same [sdk] configuration as the Dag processor.

For a coordinator author, one hook now starts an artifact for the probe and the task, another says which files a task may run, and a coordinator that is not subprocess based raises one error before its runtime starts:

class MyCoordinator(SubprocessCoordinator):
    def _read_task_handler_candidate(self, path, *, rel_path): ...  # None: not one of my artifacts; error: it cannot run
    def _build_task_handler_command(self, *, path): ...  # was _build_parse_task_handler_command; the task uses it too


raise TaskHandlerArtifactError("...")  # from execute_task, before the runtime starts
  • Failing in the task runtime. supervise_task catches TaskHandlerArtifactError, starts the task instance, writes the reason to the task log and the supervisor log, reports the try as up for retry or failed with the reason (redacted, cut to 500 characters) as its state reason, uploads the log and exits with 1. A task cleared while it was queued reports no state, and a task already running elsewhere raises as before. Task callbacks do not run, as for every coordinator task failure (Lang-SDK coordinator tasks do not run their task callbacks #74063). A runtime that fails to start is not covered, as before.
  • Running exactly one file. A task runs the file its workload references, in the bundle the reference names (or the task's own bundle at the version the run uses, as before). Without a reference it runs its own Dag file in its own bundle, which is how a task of a Dag defined in a Lang SDK runs. task_handler_bundle_name is no longer read at execution.
  • Checks before the runtime starts. The file must be in the bundle, _read_task_handler_candidate must return a candidate without an error (so a stub task's Python file fails, and so does a Go bundle without the executable bit), _build_task_handler_command must build the command (the integrity check), and the schema version, when there is one, must be one the supervisor knows. Each failure is a TaskHandlerArtifactError that gives the reason and names the file and its Dag bundle (the unbound-stub reason names the task, its Dag file and its queue instead). A file the worker cannot open gives the error as its reason, and a schema version the worker's Task SDK does not know says so instead of the migrator's text. A coordinator without a hook fails its tasks naming the hooks it lacks.
  • The scan is deleted. _build_execute_task_command, the Go and TypeScript _Bundle.find, the Java _JarInfo and ResolvedBundle are gone. The Java classpath uses the sorted walk_files (same JARs, same order). The TypeScript reader no longer requires task_handlers in the bundle metadata, and does not read it.
  • Java main_class now also runs the probe, so a handler JAR without a manifest Main-Class is listed and probed when it is set. Without it the bound JAR's manifest decides, which ends the arbitrary pick between JARs that declare a Main-Class (Coordinator: detect ambiguous entrypoints and duplicate dag_ids at import time #71134). With it every handler JAR in the bundle reports the same handlers, so keep one per bundle. The digest recorded for a JAR covers main_class when it is set (it is the manifest's digest when unset), so changing main_class makes the Dag processor probe the JARs again. A thin JAR takes its schema version from the first JAR with one, in sorted path order.
  • e2e. The Lang-SDK compose modes skip airflow dags reserialize, which writes stub Dags without bindings and would fail a task queued before the Dag processor binds it. The Dag processor commits bindings with the Dag rows, and trigger_dag waits for the Dag.
  • Docs. The language pages, the Language SDK guide, the bundle specs, the contributor guide for a new Language SDK, the READMEs, the Java and new SDK agent skills, the harness text and ADRs 0011 to 0013 say a worker runs the bound artifact and nothing else. The Java page lists the manifest attributes a Maven-built handler JAR needs (Airflow-Cache-Digest, Main-Class, Airflow-Supervisor-Schema-Version), with a digest that changes on every build; the Gradle plugin writes them. Upgrade the Dag processor, the scheduler and the workers together: a task queued without its artifact by an older scheduler fails with the reason, and retries apply.
Messages

These strings are kept stable for the later PRs of the stack, which assert on them. {origin} is Task handler artifact for a file a workload references, or Dag file for the task's own Dag file. Two log events go to the task log:

  • Running a Lang-SDK artifact, with origin, bundle_name, bundle_version and path.
  • Cannot run the task's Lang-SDK artifact, with reason. The supervisor log gets it too, as a warning with ti_id and reason.

Every final reason, as TaskHandlerArtifactError raises it:

  • A stub task with no binding whose Dag file is not an artifact of the coordinator, filled in with the task, Dag, Dag file, queue and coordinator class name: Task 'load' of Dag 'etl' has no task handler artifact, and its Dag file 'etl.py' is not an artifact that ExecutableCoordinator runs. Queue 'golang' routes it to a Lang-SDK coordinator, so it must be a @task.stub task the Dag processor bound to an artifact, or a task of a Dag defined in a Lang SDK. Check the import errors of 'etl.py', and that the scheduler has the same [sdk] configuration as the Dag processor.
  • A referenced file that is not in its bundle: Task handler artifact 'bin/etl' is not a file in Dag bundle 'go-task-handlers' at version 'abc'. The at version 'abc' part is left out when the bundle has no version.
  • A referenced file of another coordinator: Task handler artifact 'bin/etl' in Dag bundle 'go-task-handlers' is not an artifact that ExecutableCoordinator runs. A newer parse may route this task to another coordinator.
  • A file that cannot run: {origin} 'bin/etl' in Dag bundle 'go-task-handlers' cannot run: {reason}. The reason is one of:
    • The error from a file or a parent directory the worker cannot open or search, or one that _read_task_handler_candidate raises, such as [Errno 13] Permission denied: '/path/to/bin/etl'.
    • The reason the coordinator gives for the file, such as bin/etl is not executable. Use a Dag bundle that keeps the executable bit; object-store Dag bundles such as S3DagBundle drop it.
    • The error of _build_task_handler_command, such as an integrity check that fails.
    • An unknown supervisor schema version: uses supervisor schema version '1999-01-01', which this worker's Task SDK does not support.
  • A bundle that cannot be read: {origin} 'bin/etl' cannot run: Dag bundle 'go-task-handlers' cannot be read: {reason}. The reason is the error of the bundle, or it resolved to /path/to/bundle, which does not exist.
  • A coordinator without the hooks: MyCoordinator cannot run tasks: implement _read_task_handler_candidate and _build_task_handler_command. It names only the hooks that are missing.

The Dag processor's probe is not covered: it reports an unknown supervisor schema version as before, with the migrator's text in the import error of the stub Dag.

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/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/dag_processing/ -q
uv run --project airflow-core --with-editable shared/secrets_masker pytest airflow-core/tests/unit/executors/ -q
AIRFLOW_LANG_SDK_REAL_PROBE_TESTS=1 uv run --project airflow-core --with-editable shared/secrets_masker pytest airflow-core/tests/unit/dag_processing/test_task_handler_processor_go.py airflow-core/tests/unit/dag_processing/test_task_handler_processor_ts.py airflow-core/tests/unit/dag_processing/test_task_handler_processor_java.py -v
prek run mypy-airflow-e2e-tests generate-supervisor-schemas-snapshot check-go-sdk-generated-drift --all-files

Ran:

  • task-sdk coordinators/, supervisor and coordinator tests: 677 passed, 1 skipped (non-Linux only).
  • airflow-core dag_processing/: 798 passed, 17 skipped (8 real probe tests, which need AIRFLOW_LANG_SDK_REAL_PROBE_TESTS=1, and 9 that need FabAuthManager, a system test module or a Postgres or MySQL backend). executors/: 261 passed, 1 skipped (non-fork spawning only). The new executors/test_lang_sdk_unbound_task.py runs supervise_task against the in-process Execution API for a stub task queued without a binding: the try is FAILED, or UP_FOR_RETRY with a retry left, with the reason as its retry_reason and in the task log.
  • The real probes with Node 22.22.0, pnpm, JDK 11 and Go 1.25: 8 passed (Go 3, TypeScript 3, Java 2).
  • prek run --from-ref <parent> --stage pre-commit passed. mypy-airflow-core, mypy-task-sdk and mypy-airflow-e2e-tests passed with --all-files. The schema snapshot and the Go SDK generated files are unchanged, since no supervisor message changed.
  • The compose e2e modes and the Kubernetes test ran once, on the last layer of the stack, and passed (results in Add e2e tests for checking stub tasks against their task handlers #74140 and Read the k8s Java task handler bundle from S3 and test failure cases #74141). Their stub tasks ran the artifacts they are bound to, and the tasks without an artifact (an unbound stub task in compose, a Python task on the golang queue on Kubernetes) failed with this PR's reason as their state reason.

The tests of each behavior change fail without it, checked by reverting the code of its commit: the supervisor and coordinator failure tests (23 fail), the renamed hook tests (20), the Java main_class probe and digest tests (6, and the thin JAR schema version test fails against the old lookup), the by-path tests of the three coordinators, including the unopenable file and unknown schema version tests (27, and an import error in the base class tests), the core unbound stub test (2) and the TypeScript reader tests that accept metadata with and without task_handlers (2). The by-path tests fail against the old scan whatever order the filesystem lists files in. The other new tests guard behavior a commit keeps, such as the classpath order and symlink loop moved from the deleted scan tests and the cases the scan already failed the same way, and pass before and after it.


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

airflow dags reserialize writes stub Dags without task handler bindings. Once a worker fails a stub task queued without its binding, a task queued before the Dag processor's own parse binds it would fail. The go_sdk, ts_sdk and java_sdk modes now leave the Dags to the Dag processor, which commits the bindings in the same transaction as the Dag rows. AirflowClient.trigger_dag already waits for the Dag to exist before it unpauses and triggers it.
…lved

A failure before the runtime starts, such as a bundle that cannot be read or a referenced file that is missing, raised out of supervise_task before the task instance was started, so the try stayed queued, the task log was empty, and the scheduler later reported it as killed externally. A coordinator now raises TaskHandlerArtifactError from execute_task for such a failure. supervise_task catches it, starts the task instance, writes the reason to the task log and the supervisor log, reports the try as up for retry or failed with the reason (redacted and cut to 500 characters) as its state reason, uploads the task log and returns 1. A task cleared while it was queued is left as the server set it, and no state is reported.

SubprocessCoordinator raises it for a bundle that cannot be read, a referenced file that is not in the bundle, and a command that cannot be built. The scan still runs in this commit.
_build_parse_task_handler_command builds the command that runs an artifact. A task runs the same artifact the same way, so the hook is now _build_task_handler_command and the next commits use it for both the probe and the task. Only the name and the docstring change. A SubprocessCoordinator without the hook now fails with "does not build task handler commands".
JavaCoordinator ran main_class only for tasks, so the probe that reports a JAR's task handlers used the Main-Class of the JAR's own manifest, and a JAR with an Airflow-Cache-Digest but no Main-Class was reported as unable to run even when main_class was set. The probe now runs main_class when it is set, and such a JAR is listed and probed. Without main_class the bound JAR's manifest decides. With it every handler JAR in the bundle is probed through the same class and reports the same handlers, so they collide, and the docstring says to keep one handler JAR per bundle. The cache digest recorded for a JAR covers main_class when it is set, so changing it probes the JAR again, and it is the manifest's digest when it is unset.

A thin JAR takes its schema version from the first JAR with one, in sorted walk order, instead of the first JAR that declares the matching Main-Class, which no longer exists when main_class names a class no manifest declares.
A worker scanned a bundle for an artifact that declares the task's Dag id, so a task could run another artifact than the one the Dag processor bound it to, and one that could not be found failed with no reason. A task now runs one file, picked the way a Python task picks its Dag file: the artifact its workload references, in the bundle the reference names, or without a reference its own Dag file in its own bundle, which is how a task of a Dag defined in a Lang SDK runs. task_handler_bundle_name is no longer read at execution.

The file must be in the bundle, and must be an artifact the coordinator runs (checked with _read_task_handler_candidate, so a stub task's Python file is not) that can run. The command and supervisor schema version come from _build_task_handler_command, the version is resolved before the runtime starts when there is one, and the task log says which artifact runs. Each failure is a TaskHandlerArtifactError whose message gives the reason and names the file and its Dag bundle. For a stub task nobody bound it names the task, its Dag file and its queue instead. A file the worker cannot open gives the error as its reason, and a supervisor schema version the worker's Task SDK does not know gives "uses supervisor schema version X, which this worker's Task SDK does not support". A subclass of SubprocessCoordinator without a hook fails its tasks with a message naming the hooks it lacks. _build_execute_task_command is no longer called.

A core test runs supervise_task against the in-process Execution API for a stub task queued without a binding, and checks its state, state reason and task log.
No code searches a bundle for an artifact any more: a task runs the file it points at. This removes SubprocessCoordinator._build_execute_task_command and its implementations, the Go bundle walk and _Bundle.find, the TypeScript _Bundle.find, the Java _JarInfo and its Main-Class and schema version discovery, the JAR walk that duplicated walk_files, and ResolvedBundle with validate_schema_version, which had no user left. The Java classpath is now collected with the sorted walk_files, which lists the same JARs in the same order.

The TypeScript bundle reader no longer requires task_handlers in the embedded metadata, and BundleMetadata.dag_ids is gone, so a bundle reads with or without the field, and the field is not read.
A worker no longer searches a bundle for the artifact of a task, so the text that said it does is changed. The Language SDK guide and the Go, Java and TypeScript pages say the Dag processor lists the artifacts of task_handler_bundle_name and binds each stub task to one, a worker runs only that artifact, and a stub task queued without a binding fails with the reason in its task log. The Go executable bit is required of the file a worker runs, and the Java main_class option applies to the probe too, with one handler JAR per Dag bundle when it is set and a change to it probing the JARs again. The Java page lists the manifest attributes a Maven-built handler JAR needs, with a cache digest that changes on every build, and the Go page and README no longer say the coordinator reads the Dag ids of a bundle. The two bundle specs describe listing by trailer or layout line and running the bound file, and an unknown schema version fails the task. The contributor guide for a new Language SDK has one _build_task_handler_command section for the probe and the task, says _read_task_handler_candidate also decides what a task may run, and names TaskHandlerArtifactError for coordinators that do not extend SubprocessCoordinator. The READMEs, the Java and new SDK agent skills, the Go example Dag, the compose and Kubernetes harness text and the e2e docstrings follow.
ADR-0013 Flow 3 shows a worker picking one file the way a Python task picks its Dag file, checking that it is in its bundle and is an artifact of the coordinator that can run, starting it with the hook the probe uses, and resolving its schema version before the runtime starts. Failure handling says a stub task without a binding, and an artifact the worker cannot run, fail the try with the reason in its task log and as its state reason, that retries apply, and that callbacks do not run. The Consequences say a task of a Dag defined in a Lang SDK runs its own Dag file, that task_handler_bundle_name matters only to the Dag processor, that a change to Java's main_class probes the JARs again, that a stub task runs only once bound including after an upgrade, and that dag.test() with an executor does the same. ADR-0012 names one _build_task_handler_command for execute_task and parse_task_handler, and ADR-0011 uses the name in its diagram.

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant