Skip to content

Respect run on latest version when clearing a running Dag run - #71425

Merged
dheerajturaga merged 3 commits into
apache:mainfrom
astronomer:fix-clear-run-on-latest-running-dr
Oct 7, 2026
Merged

dheerajturaga merged 3 commits into
apache:mainfrom
astronomer:fix-clear-run-on-latest-running-dr

Conversation

@ephraimbuddy

Copy link
Copy Markdown
Contributor

run_on_latest_version should have the same meaning for running and queued Dag runs, and when callers preserve a finished run's state. Otherwise cleared running tasks can retry against stale code while the Dag run still points at an older serialized Dag or bundle.

Keeping the run, cleared task instances, integrity checks, and bundle selection aligned preserves the user's explicit rerun choice without rewriting unrelated task-version history.


Was generative AI tooling used to co-author this PR?
  • Yes — Codex (5.6 Sol)

Generated-by: Codex (5.6 Sol) following the guidelines

@pierrejeambrun pierrejeambrun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM, but would love another pair of eyes.

Just one nit

Comment thread airflow-core/src/airflow/models/taskinstance.py Outdated
@dstandish

Copy link
Copy Markdown
Contributor

Hi Ephraim.

Want to flag a design concern

The field created_dag_version_id has existed as a column since AIP-65 DAG versioning (#42913/#43735). #49097 added the ORM relationship (created_dag_version) on top of it, with a docstring restating what the name already communicated: "the dag version column that was in effect at dag run creation time."

Documented or not, the name itself was always the contract — a write-once historical fact, not a live pointer. Per a conversation with @jedcunningham (the original author of AIP-65), that was intentional: bundle-versioned (pinned) runs were always meant to execute the version pinned at creation, with no path to move off it. "Run with the latest code" was supposed to be achieved by not using a versioned bundle — unpinned bundles already always resolve to the latest DagVersion (DBDagBag._version_from_dag_run).

Starting with #54984 and extended through #59764 and #65835/#66901, created_dag_version_id started being mutated after creation to support run_on_latest_version for pinned runs — a capability the original design didn't intend to exist for that case. This PR extends the same mutation to running/queued runs, going further in that direction.

That mutation has already produced concrete, reachable bugs elsewhere in the codebase, because other code correctly assumed the field's original (immutable) contract:

#71454 — DagRun.dag_versions silently drops versions still in use by task instances that weren't part of a given clear.
#71455 — only_new clear can report zero new tasks when there genuinely are some, for the same reason.
I've filed #71453 to track the underlying contract problem across all of this.

Given the original design intent, I think this is worth pausing on: should run_on_latest_version apply to bundle-versioned/pinned runs at all, rather than being extended further here? If the answer is yes and this is a deliberate, considered departure from the original pinning guarantee, it should probably written down somewhere (and the field's contract/docs updated to match), rather than continuing to build on a field whose name and documentation still say something the code no longer does.

@ephraimbuddy

ephraimbuddy commented Aug 14, 2026 •

Copy link
Copy Markdown
Contributor Author

Hi Ephraim.

Want to flag a design concern

The field created_dag_version_id has existed as a column since AIP-65 DAG versioning (#42913/#43735). #49097 added the ORM relationship (created_dag_version) on top of it, with a docstring restating what the name already communicated: "the dag version column that was in effect at dag run creation time."

Documented or not, the name itself was always the contract — a write-once historical fact, not a live pointer. Per a conversation with @jedcunningham (the original author of AIP-65), that was intentional: bundle-versioned (pinned) runs were always meant to execute the version pinned at creation, with no path to move off it. "Run with the latest code" was supposed to be achieved by not using a versioned bundle — unpinned bundles already always resolve to the latest DagVersion (DBDagBag._version_from_dag_run).

Starting with #54984 and extended through #59764 and #65835/#66901, created_dag_version_id started being mutated after creation to support run_on_latest_version for pinned runs — a capability the original design didn't intend to exist for that case. This PR extends the same mutation to running/queued runs, going further in that direction.

That mutation has already produced concrete, reachable bugs elsewhere in the codebase, because other code correctly assumed the field's original (immutable) contract:

#71454 — DagRun.dag_versions silently drops versions still in use by task instances that weren't part of a given clear. #71455 — only_new clear can report zero new tasks when there genuinely are some, for the same reason. I've filed #71453 to track the underlying contract problem across all of this.

Given the original design intent, I think this is worth pausing on: should run_on_latest_version apply to bundle-versioned/pinned runs at all, rather than being extended further here? If the answer is yes and this is a deliberate, considered departure from the original pinning guarantee, it should probably written down somewhere (and the field's contract/docs updated to match), rather than continuing to build on a field whose name and documentation still say something the code no longer does.

Thanks for raising this. I agree that the name and documentation of created_dag_version_id no longer match its behavior, and that should be resolved.

I don’t think that should block this PR, though. #71425 does not introduce the mutation for queued/running Dag runs: main already updates created_dag_version_id and bundle_version in that path following #65835, which was backported by #66901. This PR makes the existing behavior consistent by also moving a running TI that becomes RESTARTING, calling verify_integrity, and applying the update when dag_run_state=False.

The #71455 repro also does not reproduce on the finished-run path as written. That path calls verify_integrity, which creates the TI for task C during the first clear, so the later only_new result is correctly empty. There is a reachable queued/running variant because main currently moves the run pointer without verify_integrity; #71425 fixes that inconsistency.

I agree #71454 is valid. DagRun.dag_versions should report versions represented by the run’s TI/TIH rows rather than assume that a versioned run only contains its current pointer. Likewise, _get_new_task_ids should determine missing tasks from the run’s actual TI rows.

On the broader design question, run_on_latest_version is now a released user-facing capability across the clear UI/API, backfills, configuration, and the Dag-level parameter. The API still marks it experimental, but removing it from versioned bundles would nevertheless be a separate user-visible design decision. Disabling versioning is not equivalent: users can reasonably want runs pinned by default while deliberately rerunning a failed task against fixed code.

From the consumers I checked, the execution paths use created_dag_version_id as the version the run currently resolves to. I have not found a consumer that requires immutable creation-time provenance. Adding another column would also have unclear upgrade semantics because _update_dagrun_to_latest_version can rewrite all current TI version IDs without creating TI history, meaning the original value may already be unrecoverable.

My preference is therefore to land #71425 to restore the existing invariants, fix #71454 and _get_new_task_ids from the TI/TIH source of truth, and use #71453 to settle the documentation and possible physical rename separately.


Drafted-by: Codex (5.6 Sol); reviewed by @ephraimbuddy before posting

@bbovenzi

bbovenzi commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor

Let's rebase and get this in. So what we can then get #73915 merged too

@gang-zh

gang-zh commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

Looking forward to this change, can we backport it to 3.3-x

@dheerajturaga dheerajturaga left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This looks good. Lets merge this asap to get it in 3.4.0 . I was unable to rebase given I dont have write access to this branch. (need someone from astronomer here)
@bbovenzi , @ephraimbuddy Can you rebase? we should be able to merge

@dheerajturaga

Copy link
Copy Markdown
Member

Looking forward to this change, can we backport it to 3.3-x

next release is 3.4 . 3.3.x is not planned

ephraimbuddy and others added 3 commits October 7, 2026 10:49
run_on_latest_version should have the same meaning for running and queued Dag runs, and when callers preserve a finished run's state. Otherwise cleared running tasks can retry against stale code while the Dag run still points at an older serialized Dag or bundle.

Keeping the run, cleared task instances, integrity checks, and bundle selection aligned preserves the user's explicit rerun choice without rewriting unrelated task-version history.
Co-authored-by: Pierre Jeambrun <pierrejbrun@gmail.com>
@ephraimbuddy
ephraimbuddy force-pushed the fix-clear-run-on-latest-running-dr branch from 4c56b59 to e61db15 Compare October 7, 2026 09:58

@dheerajturaga dheerajturaga left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I re-read this after the rebase and the new third commit (e61db152). Moving the version choice off the RESTARTING row and into complete_restart is an improvement: the terminated attempt keeps the version it actually ran, and only the successor follows the run. My earlier approval stands. I left one minor, non-blocking observation inline on taskinstance.py:1244, which can be handled in a follow-up.

Worth a second look from

The third commit also changes how verify_integrity handles mapped tasks. These people know this area best:

  • @dstandish — wrote 8 of the recent commits to models/taskinstance.py / models/dagrun.py, and raised the created_dag_version_id contract question on this PR
  • @uranusjr — CODEOWNERS owner for the mapped-task code (models/mappedoperator.py, models/expandinput.py)

None of them have been notified — asking any of them for an
extra pass is the maintainer's call, and optional.


This review was drafted by an AI-assisted tool and
confirmed by an Apache Airflow maintainer. The findings
below are observations, not blockers; an Apache Airflow
maintainer — a real person — will take the next look at the
PR. If you think a finding is mis-applied, please reply on
the PR and a maintainer will weigh in.

More on how Apache Airflow handles maintainer review:
Contributing guide.

@dheerajturaga

Copy link
Copy Markdown
Member

On complete_restart in airflow-core/src/airflow/models/taskinstance.py:

if dag_version_id := DBDagBag._version_from_dag_run(self.dag_run, session=session):
successor.dag_version_id = dag_version_id

Minor, and fine to handle in a follow-up PR; it doesn't block this one.

When a run has been moved to a newer version, the successor now gets that version's dag_version_id. Both callers of complete_restart set self.task from get_dag_for_run(ti.dag_run), so max_tries follows the new version too.

The other task-derived columns don't follow it. prepare_db_for_next_try copies queue, pool, pool_slots, priority_weight, executor, executor_config, run_as_user and operator from the terminated attempt's row. schedule_tis only calls refresh_from_task for UP_FOR_RETRY task instances, so a successor in state None keeps the old values. If the new version moves the task to a different executor, queue or pool, the restarted attempt runs the new code with the old routing.

The non-running branch of clear_task_instances already handles this with ti.refresh_from_task(task, dag_run=dr). Something like this would make the two paths match:

if self.task is not None:
    successor.refresh_from_task(self.task, dag_run=self.dag_run)
    successor.max_tries = self.try_number + self.task.retries

Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting

@dheerajturaga
dheerajturaga merged commit 96582ad into apache:main Oct 7, 2026
149 of 150 checks passed
dheerajturaga added a commit to dheerajturaga/airflow that referenced this pull request Oct 7, 2026
The rebase onto apache#71425 dropped this PR's clear_task_instances change, and apache#71425 already tests repinning a run after a bundle-only update, so the test no longer exercises anything this PR changes.
dheerajturaga added a commit to dheerajturaga/airflow that referenced this pull request Oct 7, 2026
The rebase onto apache#71425 dropped this PR's clear_task_instances change, and apache#71425 already tests repinning a run after a bundle-only update, so the test no longer exercises anything this PR changes.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants