Skip to content

Record probed artifacts and task handler bindings from a Dag parse - #74031

Draft
jason810496 wants to merge 12 commits into
jason/core-taskhandler-refactor/09-known-artifacts-push-downfrom
jason/core-taskhandler-refactor/10-task-handler-bindings-reconcile
Draft

jason810496 wants to merge 12 commits into
jason/core-taskhandler-refactor/09-known-artifacts-push-downfrom
jason/core-taskhandler-refactor/10-task-handler-bindings-reconcile

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 processor has to persist which artifact runs each stub task, so a worker stops searching for it, and each probed artifact's answer, so no other parse probes that artifact again while it is unchanged (ADR-0013 "Flow 1"). This PR adds both to the parse result, writes them into the two tables from #73971, and reclaims artifact rows nothing uses. Nothing produces them yet (a later layer probes and validates), so in production the bindings are always None, probed_artifacts is empty, and nothing is written.

What changes

DagFileParsingResult gains the resolved bindings and the probed artifacts (TaskHandlerArtifact from the previous PR):

class TaskHandlerBinding(BaseModel):
    dag_id: str
    task_id: str
    artifact_bundle_name: str
    artifact_rel_path: str = Field(max_length=2000)


class DagFileParsingResult(BaseModel):
    ...
    task_handler_bindings: list[TaskHandlerBinding] | None = None
    probed_artifacts: list[TaskHandlerArtifact] = Field(default_factory=list)
  • probed_artifacts holds every artifact the parse probed, with its answer, even when the bindings are None, so a Dag that fails validation is not probed again on every parse. It is the only thing that writes an answer.
  • handle_parsing_result records them first, in a transaction of their own, and skips that when the list is empty. Each is upserted by (bundle_name, relative_fileloc_hash) with its size, digest, answer and last_probed_at, in key order (ON CONFLICT DO UPDATE on Postgres and SQLite, ON DUPLICATE KEY UPDATE on MySQL). An artifact outside the Dag file's scope (the bundles the previous PR sends the file: the named bundles of its team, and its own Dag bundle when a coordinator falls back to it) is dropped with a warning.
  • Writing them apart from the persist releases their exclusive locks before the reconcile takes shared ones, so two Dag processors cannot deadlock on crossed artifacts, and keeps an answer when the persist fails. A failed write is logged and the persist still runs. persist_probed_task_handler_artifacts() can be overridden like persist_parsing_result().
  • None leaves the recorded bindings alone. A list replaces the bindings of each Dag in serialized_dags, so [] deletes the rows of this result's Dags and nothing else. The max_length matches the column, so an oversized path fails where the child builds the result, not in the manager's transaction.
  • The reconcile runs inside update_dag_parsing_results_in_db, after the Dag rows and before the serialized Dags, and is retried with them on OperationalError. It runs in a savepoint: on an IntegrityError or DataError only the bindings are rolled back, with a warning naming the file, and the serialized Dags, code and import errors are still written. persist_parsing_result passes the bindings, the dispatched file path and the file's scope; bindings without a path or a scope raise ValueError. Because the bindings are written first, a throttled or failed serialized-Dag write can leave them briefly ahead of the serialized Dag; nothing reads the bindings yet, and the scheduler-read layer decides how to handle that.
  • Handler rows are reconciled by the result's Dag ids, never by file: rows whose task is gone are deleted, the rest inserted or updated. A Dag that moved files keeps its rows and gets the new path and dag_relative_fileloc_hash. Rows of Dags the result lacks stay. A Dag removed from its file is only marked stale, so its rows, which keep their artifacts from the sweep, stay until the dag row is deleted, and then go by cascade. The Dag rows are already locked FOR UPDATE for the persist, so these writes need no upsert.
  • The reconcile never writes the artifact table. It reads the bound artifacts' ids under a shared lock (FOR KEY SHARE / FOR SHARE, in key order) held until commit. The read also selects size_bytes, which is not in the unique index, so MySQL locks the row the sweep checks and not only the index entry; without it the SKIP LOCKED test fails on MySQL. A Dag bound to an artifact outside the scope, or to one with no row (swept concurrently, or its probed-artifact write failed), keeps its rows, with a warning; the next parse finds no recorded answer and probes again. A steady-state parse with bindings issues two SELECTs inside a savepoint and no writes.
  • A binding for a Dag the parse did not produce is dropped with a warning. One for a Dag rejected for another team's plugin class is dropped quietly, since that Dag already gets an import error. A Dag that binds one task twice keeps its rows as they are, with a warning; the parse reports that case as an import error.
  • The Dag processor deletes unreferenced artifact rows on [scheduler] parsing_cleanup_interval, in batches of 1000, each in its own transaction. FOR UPDATE SKIP LOCKED passes over rows a parse is binding. A row goes only once its last_probed_at is older than five cleanup intervals, so a probed artifact that binds nothing keeps its answer for the files parsed next. delete_unreferenced_task_handler_artifacts() can be overridden like cleanup_stale_bundle_versions().
  • This PR carries its own migration (0143, one last_probed_at index for the sweep) because the tables' migration ships in Add Lang-SDK task handler binding tables #73971, a separately reviewed PR.
  • AddTaskHandlerBindingsToDagFileParsingResult and AddProbedArtifactsToDagFileParsingResult in the in-progress 2026-10-30 version drop the fields for a runtime pinned to 2026-06-16, and upgrade that runtime's result with no bindings and no probed artifacts. They follow main's AddDagDefinitionsToDagFileParsingResult in that version. The snapshot, ts-sdk/src/generated/supervisor.ts, the vendored java-sdk/sdk/schema/schema.json and the Go genmodels are regenerated, since check-go-sdk-generated-drift (Go SDK: bring the coordinator-protocol models back to the supervisor schema #73963) fails when the Go genmodels lag the snapshot.
  • ADR-0013: the binding and probed_artifacts sketches, the probed-artifact write ahead of the persist, the missing and out-of-scope rule, [] narrowed to this result's Dags, a stale Dag's rows kept until its dag row is deleted, and "or by db clean" dropped (the artifact table is excluded from db clean).

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/execution_time/schema task-sdk/tests/task_sdk/execution_time/test_coordinator.py -q
breeze run --backend postgres pytest airflow-core/tests/unit/dag_processing/test_collection.py::TestRecordProbedTaskHandlerArtifacts airflow-core/tests/unit/dag_processing/test_collection.py::TestTaskHandlerBindingReconcile airflow-core/tests/unit/dag_processing/test_manager.py::TestTaskHandlerBindings airflow-core/tests/unit/dag_processing/test_manager.py::TestTaskHandlerArtifactSweep -q

Ran:

  • SQLite: airflow-core/tests/unit/dag_processing/ and test_lang_sdk_task_handler.py, 678 passed, 13 skipped (Postgres-only budgets, the Postgres/MySQL-only SKIP LOCKED test, opt-in real probes). test_db_cleanup.py: 98 passed, 3 skipped. Supervisor schema and coordinator tests: 93 passed.
  • Postgres and MySQL (breeze), on a schema built from the migrations: the schema/model sync test, TestRecordProbedTaskHandlerArtifacts, TestTaskHandlerBindingReconcile, the rejected-Dag binding test, TestKnownTaskHandlerArtifacts, TestTaskHandlerBindings, TestTaskHandlerArtifactSweep (with the SKIP LOCKED test) and the model tests, 65 passed each, and again after 0143 was downgraded (index gone) and upgraded. The Postgres statement-budget tests passed unchanged.
  • The savepoint test breaks a foreign key for a real IntegrityError and injects a DataError, on all three backends; without the savepoint both fail the persist.
  • The race test commits an artifact from a second, separate session before the upsert, which then takes the conflict path. The re-probe test reads the answer back after that path, so the JSON value binds as the column type in the update clause on all three backends.
  • prek run --from-ref <parent> --stage pre-commit passed, including check-supervisor-schemas-versions without SKIP, mypy-airflow-core, mypy-task-sdk and check-core-imports. prek run migration-round-trip --hook-stage manual, compile-ts-sdk and ./gradlew -p java-sdk :sdk:test --tests '*TaskHandlerParseTest*' passed.
  • Without AddProbedArtifactsToDagFileParsingResult in the bundle, a downgrade to 2026-06-16 keeps the field, which the migrator test catches.

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

@jason810496
jason810496 added this pull request to stack #73978 October 1, 2026 13:34
@jason810496 jason810496 changed the title jason/core taskhandler refactor/10 task handler bindings reconcile Record probed artifacts and task handler bindings from a Dag parse Oct 1, 2026
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/10-task-handler-bindings-reconcile branch from 568fa9f to fead951 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/10-task-handler-bindings-reconcile branch from fead951 to a937384 Compare October 2, 2026 00:54
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/09-known-artifacts-push-down branch 2 times, most recently from 7b0b329 to 4ca5e76 Compare October 2, 2026 05:48
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/10-task-handler-bindings-reconcile branch from a937384 to 391143b Compare October 2, 2026 05:48
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/10-task-handler-bindings-reconcile branch from 391143b to dc3fd66 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 manager has to persist which artifact runs each stub task, so a worker no longer searches for it (ADR-0013), and the answers the parse probed, so no other parse probes those artifacts again. None tells a parse that did not evaluate task handlers apart from one that found none, so a failed or skipped evaluation cannot wipe the recorded bindings. The answers travel apart from the bindings because they are worth keeping when a Dag fails validation and has none, and a binding names only its artifact, whose fingerprint and answer belong to the artifact rather than to the stub task. DagFileParsingResult already existed at 2026-06-16, so each new field needs a VersionChange: a runtime pinned there is upgraded with the bindings unset, which means "leave them alone", and no probed artifacts, which records nothing.
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 DagFileParsingResult must carry task_handler_bindings and probed_artifacts on the Go side too. go-jsonschema moves some type blocks, so the diff is larger than the new fields.
A handler row belongs to a Dag, not to the file that last wrote it, so the rows are reconciled by the result's Dag ids: a Dag that moved files keeps its rows, and a parse that yields fewer Dags cannot delete rows it did not evaluate. The Dag rows are already locked for the persist, so handler rows need no upsert. Artifact rows are shared across files and Dag processors and hold answers that other files trust, so only a probe writes them, only within the Dag file's scope, and in a transaction of its own before the persist: its exclusive locks are released before the reconcile takes shared ones, so two Dag processors cannot deadlock on crossed artifacts, and an answer survives a failed persist. The reconcile only looks artifacts up, under a shared lock that keeps the orphan sweep from deleting one before its handler rows are in; a Dag whose artifact has no row or lies outside the scope keeps its rows, and the next parse probes again.
The orphan sweep deletes unreferenced artifact rows only once last_probed_at is past a grace period, and a batched cleanup must filter on indexed columns. It is a separate migration because the tables' migration already ships in an earlier, separately reviewed PR.
A parse never evicts artifact rows, because one artifact serves Dags in many files, so the Dag processor reclaims them itself; db clean leaves the table alone. The grace period keeps rows for probed artifacts that bind nothing, which later parses reuse instead of probing again. Batches commit on their own so the sweep never holds locks across the whole table, and SKIP LOCKED keeps it from waiting on, or deleting under, a parse that is binding an artifact.
A binding names only its artifact, and the probed answers travel in a field of their own, because an answer belongs to the artifact and is worth keeping even when a Dag fails validation. Writing them in their own transaction before the persist takes the upsert out of the reconcile, which then only looks artifacts up, so the ADR has to say what happens to a Dag whose artifact is missing or outside the file's scope. An empty list deletes only the rows of the Dags in the result, since the reconcile is keyed by Dag id rather than by file. db clean excludes the artifact table, so the orphan sweep is the only reclaim path.
create_session() hands back the thread's scoped session, so the session standing in for a second Dag processor shared the test's transaction: on Postgres the sweep then deleted a row the "parse" had locked, because the lock was its own. A non-scoped session is a separate transaction, as a second processor would be.
A Dag rejected for another team's plugin class already gets an import error on every parse, so warning about its dropped bindings as well was noise that buried the case worth a warning: a binding for a Dag the parse never produced, which points at a bug in the parse.
The sweep's interval check now reads like the bundle-version cleanup next to it, and its deleted-row count no longer hides a missing rowcount behind a default that can never apply. The parsing_cleanup_interval description said the artifact deletion in a clause that did not parse cleanly.
TaskHandlerParam no longer has required. A model dump no longer carries it, so the expected result bodies must drop it, and the constructors would otherwise pass a keyword pydantic only ignores.

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