Skip to content

Read the k8s Java task handler bundle from S3 and test failure cases - #74141

Draft
jason810496 wants to merge 172 commits into
jason/core-taskhandler-refactor/19-task-handler-compose-e2efrom
jason/core-taskhandler-refactor/20-task-handler-k8s-tests
Draft

jason810496 wants to merge 172 commits into
jason/core-taskhandler-refactor/19-task-handler-compose-e2efrom
jason/core-taskhandler-refactor/20-task-handler-k8s-tests

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 earlier PRs of the stack have the Dag processor bind each stub task to an artifact, the scheduler send it with the task, and the worker run exactly that file. The Kubernetes test harness still copied the Java JAR from S3 into an emptyDir in every Java task pod and in the Dag processor pod, where nothing refreshed it, which ADR-0013 calls a second addressing layer. It also ran only Dags that succeed. This PR reads the Java task handler bundle from S3 as a Dag bundle, and adds the failure cases the Kubernetes job can check: a Dag that does not import, and a task that no artifact is bound to.

What changes

The Java task handler bundle is the java-artifacts bucket itself, and the Java task pods run on the image --java-image names:

# kubernetes-tests/lang_sdk/config/values.yaml
dagProcessor:
  dagBundleConfigList:
    - name: java-task-handlers  # was a LocalDagBundle over an emptyDir that an init container filled
      classpath: "airflow.providers.amazon.aws.bundles.s3.S3DagBundle"
      kwargs: {bucket_name: "java-artifacts", aws_conn_id: "aws_localstack"}
