amoghrajesh commented on code in PR #73030:
URL: https://github.com/apache/airflow/pull/73030#discussion_r4093310934
##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py:
##########
@@ -89,6 +89,7 @@ class TaskInstanceResponse(BaseModel):
queued_by_job: JobResponse | None = Field(alias="triggerer_job")
dag_version: DagVersionResponse | None
team_name: str | None = None
+ state_reason: str | None = Field(default=None,
validation_alias="retry_reason")
Review Comment:
[after rebase comments on PR
2](https://github.com/apache/airflow/pull/73030/commits/74e7573a5606e4fa7bf83616ab1aef5106a1bb62)
##########
task-sdk/tests/task_sdk/execution_time/test_task_runner.py:
##########
@@ -1196,6 +1196,95 @@ 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",
+ retries=2,
+ retry_policy=ExceptionRetryPolicy(
+ rules=[RetryRule(exception=RuntimeError, action=RetryAction.FAIL,
reason="do not retry")]
+ ),
+ )
+ ti = create_runtime_ti(task=task)
+
+ 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",
+ retries=2,
+ retry_policy=ExceptionRetryPolicy(
+ rules=[RetryRule(exception=RuntimeError, action=RetryAction.RETRY,
reason="rate limit")]
+ ),
+ )
+ ti = create_runtime_ti(task=task, try_number=3)
+
+ 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 (3 of 3)"
Review Comment:
[after rebase comments on PR
2](https://github.com/apache/airflow/pull/73030/commits/74e7573a5606e4fa7bf83616ab1aef5106a1bb62)
--
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]