kaxil commented on code in PR #73889:
URL: https://github.com/apache/airflow/pull/73889#discussion_r4201160324


##########
airflow-core/src/airflow/models/dagrun.py:
##########
@@ -2296,6 +2296,9 @@ def schedule_tis(
                         state=TaskInstanceState.SCHEDULED,
                         scheduled_dttm=timezone.utcnow(),
                         try_number=next_try_number,
+                        # Already archived with the finished try; the new one 
must not inherit it.
+                        retry_reason=None,

Review Comment:
   A task with `start_from_trigger=True` never reaches this UPDATE: 
`defer_task()` at line 2267 sends it straight to DEFERRED and doesn't touch 
either column. So a retry that starts in the triggerer still carries the 
previous try's values. If that trigger then fails with retries left, the 
task-end-event handler in `trigger.py` archives the stale reason under the new 
try number, and `next_retry_datetime()` times the following retry from the 
stale override, which is the case this PR is fixing. Clearing both in 
`defer_task()` next to where it sets the DEFERRED state would cover it, along 
with a `start_from_trigger` case in the `schedule_tis` test.



##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -2802,6 +2802,28 @@ def 
test_schedule_tis_preserves_allocated_attempt(dag_maker, session, state, try
     assert session.get(TI, ti_id).try_number == try_number
 
 
+def test_schedule_tis_clears_retry_policy_columns(dag_maker, session):
+    """The new try must not inherit the finished try's reason."""
+    with dag_maker(session=session) as dag:
+        BashOperator(task_id="task", bash_command="echo 1")
+    dr = dag_maker.create_dagrun(session=session)
+    ti = dr.get_task_instance("task", session=session)
+    ti.refresh_from_task(dag.get_task("task"))
+    ti.state = None

Review Comment:
   This only covers a None-state TI, but the case the PR is after is 
up_for_retry moving to scheduled (a None-state TI with these set mostly comes 
from a clear, which now resets them anyway). 
`test_schedule_tis_preserves_allocated_attempt` just above already parametrizes 
None / up_for_retry / up_for_reschedule with the same setup, so seeding the two 
columns there and asserting on them would cover all three without the copy.



##########
airflow-core/tests/unit/models/test_taskinstance.py:
##########
@@ -4342,6 +4342,25 @@ def 
test_clear_task_instances_resets_context_carrier(dag_maker, session):
     assert dag_run.context_carrier["traceparent"] != original_dr_traceparent
 
 
[email protected]_test
+def test_clear_task_instances_clears_retry_policy_columns(dag_maker, session):
+    """Neither column may survive the clear."""
+    with dag_maker("test_clear_retry_policy_columns"):
+        EmptyOperator(task_id="t1")
+    dag_run = dag_maker.create_dagrun()
+    ti = dag_run.get_task_instance("t1", session=session)
+    ti.state = TaskInstanceState.FAILED
+    ti.retry_reason = "auth error, do not retry"
+    ti.retry_delay_override = 300.0
+    session.flush()
+
+    clear_task_instances([ti], session)
+    session.flush()
+
+    assert ti.retry_reason is None

Review Comment:
   These assertions read back the same in-memory `ti` that 
`clear_task_instances` just changed, so they pass as long as the two 
assignments exist. Could you also check the `TaskInstanceHistory` row for the 
old try? It should still have the seeded reason and override, and that check 
would fail if the reset ever got moved above `prepare_db_for_next_try`.



##########
airflow-core/src/airflow/models/taskinstance.py:
##########
@@ -447,6 +447,10 @@ def clear_task_instances(
                 ti.max_tries = max(ti.max_tries, previous_try_number)
             ti.state = None
             ti.external_executor_id = None
+            # retry_delay_override is the functional one: 
next_retry_datetime() prefers it over
+            # the task's own retry_delay, so a stale value would retime the 
next retry.
+            ti.retry_reason = None

Review Comment:
   The `state_reason` description in `core_api/datamodels/task_instances.py` 
still says the reason "is cleared only when the task next starts running, so a 
task waiting to be retried or re-run can still carry the reason". After this 
change a cleared task never carries it, and it's gone as soon as the next try 
is scheduled. That text needs updating (and the generated OpenAPI, UI and 
airflowctl copies regenerating), along with the "Cleared on task start 
(ti_run)" comment on the columns at line 695 and the one at 
`execution_api/routes/task_instances.py:741`.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to