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]