Skip to content

Send known task-handler artifacts to the Dag-parsing child - #74030

Draft
jason810496 wants to merge 94 commits into
jason/core-taskhandler-refactor/07-ts-sdk-task-handler-parsefrom
jason/core-taskhandler-refactor/09-known-artifacts-push-down
Draft

jason810496 wants to merge 94 commits into
jason/core-taskhandler-refactor/07-ts-sdk-task-handler-parsefrom
jason/core-taskhandler-refactor/09-known-artifacts-push-down

Conversation

@jason810496

@jason810496 jason810496 commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

Stack (bottom to top): #73970, #73971, #73972, #73973, #73974, #73975, #73976, #73977, #74030, #74031, #74032, #74033, #74067

Why

The Dag-parsing child has no database, so the task-handler artifacts the Dag processor already recorded, each with its probe answer, must arrive with the parse request (ADR-0013 "Objects on the wire"). A later layer's fast path uses an artifact's recorded answer when its size and digest still match, instead of probing it again. Nothing reads the field yet.

What changes

DagFileParseRequest gains the known artifacts:

class TaskHandlerArtifact(BaseModel):
    bundle_name: str
    relative_fileloc: str = Field(max_length=2000)
    size_bytes: int
    cache_digest: str | None = Field(max_length=128)
    task_handlers: dict[str, list[TaskHandlerDeclaration]]  # the artifact's whole probe answer


class DagFileParseRequest(BaseModel):
    ...
    known_artifacts: list[TaskHandlerArtifact] = Field(default_factory=list)

