You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Since #68336 (fix for #67657), when a bundle advances to a new commit but the Dag's serialized
content is unchanged, SerializedDagModel.write_dag repoints the latest DagVersion.bundle_version
to the new commit in place instead of creating a new DagVersion — even when that version already
has task instances
(serialized_dag.py#L704-L722).
A code-only change — e.g. editing the body of a @task function, which is not part of the
serialized Dag — therefore rewrites history of an existing DagVersion. Two user-visible side
effects follow:
Old runs show the wrong commit. A run created on commit A keeps dag_run.bundle_version = A
and correctly executes A (workloads use dag_run.bundle_version), but GET /dagRuns/{run_id}
reports its dag_versions[].bundle_version and bundle_url as commit B. The UI "view source"
link for that run points at code it did not run.
A run's reported version/commit/source link should be the commit it actually executed
(dag_run.bundle_version), not whatever the DagVersion row points at now.
Triggering at a commit Airflow has already parsed and run should work, or fail with a message
explaining that the commit was superseded in place.
This should not undo #67657's fix. Possible directions (for discussion): derive per-run
version/URL from dag_run.bundle_version; record superseded commits for a DagVersion
(e.g. a small dag_version_bundle_version history table) so lookups by commit still resolve;
or create a new DagVersion when the current one is already referenced by runs.
How to reproduce
Git-backed bundle, a Dag with a @task:
Commit A (return "GOOD"), let it parse, trigger a run.
Commit B changing only the task body (return "BUGGY"), let it parse.
dag_version still has one row, now with bundle_version = B (same dag_hash).
POST /dags/{dag_id}/dagRuns {"bundle_version": "<A>"} → 404.
GET /dags/{dag_id}/dagRuns/<run from step 1> → bundle_version: A but dag_versions[0].bundle_version: B.
Log from a local reproduction using real git commits, the Dag processor's sync_bag_to_db
path and the real REST endpoint:
commit A = 04f74a7423 (compute() returns GOOD)
v1 bundle_version=04f74a7423 dag_hash=95fd125570
POST .../dagRuns bundle_version=04f74a7423 -> HTTP 200
commit B = 607ad282da (compute() returns BUGGY; Dag structure unchanged)
v1 bundle_version=607ad282da dag_hash=95fd125570
POST .../dagRuns bundle_version=04f74a7423 -> HTTP 404
existing run from A:
dag_run.bundle_version = 04f74a7423
workload.bundle_info.version = 04f74a7423
GET dagRun dag_versions[v1].bundle_version = 607ad282da bundle_url=.../tree/607ad282da
Note: with the default [core] min_serialized_dag_update_interval (30s), commits parsed within
30s of each other are collapsed as well.
Anything else
Also relevant to AIP-109 (Dag version pinning), where rollback targets DagVersion rows.
Are you willing to submit PR?
Yes I am willing to submit a PR!
Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting
Apache Airflow version
main (reproduced at 5428b3b)
What happened
Since #68336 (fix for #67657), when a bundle advances to a new commit but the Dag's serialized
content is unchanged,
SerializedDagModel.write_dagrepoints the latestDagVersion.bundle_versionto the new commit in place instead of creating a new
DagVersion— even when that version alreadyhas task instances
(serialized_dag.py#L704-L722).
A code-only change — e.g. editing the body of a
@taskfunction, which is not part of theserialized Dag — therefore rewrites history of an existing
DagVersion. Two user-visible sideeffects follow:
dag_run.bundle_version = Aand correctly executes A (workloads use
dag_run.bundle_version), butGET /dagRuns/{run_id}reports its
dag_versions[].bundle_versionandbundle_urlas commit B. The UI "view source"link for that run points at code it did not run.
POST /dags/{dag_id}/dagRunswithbundle_version(Add the option to select bundle version parameter on dag run trigger endpoint #61550) looks up a
DagVersionrow bybundle_version(dag_version.py#L176-L178),
so for any commit that was overwritten in place it returns
404 ... does not have a version for bundle_version '<A>', even though A was a real, parsed,previously run version.
What you think should happen instead
(
dag_run.bundle_version), not whatever theDagVersionrow points at now.explaining that the commit was superseded in place.
This should not undo #67657's fix. Possible directions (for discussion): derive per-run
version/URL from
dag_run.bundle_version; record superseded commits for aDagVersion(e.g. a small
dag_version_bundle_versionhistory table) so lookups by commit still resolve;or create a new
DagVersionwhen the current one is already referenced by runs.How to reproduce
Git-backed bundle, a Dag with a
@task:return "GOOD"), let it parse, trigger a run.return "BUGGY"), let it parse.dag_versionstill has one row, now withbundle_version = B(samedag_hash).POST /dags/{dag_id}/dagRuns {"bundle_version": "<A>"}→ 404.GET /dags/{dag_id}/dagRuns/<run from step 1>→bundle_version: Abutdag_versions[0].bundle_version: B.Log from a local reproduction using real git commits, the Dag processor's
sync_bag_to_dbpath and the real REST endpoint:
Note: with the default
[core] min_serialized_dag_update_interval(30s), commits parsed within30s of each other are collapsed as well.
Anything else
Also relevant to AIP-109 (Dag version pinning), where rollback targets
DagVersionrows.Are you willing to submit PR?
Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting