diff --git a/airflow-core/docs/administration-and-deployment/dag-bundles.rst b/airflow-core/docs/administration-and-deployment/dag-bundles.rst index 354e6ecda1626..ef65b5eaa7135 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 c4bd8eceea102..5c0556974ff3b 100644 --- a/airflow-core/src/airflow/models/dagbag.py +++ b/airflow-core/src/airflow/models/dagbag.py @@ -209,7 +209,10 @@ def clear_cache(self) -> int: @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 7dae21baa2ccc..aca45da639116 100644 --- a/airflow-core/src/airflow/models/taskinstance.py +++ b/airflow-core/src/airflow/models/taskinstance.py @@ -333,6 +333,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, @@ -353,7 +374,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: """ @@ -377,7 +400,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) @@ -399,11 +425,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: @@ -435,10 +465,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: @@ -452,14 +486,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: @@ -473,6 +507,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 472c949ce284d..b8fbc5a0f966f 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 @@ -75,6 +75,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 c864f0b047c7f..8c471eaad9eb0 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 @@ -63,6 +63,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, }); @@ -92,17 +103,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 ( @@ -158,8 +158,14 @@ const ClearRunDialog = ({ dagRun, onClose, open }: Props) => { > {shouldShowBundleVersionOption ? ( setRunOnLatestVersion(Boolean(event.checked))} + title={ + runOnLatestVersionForced + ? translate("dags:runAndTaskActions.options.runOnLatestVersionForced") + : undefined + } > {translate("dags:runAndTaskActions.options.runOnLatestVersion")} 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 ab61c46e7ae28..a215486ebe076 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"; @@ -78,14 +82,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({ @@ -180,8 +190,14 @@ export const ClearGroupTaskInstanceDialog = ({ onClose, open, taskInstance }: Pr > {shouldShowRunOnLatestOption ? ( setRunOnLatestVersion(Boolean(event.checked))} + title={ + runOnLatestVersionForced + ? translate("dags:runAndTaskActions.options.runOnLatestVersionForced") + : undefined + } > {translate("dags:runAndTaskActions.options.runOnLatestVersion")} 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 0026da07395c3..57e8675704c2b 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"; @@ -104,13 +104,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({ @@ -258,8 +264,14 @@ const ClearTaskInstanceDialog = (props: Props) => { > {shouldShowRunOnLatestOption ? ( setRunOnLatestVersion(Boolean(event.checked))} + title={ + runOnLatestVersionForced + ? translate("dags:runAndTaskActions.options.runOnLatestVersionForced") + : undefined + } > {translate("dags:runAndTaskActions.options.runOnLatestVersion")} 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 1f3d6ba4d3c63..cc09d604e6548 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 fa183b37d0c06..21d7591e84c69 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 6c633fe8a2a6d..947fb3cce8bf6 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 @@ -2031,6 +2031,23 @@ def test_clear_dag_run(self, test_client, session): 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 9972191e1ebce..407b51bbec23e 100644 --- a/airflow-core/tests/unit/models/test_cleartasks.py +++ b/airflow-core/tests/unit/models/test_cleartasks.py @@ -21,7 +21,7 @@ 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 @@ def test_clear_task_instances_with_run_on_latest_version(self, run_on_latest_ver 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(