config:
  sdk:
    # "java-sdk" extra, was "worker_container_repository": "lang-sdk-java-worker" and "worker_container_tag": "latest"
    coordinators: '{..., "java-sdk": {..., "extra": {..., "worker_container_repository": "{{ .Values.images.airflow.repository }}", "worker_container_tag": "{{ .Values.images.airflow.tag }}"}}}'
  • Java from S3. The Dag processor refreshes java-task-handlers like any Dag bundle, and the Java task pod downloads it when the task starts, with AIRFLOW_CONN_AWS_LOCALSTACK on its base container. The Java init container, its volumes and the Dag processor's stage-java-jar init container are gone.
  • Go stays staged. An S3 download drops the executable bit that the Go coordinator requires, so the Go bundle is still staged into an emptyDir by an init container. stage_artifacts.py is Go only and always restores the bit, so STAGE_CHMOD_EXEC is gone.
  • --java-image. The Java coordinator's extra takes its image from images.airflow, which breeze sets to --java-image, so the Java task pods use it too. It used to name lang-sdk-java-worker:latest, so a custom image reached the Airflow components but not the container that runs the task. The default is unchanged, and the Go and Python task pods keep the plain prod image.
  • Import check. Each test waits until the Dag processor has imported its Dag and fails with the import errors of the stub Dag bundle, instead of ending in a RetryError or a timeout while the trigger returns 404. The wait takes up to 600 s, and each test's execution timeout adds it to the time the run needs, so pytest does not time the test out before the wait reports.
  • A task without an artifact. lang_sdk_misrouted.py is a plain Python task on the routed golang queue. No artifact is bound to it, the scheduler queues it, and the worker fails it. The new test waits for failed and checks that the task's state reason is the worker's reason (Task 'python_task_on_golang_queue' of Dag 'lang_sdk_misrouted' has no task handler artifact, and its Dag file 'lang_sdk_misrouted.py' is not an artifact that ExecutableCoordinator runs. ...). Breeze uploads every Dag file of the harness for it, and the Go pod gets the S3 connection, because a task without an artifact reads its own Dag file from the stub Dags' S3 bundle.
  • CI and docs. The tests-kubernetes-lang-sdk job, the setup-lang-sdk-test help and the README run the tests with -k TestLangSdkCoordinatorExecutor. The README describes the Java S3 bundle, the staged Go bundle and the new Dag. The breeze command images are regenerated.
  • Go build. The container build of the Go bundle no longer sets USER. The Go SDK stopped calling user.Current() when the Go Edge Worker was removed (Remove the Go Edge Worker #71874), and the bundle, go and the packer do not link os/user. HOME stays, for the Go build cache.

How to test

uv run --project dev/breeze --locked pytest dev/breeze/tests/test_kubernetes_lang_sdk_commands.py -q
uv run --project kubernetes-tests pytest kubernetes-tests/tests/kubernetes_tests/test_lang_sdk_coordinator_executor.py --collect-only -q
prek run update-breeze-cmd-output mypy-dev mypy-kubernetes-tests --all-files
helm template airflow chart --namespace airflow --set defaultAirflowRepository=ghcr.io/apache/airflow/main/prod/python3.10-kubernetes --set defaultAirflowTag=latest --set executor=KubernetesExecutor --set config.api_auth.jwt_secret=foo --set images.airflow.repository=example/jre --set images.airflow.tag=t1 --set config.kubernetes_executor.worker_container_repository=ghcr.io/apache/airflow/main/prod/python3.10-kubernetes --set config.kubernetes_executor.worker_container_tag=latest -f kubernetes-tests/lang_sdk/config/values.yaml
breeze k8s setup-lang-sdk-test
RUN_LANG_SDK_K8S_TESTS=true breeze k8s tests --executor KubernetesExecutor -- -k TestLangSdkCoordinatorExecutor

Ran:

  • breeze test_kubernetes_lang_sdk_commands.py: 28 passed (25 before, 3 new). The kubernetes-tests module collects 2 tests. prek run --from-ref <parent> --to-ref HEAD --stage pre-commit passed, and mypy-dev, mypy-kubernetes-tests, mypy-airflow-core and mypy-task-sdk passed with --all-files. The breeze command images (output_k8s and output_k8s_setup-lang-sdk-test) are regenerated, and a second update-breeze-cmd-output run changes nothing. The schema snapshot and the Go SDK generated files are unchanged.
  • helm template with the lang-SDK values and images.airflow set to example/jre:t1: the Dag processor has the init containers wait-for-airflow-migrations and stage-go-bundle (no stage-java-jar, and no STAGE_CHMOD_EXEC anywhere), java-task-handlers is an S3DagBundle over java-artifacts, the java-sdk extra is example/jre and t1, and [kubernetes_executor] keeps the plain prod image. With lang-sdk-java-worker:latest (the harness default) the render is the base commit's apart from those changes.
  • The task pods built with PodGenerator.construct_pod the way the executor builds them, from that render: the default and Go pods run the plain prod image (the Go pod has the stage-go-bundle init container and the S3 connection on base), and the Java pod runs example/jre:t1, has no init container and has AIRFLOW_CONN_AWS_LOCALSTACK on base.
  • Checked without a cluster, in scratch tests that are not part of this PR: a plain Python task on a routed queue run through supervise_task and the in-process Execution API (the executors/test_lang_sdk_unbound_task.py setup with @task in place of @task.stub) is failed with the worker's reason as its retry_reason (2 passed, and up_for_retry with a retry left). The reason in the new test is the start of the message that ExecutableCoordinator builds for this task. stage_artifacts.py makes every staged file executable (fake bundle), and the import check returns on an imported Dag and fails with the import errors otherwise (fake session, 3 passed).
  • The import wait replayed on a clock 100 times faster under each test's own execution_timeout, with the flags breeze passes (--timeouts-order=moi) and a Dag that never imports (scratch test): both tests fail with the import errors, while an execution timeout equal to the wait ends the replay in a bare Timeout >600.0s.
  • A local kind run on Kubernetes v1.35.0 (aarch64 host), with breeze k8s deploy-airflow --executor KubernetesExecutor --multi-namespace-mode, then setup-lang-sdk-test with LANG_SDK_NATIVE_TOOLCHAIN=true and JDK 17, then the test command above: 2 passed. The Dag processor read lang-sdk-dags and java-task-handlers from localstack and bound the 4 stub tasks, lang_sdk_combined succeeded with all 6 tasks on the first try, and python_task_on_golang_queue failed with the worker's reason as its state reason. The task pods' images were not asserted. The CI job tests-kubernetes-lang-sdk, which this PR triggers because it changes kubernetes-tests/, also covers the amd64 runner and the import wait on CI hardware.
  • airflow dag-processor run locally with the [sdk] and Dag bundle configuration rendered from these values, local copies standing in for the S3 buckets: both Dag files import and the same 4 stub tasks are bound. With the Go binary left at mode 644, as an S3 download leaves it, lang_sdk_combined.py fails to import because the binary is not executable, which is why Go stays staged.

The breeze tests of each change fail without it, checked by reverting its code: the two upload tests fail against the single-Dag upload, and the Go container test fails against the old USER setting. The rendered values show what the YAML changes do: without the images.airflow template the java-sdk extra stays lang-sdk-java-worker under a custom image, and without the S3 bundle change the render still has stage-java-jar and a LocalDagBundle for java-task-handlers. The import check and the failure test fail on a cluster only: the first when the Dag processor cannot import the file, the second when the worker does not record its reason, or when the Go pod has no S3 connection (the reason is then a bundle initialization error).


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 unit tests that pack a real Go, TypeScript or Java example and probe
it run nowhere in CI: they need AIRFLOW_LANG_SDK_REAL_PROBE_TESTS=1 and
a toolchain, and no job sets the variable. Run each one as a step of the
compose e2e job of its language, which has the toolchain and its caches
set up, after the e2e tests and also when they fail. The Go cache key
now covers the test bundle's go.sum as well.

A job that lost its toolchain would still pass, because those tests skip
when it is missing. When CI is set, fail instead of skipping.
The fixture said the Go SDK reads the current user from USER and HOME
when it is built without cgo. It does not: the SDK stopped calling
user.Current() when the Go Edge Worker was removed, and the bundles do
not link os/user. The four tests pass without USER set.
The three Lang-SDK e2e jobs ran only for changes to their SDK, its
compose override and its test directory, plus _subprocess.py and its own
coordinator package. A change to the Dag processor's task handler
modules, the binding model, the task runtime's coordinator, a top-level
coordinator module or the shared e2e harness skipped them, though they
are the only end-to-end check of that code. Match those files for all
three jobs, each job's own real probe test, and the Go test bundle for
the Go job.
The contributor guide for the e2e tests says nothing about the go_sdk,
ts_sdk and java_sdk modes. Explain the wait for the Dag processor before
the first test, where a Dag file that must fail goes and why it should
fail through its stub tasks, what a new mode records for the wait and
the selective checks, and how to run the real probe tests that the same
CI job runs.
The readiness gate and the import-error checks looked at every file of the Dags folder, so a stock example Dag that cannot import in the image (example_event_driven.py without the Kafka provider) would time out the gate and fail the checks in all three Lang-SDK modes. They now check only the mode's own Lang-SDK Dag files.
A Dag file that fails the stub check gets an import error, and its Dags
are stale. The trigger then returns 404, which the test session retries,
so the test ended in a RetryError or a timeout and never showed the
import error.

Wait until the Dag has imported before triggering it, and fail with the
import errors of the stub Dag bundle when it does not.
The Java coordinator's extra names lang-sdk-java-worker:latest, and the
executor takes the task pod's image from it. A custom --java-image
reached the Airflow components and the pod's init container, but not the
container that runs the task. Render the extra from images.airflow,
which breeze already sets to --java-image, so both use the same image.
The default stays lang-sdk-java-worker:latest, and the other task pods
keep the plain prod image.
The Java task handler bundle was a LocalDagBundle over an emptyDir that an
init container filled from the java-artifacts bucket, in each Java task
pod and once in the Dag processor pod, where nothing refreshed it. A JAR
needs no executable bit, so make java-task-handlers an S3DagBundle over
the bucket instead, as ADR-0013 does: the Dag processor refreshes it and
the Java task pod downloads it when the task starts, with the connection
on its base container.

Go stays staged, because an S3 download drops the executable bit that the
Go coordinator requires. stage_artifacts.py is Go only and always
restores the bit.
…n on k8s

The scheduler queues a Python task on a queue that is routed to a
coordinator without an artifact, and the worker fails it with the
reason. The k8s test only ran Dags that succeed.

Add lang_sdk_misrouted.py, a plain Python task on the golang queue, and
a test that waits for it to fail and reads the reason from its state
reason. Breeze uploads every Dag file of the harness for it, and the Go
pod gets the connection of the S3 Dag bundle, because a task without an
artifact reads its own Dag file from its Dag bundle. CI, the breeze
help and the README run the tests by class name.
The comment said the Go SDK calls user.Current() at init, so USER and
HOME must be set. The SDK stopped doing that when the Go Edge Worker was
removed (#71874). Neither the example bundle, the go command nor the
packer links os/user, and the pack works with USER unset.

HOME stays: the container user has no home, and it points the Go build
cache at the persistent cache dir.
@jason810496
jason810496 added this pull request to stack #73978 October 3, 2026 07:56
@jason810496 jason810496 changed the title jason/core taskhandler refactor/20 task handler k8s tests Read the k8s Java task handler bundle from S3 and test failure cases Oct 3, 2026
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/19-task-handler-compose-e2e branch from eee7d4f to 8e11680 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/19-task-handler-compose-e2e branch from 8e11680 to b119b49 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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant