This is an automated email from the ASF dual-hosted git repository. ashb pushed a commit to branch store-historic-ti-ownership-data in repository https://gitbox.apache.org/repos/asf/airflow.git
commit 08bd03ad9d7ec85e8670943625083f62711b211b Author: Ash Berlin-Taylor <[email protected]> AuthorDate: Tue Oct 6 11:40:09 2026 +0100 fixup! Keep retired task attempts and their data under the attempt UUID --- airflow-core/docs/howto/usage-cli.rst | 2 +- .../api_fastapi/core_api/routes/ui/gantt.py | 2 +- .../execution_api/routes/task_instances.py | 6 +- .../execution_api/routes/task_state_store.py | 2 +- .../api_fastapi/execution_api/routes/xcoms.py | 8 +- .../airflow/api_fastapi/execution_api/security.py | 10 +-- .../api_fastapi/execution_api/versions/__init__.py | 4 +- .../execution_api/versions/v2026_10_30.py | 2 +- .../0142_3_4_0_unify_task_attempt_ownership.py | 4 - airflow-core/src/airflow/models/taskinstance.py | 26 +++---- airflow-core/src/airflow/models/xcom.py | 2 +- .../core_api/routes/public/test_hitl.py | 9 ++- .../core_api/routes/public/test_task_instances.py | 2 +- .../core_api/routes/ui/test_dashboard.py | 2 +- .../api_fastapi/execution_api/test_security.py | 52 ++++++------- .../versions/head/test_task_instances.py | 2 +- .../versions/head/test_task_state_store.py | 18 +++++ .../execution_api/versions/head/test_xcoms.py | 2 +- .../versions/v2026_10_30/test_task_instances.py | 6 +- .../test_0142_unify_task_attempt_ownership.py | 4 +- airflow-core/tests/unit/models/test_cleartasks.py | 8 +- airflow-core/tests/unit/models/test_dagrun.py | 6 +- .../tests/unit/models/test_mappedoperator.py | 4 +- .../tests/unit/models/test_taskinstance.py | 88 +++++++++++----------- airflow-core/tests/unit/models/test_trigger.py | 4 +- airflow-core/tests/unit/utils/test_db_cleanup.py | 2 +- .../tests_common/test_utils/attempt_ownership.py | 2 +- .../unit/common/ai/plugins/test_hitl_review.py | 10 +-- .../task_sdk/execution_time/test_supervisor.py | 2 +- 29 files changed, 153 insertions(+), 138 deletions(-) diff --git a/airflow-core/docs/howto/usage-cli.rst b/airflow-core/docs/howto/usage-cli.rst index fdb40c82836..a45b30b2ad0 100644 --- a/airflow-core/docs/howto/usage-cli.rst +++ b/airflow-core/docs/howto/usage-cli.rst @@ -233,7 +233,7 @@ You can optionally provide a list of tables to perform deletes on with ``--table ``xcom_v1`` and ``xcom_v2``, all of which the ``--dry-run`` output lists. ``--tables xcom`` selects both XCom stores. ``--tables task_instance_history`` selects only - retired attempts from ``task_instance``. Selecting ``xcom_v2`` also removes legacy values + older tries from ``task_instance``. Selecting ``xcom_v2`` also removes legacy values shadowed by the selected v2 rows, so those values cannot reappear after cleanup. That list is maintained per table rather than derived from the schema, so it does not cover every diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/gantt.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/gantt.py index 119fb208038..ba4778b4a9c 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/gantt.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/gantt.py @@ -60,7 +60,7 @@ def get_gantt_data( session: SessionDep, ) -> GanttResponse: """Get all task instance tries for Gantt chart.""" - # Pending retries retain timing for backoff; only the retired attempt belongs on the chart. + # Pending retries retain timing for backoff; only the archived attempt belongs on the chart. current_tis = select( TaskInstance.task_id.label("task_id"), TaskInstance.task_display_name.label("task_display_name"), # type: ignore[attr-defined] diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py index b05dd9e3bdc..1525be1c4bc 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py @@ -88,7 +88,7 @@ from airflow.api_fastapi.execution_api.services.task_instances import ( client_supports_arg_bindings, get_arg_bindings, ) -from airflow.api_fastapi.execution_api.versions.v2026_10_30 import IdentifyRetiredTaskStateUpdates +from airflow.api_fastapi.execution_api.versions.v2026_10_30 import IdentifyArchivedTaskStateUpdates from airflow.configuration import conf from airflow.exceptions import InvalidPartitionKeyError, TaskNotFound from airflow.models.asset import AssetActive @@ -472,7 +472,7 @@ def ti_update_state( _raise_ti_not_in_live_table(task_instance_id, archived_in_history=False) if working_set is None: _raise_ti_not_in_live_table( - task_instance_id, archived_in_history=IdentifyRetiredTaskStateUpdates.is_applied + task_instance_id, archived_in_history=IdentifyArchivedTaskStateUpdates.is_applied ) # TIStateUpdate can include terminal and intermediate states. This idempotency check handles @@ -1084,7 +1084,7 @@ async def ti_heartbeat( except NoResultFound: _raise_ti_not_in_live_table(task_instance_id, archived_in_history=False) if working_set is None: - # A retired attempt was likely cleared while running, so return 410 Gone + # An archived attempt was likely cleared while running, so return 410 Gone # instead of 404 Not Found to give the client a more specific signal. _raise_ti_not_in_live_table(task_instance_id, archived_in_history=True) diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_state_store.py b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_state_store.py index e1dfe7fcd97..beda54f23f0 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_state_store.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_state_store.py @@ -47,7 +47,7 @@ router = VersionedAPIRouter( def _get_task_scope_for_ti(task_instance_id: UUID, session: Session) -> TaskScope: ti = session.get(TI, task_instance_id) - if ti is None: + if ti is None or ti.working_set is not True: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={ diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py b/airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py index 542f95e4cb4..c9d58fc9a05 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py @@ -35,7 +35,7 @@ from airflow.api_fastapi.execution_api.datamodels.xcom import ( XComSequenceSliceResponse, ) from airflow.api_fastapi.execution_api.security import CurrentTIToken -from airflow.api_fastapi.execution_api.versions.v2026_10_30 import IdentifyRetiredTaskStateUpdates +from airflow.api_fastapi.execution_api.versions.v2026_10_30 import IdentifyArchivedTaskStateUpdates from airflow.models.taskinstance import TaskInstance from airflow.models.xcom import XCOM_RETURN_KEY, XComModel, xcom_entity from airflow.utils.db import get_query_count @@ -350,14 +350,14 @@ def get_xcom( and not params.include_prior_dates and params.offset is None ): - identify_retired = IdentifyRetiredTaskStateUpdates.is_applied + identify_archived = IdentifyArchivedTaskStateUpdates.is_applied raise HTTPException( - status_code=status.HTTP_410_GONE if identify_retired else status.HTTP_404_NOT_FOUND, + status_code=status.HTTP_410_GONE if identify_archived else status.HTTP_404_NOT_FOUND, detail={ "reason": "not_found", "message": ( "Task Instance not found in the working set; its attempt has been archived" - if identify_retired + if identify_archived else "Task Instance not found" ), }, diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/security.py b/airflow-core/src/airflow/api_fastapi/execution_api/security.py index 806c7fff4f5..6efd453c4f8 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/security.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/security.py @@ -43,7 +43,7 @@ Enforcement flow: - ``ti:self`` scope — checks that the JWT ``sub`` matches the ``{task_instance_id}`` path parameter. - Mutating task requests — checks that the attempt UUID still exists - in the live TI table. Already-admitted requests may finish after retirement. + in the live TI table. Already-admitted requests may finish after archival. 3. ``ExecutionAPIRoute`` precomputes ``allowed_token_types`` from ``token:*`` Security scopes at route registration time. Routes without explicit ``token:*`` scopes default to execution-only. @@ -231,9 +231,9 @@ async def require_auth( and not getattr(route, _SKIP_AUTO_TI_ATTEMPT_LIVE, False) ): # The versions package imports routes, which depend on this module. - from airflow.api_fastapi.execution_api.versions.v2026_10_30 import IdentifyRetiredTaskStateUpdates + from airflow.api_fastapi.execution_api.versions.v2026_10_30 import IdentifyArchivedTaskStateUpdates - if IdentifyRetiredTaskStateUpdates.is_applied: + if IdentifyArchivedTaskStateUpdates.is_applied: await _require_live_attempt(token, allow_callback="task_instance_id" not in request.path_params) request.scope[_REQUEST_SCOPE_LIVE_ATTEMPT_KEY] = True @@ -244,9 +244,9 @@ async def _require_live_attempt(token: TIToken, *, allow_callback: bool) -> None """ Reject mutations from an attempt whose UUID is no longer in the working set. - This is an admission check, not a lock: retirement may race with an already + This is an admission check, not a lock: archival may race with an already admitted request. Use a fresh session so a prior transaction's snapshot - cannot keep a retired UUID visible. Historical attempts never grant access. + cannot keep an archived UUID visible. Historical attempts never grant access. """ async with create_session_async() as session: attempt = ( diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py b/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py index 3f5050b19b1..44df7a51478 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py @@ -59,7 +59,7 @@ from airflow.api_fastapi.execution_api.versions.v2026_10_30 import ( AddMultiTeamToTIRunContext, AddStoppedTaskReport, AddTerminalStateRetryReasonField, - IdentifyRetiredTaskStateUpdates, + IdentifyArchivedTaskStateUpdates, ) bundle = VersionBundle( @@ -72,7 +72,7 @@ bundle = VersionBundle( AddTerminalStateRetryReasonField, AddMultiTeamToTIRunContext, AddStoppedTaskReport, - IdentifyRetiredTaskStateUpdates, + IdentifyArchivedTaskStateUpdates, ), Version( "2026-06-30", diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_10_30.py b/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_10_30.py index 8ca10d7c0bd..5739bbd787c 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_10_30.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_10_30.py @@ -45,7 +45,7 @@ class AddStoppedTaskReport(VersionChange): ) -class IdentifyRetiredTaskStateUpdates(VersionChangeWithSideEffects): +class IdentifyArchivedTaskStateUpdates(VersionChangeWithSideEffects): """Reject every mutation from an archived attempt with 410; older clients keep each endpoint's own response.""" description = __doc__ diff --git a/airflow-core/src/airflow/migrations/versions/0142_3_4_0_unify_task_attempt_ownership.py b/airflow-core/src/airflow/migrations/versions/0142_3_4_0_unify_task_attempt_ownership.py index 0b39848899e..1ded51930f8 100644 --- a/airflow-core/src/airflow/migrations/versions/0142_3_4_0_unify_task_attempt_ownership.py +++ b/airflow-core/src/airflow/migrations/versions/0142_3_4_0_unify_task_attempt_ownership.py @@ -256,9 +256,6 @@ def upgrade(): batch.drop_constraint("task_instance_composite_key", type_="unique") batch.create_unique_constraint("task_instance_current_key", [*_COORDINATES, "working_set"]) batch.create_unique_constraint("task_instance_try_key", [*_COORDINATES, "try_number"]) - batch.create_check_constraint( - "ti_working_set_true_or_null", "working_set IS NULL OR working_set = TRUE" - ) ti = sa.table( "task_instance", sa.column("id"), @@ -398,7 +395,6 @@ def downgrade(): with op.batch_alter_table("task_instance") as batch: batch.drop_constraint("task_instance_current_key", type_="unique") batch.drop_constraint("task_instance_try_key", type_="unique") - batch.drop_constraint("ti_working_set_true_or_null", type_="check") batch.create_unique_constraint("task_instance_composite_key", list(_COORDINATES)) for column in ("working_set", "archived_reason"): batch.drop_column(column) diff --git a/airflow-core/src/airflow/models/taskinstance.py b/airflow-core/src/airflow/models/taskinstance.py index 1aef66e04c3..6aa539209e9 100644 --- a/airflow-core/src/airflow/models/taskinstance.py +++ b/airflow-core/src/airflow/models/taskinstance.py @@ -37,7 +37,6 @@ from opentelemetry import trace from sqlalchemy import ( JSON, Boolean, - CheckConstraint, Float, ForeignKey, ForeignKeyConstraint, @@ -422,7 +421,7 @@ def clear_task_instances( for original in tis: ti = original if original in session else session.get(TaskInstance, original.id) if ti is None or ti.working_set is not True: - raise ValueError("A retired task instance cannot be cleared") + raise ValueError("An archived task instance cannot be cleared") if ti.state in (TaskInstanceState.RUNNING, TaskInstanceState.RESTARTING): if prevent_running_task: raise AirflowClearRunningTaskException( @@ -667,10 +666,10 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload): a TI with mapped tasks that expanded to an empty list (state=skipped). Every try of a task is its own row with its own UUID. Only the latest try is live (``working_set`` is - true); earlier tries are retired (``working_set`` is NULL) and kept as history. + true); earlier tries are archived (``working_set`` is NULL) and kept as history. ORM queries see only live rows by default: a session hook adds ``working_set IS TRUE`` to every ORM - select, update and delete, including joins to this model. To include retired rows, set the execution + select, update and delete, including joins to this model. To include archived rows, set the execution option on the statement:: session.scalars(select(TaskInstance).where(...).execution_options(include_all_attempts=True)) @@ -678,9 +677,9 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload): The default does not apply to: * primary key lookups (``Session.get``, ``merge``, ``refresh``), which return the row with that UUID - whether or not it is retired; + whether or not it is archived; * relationship loads, which follow the join the relationship defines (``DagRun.task_instances`` is live, - ``DagRun.historical_task_instances`` is retired); + ``DagRun.historical_task_instances`` is archived); * Core statements on ``TaskInstance.__table__``, and an ``exists()`` that does not name this model in a FROM clause. These see every row unless they filter ``working_set`` themselves. """ @@ -790,7 +789,6 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload): UniqueConstraint( "dag_id", "task_id", "run_id", "map_index", "try_number", name="task_instance_try_key" ), - CheckConstraint("working_set IS NULL OR working_set = TRUE", name="ti_working_set_true_or_null"), ForeignKeyConstraint( [trigger_id], ["trigger.id"], @@ -1152,13 +1150,13 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload): # is the task still in the retry waiting period? return self.state == TaskInstanceState.UP_FOR_RETRY and not self.ready_for_retry() - def retire(self, *, reason: str, session: Session) -> None: + def archive(self, *, reason: str, session: Session) -> None: """Remove this attempt from the working set while retaining its UUID and children.""" current = session.scalar( select(TaskInstance.working_set).where(TaskInstance.id == self.id).with_for_update() ) if current is not True: - raise ValueError("A retired task instance cannot be retired again") + raise ValueError("An archived task instance cannot be archived again") if self.state not in State.finished: self.state = TaskInstanceState.FAILED if self.end_date is None: @@ -1179,7 +1177,7 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload): map_index: int | None = None, session: Session, ) -> None: - """Delete every attempt, current and retired, of a task; all map indexes if ``map_index`` is None.""" + """Delete every attempt, current and archived, of a task; all map indexes if ``map_index`` is None.""" statement = ( delete(cls) .where(cls.dag_id == dag_id, cls.run_id == run_id, cls.task_id == task_id) @@ -1213,10 +1211,10 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload): return {map_index: last_try for map_index, last_try in session.execute(statement)} def prepare_db_for_next_try(self, session: Session) -> TaskInstance: - """Retire this UUID and return its successor in the caller's transaction.""" + """Archive this UUID and return its successor in the caller's transaction.""" successor_state = self.state dag_version_id = self.dag_version_id - self.retire(reason="retry", session=session) + self.archive(reason="retry", session=session) values = { attribute.columns[0].name: getattr(self, attribute.key) for attribute in inspect(TaskInstance).column_attrs @@ -2993,10 +2991,10 @@ def _is_primary_key_lookup(statement) -> bool: @sqlalchemy_event.listens_for(Session, "do_orm_execute") def _restrict_to_current_attempts(state: ORMExecuteState) -> None: """ - Hide retired attempts from ORM queries over task instances unless ``include_all_attempts`` is set. + Hide archived attempts from ORM queries over task instances unless ``include_all_attempts`` is set. Primary key lookups (``Session.get``, ``merge`` and ``refresh``) are exempt: asking for an attempt by - its UUID returns it whether or not it has been retired. + its UUID returns it whether or not it has been archived. """ if ( state.is_column_load diff --git a/airflow-core/src/airflow/models/xcom.py b/airflow-core/src/airflow/models/xcom.py index 71e683e9295..e64789cb8ce 100644 --- a/airflow-core/src/airflow/models/xcom.py +++ b/airflow-core/src/airflow/models/xcom.py @@ -340,7 +340,7 @@ class _XComOperations: specified DAG run are returned. If *True*, all matching XComs are returned regardless of the run it belongs to. :param limit: Limiting returning XComs - :param try_number: Read the XComs of this public try, current or retired, instead of + :param try_number: Read the XComs of this public try, current or archived, instead of the current attempt. """ if key is not None and not key: diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py index 2f8286c6609..a0dec74e2d9 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py @@ -329,7 +329,7 @@ def sample_update_payload() -> dict[str, Any]: class TestUpdateHITLDetailEndpoint: - def test_response_rejects_attempt_retired_after_lookup( + def test_response_rejects_attempt_archived_after_lookup( self, test_client, sample_ti, @@ -341,7 +341,7 @@ class TestUpdateHITLDetailEndpoint: ): get_task_instance = hitl_routes._get_task_instance_with_hitl_detail - def retire_after_lookup(*args, **kwargs): + def archive_after_lookup(*args, **kwargs): task_instance = get_task_instance(*args, **kwargs) with Session(bind=session.get_bind()) as other_session: current = other_session.get(TIModel, task_instance.id) @@ -350,7 +350,10 @@ class TestUpdateHITLDetailEndpoint: return task_instance mocker.patch.object( - hitl_routes, "_get_task_instance_with_hitl_detail", autospec=True, side_effect=retire_after_lookup + hitl_routes, + "_get_task_instance_with_hitl_detail", + autospec=True, + side_effect=archive_after_lookup, ) response = test_client.patch(f"{sample_ti_url_identifier}/hitlDetails", json=sample_update_payload) diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py index b0daa900181..7efd8fe3504 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py @@ -6433,7 +6433,7 @@ class TestBulkTaskInstances(TestTaskInstanceEndpoint): WILDCARD_ENDPOINT = "/dags/~/dagRuns/~/taskInstances" @pytest.mark.parametrize("delete_mode", ["single", "bulk-exact", "bulk-all"]) - def test_delete_removes_current_and_retired_tries(self, test_client, session, delete_mode): + def test_delete_removes_current_and_archived_tries(self, test_client, session, delete_mode): current_tis = self.create_task_instances( session, task_instances=[{"task_id": self.TASK_ID, "state": State.SUCCESS, "map_indexes": (0, 1, 2)}], diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dashboard.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dashboard.py index e7948a2049b..6747bbedce9 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dashboard.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_dashboard.py @@ -248,7 +248,7 @@ def make_multiple_dags(dag_maker, session): class TestHistoricalMetricsDataEndpoint: @pytest.mark.usefixtures("freeze_time_for_dagruns", "make_dag_runs") - def test_retired_task_instance_is_excluded_from_state_counts(self, test_client, session): + def test_archived_task_instance_is_excluded_from_state_counts(self, test_client, session): ti = session.scalar(select(TaskInstance).where(TaskInstance.state == TaskInstanceState.SUCCESS)) successor = ti.prepare_db_for_next_try(session) successor.state = TaskInstanceState.RUNNING diff --git a/airflow-core/tests/unit/api_fastapi/execution_api/test_security.py b/airflow-core/tests/unit/api_fastapi/execution_api/test_security.py index e5056e8bd6d..39b01143093 100644 --- a/airflow-core/tests/unit/api_fastapi/execution_api/test_security.py +++ b/airflow-core/tests/unit/api_fastapi/execution_api/test_security.py @@ -319,24 +319,24 @@ class TestAttemptLiveness: monkeypatch.setitem(exec_app.dependency_overrides, _jwt_bearer, authenticated_token) return ti, token - @pytest.mark.parametrize("retirement", ["retry", "delete"]) - def test_retired_attempt_cannot_replace_or_delete_xcom(self, client, caller, session, retirement): + @pytest.mark.parametrize("archival", ["retry", "delete"]) + def test_archived_attempt_cannot_replace_or_delete_xcom(self, client, caller, session, archival): ti, token = caller path = f"/execution/xcoms/{ti.dag_id}/{ti.run_id}/{ti.task_id}/result" assert client.post(path, json="original").status_code == 201 old_id = ti.id - if retirement == "retry": + if archival == "retry": successor = ti.prepare_db_for_next_try(session) successor.state = TaskInstanceState.UP_FOR_RETRY else: session.delete(ti) session.commit() - expected_status = 410 if retirement == "retry" else 404 + expected_status = 410 if archival == "retry" else 404 assert client.post(path, json="stale").status_code == expected_status assert client.delete(path).status_code == expected_status - if retirement == "retry": + if archival == "retry": assert ( session.scalar( select(TaskInstance.id) @@ -364,7 +364,7 @@ class TestAttemptLiveness: finally: Variable.delete(key) - RETIRED_ROUTES = [ + ARCHIVED_ROUTES = [ ("put", "/execution/store/ti/{ti_id}/key", {"value": "stale"}), ("delete", "/execution/store/ti/{ti_id}/key", None), ("delete", "/execution/store/ti/{ti_id}", None), @@ -382,9 +382,9 @@ class TestAttemptLiveness: @pytest.mark.parametrize( ("method", "path", "body"), - [pytest.param(*route, id=f"{route[0]}:{route[1]}") for route in RETIRED_ROUTES], + [pytest.param(*route, id=f"{route[0]}:{route[1]}") for route in ARCHIVED_ROUTES], ) - def test_retired_attempt_rejected_before_mutation(self, client, caller, session, method, path, body): + def test_archived_attempt_rejected_before_mutation(self, client, caller, session, method, path, body): ti, token = caller path = path.format(ti_id=token.id, dag_id=ti.dag_id, run_id=ti.run_id, task_id=ti.task_id) ti.prepare_db_for_next_try(session) @@ -448,26 +448,26 @@ class TestAttemptLiveness: finally: Variable.delete(key) - @pytest.mark.parametrize("delete_rejected", [False, True], ids=["accepted", "retired-token-rejected"]) + @pytest.mark.parametrize("delete_rejected", [False, True], ids=["accepted", "archived-token-rejected"]) def test_trusted_in_process_xcom_uses_current_owner_without_attempt_token( self, caller, create_task_instance, session, monkeypatch, delete_rejected ): - retired, _ = caller + archived, _ = caller XComModel.set_for_attempt( - task_instance_id=retired.id, key="key", value="retired", serialize=False, session=session + task_instance_id=archived.id, key="key", value="archived", serialize=False, session=session ) - current = retired.prepare_db_for_next_try(session) + current = archived.prepare_db_for_next_try(session) other = create_task_instance(dag_id="other_dag", task_id="other_task") session.commit() - path = f"/xcoms/{retired.dag_id}/{retired.run_id}/{retired.task_id}/key" + path = f"/xcoms/{archived.dag_id}/{archived.run_id}/{archived.task_id}/key" with TestClient(InProcessExecutionAPI().app) as client: other_attempt = {"X-Airflow-In-Process-Attempt-Id": str(other.id)} assert client.post(path, json="other", headers=other_attempt).status_code == 201 assert client.delete(path, headers=other_attempt).status_code == 200 - retired_header = {"X-Airflow-In-Process-Attempt-Id": str(retired.id)} - assert client.post(path, json="stale", headers=retired_header).status_code == 410 - assert client.delete(path, headers=retired_header).status_code == 410 + archived_header = {"X-Airflow-In-Process-Attempt-Id": str(archived.id)} + assert client.post(path, json="stale", headers=archived_header).status_code == 410 + assert client.delete(path, headers=archived_header).status_code == 410 assert client.post(path, json="current").status_code == 201 session.expire_all() assert XComModelV2.get_for_attempt(current.id, "key", session=session).value == "current" @@ -480,7 +480,7 @@ class TestAttemptLiveness: def delete(self, path, *, params): headers = ( - {"X-Airflow-In-Process-Attempt-Id": str(retired.id)} if delete_rejected else None + {"X-Airflow-In-Process-Attempt-Id": str(archived.id)} if delete_rejected else None ) response = client.delete(f"/{path}", params=params, headers=headers) if response.status_code != 200: @@ -501,10 +501,10 @@ class TestAttemptLiveness: with patch.object(BaseXCom, "purge", autospec=True) as purge: kwargs = { "key": "key", - "dag_id": retired.dag_id, - "run_id": retired.run_id, - "task_id": retired.task_id, - "map_index": retired.map_index, + "dag_id": archived.dag_id, + "run_id": archived.run_id, + "task_id": archived.task_id, + "map_index": archived.map_index, } if delete_rejected: with pytest.raises(RuntimeError, match="DELETE rejected: 410"): @@ -519,10 +519,10 @@ class TestAttemptLiveness: assert (current_row.value if current_row is not None else None) == ( "current" if delete_rejected else None ) - assert XComModelV2.get_for_attempt(retired.id, "key", session=session).value == "retired" + assert XComModelV2.get_for_attempt(archived.id, "key", session=session).value == "archived" @pytest.mark.parametrize("operation", ["xcom_write", "xcom_delete", "rtif_write", "variable_write"]) - def test_admitted_child_mutation_stays_with_retired_uuid( + def test_admitted_child_mutation_stays_with_archived_uuid( self, client, caller, session, monkeypatch, request, operation ): attempt, _ = caller @@ -621,12 +621,12 @@ class TestAttemptLiveness: assert RenderedTaskInstanceFields.get_for_attempt(old_id, session=session) is None def test_public_execution_api_ignores_in_process_attempt_header(self, client, caller, session): - retired, _ = caller - successor = retired.prepare_db_for_next_try(session) + archived, _ = caller + successor = archived.prepare_db_for_next_try(session) session.commit() response = client.post( - f"/execution/xcoms/{retired.dag_id}/{retired.run_id}/{retired.task_id}/value", + f"/execution/xcoms/{archived.dag_id}/{archived.run_id}/{archived.task_id}/value", headers={"X-Airflow-In-Process-Attempt-Id": str(successor.id)}, json="spoofed", ) diff --git a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py index 6c8039f84fd..62c69f8b47b 100644 --- a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py +++ b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py @@ -2566,7 +2566,7 @@ class TestTIUpdateState: def test_ti_update_state_retry_policy_overrides_persisted_in_history( self, client, session, create_task_instance ): - """The retired attempt keeps its retry policy and rendered-index values.""" + """The archived attempt keeps its retry policy and rendered-index values.""" ti = create_task_instance( task_id="test_retry_policy_override_history", state=State.RUNNING, diff --git a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_state_store.py b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_state_store.py index 671ceffbb4a..2bd8fd68485 100644 --- a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_state_store.py +++ b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_state_store.py @@ -34,6 +34,7 @@ from airflow.api_fastapi.execution_api.security import _jwt_bearer from airflow.models.dagrun import DagRun from airflow.models.task_state_store import TaskStateStoreModel from airflow.utils.session import create_session +from airflow.utils.state import TaskInstanceState if TYPE_CHECKING: from tests_common.pytest_plugin import CreateTaskInstance @@ -87,6 +88,23 @@ class TestGetTaskState: assert response.status_code == 200 assert response.json() == {"value": "spark_001"} + @pytest.mark.parametrize("method", ["get", "put", "delete"]) + def test_archived_attempt_returns_404( + self, client: TestClient, create_task_instance: CreateTaskInstance, session, method + ): + ti = create_task_instance(state=TaskInstanceState.RUNNING) + session.commit() + old_id = ti.id + ti.prepare_db_for_next_try(session) + session.commit() + client.headers["Airflow-API-Version"] = "2026-06-30" + kwargs = {"json": {"value": "stale"}} if method == "put" else {} + + response = getattr(client, method)(_api_url(old_id, "job_id"), **kwargs) + + assert response.status_code == 404 + assert "Task instance" in response.json()["detail"]["message"] + class TestPutTaskState: def test_put_creates_row(self, client: TestClient, create_task_instance: CreateTaskInstance): diff --git a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_xcoms.py b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_xcoms.py index c762d29f513..e1938eb7ec7 100644 --- a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_xcoms.py +++ b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_xcoms.py @@ -402,7 +402,7 @@ class TestXComsGetEndpoint: ("version", "expected_status"), [("2025-04-11", 404), ("2026-06-30", 404), ("2026-10-30", 410)], ) - def test_retired_attempt_read_response_by_version( + def test_archived_attempt_read_response_by_version( self, client, exec_app, monkeypatch, create_task_instance, session, version, expected_status ): attempt = create_task_instance(state=TaskInstanceState.RUNNING) diff --git a/airflow-core/tests/unit/api_fastapi/execution_api/versions/v2026_10_30/test_task_instances.py b/airflow-core/tests/unit/api_fastapi/execution_api/versions/v2026_10_30/test_task_instances.py index 1dd49b684ee..3da74e25056 100644 --- a/airflow-core/tests/unit/api_fastapi/execution_api/versions/v2026_10_30/test_task_instances.py +++ b/airflow-core/tests/unit/api_fastapi/execution_api/versions/v2026_10_30/test_task_instances.py @@ -270,7 +270,7 @@ class TestArgBindingsFieldBackwardCompat: @pytest.fixture -def retired_attempt(client, exec_app, monkeypatch, create_task_instance, session): +def archived_attempt(client, exec_app, monkeypatch, create_task_instance, session): """A running attempt, authenticated with a signed token, that has since been retried.""" ti = create_task_instance(state=State.RUNNING) session.commit() @@ -291,9 +291,9 @@ def retired_attempt(client, exec_app, monkeypatch, create_task_instance, session ], ) def test_mutations_after_retry_response_by_version( - client, retired_attempt, session, version, rtif_status, heartbeat_status, xcom_status, xcom_kept + client, archived_attempt, session, version, rtif_status, heartbeat_status, xcom_status, xcom_kept ): - ti, successor = retired_attempt + ti, successor = archived_attempt client.headers["Airflow-API-Version"] = version rtif = client.put(f"/execution/task-instances/{ti.id}/rtif", json={"field": "late"}) diff --git a/airflow-core/tests/unit/migrations/test_0142_unify_task_attempt_ownership.py b/airflow-core/tests/unit/migrations/test_0142_unify_task_attempt_ownership.py index 019376fb79a..4fa7aab1ce2 100644 --- a/airflow-core/tests/unit/migrations/test_0142_unify_task_attempt_ownership.py +++ b/airflow-core/tests/unit/migrations/test_0142_unify_task_attempt_ownership.py @@ -582,9 +582,9 @@ def test_upgrade_enforces_current_and_public_try_uniqueness(populated_predecesso pool_slots=1, try_number=3, ) - for values in ({}, {"working_set": False}, {"working_set": None, "try_number": 1}): + for values in ({}, {"working_set": None, "try_number": 1}): with ( - pytest.raises((IntegrityError, OperationalError), match="(?i)unique|check constraint|duplicate"), + pytest.raises((IntegrityError, OperationalError), match="(?i)unique|duplicate"), connection.begin_nested(), ): connection.execute(ti.insert().values(**(row | values), id=uuid4())) diff --git a/airflow-core/tests/unit/models/test_cleartasks.py b/airflow-core/tests/unit/models/test_cleartasks.py index c66e0c26579..5bbe5ad778b 100644 --- a/airflow-core/tests/unit/models/test_cleartasks.py +++ b/airflow-core/tests/unit/models/test_cleartasks.py @@ -80,7 +80,7 @@ class TestClearTasks: ), ], ) - def test_clear_detached_attempt_retires_its_uuid_before_inserting_successor( + def test_clear_detached_attempt_archives_its_uuid_before_inserting_successor( self, create_task_instance, session, mocker, state, dag_metadata_missing, archived_state ): attempt = create_task_instance(state=state, session=session) @@ -105,21 +105,21 @@ class TestClearTasks: assert successor.id != old_id assert (successor.try_number, successor.working_set, successor.state) == (2, True, None) - @pytest.mark.parametrize(("non_current", "expected_rows"), [("deleted", 0), ("retired", 2)]) + @pytest.mark.parametrize(("non_current", "expected_rows"), [("deleted", 0), ("archived", 2)]) def test_clear_rejects_non_current_attempt_without_allocating_successor( self, create_task_instance, session, non_current, expected_rows ): attempt = create_task_instance(state=TaskInstanceState.SUCCESS, session=session) attempt.try_number = 1 session.commit() - if non_current == "retired": + if non_current == "archived": attempt.prepare_db_for_next_try(session) else: session.expunge(attempt) session.delete(session.get(TaskInstance, attempt.id)) session.commit() - with pytest.raises(ValueError, match="retired task instance cannot be cleared"): + with pytest.raises(ValueError, match="archived task instance cannot be cleared"): clear_task_instances([attempt], session=session) assert ( diff --git a/airflow-core/tests/unit/models/test_dagrun.py b/airflow-core/tests/unit/models/test_dagrun.py index 926524c9cbe..093ba27d5f3 100644 --- a/airflow-core/tests/unit/models/test_dagrun.py +++ b/airflow-core/tests/unit/models/test_dagrun.py @@ -1366,14 +1366,14 @@ class TestDagRun: dm = session.scalar(select(DagModel).options(joinedload(DagModel.dag_versions))) assert dag_run.dag_versions[0].id == dm.dag_versions[0].id - def test_retired_task_instance_version_relationship_and_indexed_lookup(self, dag_maker, session): - with dag_maker("retired_version_lookup", session=session): + def test_archived_task_instance_version_relationship_and_indexed_lookup(self, dag_maker, session): + with dag_maker("archived_version_lookup", session=session): EmptyOperator(task_id="task") dag_run = dag_maker.create_dagrun() ti = session.merge(dag_run.get_task_instance("task")) version_id = ti.dag_version_id ti.state = TaskInstanceState.SUCCESS - ti.retire(reason="retry", session=session) + ti.archive(reason="retry", session=session) session.expire(ti, ["dag_version"]) assert ti.dag_version.id == version_id diff --git a/airflow-core/tests/unit/models/test_mappedoperator.py b/airflow-core/tests/unit/models/test_mappedoperator.py index 7834cf39fee..6b48d92d854 100644 --- a/airflow-core/tests/unit/models/test_mappedoperator.py +++ b/airflow-core/tests/unit/models/test_mappedoperator.py @@ -261,7 +261,7 @@ def test_missing_mapped_index_uses_retained_max_try(dag_maker, session): expand_mapped_task_instances(dag.task_dict[mapped.task_id], run.run_id, session=session) removed = run.get_task_instance(mapped.task_id, map_index=1, session=session) removed.try_number = 3 - removed.retire(reason="retry", session=session) + removed.archive(reason="retry", session=session) new_tis = run._revise_map_indexes_if_mapped( dag.task_dict[mapped.task_id], dag_version_id=run.created_dag_version_id, session=session @@ -1877,7 +1877,7 @@ def test_placeholder_promotion_keeps_legacy_owner_and_avoids_historical_try_coll historical.try_number = 3 session.add(historical) session.flush() - historical.retire(reason="retry", session=session) + historical.archive(reason="retry", session=session) session.add( LegacyTaskDataOwner( dag_id=run.dag_id, diff --git a/airflow-core/tests/unit/models/test_taskinstance.py b/airflow-core/tests/unit/models/test_taskinstance.py index 80083b4d80d..0820c2f7881 100644 --- a/airflow-core/tests/unit/models/test_taskinstance.py +++ b/airflow-core/tests/unit/models/test_taskinstance.py @@ -857,15 +857,15 @@ class TestTaskInstance: # Clear the task instance. dag.clear() - retired = ti - ti = retired.dag_run.get_task_instance( - retired.task_id, map_index=retired.map_index, session=dag_maker.session + archived = ti + ti = archived.dag_run.get_task_instance( + archived.task_id, map_index=archived.map_index, session=dag_maker.session ) - assert ti.id != retired.id + assert ti.id != archived.id assert ti.state == State.NONE assert ti.try_number == 1 - # The reschedules stay with the retired attempt and none carry over to the new one. - assert len(task_reschedules_for_ti(retired)) == 1 + # The reschedules stay with the archived attempt and none carry over to the new one. + assert len(task_reschedules_for_ti(archived)) == 1 assert not task_reschedules_for_ti(ti) @pytest.mark.usefixtures("test_pool") @@ -925,15 +925,15 @@ class TestTaskInstance: # Clear the task instance. dag.clear() - retired = ti - ti = retired.dag_run.get_task_instance( - retired.task_id, map_index=retired.map_index, session=dag_maker.session + archived = ti + ti = archived.dag_run.get_task_instance( + archived.task_id, map_index=archived.map_index, session=dag_maker.session ) - assert ti.id != retired.id + assert ti.id != archived.id assert ti.state == State.NONE assert ti.try_number == 1 - # The reschedules stay with the retired attempt and none carry over to the new one. - assert len(task_reschedules_for_ti(retired)) == 1 + # The reschedules stay with the archived attempt and none carry over to the new one. + assert len(task_reschedules_for_ti(archived)) == 1 assert not task_reschedules_for_ti(ti) def test_depends_on_past_catchup_true(self, dag_maker): @@ -2798,7 +2798,7 @@ class TestTaskInstance: session.flush() with time_machine.travel(archive_time, tick=False): - ti.retire(reason="retry", session=session) + ti.archive(reason="retry", session=session) session.flush() tih = session.scalars( @@ -2891,12 +2891,12 @@ class TestTaskInstance: [ pytest.param(None, {-1: 2, 0: 5, 1: 6}, id="all-indexes"), pytest.param((0,), {0: 5}, id="selected-current-index"), - pytest.param((1,), {1: 6}, id="selected-retired-index"), + pytest.param((1,), {1: 6}, id="selected-archived-index"), pytest.param((), {}, id="empty-indexes"), pytest.param((99,), {}, id="missing-index"), ], ) - def test_get_last_try_numbers_includes_current_and_retired_attempts( + def test_get_last_try_numbers_includes_current_and_archived_attempts( self, ownership_session, map_indexes, expected ): session = ownership_session @@ -2930,7 +2930,7 @@ class TestTaskInstance: assert result == expected @pytest.mark.execution_timeout(10) - def test_retirement_preserves_attempt_and_children_and_returns_successor(self, ownership_session): + def test_archival_preserves_attempt_and_children_and_returns_successor(self, ownership_session): session = ownership_session attempt = session.get(TaskInstance, CURRENT_ID) original_try = attempt.try_number @@ -2966,7 +2966,7 @@ class TestTaskInstance: XComModel.set_for_attempt(task_instance_id=attempt.id, key="late", value=1, session=session) assert XComModelV2.get_for_attempt(successor.id, "late", session=session) is None - with pytest.raises(ValueError, match="retired"): + with pytest.raises(ValueError, match="archived"): attempt.prepare_db_for_next_try(session) assert ( @@ -2979,7 +2979,7 @@ class TestTaskInstance: ) assert successor.working_set is True - def test_retirement_carries_the_note_to_the_successor(self, ownership_session): + def test_archival_carries_the_note_to_the_successor(self, ownership_session): session = ownership_session attempt = session.get(TaskInstance, CURRENT_ID) attempt.note = "needs a look" @@ -3012,7 +3012,7 @@ class TestTaskInstance: pytest.param(5, False, id="other-map-index"), ], ) - def test_delete_attempts_removes_current_and_retired_attempts_with_their_data( + def test_delete_attempts_removes_current_and_archived_attempts_with_their_data( self, ownership_session, map_index, deleted ): session = ownership_session @@ -3040,7 +3040,7 @@ class TestTaskInstance: assert set(session.scalars(select(XComModelV2.task_instance_id))) == expected assert set(session.scalars(select(RenderedTaskInstanceFields.task_instance_id))) == expected - def test_orm_statements_ignore_retired_attempts_unless_asked(self, ownership_session): + def test_orm_statements_ignore_archived_attempts_unless_asked(self, ownership_session): session = ownership_session session.expunge_all() include_all_attempts = {"include_all_attempts": True} @@ -3058,19 +3058,19 @@ class TestTaskInstance: .execution_options(**include_all_attempts) ) - def test_primary_key_lookups_see_retired_attempts(self, ownership_session): + def test_primary_key_lookups_see_archived_attempts(self, ownership_session): session = ownership_session session.expunge_all() - retired = session.get(TaskInstance, HISTORY_ID) + archived = session.get(TaskInstance, HISTORY_ID) - assert retired is not None - assert retired.working_set is None - retired.state = TaskInstanceState.SKIPPED - merged = session.merge(retired) - assert merged is retired + assert archived is not None + assert archived.working_set is None + archived.state = TaskInstanceState.SKIPPED + merged = session.merge(archived) + assert merged is archived - def test_bulk_update_ignores_retired_attempts(self, ownership_session): + def test_bulk_update_ignores_archived_attempts(self, ownership_session): session = ownership_session updated = session.execute(update(TaskInstance).values(pid=99)).rowcount @@ -3085,17 +3085,17 @@ class TestTaskInstance: assert pids[CURRENT_ID] == 99 assert pids[HISTORY_ID] != 99 - def test_retired_attempt_loads_through_relationships_and_refresh(self, ownership_session): + def test_archived_attempt_loads_through_relationships_and_refresh(self, ownership_session): session = ownership_session - retired = session.get(TaskInstance, HISTORY_ID, execution_options={"include_all_attempts": True}) - retired.note = "kept" + archived = session.get(TaskInstance, HISTORY_ID, execution_options={"include_all_attempts": True}) + archived.note = "kept" session.flush() session.expire_all() note = session.scalar(select(TaskInstanceNote).where(TaskInstanceNote.ti_id == HISTORY_ID)) assert note.task_instance.id == HISTORY_ID - retired.refresh_from_db(session=session) - assert retired.working_set is None + archived.refresh_from_db(session=session) + assert archived.working_set is None def test_filter_for_tis_selects_only_the_current_attempt(self, ownership_session): session = ownership_session @@ -3186,9 +3186,9 @@ class TestTaskInstance: @pytest.mark.parametrize( ("delete_method", "deleted_attempt"), [ - pytest.param("orm", "retired", id="retired-orm"), + pytest.param("orm", "archived", id="archived-orm"), pytest.param("orm", "current", id="current-orm"), - pytest.param("bulk", "retired", id="retired-bulk"), + pytest.param("bulk", "archived", id="archived-bulk"), pytest.param("bulk", "current", id="current-bulk"), pytest.param("dagrun", None, id="dagrun"), ], @@ -3197,9 +3197,9 @@ class TestTaskInstance: self, ownership_session, delete_method, deleted_attempt ): session = ownership_session - retired = session.get(TaskInstance, CURRENT_ID) - current = retired.prepare_db_for_next_try(session) - for attempt in (retired, current): + archived = session.get(TaskInstance, CURRENT_ID) + current = archived.prepare_db_for_next_try(session) + for attempt in (archived, current): XComModel.set_for_attempt(task_instance_id=attempt.id, key="deletion", value=1, session=session) RenderedTaskInstanceFields.set_for_attempt( task_instance_id=attempt.id, rendered_fields={"owner": str(attempt.id)}, session=session @@ -3207,7 +3207,7 @@ class TestTaskInstance: session.flush() if delete_method == "dagrun": - session.delete(retired.dag_run) + session.delete(archived.dag_run) session.flush() assert session.scalar(sa.select(sa.func.count()).select_from(TaskInstance)) == 0 @@ -3222,7 +3222,7 @@ class TestTaskInstance: assert session.scalar(sa.text(f"SELECT count(*) FROM {name}")) == 0 return - target, retained = (retired, current) if deleted_attempt == "retired" else (current, retired) + target, retained = (archived, current) if deleted_attempt == "archived" else (current, archived) if delete_method == "orm": session.delete(target) else: @@ -4156,8 +4156,8 @@ def test__refresh_from_db_should_not_increment_try_number(dag_maker, session): assert ti.try_number == 1 # stays 1 [email protected]("retired", [False, True], ids=["current", "historical"]) -def test_delete_dagversion_restricted_when_taskinstance_exists(dag_maker, session, retired): [email protected]("archived", [False, True], ids=["current", "historical"]) +def test_delete_dagversion_restricted_when_taskinstance_exists(dag_maker, session, archived): """ Ensure that deleting a DagVersion with existing TaskInstance references is restricted (ON DELETE RESTRICT). """ @@ -4171,9 +4171,9 @@ def test_delete_dagversion_restricted_when_taskinstance_exists(dag_maker, sessio ti = session.scalars(select(TaskInstance).where(TaskInstance.dag_version_id == version.id)).first() assert ti is not None - if retired: + if archived: ti.state = TaskInstanceState.SUCCESS - ti.retire(reason="retry", session=session) + ti.archive(reason="retry", session=session) assert ti.dag_version_id == version.id session.delete(version) diff --git a/airflow-core/tests/unit/models/test_trigger.py b/airflow-core/tests/unit/models/test_trigger.py index dd24708e108..a219e536f08 100644 --- a/airflow-core/tests/unit/models/test_trigger.py +++ b/airflow-core/tests/unit/models/test_trigger.py @@ -806,11 +806,11 @@ def test_queue_column_max_len_matches_ti_column_max_len() -> None: @pytest.mark.need_serialized_dag -def test_get_sorted_triggers_ignores_retired_task_instance(session, create_task_instance): +def test_get_sorted_triggers_ignores_archived_task_instance(session, create_task_instance): trigger = Trigger(classpath="airflow.triggers.testing.SuccessTrigger", kwargs={}) session.add(trigger) session.flush() - task_instance = create_task_instance(task_id="retired_trigger_owner") + task_instance = create_task_instance(task_id="archived_trigger_owner") task_instance.trigger_id = trigger.id task_instance.prepare_db_for_next_try(session) session.commit() diff --git a/airflow-core/tests/unit/utils/test_db_cleanup.py b/airflow-core/tests/unit/utils/test_db_cleanup.py index 713389d7b4c..eb650f895af 100644 --- a/airflow-core/tests/unit/utils/test_db_cleanup.py +++ b/airflow-core/tests/unit/utils/test_db_cleanup.py @@ -259,7 +259,7 @@ class TestDBCleanup: *attempt_archives, } - def test_task_instance_history_alias_cleans_only_retired_attempts(self, ownership_session): + def test_task_instance_history_alias_cleans_only_older_tries(self, ownership_session): session = ownership_session session.execute( sa.update(TaskInstance).values(start_date=NOW).execution_options(include_all_attempts=True) diff --git a/devel-common/src/tests_common/test_utils/attempt_ownership.py b/devel-common/src/tests_common/test_utils/attempt_ownership.py index 22e8e939e2f..add1dd9dac8 100644 --- a/devel-common/src/tests_common/test_utils/attempt_ownership.py +++ b/devel-common/src/tests_common/test_utils/attempt_ownership.py @@ -66,7 +66,7 @@ def _clear(session): @pytest.fixture def ownership_session(session): - """Hold the rows migration 0142 leaves behind for one retired and one current attempt.""" + """Hold the rows migration 0142 leaves behind for one archived and one current attempt.""" from airflow.models.dagrun import DagRun from airflow.models.hitl import HITLDetail from airflow.models.renderedtifields import LegacyRenderedTaskInstanceFields diff --git a/providers/common/ai/tests/unit/common/ai/plugins/test_hitl_review.py b/providers/common/ai/tests/unit/common/ai/plugins/test_hitl_review.py index 0a5d099a444..a96c2108abe 100644 --- a/providers/common/ai/tests/unit/common/ai/plugins/test_hitl_review.py +++ b/providers/common/ai/tests/unit/common/ai/plugins/test_hitl_review.py @@ -687,13 +687,13 @@ class TestIsTaskCompleted: @pytest.mark.skipif(not AIRFLOW_V_3_4_PLUS, reason="Attempt ownership starts in Airflow 3.4") @pytest.mark.parametrize( - ("old_state", "retired_state"), + ("old_state", "archived_state"), [ - pytest.param("success", "success", id="retired-success"), - pytest.param("running", "failed", id="retired-failed-retry"), + pytest.param("success", "success", id="archived-success"), + pytest.param("running", "failed", id="archived-failed-retry"), ], ) - def test_ignores_historical_attempt_state(self, session, dag_maker, old_state, retired_state): + def test_ignores_historical_attempt_state(self, session, dag_maker, old_state, archived_state): from airflow.utils.state import TaskInstanceState with dag_maker("d", schedule=None, start_date=logical_date, serialized=True): @@ -702,7 +702,7 @@ class TestIsTaskCompleted: old = run.get_task_instance("t", session=session) old.state = TaskInstanceState(old_state) current = old.prepare_db_for_next_try(session) - assert old.state == TaskInstanceState(retired_state) + assert old.state == TaskInstanceState(archived_state) current.state = TaskInstanceState.RUNNING session.flush() diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py index b13b064c9b9..2953bb49a92 100644 --- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py +++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py @@ -4250,7 +4250,7 @@ class TestHandleRequest: proc.client.task_instances.finish.assert_called_once() @pytest.mark.parametrize("status_code", [404, 409, 410]) - def test_server_termination_acknowledgement_after_retirement(self, watched_subprocess, status_code): + def test_server_termination_acknowledgement_after_archival(self, watched_subprocess, status_code): proc, _ = watched_subprocess proc._exit_code = -signal.SIGTERM proc._terminal_state = SERVER_TERMINATED