With this config, a Dag file in the dags-a bundle gets the rows of java-task-handlers (named) and of dags-a (because go-sdk names no bundle, so it reads the task's own Dag bundle). With [core] multi_team, it gets the java-task-handlers rows only if that bundle's team equals the team of dags-a:

[sdk]
coordinators = {
    "jdk-17": {"classpath": "...JavaCoordinator", "kwargs": {"task_handler_bundle_name": "java-task-handlers"}},
    "go-sdk": {"classpath": "...ExecutableCoordinator"}
}
queue_to_coordinator = {"java": "jdk-17", "golang": "go-sdk"}
  • The manager reads lang_sdk_task_handler_artifact, including the task_handlers answer, once per parsing loop and only when it starts a child. That makes rows persisted by earlier children visible to later ones. The query covers the task_handler_bundle_name of each coordinator a queue routes to, plus the Dag bundles it parses when one of them names none, and filters on the leading column of the (bundle_name, relative_fileloc_hash) unique index. Rows come back ordered by bundle and path, so a child's payload is reproducible. There is no query when no queue routes to a coordinator.
  • Each child gets the rows of the named bundles whose team equals its Dag bundle's team, plus its own Dag bundle's rows when a coordinator falls back to it, never another Dag bundle's. A recorded answer is trusted by every Dag file that reads it, so it must not cross a team. Without [core] multi_team every team is None and every named bundle is in scope; with it, a team-less Dag bundle sees only team-less named bundles. Files without stub tasks get their scope's rows too.
  • The scope is computed in one manager method, which the next layer reuses for the bundles a parse may record answers in.
  • cache_digest is None for an artifact that stores no digest. It never matches, so such an artifact is probed on every parse. The max_lengths match the columns.
  • CoordinatorManager.get_task_handler_bundle_names() maps each coordinator a queue routes to, by key, to its bundle name, or None for the Dag's own bundle. It reads the specs from_config validated and builds no coordinator. A coordinator no queue routes to runs no stub task, so its bundle is not read, and from_config does not validate its bundle name.
  • An invalid [sdk] coordinators value is logged once, and the Dag files are parsed with no known artifacts, so Python-only parsing keeps running.
  • A stored row that fails validation is skipped with a warning, so the other rows still go out and that artifact is probed again. A failed read is logged, and the children started in that loop get no known artifacts, so they probe.
  • get_known_task_handler_artifacts(bundle_names) reads the database by default and can be overridden to source the rows elsewhere, like fetch_callbacks.
  • AddKnownArtifactsToDagFileParseRequest in the in-progress 2026-10-30 version drops the field for a runtime pinned to 2026-06-16. The snapshot, ts-sdk/src/generated/supervisor.ts and the vendored java-sdk/sdk/schema/schema.json are regenerated.
  • The Go genmodels are regenerated too, since check-go-sdk-generated-drift (Go SDK: bring the coordinator-protocol models back to the supervisor schema #73963) fails when they lag the snapshot.
  • ADR-0013 now describes the fallback bundles, the per-child and per-team scope, and TaskHandlerArtifact with its answer, and its probe sketch no longer sends Dag ids.

How to test

uv run --project airflow-core --with-editable shared/secrets_masker pytest airflow-core/tests/unit/dag_processing/ airflow-core/tests/unit/models/test_lang_sdk_task_handler.py -q
uv run --project task-sdk pytest task-sdk/tests/task_sdk/coordinators task-sdk/tests/task_sdk/execution_time/schema task-sdk/tests/task_sdk/execution_time/test_coordinator.py -q
uv run --project task-sdk --with-editable shared/secrets_masker pytest task-sdk/tests/task_sdk/execution_time/test_supervisor.py -q
breeze run --backend postgres pytest airflow-core/tests/unit/dag_processing/test_manager.py -k "KnownTaskHandlerArtifacts or serialize_parse_request or callback_queue"

Ran:

  • airflow-core/tests/unit/dag_processing/ and test_lang_sdk_task_handler.py: 640 passed, 12 skipped (Postgres-only budgets and opt-in real probes).
  • task-sdk: coordinators, schema and coordinator tests 327 passed; test_supervisor.py 332 passed, 1 skipped (non-Linux only).
  • The manager tests above on Postgres and MySQL through breeze: 23 passed each.
  • ./gradlew -p java-sdk :sdk:test --tests '*TaskHandlerParseTest*' and compile-ts-sdk passed against the regenerated schema.
  • prek run --from-ref <parent> --stage pre-commit passed, including check-supervisor-schemas-versions without SKIP, mypy-airflow-core and mypy-task-sdk.

The migrator test fails with AddKnownArtifactsToDagFileParseRequest removed from the bundle.


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

@boring-cyborg boring-cyborg Bot added area:coordinator Coordinator: The interface to spawn Lang-SDK subprocesses area:DAG-processing area:java-sdk area:task-sdk area:ts-sdk labels Oct 1, 2026
@jason810496
jason810496 added this pull request to stack #73978 October 1, 2026 13:34
@jason810496 jason810496 changed the title jason/core taskhandler refactor/09 known artifacts push down Send known task-handler artifacts to the Dag-parsing child Oct 1, 2026
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/08-java-sdk-task-handler-parse branch from 03405f9 to b6456cd Compare October 1, 2026 14:44
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/09-known-artifacts-push-down branch from a19f24f to 34aabf6 Compare October 1, 2026 14:44
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/08-java-sdk-task-handler-parse branch from b6456cd to 949c919 Compare October 2, 2026 00:54
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/09-known-artifacts-push-down branch from 34aabf6 to 7b0b329 Compare October 2, 2026 00:54
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.
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/08-java-sdk-task-handler-parse branch from 949c919 to f508e20 Compare October 2, 2026 05:48
An unreadable file or a JAR without a manifest was reported as having no Main-Class, which points the Dag author at the wrong fix. The two cases now have their own messages. Manifest values are now decoded after their continuation lines are joined, because the 72-byte fold can split a multi-byte character. Execution and the probe share one command builder, so the two cannot drift.
The branch that copies the Task SDK snapshot had only a manual run behind it, and its version reader had a noun name that reads as an attribute.
An unspecced patch of _should_use_exec accepts any call, and the noun or adjective helper names read as attributes.
Composite and included builds put directories on the runtime classpath, and a classpath file can be something other than a JAR. Neither path was tested, so a change there could leave a stale digest unnoticed.
The parse request no longer names Dags, and the Dag processor reuses an artifact's answer for every Dag file that uses the artifact, so the answer must not depend on what was asked. The vendored schema is resynced in the same commit because the request model generated from it no longer has dag_ids.
The supervisor schema dropped TaskHandlerParam.required, so a reply that still sends it no longer matches the schema the Dag processor decodes with. Flat params stay positional and a TaskInput stays named, which now means a field nothing fills only warns. The vendored schema is synced from the monorepo snapshot with sync-java-sdk-supervisor-schema.
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/08-java-sdk-task-handler-parse branch from f508e20 to f599ee6 Compare October 2, 2026 11:36
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/09-known-artifacts-push-down branch from 4ca5e76 to 7e11dfa Compare October 2, 2026 11:36
The Dag processor manager must know which bundles can hold task-handler artifacts before it dispatches a parse, and it must not build coordinators to find out, since that imports and constructs every runtime. The validated specs already carry the answer. Only the coordinators a queue routes to are listed, since only they run stub tasks.
The Dag-parsing child has no database, so the artifacts the fast path compares against must arrive with the parse request. Each carries its whole probe answer next to its fingerprint: a match then lets the child validate its stub tasks against that answer without starting a runtime, and an answer is only ever trusted under the fingerprint it was probed with. DagFileParseRequest already existed at 2026-06-16, so a runtime pinned there must not receive the new field.
check-ts-sdk-supervisor-schema fails whenever supervisor.ts lags the snapshot, and CI runs it over all files on any ts-sdk change. The shifted numeric suffixes are the generator's own renumbering; no hand-written TS code refers to them.
The vendored copy tracks the monorepo snapshot while both declare the same api_version, and sync-java-sdk-supervisor-schema rewrites it on the next Java SDK change otherwise.
check-go-sdk-generated-drift regenerates the Go SDK's coordinator-protocol models from the supervisor schema snapshot and fails on any difference, so DagFileParseRequest must carry known_artifacts on the Go side too. go-jsonschema moves some type blocks, so the diff is larger than the new field.
The fast path can only skip a probe if the child knows what the Dag processor already recorded, answers included. Reading once per loop, not once per pass, lets children started later see rows persisted by earlier ones, so a cold start does not re-probe every artifact for every file. An answer one file recorded is trusted by every file that reads it, so a child only gets the named bundles of its own team, plus its own Dag bundle when a coordinator falls back to it; never another Dag bundle, which keeps the request small. A broken [sdk] coordinators value must not stop Python-only parsing.
ADR-0013 predates the fallback to the task's own Dag bundle, so it said the manager reads only the named bundles and sends every child the same list. A fallback coordinator makes Dag bundles hold artifacts too, and no child can use another Dag bundle's rows. Each row now carries the artifact's whole answer, which every file reading it trusts, so a child only sees the named bundles of its own team, and the ADR's probe sketch still asked for Dag ids the runtimes no longer filter by. The message is named TaskHandlerArtifact, like the other task-handler parse messages.
The other tests that start children all configure a coordinator that falls back to the Dag's own bundle, so dropping that condition on the way to the child would have gone unnoticed, and an override of get_known_task_handler_artifacts may return bundles it was not asked for.
TaskHandlerParam no longer has required. A model dump no longer carries it, so the expected request bodies must drop it, and the constructors would otherwise pass a keyword pydantic only ignores.
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/09-known-artifacts-push-down branch from 7e11dfa to 00586db Compare October 2, 2026 14:57
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/08-java-sdk-task-handler-parse branch from f599ee6 to 94a8e22 Compare October 5, 2026 18:12
Base automatically changed from jason/core-taskhandler-refactor/08-java-sdk-task-handler-parse to jason/core-taskhandler-refactor/07-ts-sdk-task-handler-parse 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/07-ts-sdk-task-handler-parse branch from b52f31f to 9c284fc 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

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant