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 09ad0e53fd1a49917735500a0d260d39db4d3ac7 Author: Ash Berlin-Taylor <[email protected]> AuthorDate: Mon Oct 5 14:12:46 2026 +0100 fixup! Keep retired task attempts and their data under the attempt UUID --- airflow-core/src/airflow/models/taskinstance.py | 2 +- airflow-core/src/airflow/utils/db_cleanup.py | 6 +++++- airflow-core/tests/unit/models/test_taskinstance.py | 15 +++++++++++++++ airflow-core/tests/unit/utils/test_db_cleanup.py | 16 ++++++++++++++++ .../ai/tests/unit/common/ai/plugins/test_hitl_review.py | 1 + 5 files changed, 38 insertions(+), 2 deletions(-) diff --git a/airflow-core/src/airflow/models/taskinstance.py b/airflow-core/src/airflow/models/taskinstance.py index e6a0055282f..45d25f92545 100644 --- a/airflow-core/src/airflow/models/taskinstance.py +++ b/airflow-core/src/airflow/models/taskinstance.py @@ -2142,7 +2142,7 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload): fail_fast=fail_fast, ) - _log_state(task_instance=self) + _log_state(task_instance=ti) if not test_mode: TaskInstance.save_to_db(ti, session=session) diff --git a/airflow-core/src/airflow/utils/db_cleanup.py b/airflow-core/src/airflow/utils/db_cleanup.py index 86a7f5b22f0..92cebbd236e 100644 --- a/airflow-core/src/airflow/utils/db_cleanup.py +++ b/airflow-core/src/airflow/utils/db_cleanup.py @@ -982,7 +982,11 @@ def _get_archived_table_names(table_names: list[str] | None, session: Session) - inspector = inspect(session.bind) _, effective_config_dict = _effective_table_names(table_names=table_names) schemas = {config.schema_name for config in effective_config_dict.values()} - archive_sources = {config.bare_table_name for config in effective_config_dict.values()} + archive_sources = { + config.bare_table_name + for name, config in effective_config_dict.items() + if name != "task_instance_history" + } if "xcom_v2" in archive_sources: archive_sources.add("xcom_v1") if "xcom_v1" in archive_sources or "xcom_v2" in archive_sources: diff --git a/airflow-core/tests/unit/models/test_taskinstance.py b/airflow-core/tests/unit/models/test_taskinstance.py index 31fbaa6aea2..016fca7690c 100644 --- a/airflow-core/tests/unit/models/test_taskinstance.py +++ b/airflow-core/tests/unit/models/test_taskinstance.py @@ -2456,6 +2456,21 @@ class TestTaskInstance: ti.handle_failure("test queued ti", test_mode=True) assert ti.state == State.UP_FOR_RETRY + @mock.patch("airflow.models.taskinstance._log_state", autospec=True) + def test_handle_failure_logs_the_attempt_that_continues(self, mock_log_state, dag_maker, session): + with dag_maker(): + task = EmptyOperator(task_id="mytask", retries=1) + dr = dag_maker.create_dagrun() + ti = dr.get_task_instance(task.task_id, session=session) + ti.state = State.RUNNING + session.flush() + + retried = ti.handle_failure("test retry", session=session) + + assert retried.id != ti.id + mock_log_state.assert_called_once_with(task_instance=retried) + assert retried.state == State.UP_FOR_RETRY + @patch("airflow._shared.observability.metrics.stats._get_backend") def test_handle_failure_no_task(self, mock_get_backend, dag_maker): """ diff --git a/airflow-core/tests/unit/utils/test_db_cleanup.py b/airflow-core/tests/unit/utils/test_db_cleanup.py index 0cbf31d58a5..a8a4bfd91df 100644 --- a/airflow-core/tests/unit/utils/test_db_cleanup.py +++ b/airflow-core/tests/unit/utils/test_db_cleanup.py @@ -243,6 +243,22 @@ class TestDBCleanup: inspect_mock.return_value.get_table_names.return_value = [archive] assert _get_archived_table_names([requested], session) == [archive] + def test_history_alias_selects_only_history_archives(self): + history_archive = f"{ARCHIVE_TABLE_PREFIX}task_instance_history__20260929" + attempt_archives = [ + f"{ARCHIVE_TABLE_PREFIX}{name}__20261005" + for name in ("task_instance", "xcom_v2", "rtif_v2", "task_instance_note", "hitl_detail") + ] + + with create_session() as session: + with patch("airflow.utils.db_cleanup.inspect", autospec=True) as inspect_mock: + inspect_mock.return_value.get_table_names.return_value = [history_archive, *attempt_archives] + assert _get_archived_table_names(["task_instance_history"], session) == [history_archive] + assert set(_get_archived_table_names(["task_instance"], session)) == { + history_archive, + *attempt_archives, + } + def test_task_instance_history_alias_cleans_only_retired_attempts(self, ownership_session): session = ownership_session session.execute(sa.update(TaskInstance).values(start_date=NOW)) 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 0bd19ac6718..0a5d099a444 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 @@ -622,6 +622,7 @@ class TestWriteXcom: session, dag_id="d", run_id="r", task_id="t", map_index=-1, key=XCOM_AGENT_SESSION ) assert result == expected + _clear_db() @pytest.mark.skipif(not AIRFLOW_V_3_4_PLUS, reason="Attempt ownership starts in Airflow 3.4") def test_write_uses_current_attempt_after_retry(self, session, dag_maker):
