amoghrajesh commented on code in PR #73030:
URL: https://github.com/apache/airflow/pull/73030#discussion_r4060007392
##########
task-sdk/src/airflow/sdk/execution_time/task_runner.py:
##########
@@ -1931,9 +1932,15 @@ def _finalize_task_failure(
if retry_reason is not None:
retry_kwargs["retry_reason"] = retry_reason[:500]
return RetryTask(**retry_kwargs), TaskInstanceState.UP_FOR_RETRY
+ if retry_reason is not None and ti._ti_context_from_server is not None:
+ max_tries = ti._ti_context_from_server.max_tries
+ retry_reason = f"{retry_reason}; retries exhausted ({ti.try_number} of
{max_tries})"
Review Comment:
Fixed on #73027: the denominator is now `max_tries + 1`, matching `Starting
attempt %s of %s` in core and the alert template in this same file. The
`max_tries == 0` case you flagged is handled too — the suffix is skipped
entirely rather than rendering "(1 of 0)", since there was no budget to exhaust.
Note this branch is still based on `main` and carries the pre-fix commits,
so the diff here shows the old form until I restack it onto #73027.
---
Drafted-by: Claude Opus 5 (1M context); reviewed by @amoghrajesh before
posting
##########
task-sdk/tests/task_sdk/execution_time/test_task_runner.py:
##########
@@ -1196,6 +1196,65 @@ def execute(self, context):
assert counted.count("operator_failures") == 1
+def test_retry_policy_fail_persists_reason(create_runtime_ti,
mock_supervisor_comms):
+ class _AlwaysFails(BaseOperator):
+ def execute(self, context):
+ raise RuntimeError("boom")
+
+ task = _AlwaysFails(
+ task_id="fail_with_reason",
+ retry_policy=ExceptionRetryPolicy(
+ rules=[RetryRule(exception=RuntimeError, action=RetryAction.FAIL,
reason="do not retry")]
+ ),
+ )
+ ti = create_runtime_ti(task=task, should_retry=True)
+
+ state, msg, error = run(ti, ti.get_template_context(), mock.MagicMock())
+
+ assert state == TaskInstanceState.FAILED
+ assert isinstance(msg, TaskState)
+ assert msg.retry_reason == "do not retry"
+
+
+def
test_retry_policy_retry_exhausted_persists_combined_reason(create_runtime_ti,
mock_supervisor_comms):
+ """A policy-chosen RETRY that hits an exhausted budget still fails, with
both reasons recorded."""
+
+ class _AlwaysFails(BaseOperator):
+ def execute(self, context):
+ raise RuntimeError("boom")
+
+ task = _AlwaysFails(
+ task_id="retry_exhausted",
+ retry_policy=ExceptionRetryPolicy(
+ rules=[RetryRule(exception=RuntimeError, action=RetryAction.RETRY,
reason="rate limit")]
+ ),
+ )
+ ti = create_runtime_ti(task=task, try_number=2, max_tries=2,
should_retry=False)
+
+ state, msg, error = run(ti, ti.get_template_context(), mock.MagicMock())
+
+ assert state == TaskInstanceState.FAILED
+ assert isinstance(msg, TaskState)
+ assert msg.retry_reason == "rate limit; retries exhausted (2 of 2)"
Review Comment:
Fixed on #73027, and taken a bit further than the suggestion: rather than
hardcoding `try_number=3, max_tries=2`, the operator now carries `retries=2`
and the fixture derives `max_tries` and `should_retry` itself. That keeps the
triple inside the reachable space by construction rather than by convention, so
it can't drift back out. Expected string is `(3 of 3)`.
Not visible in this diff yet — this branch still needs restacking onto
#73027.
---
Drafted-by: Claude Opus 5 (1M context); reviewed by @amoghrajesh before
posting
##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py:
##########
@@ -670,7 +670,11 @@ def _create_ti_state_update_query_and_update_state(
query = query.values(state=updated_state, next_method=None,
next_kwargs=None)
if updated_state == TaskInstanceState.FAILED:
- # This is the only case needs extra handling for
TITerminalStatePayload
+ if isinstance(ti_patch_payload, TITerminalStatePayload) and
ti_patch_payload.retry_reason:
+ failed_retry_reason: str | None =
ti_patch_payload.retry_reason[:500]
Review Comment:
Fixed on #73027. `_finalize_task_failure` now caps the base reason at `500 -
len(suffix)` before appending, so the suffix survives any reason length; the
`[:500]` here stays as a backstop for the plain FAIL path, which has no suffix.
Pinned by a test using a 600-character reason that asserts both `len == 500`
and `endswith("; retries exhausted (3 of 3)")` — I checked it fails on the old
form.
Not visible in this diff yet — this branch still needs restacking onto
#73027.
---
Drafted-by: Claude Opus 5 (1M context); reviewed by @amoghrajesh before
posting
##########
airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py:
##########
@@ -2446,6 +2446,72 @@ def test_ti_update_state_to_failed_table_check(self,
client, session, create_tas
assert ti.next_kwargs is None
assert ti.duration == 3600.00
+ def test_ti_update_state_to_failed_persists_retry_reason(self, client,
session, create_task_instance):
Review Comment:
Superseded by your own follow-up on #73027 — you measured that cadwyn
doesn't reach through the `TIStateUpdate` discriminated union, so the route
validates against the head models whatever version is pinned, and the boundary
test would fail today if written. Leaving it out on that basis. The
supervisor-side gate does bind and is covered by
`TestRealBundleRetryReasonUpgrade` in `test_migrator.py`.
---
Drafted-by: Claude Opus 5 (1M context); reviewed by @amoghrajesh before
posting
--
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]