This is an automated email from the ASF dual-hosted git repository.
potiuk 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 3b29a4c3b3d Record TaskInstanceHistory when a task fails from a
non-running state (#67372)
3b29a4c3b3d is described below
commit 3b29a4c3b3d26142b80dad1c10481f449da5ecaa
Author: Pradeep Kalluri <[email protected]>
AuthorDate: Mon Sep 21 22:52:39 2026 +0100
Record TaskInstanceHistory when a task fails from a non-running state
(#67372)
When a task instance failed from a non-RUNNING state (QUEUED or SCHEDULED
after
being killed externally by the executor), handle_failure() set UP_FOR_RETRY
without
calling prepare_db_for_next_try(), so no TaskInstanceHistory row was
written and the
attempt lost its hostname and execution metadata. The UI then had no log
entry for
that try.
The condition now excludes only RESTARTING, whose history is already
recorded when a
running task is cleared.
closes: #65366
closes: #67238
---
airflow-core/newsfragments/67372.bugfix.rst | 1 +
airflow-core/src/airflow/models/taskinstance.py | 12 +++--
airflow-core/tests/unit/jobs/test_scheduler_job.py | 51 ++++++++++++++++++++++
3 files changed, 60 insertions(+), 4 deletions(-)
diff --git a/airflow-core/newsfragments/67372.bugfix.rst
b/airflow-core/newsfragments/67372.bugfix.rst
new file mode 100644
index 00000000000..3cd5c2f6368
--- /dev/null
+++ b/airflow-core/newsfragments/67372.bugfix.rst
@@ -0,0 +1 @@
+Fix TaskInstanceHistory missing host/metadata when task is retried from a
non-RUNNING state (e.g. killed externally or failed in queued/deferred status)
diff --git a/airflow-core/src/airflow/models/taskinstance.py
b/airflow-core/src/airflow/models/taskinstance.py
index 5c8559d11f4..d9c3f8cab95 100644
--- a/airflow-core/src/airflow/models/taskinstance.py
+++ b/airflow-core/src/airflow/models/taskinstance.py
@@ -1933,10 +1933,14 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload):
if task and fail_fast:
_stop_remaining_tasks(task_instance=ti, session=session)
else:
- if ti.state == TaskInstanceState.RUNNING:
- # If the task instance is in the running state, it means it
raised an exception and
- # about to retry so we record the task instance history. For
other states, the task
- # instance was cleared and already recorded in the task
instance history.
+ if ti.state != TaskInstanceState.RESTARTING:
+ # Record the current attempt and prepare the TI for its next
try.
+ # Covers every path eligible for retry reaching
handle_failure():
+ # - RUNNING: task raised an exception during execution (normal
failure)
+ # - QUEUED/SCHEDULED: executor killed the task externally
before
+ # it could start (e.g. pod OOMKilled in KubernetesExecutor)
+ # RESTARTING is excluded: the task was cleared via the UI/API
while running;
+ # prepare_db_for_next_try() was already called during that
clear operation.
ti.prepare_db_for_next_try(session)
ti.state = State.UP_FOR_RETRY
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index d3ed94e70e8..05d6a76c7cc 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -5427,6 +5427,57 @@ class TestSchedulerJob:
is not None
)
+ @pytest.mark.parametrize("ti_state", [TaskInstanceState.QUEUED,
TaskInstanceState.SCHEDULED])
+ def test_process_executor_events_queued_ti_retry_preserves_history(self,
ti_state, dag_maker, session):
+ """
+ Regression test for #65366 / #67238.
+
+ When an executor reports FAILED for a TI that is QUEUED or SCHEDULED
+ (killed externally before it could start), the scheduler calls
handle_failure()
+ which must call prepare_db_for_next_try() so that TaskInstanceHistory
is recorded
+ with the correct hostname and start_date.
+ """
+ dag_id = "test_queued_ti_retry_history"
+ task_id = "dummy"
+ hostname = "worker-node-42"
+
+ with dag_maker(dag_id=dag_id, fileloc="/test_path/"):
+ task = EmptyOperator(task_id=task_id, retries=2)
+
+ dr = dag_maker.create_dagrun()
+ ti = dr.get_task_instance(task.task_id, session=session)
+ ti.state = ti_state
+ ti.hostname = hostname
+ ti.start_date = DEFAULT_DATE
+ ti.try_number = 1
+ ti.max_tries = 2
+ session.merge(ti)
+ session.commit()
+
+ old_ti_id = ti.id
+
+ executor = MockExecutor(do_update=False)
+ executor.event_buffer[ti.key] = TaskInstanceState.FAILED, None
+
+ scheduler_job = Job()
+ self.job_runner = SchedulerJobRunner(job=scheduler_job,
executors=[executor])
+ self.job_runner._process_executor_events(executor=executor,
session=session)
+
+ session.expire_all()
+ ti.refresh_from_db(session=session)
+
+ assert ti.state == State.UP_FOR_RETRY
+ assert ti.id != old_ti_id, "prepare_db_for_next_try must assign a new
UUID"
+
+ from airflow.models.taskinstancehistory import TaskInstanceHistory
+
+ tih = session.scalar(
+
select(TaskInstanceHistory).where(TaskInstanceHistory.task_instance_id ==
old_ti_id)
+ )
+ assert tih is not None, "TaskInstanceHistory must be created for
non-RUNNING retry"
+ assert tih.hostname == hostname
+ assert tih.start_date == DEFAULT_DATE
+
def test_adopt_or_reset_orphaned_tasks_external_triggered_dag(self,
dag_maker, session):
dag_id = "test_reset_orphaned_tasks_external_triggered_dag"
with dag_maker(dag_id=dag_id, schedule="@daily"):