Skip to content

Add e2e tests for checking stub tasks against their task handlers - #74140

Draft
jason810496 wants to merge 12 commits into
jason/core-taskhandler-refactor/14-task-handler-parse-validationfrom
jason/core-taskhandler-refactor/19-task-handler-compose-e2e
Draft

jason810496 wants to merge 12 commits into
jason/core-taskhandler-refactor/14-task-handler-parse-validationfrom
jason/core-taskhandler-refactor/19-task-handler-compose-e2e

Conversation

@jason810496

@jason810496 jason810496 commented Oct 3, 2026 •

Copy link
Copy Markdown
Member

Stack (bottom to top), on native #74037 (stack #74170): #73973, #73974, #73975, #74317, #73976, #74067, #74135, #74140, #74141, #73971, #74030, #74031, #74032, #74136, #74137, #74138, #74139

Adds end-to-end tests of #74135's check against real Go, Java and TypeScript artifacts. A Go and a Java fixture have deliberately wrong stub tasks, and the e2e tests of all three Lang-SDK modes read the Dag processor's parse logs and the REST API to see what it did.

# airflow-e2e-tests/go-test-bundle/dags/go_task_handler_failures.py
@task.stub(queue=_QUEUE)
def not_registered(): ...

@task.stub(queue=_QUEUE)
def takes_two_numbers(first: int, second: int, third: int): ...

@dag(dag_id="go_task_handler_failures")
def go_task_handler_failures():
    not_registered()
    takes_two_numbers(1, 2, 3)
Stub tasks in go_task_handler_failures.py do not match their task handlers:
- Dag 'go_task_handler_failures', task 'not_registered': 'handlers' in Dag bundle 'go-test-task-handlers' registers no task handler for it
- Dag 'go_task_handler_failures', task 'takes_two_numbers' ('handlers' in Dag bundle 'go-test-task-handlers'): passes 3 arguments, the task handler takes 2
  • The Dag processor gets the worker's artifacts and runtimes in the go_sdk, java_sdk and ts_sdk compose overrides, so it can run the check.
  • Native's Java and Node Dag importers register in every Dag bundle, so they would import every mounted JAR and TypeScript bundle as a Dag file, hundreds of them through the Scala example's dependency JARs. The harness writes an .airflowignore holding * into each artifact directory to keep them out. The check is unaffected, because the coordinator's scan reads the filesystem directly. A real deployment with several coordinators of one class sets [sdk] dag_bundle_to_coordinator instead.
  • airflow dags reserialize is skipped in the Lang-SDK modes, so the Dag processor's own parse, with the check, is the first one. A readiness gate then waits until every Lang-SDK Dag file has a probe record and the import errors are exactly the expected ones, so no test races that parse.
  • The tests check that every Lang-SDK Dag file gets a probe record, no artifact is probed twice in one parse, a later parse probes nothing new, only the fixtures fail to import, with the exact expected text, and Go's name mismatch is only a warning.
  • The real-probe unit tests now run in each Lang-SDK e2e job, and in CI they fail instead of skipping when the toolchain is missing.
  • The three Lang-SDK e2e jobs also run for changes to the Dag processor's task handler modules, the coordinators and the shared e2e harness. dag_processing/processor.py stays out on purpose: it is only the check's call site, its unit tests cover it, and it changes too often to trigger three compose jobs each time.

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

Generated-by: Claude Code (Sonnet 5) following the guidelines

The Go e2e checks added later run `go` too, and must do it the way the
harness packs: with the host toolchain in CI and in the pinned builder
image locally. Move that choice into `run_go`, and have the one pack
function take the module, package and output so the example bundle and
the test bundles pack the same way.
The e2e check for stub task problems needs a handler the user-facing
Go example must not carry: one whose stub tasks are missing and
mismatched on purpose. Put it in a Go module of its own, packed into
its own bundle, so the example's bindings, parse warnings and probe
counts stay as the example defines them. Nothing deploys it yet.
The Go e2e test bundle is a Go module of its own that resolves go-sdk
through a replace, as the Kubernetes example module does. A dependency
bump in go-sdk leaves its go.mod stale and breaks the Go SDK e2e job on
every pull request, and a bump of the go directive leaves it on another
toolchain. Register the module in the two hooks that guard that for the
example module, and have the tidy hook name the CI job each module
breaks.
The Dag processor now runs each artifact to check the stub tasks of
the Python Dags, so the e2e harnesses must mount the artifacts and
install the runtimes there too, not only on the worker. Native's Java
and Node Dag importers register in every bundle, so without an
.airflowignore they would also try to import these artifacts as Dag
files; a bare "*" keeps the Dag processor's bundles to the artifacts
alone, since the check reads the filesystem directly and ignores it.
Pack the test bundle next to the example, deploy it as a Dag bundle of
its own served by a coordinator of its own, and copy its Dag file. The
example's coordinator, Dag bundle and Dag files stay as they are.
takesTwoNumbers registers through the generated annotation processor
path (Bundle.register(Class)) rather than TestBundleBuilder's
hand-written Task classes, so its TaskParams carries flat positional
parameters: the Dag processor sees a two-argument task handler and a
three-argument call. not_registered has no handler at all. Both land
in one Dag file so the check reports them as one import error.
…DK e2e tests run

The Lang-SDK modes skip "airflow dags reserialize", so the first test
can trigger a Dag before the Dag processor has checked its stub tasks
against the task handlers, and a probe that fails is only a warning,
so a file whose artifact cannot be probed gets no import error either.
Add the helpers to read the parse logs and the REST API, and a session
fixture that waits until every Lang-SDK Dag file has a probe record and
the import errors are exactly the ones that fail by design. It asks
once for a reparse of each file still missing a probe record, and fails
with what is missing and the text of any unexpected import error.

Only each mode's own Lang-SDK Dag files are checked: a stock example
Dag (such as one needing the Kafka provider) may fail to import in CI
for reasons that have nothing to do with this check.
A reserialize imports the stub Dag files without the Dag processor's
stub-task check having run on them, so the first rows for a Lang-SDK
Dag file would carry neither a probe record nor its import errors.
Leaving the files to the Dag processor makes its own parse the first
one; AirflowClient.trigger_dag already waits for a Dag to exist before
it unpauses and triggers it.
…orts

Every stub Dag file of each mode gets a probe record, with no artifact
probed twice in one parse, and a later parse probes nothing an earlier
one did not: nothing records the answer, so every parse checks
again. The Go and Java fixtures each produce one import error with
both their problems as separate lines, and no other Lang-SDK Dag file
of the mode fails to import. Go's named-binding mismatch is only a
warning in the parse log, and the example's stub tasks still serialize
and run on the routed queue.

Java has no named-binding mismatch fixture to observe (its existing
examples all bind exactly), so it gets four of the five assertions.
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 three Lang-SDK e2e jobs ran only for changes to their own 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, a top-level coordinator module, the task runtime's
coordinator 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.

processor.py (the check's call site) stays out of the shared list:
the Dag processor's own unit tests cover it, and it is hot enough that adding it here
would cost three compose jobs on every change unrelated to task
handlers.
The contributor guide for the e2e tests said 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.
@jason810496
jason810496 removed this pull request from stack #73978 October 6, 2026 02:27
@jason810496
jason810496 changed the base branch from jason/core-taskhandler-refactor/18-bundle-metadata-without-dags to jason/core-taskhandler-refactor/14-task-handler-parse-validation October 6, 2026 02:30
@jason810496
jason810496 added this pull request to stack #74318 October 6, 2026 02:35
@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