Let a coordinator parse the Dag files of the bundles it serves - #74042
Open
jason810496 wants to merge 7 commits into
Open
jason810496 wants to merge 7 commits into
jason810496 wants to merge 7 commits into
Conversation
This was referenced Oct 1, 2026
jason810496
added this pull request to stack #74044
October 1, 2026 15:18
kaxil
reviewed
Oct 1, 2026
jason810496
force-pushed
the
jason/lang-sdk-e2e/03c-coordinator-dag-parsing
branch
2 times, most recently
from
October 2, 2026 04:07
107f1eb to
e3fc440
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/03c-coordinator-dag-parsing
branch
from
October 2, 2026 05:44
e3fc440 to
3018722
Compare
1 task done
pierrejeambrun
force-pushed
the
jason/lang-sdk-e2e/03c-coordinator-dag-parsing
branch
from
October 2, 2026 09:00
3018722 to
d21aea9
Compare
1 task done
jason810496
force-pushed
the
jason/lang-sdk-e2e/03c-coordinator-dag-parsing
branch
from
October 2, 2026 11:41
d21aea9 to
42c00d4
Compare
jason810496
marked this pull request as ready for review
October 2, 2026 12:28
jason810496
requested review from
amoghrajesh,
ashb and
uranusjr
as code owners
October 2, 2026 12:28
jason810496
requested review from
bugraoz93,
gopidesupavan,
jscheffl and
potiuk
as code owners
October 2, 2026 12:28
jason810496
commented
Oct 2, 2026
jason810496
left a comment
Member
Author
There was a problem hiding this comment.
Thanks for the reivew.
jason810496
force-pushed
the
jason/lang-sdk-e2e/03c-coordinator-dag-parsing
branch
from
October 2, 2026 14:07
42c00d4 to
63a5cd3
Compare
A coordinator opts in to parsing native Dags by returning a Dag importer from get_dag_importer(). It serves the Dag bundle its dag_bundle_name names, or every bundle when that is unset. Two coordinators claiming one extension in a bundle is a config error. The jars_root, bundles_root and executables_root kwargs are removed. A coordinator now always finds its artifacts in a Dag bundle: the one dag_bundle_name names, or the task's own bundle. SubprocessCoordinator.parse_dag execs the runtime that parses a file. The runtime inherits only the standard streams and connects back to the given addresses. Like the log levels, a runtime now gets the resolved [api] base_url, [operators] default_deferrable and [triggerer] queues_enabled in its environment, with the fallbacks Python uses.
importers/base.py and execution_time/coordinator.py do not import each other at load time, so base.py now imports the coordinator module at the top. The imports back in coordinator.py stay lazy, with a comment naming the cycle. Rename _normalize_extensions to normalize_extensions, since coordinator.py uses it too, and let reset_coordinator_manager reuse reset_importer_registry. Co-Authored-By: Claude <noreply@anthropic.com>
The claim check ran in CoordinatorManager.from_config, which every task hits through for_queue, so one conflicting pair of coordinators stopped all tasks from starting. It now runs in for_bundle, before anything is built, and only for the bundle being parsed. for_bundle selects the coordinators to build from their specs alone, so a coordinator that cannot parse the bundle is never built. A class that fails to import is logged and skipped, and an importer whose supported_extensions is not a class-level list is rejected. Co-Authored-By: Claude <noreply@anthropic.com>
for_bundle selects coordinators from get_parsed_bundles, so serves_bundle was a second answer to the same question that could disagree with it. Remove it. get_dag_importer is called only on coordinators whose get_dag_importer_class is set, so it now always returns an importer and the base raises NotImplementedError. ADR-0010 describes the class-level lookups. Co-Authored-By: Claude <noreply@anthropic.com>
Reword the docs so dag_bundle_name reads as optional with one coordinator per class, and needed when one class is configured more than once. Co-Authored-By: Claude <noreply@anthropic.com>
Make get_dag_importer the only hook a coordinator implements to parse native Dags: it returns the Dag importer, or None, the default, to parse none. get_dag_importer_class and get_parsed_bundles are removed. for_bundle selects coordinators by the dag_bundle_name kwarg in their config, unset meaning every bundle, builds them, and returns their importers keyed by coordinator. One that cannot be built, or returns None, is skipped. Two importers claiming one extension in the bundle still raise. Rename normalize_extensions back to _normalize_extensions, since the coordinator module no longer uses it. Co-Authored-By: Claude <noreply@anthropic.com>
dag_bundle_name lets Dag processing pick which of several coordinators of one class starts the runtime that parses an artifact. With a single coordinator it can be omitted. Reword the Java, TypeScript and Go docs and READMEs to say so, without tying it to the Python stub Dag. Co-Authored-By: Claude <noreply@anthropic.com>
jason810496
force-pushed
the
jason/lang-sdk-e2e/03c-coordinator-dag-parsing
branch
from
October 2, 2026 15:09
63a5cd3 to
b830155
Compare
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Stack (bottom to top): #74040, #74041, #74042, #74035, #74043, #74036, #74037, #73841, #73845, #73846, #73847
Why
To parse native Dags, a coordinator has to say which Dag bundles it parses, and the Dag importer registry has to use its importer for them.
A coordinator's artifacts now always come from a Dag bundle, so the explicit filesystem roots go away.
What changes
BaseCoordinator.get_dag_importer()returns the Dag importer that parses a coordinator's native Dag files. It returnsNoneby default, so a coordinator opts in.dag_bundle_namekwarg names, or every bundle when that is unset.DagImporterRegistry.from_config(bundle_name)registers the importers of the coordinators that parse the bundle.CoordinatorManager.for_bundlereturns the Dag importers of the coordinators that parse the bundle, and rejects two that claim one extension there. Give each its owndag_bundle_name. Routing a task to a coordinator's queue does not run this check.[api] base_url,[operators] default_deferrableand[triggerer] queues_enabledin its environment, like the log levels.A coordinator opts in by handing out its Dag importer:
Was generative AI tooling used to co-author this PR?