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

ashb 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 79de1d48039 Fix `refresh_from_task` not applied on retries (#65932)
79de1d48039 is described below

commit 79de1d4803941e7ba4a6361529cc7079bb36373e
Author: Jorge Rocamora <[email protected]>
AuthorDate: Sat Oct 3 11:19:13 2026 +0200

    Fix `refresh_from_task` not applied on retries (#65932)
    
    task_instance_mutation_hook and dynamic PriorityWeightStrategy
    subclasses are not re-applied on natural retries: refresh_from_task
    is missing from the UP_FOR_RETRY → SCHEDULED transition.
    
    This adds it to DagRun.schedule_tis against the new computed
    try_number. First attempts and UP_FOR_RESCHEDULE are unaffected,
    per-TI UPDATEs only fire when field values actually change.
---
 .../cluster-policies.rst                           | 11 ++--
 airflow-core/src/airflow/models/dagrun.py          |  4 ++
 airflow-core/tests/unit/models/test_dagrun.py      | 68 ++++++++++++++++++++++
 3 files changed, 79 insertions(+), 4 deletions(-)

diff --git 
a/airflow-core/docs/administration-and-deployment/cluster-policies.rst 
b/airflow-core/docs/administration-and-deployment/cluster-policies.rst
index 1a09f69916b..c08a6bebfa0 100644
--- a/airflow-core/docs/administration-and-deployment/cluster-policies.rst
+++ b/airflow-core/docs/administration-and-deployment/cluster-policies.rst
@@ -41,10 +41,13 @@ There are three main types of cluster policy:
   ``task_instance_mutation_hook`` applies not to a task but to the instance of 
a task that
   relates to a particular DagRun. It is executed scheduler-side while task 
instances are created or
   reconciled (not in the Dag file processor, and not on the worker). The 
policy is only applied to the
-  currently executed run (i.e. instance) of that task. The ``dag_run`` 
argument lets the policy route on
-  run configuration (``dag_run.conf``); it may be ``None`` in early 
task-instance construction, and a hook
-  that only declares ``task_instance`` keeps working unchanged. Note that 
``dag_run.conf`` is only populated
-  for manually triggered or API-triggered runs; scheduled runs carry an empty 
``conf``.
+  currently executed run (i.e. instance) of that task. The policy can run more 
than once for the same task
+  instance, for example when it is created or reconciled, so implementations 
should be idempotent. It also
+  runs when a retry is scheduled, with ``task_instance.try_number`` set to the 
attempt about to run. The
+  ``dag_run`` argument lets the policy route on run configuration 
(``dag_run.conf``); it may be ``None`` in
+  early task-instance construction, and a hook that only declares 
``task_instance`` keeps working
+  unchanged. Note that ``dag_run.conf`` is only populated for manually 
triggered or API-triggered runs;
+  scheduled runs carry an empty ``conf``.
 
 .. warning::
 
diff --git a/airflow-core/src/airflow/models/dagrun.py 
b/airflow-core/src/airflow/models/dagrun.py
index c2b165d49cc..fc26cbccd16 100644
--- a/airflow-core/src/airflow/models/dagrun.py
+++ b/airflow-core/src/airflow/models/dagrun.py
@@ -2257,6 +2257,10 @@ class DagRun(Base, LoggingMixin):
         debug_try_number_check = self.log.isEnabledFor(logging.DEBUG)
         expected_try_number_by_ti_id: dict[UUID, tuple[int, int, str | None]] 
= {}
         for ti in schedulable_tis:
+            if ti.state == TaskInstanceState.UP_FOR_RETRY:
+                if TYPE_CHECKING:
+                    assert ti.task
+                ti.refresh_from_task(ti.task, dag_run=self)
             if not ti.is_schedulable:
                 empty_ti_ids.append(ti.id)
             # The defer_task method will check "start_trigger_args" to see 
whether the operator
diff --git a/airflow-core/tests/unit/models/test_dagrun.py 
b/airflow-core/tests/unit/models/test_dagrun.py
index 07a3151181b..e897a59f35e 100644
--- a/airflow-core/tests/unit/models/test_dagrun.py
+++ b/airflow-core/tests/unit/models/test_dagrun.py
@@ -95,6 +95,7 @@ from tests_common.test_utils.mapping import 
expand_mapped_task, push_mapped_leng
 from tests_common.test_utils.mock_operators import MockOperator
 from tests_common.test_utils.taskinstance import create_task_instance, 
run_task_instance
 from unit.models import DEFAULT_DATE as _DEFAULT_DATE
+from unit.plugins.priority_weight_strategy import DecreasingPriorityStrategy, 
TestPriorityWeightStrategyPlugin
 
 if TYPE_CHECKING:
     from airflow.serialization.definitions.dag import SerializedDAG
@@ -2832,6 +2833,73 @@ def 
test_schedule_tis_up_for_reschedule_does_not_increment_try_number(dag_maker,
     assert refreshed_ti.try_number == 3
 
 
[email protected]_plugin_manager(plugins=[TestPriorityWeightStrategyPlugin])
+def test_schedule_tis_refreshes_task_instance_only_on_retry(dag_maker, 
session):
+    with dag_maker(session=session) as dag:
+        for task_id in ("first_attempt", "retry", "reschedule"):
+            BashOperator(task_id=task_id, bash_command="echo 1", 
weight_rule=DecreasingPriorityStrategy())
+
+    dr = dag_maker.create_dagrun(session=session)
+    tis = {ti.task_id: ti for ti in dr.get_task_instances(session=session)}
+    for task_id, state, try_number in (
+        ("first_attempt", None, 0),
+        ("retry", TaskInstanceState.UP_FOR_RETRY, 2),
+        ("reschedule", TaskInstanceState.UP_FOR_RESCHEDULE, 3),
+    ):
+        tis[task_id].refresh_from_task(dag.get_task(task_id))
+        tis[task_id].state = state
+        tis[task_id].try_number = try_number
+    session.commit()
+
+    hook_calls = []
+
+    def route_to_retry_queue(task_instance, dag_run=None):
+        hook_calls.append((task_instance.task_id, task_instance.try_number, 
dag_run))
+        task_instance.queue = "retry_queue"
+
+    with _registered_mutation_hook(route_to_retry_queue):
+        assert dr.schedule_tis(tis.values(), session=session) == 3
+    session.commit()
+
+    assert hook_calls == [("retry", 2, dr)]
+    session.expire_all()
+    retry_ti = dr.get_task_instance("retry", session=session)
+    assert retry_ti.state == TaskInstanceState.SCHEDULED
+    assert retry_ti.try_number == 2
+    assert retry_ti.queue == "retry_queue"
+    assert retry_ti.priority_weight == 2
+
+
+def test_schedule_tis_refreshes_a_retry_that_defers_from_trigger(dag_maker, 
session):
+    with dag_maker(session=session):
+        task = MockOperator(task_id="task")
+        task.start_from_trigger = True
+        task.start_trigger_args = StartTriggerArgs(
+            trigger_cls="airflow.triggers.testing.SuccessTrigger",
+            next_method="execute_complete",
+        )
+
+    dr = dag_maker.create_dagrun(session=session)
+    ti = dag_maker.create_ti("task", dag_run=dr)
+    ti.state = TaskInstanceState.UP_FOR_RETRY
+    ti.try_number = 2
+    session.commit()
+
+    hook_calls = []
+
+    def route_to_retry_queue(task_instance, dag_run=None):
+        hook_calls.append((task_instance.task_id, task_instance.try_number))
+        task_instance.queue = "retry_queue"
+
+    with _registered_mutation_hook(route_to_retry_queue):
+        dr.schedule_tis((ti,), session=session)
+    session.commit()
+
+    assert hook_calls == [("task", 2)]
+    session.expire_all()
+    assert (ti.state, ti.queue) == (TaskInstanceState.DEFERRED, "retry_queue")
+
+
 def test_schedule_tis_empty_operator_is_noop_if_ti_already_running(dag_maker, 
session):
     with dag_maker(session=session) as dag:
         EmptyOperator(task_id="empty_task")

Reply via email to