amoghrajesh commented on code in PR #73027:
URL: https://github.com/apache/airflow/pull/73027#discussion_r4015149435
##########
task-sdk/tests/task_sdk/api/test_client.py:
##########
@@ -432,14 +432,19 @@ def handle_request(request: httpx.Request) ->
httpx.Response:
assert actual_body["end_date"] == "2024-10-31T12:00:00Z"
assert actual_body["state"] == state
assert actual_body["rendered_map_index"] == "test"
+ assert actual_body["retry_reason"] == "auth error, do not
retry"
Review Comment:
Handled in [comments from
kaxil](https://github.com/apache/airflow/pull/73027/commits/a6f747cd5c9a5723f264951d076ec6d26630f394)
##########
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]
+ query = query.values(retry_reason=failed_retry_reason)
+ if ti is not None:
+ ti.retry_reason = failed_retry_reason
Review Comment:
Handled in [comments from
kaxil](https://github.com/apache/airflow/pull/73027/commits/a6f747cd5c9a5723f264951d076ec6d26630f394)
##########
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:
Handled in [comments from
kaxil](https://github.com/apache/airflow/pull/73027/commits/a6f747cd5c9a5723f264951d076ec6d26630f394)
##########
task-sdk/src/airflow/sdk/execution_time/comms.py:
##########
@@ -869,6 +869,7 @@ class TaskState(BaseModel):
end_date: datetime | None = None
type: Literal["TaskState"] = "TaskState"
rendered_map_index: str | None = None
+ retry_reason: str | None = None
Review Comment:
Handled in [comments from
kaxil](https://github.com/apache/airflow/pull/73027/commits/a6f747cd5c9a5723f264951d076ec6d26630f394)
##########
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)
Review Comment:
Handled in [comments from
kaxil](https://github.com/apache/airflow/pull/73027/commits/a6f747cd5c9a5723f264951d076ec6d26630f394)
##########
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:
Handled in [comments from
kaxil](https://github.com/apache/airflow/pull/73027/commits/a6f747cd5c9a5723f264951d076ec6d26630f394)
--
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]