Skip to content

Check stub tasks against task handlers when a Dag file is parsed - #74135

Draft
jason810496 wants to merge 6 commits into
jason/core-taskhandler-refactor/13-stub-task-checksfrom
jason/core-taskhandler-refactor/14-task-handler-parse-validation
Draft

jason810496 wants to merge 6 commits into
jason/core-taskhandler-refactor/13-stub-task-checksfrom
jason/core-taskhandler-refactor/14-task-handler-parse-validation

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

A stub task whose artifact registers no handler for it, or one that cannot take its arguments, used to serialize cleanly and fail only on a worker. The Dag processor now probes the artifact a worker would run (#73974) while it parses a Dag file, and checks the file's stub tasks against the answer (#74067).

Stub tasks in etl.py do not match their task handlers:
- Dag 'etl', task 'load': 'bin/etl' in Dag bundle 'go-task-handlers' registers no task handler for it
- Dag 'etl', task 'transform' ('bin/etl' in Dag bundle 'go-task-handlers'): passes 2 arguments, the task handler takes 3

For SDK authors, the coordinator registry can now name the coordinator a queue routes to without building it:

manager.get_coordinator_key("golang")  # -> "go-sdk", or None if no coordinator serves the queue
  • Only a mismatch a probe answer proves fails the import. Everything else is a parse-log warning, such as a name mismatch under named binding, a missing artifact or artifact bundle, a bundle owned by another team, or a probe that fails or times out. A Dag processor without the artifacts therefore cannot stop Dags that run today.
  • Stub tasks are grouped by the coordinator their queue routes to and by Dag. Each group's artifact is found with the coordinator's own scan, in its task_handler_bundle_name bundle or the Dag's own bundle, and each distinct (coordinator, artifact) pair is probed once per parse. A stub task on a queue no coordinator serves is not checked.
  • Each probe is bounded by [core] dagbag_import_timeout (or the get_dagbag_import_timeout policy), and the whole check by 90% of [dag_processor] dag_file_processor_timeout, so the parse result is sent before the manager kills the process.
  • The language SDK guide gets a "What the Dag processor checks" section, and ADR-0011 Steps 1 to 4 describe the flow.

Known cost: the Dag processor now needs the same [sdk] config, artifact bundles and runtimes as a worker, and starts one runtime per (coordinator, artifact) pair on every parse of a file with stub tasks. A Dag processor started with --bundle-name must also be given the artifact bundles its coordinators name.


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

The Dag-file parse groups stub tasks by the coordinator their queue
routes to, and leaves a stub task on an unrouted queue to an external
worker. for_queue cannot tell the two apart: it returns the Python
coordinator for an unrouted queue and builds the routed one.
get_coordinator_key returns the key, or None, without building
anything.
Group a Dag file's stub tasks by the coordinator their queue routes to
and by their Dag. Find each group's artifact with the coordinator's
own scan, in the Dag bundle named by its task_handler_bundle_name or
in the Dag's own bundle, and probe each distinct (coordinator,
artifact) pair once, until the parse's deadline. Anything that leaves
a group unchecked (an unbuildable coordinator, a bundle of another
team, a missing or unreadable bundle, a coordinator with no find hook,
an artifact whose schema is too old, a failed or timed-out probe) is
logged and never an import error; only a mismatch a probe answer
proves becomes one.
A stub task whose artifact registers no handler for it, or one that
cannot take its arguments, serialized cleanly and only failed on a
worker, far from the Dag file that caused it. The parse now resolves
every stub task and adds the mismatches a probe answer proves to the
Dag file's import error. Everything else (no artifacts, a coordinator
or [sdk] failure, a failed or timed-out probe, an unexpected error in
the check itself) is logged and never stops a Dag that runs today. The
check gets 90% of dag_file_processor_timeout, counted from the
creation of the parse child, so its result is sent before the manager
kills the process.
The Go, Java and TypeScript examples are what users copy, and the e2e
tests run them, so each example Dag file must now pass the parse-time
check against the bundle built from the same example. Each real-probe
test parses the example Dag file against its own built bundle, expects
no import error, and the "Probed a task handler artifact" record
proves the parse really ran the check. The Go example's two name
mismatches stay warnings in the parse log.
Dag authors now see a missing task handler or an argument a handler
cannot take as an import error of the Python Dag file, and a name
mismatch under named binding as a warning in its parse log, so the
language SDK guide says what is checked, what is not, and what it
costs. The Go and Java pages said a flat count or type mismatch fails
the task; it now fails the Dag file's import first. The pages describe
defaulted arguments as the check counts them, which is all or none.
The Java page adds the probe's JAR/JRE/main_class specifics, and the
TypeScript page notes its check is presence only and that a bundle
packed before this change must be packed again.
Step 1 now groups stub tasks by the coordinator their queue routes to
and by their Dag, with an unrouted stub left to an external worker.
Step 2 names get_coordinator_key, which replaces for_queue because
for_queue built a Python coordinator for an unrouted queue instead of
telling the two apart. Step 3 replaces BundleScanner and its mode list
with the coordinator's own scan, which re-runs on every parse the find
hook a worker uses, so the artifact probed is the one a worker would
run. Step 4 names the probe class and call F2 introduced. Consequences
records what the probe costs on every parse, the lenient outcome rule,
and the parse's deadline; Appendix B adds the one new outcome that is
only logged. The ADR-0003 BundleScanner reference is dropped with it.
@jason810496
jason810496 removed this pull request from stack #73978 October 6, 2026 02:27
@jason810496
jason810496 added this pull request to stack #74318 October 6, 2026 02:35
@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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant