This is an automated email from the ASF dual-hosted git repository.

jedcunningham pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new f5ec36107f9 Fix cleared tasks getting stuck when a Dag run has no 
version (#71696)
f5ec36107f9 is described below

commit f5ec36107f9e02fa3af42a8c9129ad7a9ff56e85
Author: Jed Cunningham <[email protected]>
AuthorDate: Mon Aug 17 16:28:04 2026 -0600

    Fix cleared tasks getting stuck when a Dag run has no version (#71696)
    
    A Dag run can have no version of its own: runs carried over from Airflow 2
    predate the version columns, and `airflow db clean` can remove the version a
    run was created with. Clearing such a run left its tasks without a version
    too, and a task instance with no version is never enqueued — so the task sat
    in queued until it timed out, while its revived run held up new ones. The
    clear dialogs offered no way out either, since they only showed the
    run-on-latest option when there was a version to compare against.
---
 .../administration-and-deployment/dag-bundles.rst  |   4 +
 airflow-core/src/airflow/models/dagbag.py          |   5 +-
 airflow-core/src/airflow/models/taskinstance.py    |  53 +++++-
 .../airflow/ui/public/i18n/locales/en/dags.json    |   1 +
 .../ui/src/components/Clear/Run/ClearRunDialog.tsx |  30 ++--
 .../TaskInstance/ClearGroupTaskInstanceDialog.tsx  |  32 +++-
 .../Clear/TaskInstance/ClearTaskInstanceDialog.tsx |  26 ++-
 .../Clear/TaskInstance/runOnLatestVersion.test.ts  |  14 ++
 .../Clear/TaskInstance/runOnLatestVersion.ts       |  21 ++-
 .../core_api/routes/public/test_dag_run.py         |  17 ++
 airflow-core/tests/unit/models/test_cleartasks.py  | 179 ++++++++++++++++++++-
 11 files changed, 344 insertions(+), 38 deletions(-)

diff --git a/airflow-core/docs/administration-and-deployment/dag-bundles.rst 
b/airflow-core/docs/administration-and-deployment/dag-bundles.rst
index 354e6ecda16..ef65b5eaa71 100644
--- a/airflow-core/docs/administration-and-deployment/dag-bundles.rst
+++ b/airflow-core/docs/administration-and-deployment/dag-bundles.rst
@@ -276,6 +276,10 @@ The setting is resolved using the following precedence 
(highest to lowest):
 4. **Per-call-site fallback**: ``False`` for clear/rerun, ``True`` for 
backfills (preserving
    the historical default for each path)
 
+One exception: a Dag run with no version of its own — carried over from 
Airflow 2, or its version
+since removed by ``airflow db clean`` — has nothing to preserve, so clearing 
it always uses the
+latest version and bundle version regardless of the resolved setting.
+
 Global Configuration
 ~~~~~~~~~~~~~~~~~~~~
 
diff --git a/airflow-core/src/airflow/models/dagbag.py 
b/airflow-core/src/airflow/models/dagbag.py
index c4bd8eceea1..5c0556974ff 100644
--- a/airflow-core/src/airflow/models/dagbag.py
+++ b/airflow-core/src/airflow/models/dagbag.py
@@ -209,7 +209,10 @@ class DBDagBag:
 
     @staticmethod
     def _version_from_dag_run(dag_run: DagRun, *, session: Session) -> UUID | 
None:
-        if not dag_run.bundle_version:
+        # A run with no version of its own can only resolve to the latest. 
Runs carried over from
+        # Airflow 2 are like this, as are runs whose version `airflow db 
clean` has since deleted --
+        # the latter keep their bundle version, so they would otherwise 
resolve to nothing at all.
+        if not dag_run.bundle_version or not dag_run.created_dag_version_id:
             if dag_version := 
DagVersion.get_latest_version(dag_id=dag_run.dag_id, session=session):
                 return dag_version.id
 
diff --git a/airflow-core/src/airflow/models/taskinstance.py 
b/airflow-core/src/airflow/models/taskinstance.py
index 6a91d18223e..ff610f21465 100644
--- a/airflow-core/src/airflow/models/taskinstance.py
+++ b/airflow-core/src/airflow/models/taskinstance.py
@@ -352,6 +352,27 @@ def _update_dagrun_to_latest_version(
     session.flush()
 
 
+def _pin_versionless_tis_to_run_version(dag_run: DagRun, dag_version_id: UUID, 
session: Session) -> None:
+    """
+    Give the run's unfinished task instances a dag version if they have none.
+
+    Once the run is pinned the scheduler stops backfilling versions onto them, 
and one
+    without a version is never enqueued.
+    """
+    session.execute(
+        update(TaskInstance)
+        .where(
+            TaskInstance.dag_id == dag_run.dag_id,
+            TaskInstance.run_id == dag_run.run_id,
+            TaskInstance.dag_version_id.is_(None),
+            # State.unfinished holds None, which SQL IN never matches.
+            or_(TaskInstance.state.is_(None), 
TaskInstance.state.in_(State.unfinished)),
+        )
+        .values(dag_version_id=dag_version_id)
+        .execution_options(synchronize_session="evaluate")
+    )
+
+
 def clear_task_instances(
     tis: list[TaskInstance],
     session: Session,
@@ -372,7 +393,9 @@ def clear_task_instances(
     :param session: current session
     :param dag_run_state: state to set finished DagRuns to.
         If set to False, DagRuns state will not be changed.
-    :param run_on_latest_version: whether to run on latest serialized DAG and 
Bundle version
+    :param run_on_latest_version: whether to run on latest serialized DAG and 
Bundle version.
+        A run with no version of its own uses the latest either way, since 
there is nothing
+        else for it to run on; a task instance with no version joins its run's.
 
     :meta private:
     """
@@ -396,7 +419,10 @@ def clear_task_instances(
         # the task is terminated and becomes eligible for retry.
         else:
             dr = ti.dag_run
-            if run_on_latest_version:
+            # A run with no version of its own has nothing to re-run on but 
the latest, and the
+            # run loop below moves it there.
+            use_latest_version = run_on_latest_version or 
dr.created_dag_version_id is None
+            if use_latest_version:
                 ti_dag = scheduler_dagbag.get_latest_version_of_dag(ti.dag_id, 
session=session)
             else:
                 ti_dag = scheduler_dagbag.get_dag_for_run(dag_run=dr, 
session=session)
@@ -418,11 +444,15 @@ def clear_task_instances(
             ti.state = None
             ti.external_executor_id = None
             ti.clear_next_method_args()
-            # Match DagVersion to latest serialized DAG when 
run_on_latest_version.
-            if run_on_latest_version:
+            # Match DagVersion to latest serialized DAG when running on the 
latest version.
+            if use_latest_version:
                 latest_dag_version = DagVersion.get_latest_version(ti.dag_id, 
session=session)
                 if latest_dag_version is not None:
                     ti.dag_version_id = latest_dag_version.id
+            elif ti.dag_version_id is None:
+                # One without a version is never enqueued, and the run keeps 
its own, so it can
+                # only go there.
+                ti.dag_version_id = dr.created_dag_version_id
             session.merge(ti)
 
     if dag_run_state is not False and tis:
@@ -461,10 +491,14 @@ def clear_task_instances(
 
             _recalculate_dagrun_queued_at_deadlines(dr, dr.queued_at, session)
 
+            # A run with no version of its own has nothing to preserve, so the 
latest is all
+            # it can be re-run on. Runs migrated from Airflow 2 are like this, 
as are runs
+            # whose version `airflow db clean` has since deleted.
+            use_latest_version = run_on_latest_version or 
dr.created_dag_version_id is None
             if dr.state in State.finished_dr_states:
                 dr.state = dag_run_state
                 dr.start_date = timezone.utcnow()
-                if run_on_latest_version:
+                if use_latest_version:
                     dr_dag = 
scheduler_dagbag.get_latest_version_of_dag(dr.dag_id, session=session)
                     dag_version = DagVersion.get_latest_version(dr.dag_id, 
session=session)
                     if dag_version:
@@ -478,14 +512,14 @@ def clear_task_instances(
                     dr_dag = scheduler_dagbag.get_dag_for_run(dag_run=dr, 
session=session)
                 if not dr_dag:
                     log.warning("No serialized dag found for dag '%s'", 
dr.dag_id)
-                if dr_dag and not dr_dag.disable_bundle_versioning and 
run_on_latest_version:
+                if dr_dag and not dr_dag.disable_bundle_versioning and 
use_latest_version:
                     bundle_version = dr.dag_model.bundle_version
-                    if bundle_version is not None and run_on_latest_version:
+                    if bundle_version is not None:
                         dr.bundle_version = bundle_version
                 if dag_run_state == DagRunState.QUEUED:
                     dr.last_scheduling_decision = None
                     dr.start_date = None
-            elif run_on_latest_version:
+            elif use_latest_version:
                 # Queued/running DagRun: update DR to latest version/bundle 
for workloads that use it.
                 dag_version = DagVersion.get_latest_version(dr.dag_id, 
session=session)
                 if dag_version and dr.created_dag_version_id != dag_version.id:
@@ -499,6 +533,9 @@ def clear_task_instances(
                             bundle_version = dr.dag_model.bundle_version
                             if bundle_version is not None:
                                 dr.bundle_version = bundle_version
+
+            if dr.created_dag_version_id:
+                _pin_versionless_tis_to_run_version(dr, 
dr.created_dag_version_id, session)
     for ti in tis:
         ti.context_carrier = new_task_run_carrier(ti.dag_run.context_carrier)
     session.flush()
diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json 
b/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
index 753f7fe53fd..11c34b00dd5 100644
--- a/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
+++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/dags.json
@@ -77,6 +77,7 @@
       "preventRunningTasks": "Prevent rerun if task is running",
       "queueNew": "Queue up new tasks",
       "runOnLatestVersion": "Run with latest bundle version",
+      "runOnLatestVersionForced": "Always uses the latest — there's no earlier 
version to go back to",
       "upstream": "Upstream"
     }
   },
diff --git 
a/airflow-core/src/airflow/ui/src/components/Clear/Run/ClearRunDialog.tsx 
b/airflow-core/src/airflow/ui/src/components/Clear/Run/ClearRunDialog.tsx
index c4124cc8878..7be08d27af4 100644
--- a/airflow-core/src/airflow/ui/src/components/Clear/Run/ClearRunDialog.tsx
+++ b/airflow-core/src/airflow/ui/src/components/Clear/Run/ClearRunDialog.tsx
@@ -65,6 +65,17 @@ const ClearRunDialog = ({ dagRun, onClose, open }: Props) => 
{
     dagId,
   });
 
+  // Offered only where it changes the outcome. A non-versioned bundle (e.g. 
LocalDagBundle)
+  // leaves bundle_version null and resolves to the latest serialized Dag at 
run time anyway,
+  // so unless the run has no version at all the option would be a no-op there.
+  const { runOnLatestVersionForced, shouldShowRunOnLatestOption } = 
getRunOnLatestVersionState({
+    latestBundleVersion: dagDetails?.bundle_version,
+    latestDagVersionNumber: dagDetails?.latest_dag_version?.version_number,
+    selectedBundleVersion: dagRun.bundle_version,
+    selectedDagVersionNumber: dagRun.dag_versions.at(-1)?.version_number,
+    selectedVersionMissing: dagRun.dag_versions.length === 0,
+  });
+
   const { setValue: setRunOnLatestVersion, value: runOnLatestVersion } = 
useRerunWithLatestVersion({
     dagLevelConfig: dagDetails?.rerun_with_latest_version,
   });
@@ -94,17 +105,6 @@ const ClearRunDialog = ({ dagRun, onClose, open }: Props) 
=> {
     onSuccessConfirm: handleClose,
   });
 
-  // Non-versioned bundles (e.g. LocalDagBundle) always leave bundle_version 
null and
-  // resolve to the latest serialized Dag at run time, so "run on latest" is a 
no-op there.
-  // Offer it only when re-running on the latest would actually change the 
outcome:
-  // the run's Dag version differs from the latest while the bundle is 
versioned
-  // (latest bundle_version present), or the run's bundle version differs from 
the latest.
-  const { shouldShowRunOnLatestOption } = getRunOnLatestVersionState({
-    latestBundleVersion: dagDetails?.bundle_version,
-    latestDagVersionNumber: dagDetails?.latest_dag_version?.version_number,
-    selectedBundleVersion: dagRun.bundle_version,
-    selectedDagVersionNumber: dagRun.dag_versions.at(-1)?.version_number,
-  });
   const shouldShowBundleVersionOption = shouldShowRunOnLatestOption && 
!onlyNew;
 
   return (
@@ -160,8 +160,14 @@ const ClearRunDialog = ({ dagRun, onClose, open }: Props) 
=> {
           >
             {shouldShowBundleVersionOption ? (
               <Checkbox
-                checked={runOnLatestVersion}
+                checked={runOnLatestVersionForced || runOnLatestVersion}
+                disabled={runOnLatestVersionForced}
                 onCheckedChange={(event) => 
setRunOnLatestVersion(Boolean(event.checked))}
+                title={
+                  runOnLatestVersionForced
+                    ? 
translate("dags:runAndTaskActions.options.runOnLatestVersionForced")
+                    : undefined
+                }
               >
                 
{translate("dags:runAndTaskActions.options.runOnLatestVersion")}
               </Checkbox>
diff --git 
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
 
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
index 6a7e33354ae..aec9bd55c66 100644
--- 
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
+++ 
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
@@ -22,7 +22,11 @@ import { useTranslation } from "react-i18next";
 import { CgRedo } from "react-icons/cg";
 import { useParams } from "react-router-dom";
 
-import { useDagServiceGetDagDetails, useTaskInstanceServiceGetTaskInstances } 
from "openapi/queries";
+import {
+  useDagRunServiceGetDagRun,
+  useDagServiceGetDagDetails,
+  useTaskInstanceServiceGetTaskInstances,
+} from "openapi/queries";
 import type { LightGridTaskInstanceSummary, TaskInstanceResponse } from 
"openapi/requests/types.gen";
 import { ActionAccordion } from "src/components/ActionAccordion";
 import { useRerunWithLatestVersion } from 
"src/components/Clear/useRerunWithLatestVersion";
@@ -80,14 +84,20 @@ export const ClearGroupTaskInstanceDialog = ({ onClose, 
open, taskInstance }: Pr
 
   const groupTaskIds = groupTaskInstances?.task_instances.map((ti) => 
ti.task_id) ?? [];
 
-  const { dagVersionsDiffer, shouldShowRunOnLatestOption } = 
getRunOnLatestVersionState({
-    latestBundleVersion: dagDetails?.bundle_version,
-    latestDagVersionNumber: dagDetails?.latest_dag_version?.version_number,
-    selectedDagVersionNumber: taskInstance.dag_version_number,
-    // Fall back to legacy heuristic when grid summary has no version (older 
API).
-    useLatestBundleVersionAsFallback: true,
+  const { data: dagRun } = useDagRunServiceGetDagRun({ dagId, dagRunId: runId 
}, undefined, {
+    enabled: open,
   });
 
+  const { dagVersionsDiffer, runOnLatestVersionForced, 
shouldShowRunOnLatestOption } =
+    getRunOnLatestVersionState({
+      latestBundleVersion: dagDetails?.bundle_version,
+      latestDagVersionNumber: dagDetails?.latest_dag_version?.version_number,
+      selectedDagVersionNumber: taskInstance.dag_version_number,
+      selectedVersionMissing: dagRun?.dag_versions.length === 0,
+      // Fall back to legacy heuristic when grid summary has no version (older 
API).
+      useLatestBundleVersionAsFallback: true,
+    });
+
   // dagVersionsDiffer becomes the fallback so the historical "auto-check when 
versions
   // differ" heuristic still applies when neither DAG-level nor global config 
is set.
   const { setValue: setRunOnLatestVersion, value: runOnLatestVersion } = 
useRerunWithLatestVersion({
@@ -182,8 +192,14 @@ export const ClearGroupTaskInstanceDialog = ({ onClose, 
open, taskInstance }: Pr
           >
             {shouldShowRunOnLatestOption ? (
               <Checkbox
-                checked={runOnLatestVersion}
+                checked={runOnLatestVersionForced || runOnLatestVersion}
+                disabled={runOnLatestVersionForced}
                 onCheckedChange={(event) => 
setRunOnLatestVersion(Boolean(event.checked))}
+                title={
+                  runOnLatestVersionForced
+                    ? 
translate("dags:runAndTaskActions.options.runOnLatestVersionForced")
+                    : undefined
+                }
               >
                 
{translate("dags:runAndTaskActions.options.runOnLatestVersion")}
               </Checkbox>
diff --git 
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx
 
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx
index 53525e320b6..074e3f2d40a 100644
--- 
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx
+++ 
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx
@@ -21,7 +21,7 @@ import { useEffect, useMemo, useState } from "react";
 import { useTranslation } from "react-i18next";
 import { CgRedo } from "react-icons/cg";
 
-import { useDagServiceGetDagDetails } from "openapi/queries";
+import { useDagRunServiceGetDagRun, useDagServiceGetDagDetails } from 
"openapi/queries";
 import type { ClearTaskInstancesBody, TaskInstanceResponse } from 
"openapi/requests/types.gen";
 import { ActionAccordion } from "src/components/ActionAccordion";
 import { taskInstanceKey } from "src/components/ActionAccordion/columns";
@@ -110,13 +110,19 @@ const ClearTaskInstanceDialog = (props: Props) => {
     dagId,
   });
 
-  const { dagVersionsDiffer, shouldShowRunOnLatestOption } = 
getRunOnLatestVersionState({
-    latestBundleVersion: dagDetails?.bundle_version,
-    latestDagVersionNumber: dagDetails?.latest_dag_version?.version_number,
-    selectedBundleVersion: taskInstance?.dag_version?.bundle_version,
-    selectedDagVersionNumber: taskInstance?.dag_version?.version_number,
+  const { data: dagRun } = useDagRunServiceGetDagRun({ dagId, dagRunId }, 
undefined, {
+    enabled: openDialog,
   });
 
+  const { dagVersionsDiffer, runOnLatestVersionForced, 
shouldShowRunOnLatestOption } =
+    getRunOnLatestVersionState({
+      latestBundleVersion: dagDetails?.bundle_version,
+      latestDagVersionNumber: dagDetails?.latest_dag_version?.version_number,
+      selectedBundleVersion: taskInstance?.dag_version?.bundle_version,
+      selectedDagVersionNumber: taskInstance?.dag_version?.version_number,
+      selectedVersionMissing: dagRun?.dag_versions.length === 0,
+    });
+
   // dagVersionsDiffer becomes the fallback so the historical "auto-check when 
versions
   // differ" heuristic still applies when neither DAG-level nor global config 
is set.
   const { setValue: setRunOnLatestVersion, value: runOnLatestVersion } = 
useRerunWithLatestVersion({
@@ -264,8 +270,14 @@ const ClearTaskInstanceDialog = (props: Props) => {
             >
               {shouldShowRunOnLatestOption ? (
                 <Checkbox
-                  checked={runOnLatestVersion}
+                  checked={runOnLatestVersionForced || runOnLatestVersion}
+                  disabled={runOnLatestVersionForced}
                   onCheckedChange={(event) => 
setRunOnLatestVersion(Boolean(event.checked))}
+                  title={
+                    runOnLatestVersionForced
+                      ? 
translate("dags:runAndTaskActions.options.runOnLatestVersionForced")
+                      : undefined
+                  }
                 >
                   
{translate("dags:runAndTaskActions.options.runOnLatestVersion")}
                 </Checkbox>
diff --git 
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/runOnLatestVersion.test.ts
 
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/runOnLatestVersion.test.ts
index 1f3d6ba4d3c..cc09d604e65 100644
--- 
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/runOnLatestVersion.test.ts
+++ 
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/runOnLatestVersion.test.ts
@@ -137,15 +137,27 @@ describe("getRunOnLatestVersionState", () => {
       name: "does not show for group fallback when latest bundle is missing",
       useLatestBundleVersionAsFallback: true,
     },
+    {
+      expectedDagVersionsDiffer: false,
+      expectedRunOnLatestVersionForced: true,
+      expectedShouldShowRunOnLatestOption: true,
+      // A null latest bundle version pins the case that matters: the option 
is forced even
+      // on a non-versioned bundle, where it would otherwise never be offered.
+      latestBundleVersion: null,
+      name: "forces and shows the option when the selection has no Dag version 
at all",
+      selectedVersionMissing: true,
+    },
   ])(
     "$name",
     ({
       expectedDagVersionsDiffer,
+      expectedRunOnLatestVersionForced = false,
       expectedShouldShowRunOnLatestOption,
       latestBundleVersion,
       latestDagVersionNumber,
       selectedBundleVersion,
       selectedDagVersionNumber,
+      selectedVersionMissing,
       useLatestBundleVersionAsFallback,
     }) => {
       expect(
@@ -154,10 +166,12 @@ describe("getRunOnLatestVersionState", () => {
           latestDagVersionNumber,
           selectedBundleVersion,
           selectedDagVersionNumber,
+          selectedVersionMissing,
           useLatestBundleVersionAsFallback,
         }),
       ).toEqual({
         dagVersionsDiffer: expectedDagVersionsDiffer,
+        runOnLatestVersionForced: expectedRunOnLatestVersionForced,
         shouldShowRunOnLatestOption: expectedShouldShowRunOnLatestOption,
       });
     },
diff --git 
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/runOnLatestVersion.ts
 
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/runOnLatestVersion.ts
index fa183b37d0c..21d7591e84c 100644
--- 
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/runOnLatestVersion.ts
+++ 
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/runOnLatestVersion.ts
@@ -22,11 +22,26 @@ type RunOnLatestVersionParams = {
   readonly latestDagVersionNumber?: number | null;
   readonly selectedBundleVersion?: string | null;
   readonly selectedDagVersionNumber?: number | null;
+  /**
+   * True when the *run* being cleared has no Dag version at all, which is the 
case for
+   * anything carried over from Airflow 2. There is nothing to re-run it on 
but the latest
+   * version, so the backend forces that regardless of the request. Keep this 
keyed off the
+   * run: a task instance with no version of its own is given its run's 
version, not the
+   * latest, so deriving this from the task instance would promise the wrong 
thing.
+   */
+  readonly selectedVersionMissing?: boolean;
   readonly useLatestBundleVersionAsFallback?: boolean;
 };
 
 type RunOnLatestVersionState = {
   readonly dagVersionsDiffer: boolean;
+  /**
+   * Drives how the checkbox renders, not what is submitted. A clear can span 
several runs
+   * (via past/future) while the request carries one flag for all of them, so 
forcing it
+   * would pin runs the user never selected. The backend forces each 
version-less run on
+   * its own instead.
+   */
+  readonly runOnLatestVersionForced: boolean;
   readonly shouldShowRunOnLatestOption: boolean;
 };
 
@@ -38,6 +53,7 @@ export const getRunOnLatestVersionState = ({
   latestDagVersionNumber,
   selectedBundleVersion,
   selectedDagVersionNumber,
+  selectedVersionMissing = false,
   useLatestBundleVersionAsFallback = false,
 }: RunOnLatestVersionParams): RunOnLatestVersionState => {
   const dagVersionsDiffer =
@@ -55,7 +71,10 @@ export const getRunOnLatestVersionState = ({
 
   return {
     dagVersionsDiffer,
+    runOnLatestVersionForced: selectedVersionMissing,
     shouldShowRunOnLatestOption:
-      (dagVersionsDiffer && hasBundleVersion(latestBundleVersion)) || 
shouldShowForBundleVersion,
+      selectedVersionMissing ||
+      (dagVersionsDiffer && hasBundleVersion(latestBundleVersion)) ||
+      shouldShowForBundleVersion,
   };
 };
diff --git 
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py 
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py
index 9d27ddf2817..4aef320cc63 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_run.py
@@ -2147,6 +2147,23 @@ class TestClearDagRun:
             logical_date=None,
         )
 
+    @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle")
+    def test_clear_dag_run_whose_dag_version_was_deleted(self, test_client, 
session):
+        """A run that kept its bundle version after ``airflow db clean`` 
removed its Dag version."""
+        session.execute(
+            update(DagRun)
+            .where(DagRun.dag_id == DAG1_ID, DagRun.run_id == DAG1_RUN1_ID)
+            .values(created_dag_version_id=None, 
bundle_version="deleted-version")
+        )
+        session.commit()
+
+        response = test_client.post(
+            f"/dags/{DAG1_ID}/dagRuns/{DAG1_RUN1_ID}/clear",
+            json={"dry_run": False},
+        )
+        assert response.status_code == 200
+        assert response.json()["state"] == "queued"
+
     def test_should_respond_401(self, unauthenticated_test_client):
         response = unauthenticated_test_client.post(
             f"/dags/{DAG1_ID}/dagRuns/{DAG1_RUN1_ID}/clear",
diff --git a/airflow-core/tests/unit/models/test_cleartasks.py 
b/airflow-core/tests/unit/models/test_cleartasks.py
index 9972191e1eb..407b51bbec2 100644
--- a/airflow-core/tests/unit/models/test_cleartasks.py
+++ b/airflow-core/tests/unit/models/test_cleartasks.py
@@ -21,7 +21,7 @@ import datetime
 import random
 
 import pytest
-from sqlalchemy import func, select
+from sqlalchemy import func, select, update
 
 from airflow.models.dag_version import DagVersion
 from airflow.models.dagrun import DagRun
@@ -739,6 +739,183 @@ class TestClearTasks:
             for ti in dr.task_instances:
                 assert ti.dag_version_id == old_dag_version.id
 
+    def test_clear_task_instances_without_dag_version_forces_latest(self, 
dag_maker, session):
+        """A Dag run carried over from Airflow 2 has no version, so clearing 
must pin it to the latest."""
+        dag_id = "test_clear_no_dag_version"
+        dr = self._make_versionless_run(dag_maker, session, dag_id, 
DagRunState.SUCCESS)
+
+        latest_dag_version = DagVersion.get_latest_version(dr.dag_id)
+        ti0 = session.scalar(select(TI).where(TI.dag_id == dag_id))
+        assert ti0.dag_version_id is None, "Pre-condition"
+        assert ti0.dag_run.created_dag_version_id is None, "Pre-condition"
+
+        clear_task_instances([ti0], session, run_on_latest_version=False)
+        session.commit()
+
+        dr_after = session.scalar(select(DagRun).where(DagRun.dag_id == 
dag_id))
+        assert dr_after.created_dag_version_id == latest_dag_version.id
+        assert dr_after.bundle_version == latest_dag_version.bundle_version
+        assert dr_after.task_instances[0].dag_version_id == 
latest_dag_version.id
+
+    def _make_versionless_run(self, dag_maker, session, dag_id, dr_state, 
task_count=1, sibling_state=None):
+        """
+        Build a run shaped like Airflow 2 left it: no versions anywhere.
+
+        Task "0" is run for real; any further tasks are left in 
``sibling_state``.
+        """
+        with dag_maker(dag_id, start_date=DEFAULT_DATE, catchup=True, 
bundle_version="v1") as dag:
+            task0 = EmptyOperator(task_id="0")
+            for index in range(1, task_count):
+                EmptyOperator(task_id=str(index))
+        dr = dag_maker.create_dagrun(state=State.RUNNING, 
run_type=DagRunType.SCHEDULED)
+
+        ti0, *siblings = sorted(dr.task_instances, key=lambda ti: ti.task_id)
+        ti0.refresh_from_task(dag.get_task("0"))
+        run_task_instance(ti0, task0)
+        for sibling in siblings:
+            sibling.state = sibling_state
+        dr.state = dr_state
+
+        # `airflow db migrate` from Airflow 2 leaves these columns NULL. Write 
them directly so
+        # no ORM relationship syncs the old values back, then expire so the 
objects are reloaded
+        # from the database like they are in a real deployment.
+        session.flush()
+        session.execute(
+            update(DagRun).where(DagRun.id == 
dr.id).values(created_dag_version_id=None, bundle_version=None)
+        )
+        session.execute(update(TI).where(TI.dag_id == 
dag.dag_id).values(dag_version_id=None))
+        session.commit()
+        session.expire_all()
+        return dr
+
+    def 
test_clear_task_instances_pins_task_instance_restored_by_verify_integrity(self, 
dag_maker, session):
+        """
+        A task instance revived by ``verify_integrity`` is given a version too.
+
+        It comes back unfinished but unversioned, and pinning the run stops 
the scheduler
+        backfilling one, so it would never be enqueued.
+        """
+        dag_id = "test_clear_no_dag_version_restored"
+        # Task "1" was dropped from the Dag during the Airflow 2 era and later 
re-added, so
+        # verify_integrity restores it when the finished run is cleared.
+        dr = self._make_versionless_run(
+            dag_maker,
+            session,
+            dag_id,
+            DagRunState.SUCCESS,
+            task_count=2,
+            sibling_state=TaskInstanceState.REMOVED,
+        )
+        latest_dag_version = DagVersion.get_latest_version(dr.dag_id)
+        ti0 = session.scalar(select(TI).where(TI.dag_id == dag_id, TI.task_id 
== "0"))
+
+        clear_task_instances([ti0], session, run_on_latest_version=False)
+        session.commit()
+
+        restored = session.scalar(select(TI).where(TI.dag_id == dag_id, 
TI.task_id == "1"))
+        assert restored.state is None, "verify_integrity should have restored 
it"
+        assert restored.dag_version_id == latest_dag_version.id
+
+    def 
test_clear_task_instances_pins_unfinished_siblings_on_running_run(self, 
dag_maker, session):
+        """A queued/running run is pinned without verify_integrity, so its 
siblings need one too."""
+        dag_id = "test_clear_no_dag_version_running"
+        dr = self._make_versionless_run(
+            dag_maker,
+            session,
+            dag_id,
+            DagRunState.RUNNING,
+            task_count=2,
+            sibling_state=TaskInstanceState.SCHEDULED,
+        )
+        latest_dag_version = DagVersion.get_latest_version(dr.dag_id)
+        ti0 = session.scalar(select(TI).where(TI.dag_id == dag_id, TI.task_id 
== "0"))
+
+        clear_task_instances([ti0], session, run_on_latest_version=False)
+        session.commit()
+
+        dr_after = session.scalar(select(DagRun).where(DagRun.dag_id == 
dag_id))
+        assert dr_after.created_dag_version_id == latest_dag_version.id
+        sibling = session.scalar(select(TI).where(TI.dag_id == dag_id, 
TI.task_id == "1"))
+        assert sibling.dag_version_id == latest_dag_version.id
+
+    def test_clear_task_instances_keeps_run_and_task_versions_together(self, 
dag_maker, session):
+        """A run pinned to a version must not leave its cleared task instances 
on another."""
+        dag_id = "test_clear_backfilled_ti_null_run"
+        with dag_maker(dag_id, start_date=DEFAULT_DATE, catchup=True, 
bundle_version="v1") as dag:
+            task0 = EmptyOperator(task_id="0")
+        dr = dag_maker.create_dagrun(state=State.RUNNING, 
run_type=DagRunType.SCHEDULED)
+        (ti0,) = dr.task_instances
+        ti0.refresh_from_task(dag.get_task("0"))
+        run_task_instance(ti0, task0)
+        dr.state = DagRunState.SUCCESS
+        session.flush()
+
+        # The task instance keeps a version while the run loses its own, so 
the run counts as
+        # version-less and gets forced onto the latest.
+        old_dag_version = DagVersion.get_latest_version(dag_id)
+        session.execute(update(DagRun).where(DagRun.id == 
dr.id).values(created_dag_version_id=None))
+        session.commit()
+        session.expire_all()
+
+        with dag_maker(dag_id, start_date=DEFAULT_DATE, catchup=True, 
bundle_version="v2"):
+            EmptyOperator(task_id="0")
+        new_dag_version = DagVersion.get_latest_version(dag_id)
+        assert old_dag_version.id != new_dag_version.id, "Pre-condition"
+
+        ti0 = session.scalar(select(TI).where(TI.dag_id == dag_id))
+        assert ti0.dag_version_id == old_dag_version.id, "Pre-condition"
+        assert ti0.dag_run.created_dag_version_id is None, "Pre-condition"
+
+        clear_task_instances([ti0], session, run_on_latest_version=False)
+        session.commit()
+
+        dr_after = session.scalar(select(DagRun).where(DagRun.dag_id == 
dag_id))
+        ti_after = session.scalar(select(TI).where(TI.dag_id == dag_id))
+        assert dr_after.created_dag_version_id == new_dag_version.id
+        assert ti_after.dag_version_id == dr_after.created_dag_version_id, (
+            "the run and its task instance must end up on the same version"
+        )
+
+    def 
test_clear_task_instances_moves_versionless_task_to_its_run_version(self, 
dag_maker, session):
+        """A version-less task instance on a pinned run joins the run, not the 
latest version."""
+        dag_id = "test_clear_versionless_ti_pinned_run"
+        with dag_maker(dag_id, start_date=DEFAULT_DATE, catchup=True, 
bundle_version="v1") as dag:
+            task0 = EmptyOperator(task_id="0")
+            EmptyOperator(task_id="1")
+        dr = dag_maker.create_dagrun(state=State.RUNNING, 
run_type=DagRunType.SCHEDULED)
+        ti0, ti1 = sorted(dr.task_instances, key=lambda ti: ti.task_id)
+        ti0.refresh_from_task(dag.get_task("0"))
+        run_task_instance(ti0, task0)
+        ti1.state = TaskInstanceState.SUCCESS
+        dr.state = DagRunState.SUCCESS
+        session.flush()
+
+        # An Airflow 2 task instance that an earlier clear left behind: it was 
already finished, so
+        # pinning the run did not give it a version.
+        run_dag_version = DagVersion.get_latest_version(dag_id)
+        session.execute(update(TI).where(TI.dag_id == dag_id, TI.task_id == 
"1").values(dag_version_id=None))
+        session.commit()
+        session.expire_all()
+
+        with dag_maker(dag_id, start_date=DEFAULT_DATE, catchup=True, 
bundle_version="v2"):
+            EmptyOperator(task_id="0")
+            EmptyOperator(task_id="1")
+        assert DagVersion.get_latest_version(dag_id).id != run_dag_version.id, 
"Pre-condition"
+
+        ti1 = session.scalar(select(TI).where(TI.dag_id == dag_id, TI.task_id 
== "1"))
+        assert ti1.dag_version_id is None, "Pre-condition"
+        assert ti1.dag_run.created_dag_version_id == run_dag_version.id, 
"Pre-condition"
+
+        clear_task_instances([ti1], session, run_on_latest_version=False)
+        session.commit()
+
+        dr_after = session.scalar(select(DagRun).where(DagRun.dag_id == 
dag_id))
+        ti1_after = session.scalar(select(TI).where(TI.dag_id == dag_id, 
TI.task_id == "1"))
+        assert dr_after.created_dag_version_id == run_dag_version.id
+        assert ti1_after.dag_version_id == run_dag_version.id, (
+            "the run and its task instance must end up on the same version"
+        )
+
     def test_clear_subset_run_on_latest_version_only_updates_cleared_tis(self, 
dag_maker, session):
         """run_on_latest_version on a finished DR must not rewrite 
dag_version_id on TIs that were not cleared."""
         with dag_maker(

Reply via email to