kaxil commented on code in PR #73030:
URL: https://github.com/apache/airflow/pull/73030#discussion_r4073306324
##########
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:
Not reopening the suffix design, you already settled that on #73027. This is
about branch state: `git merge-base --is-ancestor 524f00dce2 6d467f3daa`
returns non-zero, so this head last picked up #73027 at `cfd075de77` and does
not contain the commit the three "Handled in 524f00dce2" replies point at. That
leaves two of this PR's own tests asserting a string the parent has since
deleted: this line expects `"rate limit; retries exhausted (3 of 3)"` and line
1268 expects `.endswith("; retries exhausted (3 of 3)")`, while `524f00dce2`
replaced both with a parametrized `pytest.param("retry_exhausted", 2, 3,
id="budget-exhausted")` expecting the bare reason. Both go red on merge-down,
and the `[: 500 - len(suffix)]` arithmetic at `task_runner.py:1943` goes dead
with them.
`retry_policies.rst` here also still carries both sentences the parent
corrected, the unconditional `; retries exhausted (N of M)` promise at line 148
and the version note at line 22; and main has rewritten that file (+393/-69)
since this branch's merge-base, moving the sentence this PR corrects to
`main:420`, so the hunk will not apply as-is. Worth landing #73027 and rebasing
before this one gets another review pass.
---
What I need before approving:
1. Land #73027 and rebase. That turns the two red assertions above green,
brings in the suffix removal and the migrator downgrade test, and clears the
current conflict with main (`mergeable` is `false` right now).
2. Redact `state_reason` before serving, on this model and
`task_instance_history.py:65`, the way `connections.py:48` does. With a test
that a `mask_secret()`-registered value does not come back in the response.
3. Make the new frontend tests capable of failing: register the i18n bundle,
and make one try-count case non-degenerate. Detail in the `Details.test.tsx`
comment.
4. Have `STATES_WITH_REASON` bind both surfaces, or drop the redundant
clause and correct the comment at 49 to 51.
5. CI actually green. Two checks have run on this head, so nothing is
verified yet.
None of the rest blocks me. The wire field is still `retry_reason` on
`TITerminalStatePayload` and `TaskState` while the REST field is now
`state_reason`; aligning them is free while 2026-10-30 is unreleased and
`API_VERSION` already points at it, and after release it costs a version change
plus the ts and java mirrors, so it is worth settling now either way.
`test_ti_update_state_to_failed_without_retry_reason` passes with both new
route lines deleted because the fixture never sets the column, and seeding a
stale value first would pin the preserve-on-absent behaviour that differs from
the retry branch. `_finalize_task_failure` logs `Retry policy decision` with
`action="fail"` after `_evaluate_retry_policy` already logged the same event
name with `action="retry"`, and the parent keeps that line so the merge-down
will not resolve it. And `retry_reason` is accepted on skipped, removed and
upstream_failed and then dropped, which a docstring saying the field is
failed-only would close
.
Clearing the columns in `clear_task_instances` is yours to do separately,
and the execution-API boundary test stays closed.
##########
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:
This is the line that makes the reason public, and it goes out with no
masking. Two siblings in this same directory redact before serving:
`connections.py:48` has a `@field_validator("password", mode="after")` calling
`redact`, and `variables.py:41` does the same for `val`. Nothing masks it
upstream either; `redact()` in `task_runner.py` is applied to rendered template
fields only (1343/1360/1387), and the worker stores `decision.reason` verbatim.
The reachable paths are the documented ones.
`LLMRetryPolicy(redact_exception=False)` is a supported opt-out
(`policies/retry.py:277`), and with it the model reads the raw exception and
can echo a credential into `reasoning`, which becomes the stored reason at
`retry.py:415`/`:424`. A hand-written `RetryPolicy` returning
`reason=f"...{exc}"` has the same shape, and that is the extension point the
docs point people at. The same string in a task log would be masked, so this is
the one surface where it is not. It also is not just the details page:
`_private_ui.yaml` picks `state_reason` up on `TaskInstanceResponse`, so it
rides the list endpoints too.
A redact validator on both this and `task_instance_history.py:65` would
match what `connections.py` already does.
##########
airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.test.tsx:
##########
@@ -0,0 +1,170 @@
+/*!
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+import "@testing-library/jest-dom";
+import { render, screen } from "@testing-library/react";
+import { describe, expect, it, vi } from "vitest";
+
+import type { TaskInstanceHistoryResponse, TaskInstanceResponse } from
"openapi/requests/types.gen";
+
+import i18n from "src/i18n/config";
+import { Wrapper } from "src/utils/Wrapper";
+
+import { Details } from "./Details";
+
+// Sibling panels each fetch their own data and are unrelated to the
state-reason
+// banner and row under test.
+vi.mock("./BlockingDeps", () => ({ BlockingDeps: () => undefined }));
+vi.mock("./ExtraLinks", () => ({ ExtraLinks: () => undefined }));
+vi.mock("./TriggererInfo", () => ({ TriggererInfo: () => undefined }));
+vi.mock("src/components/DagVersionDetails", () => ({ DagVersionDetails: () =>
undefined }));
+vi.mock("src/components/TaskTrySelect", () => ({ TaskTrySelect: () =>
undefined }));
+vi.mock("src/components/TeamName", () => ({ TeamName: () => undefined }));
+vi.mock("src/hooks/useShowTeam", () => ({ useShowTeam: () => false }));
+
+const mockTaskInstance = vi.fn<() => TaskInstanceResponse | undefined>();
+const mockTryInstance = vi.fn<() => TaskInstanceHistoryResponse | undefined>();
+
+vi.mock("openapi/queries", async () => {
+ const actual = await vi.importActual("openapi/queries");
+
+ return {
+ ...actual,
+ useTaskInstanceServiceGetMappedTaskInstance: () => ({ data:
mockTaskInstance() }),
+ useTaskInstanceServiceGetTaskInstanceTryDetails: () => ({ data:
mockTryInstance() }),
+ };
+});
+
+vi.mock("src/utils", async () => {
+ const actual = await vi.importActual("src/utils");
+
+ return { ...actual, useAutoRefresh: () => false };
+});
+
+const buildTaskInstance = (overrides: Partial<TaskInstanceResponse>):
TaskInstanceResponse =>
+ ({
+ dag_id: "test_dag",
+ dag_run_id: "run_1",
+ dag_version: null,
+ duration: null,
+ end_date: null,
+ id: "ti-id",
+ map_index: -1,
+ max_tries: 2,
+ note: null,
+ operator_name: "PythonOperator",
+ rendered_map_index: null,
+ start_date: null,
+ state: "failed",
+ state_reason: null,
+ task_display_name: "test_task",
+ task_id: "test_task",
+ trigger: null,
+ triggerer_job: null,
+ try_number: 3,
+ ...overrides,
+ }) as unknown as TaskInstanceResponse;
+
+const renderDetails = (
+ taskInstance: TaskInstanceResponse,
+ tryInstance: Partial<TaskInstanceHistoryResponse> = {},
+) => {
+ mockTaskInstance.mockReturnValue(taskInstance);
+ mockTryInstance.mockReturnValue({
+ ...taskInstance,
+ ...tryInstance,
+ });
+
+ return render(<Details />, { wrapper: Wrapper });
+};
+
+describe("Details state reason", () => {
+ it("does not render the banner when there is no reason", () => {
+ renderDetails(buildTaskInstance({ state_reason: null }));
+
+ expect(screen.queryByTestId("state-reason-alert")).not.toBeInTheDocument();
+
expect(screen.queryByText(i18n.t("common:taskInstance.stateReason"))).not.toBeInTheDocument();
+ });
+
+ it.each([
+ { state: "failed", titleKey: "failed" },
+ { state: "up_for_retry", titleKey: "upForRetry" },
+ ] as const)("titles the banner for a $state task", ({ state, titleKey }) => {
+ renderDetails(buildTaskInstance({ max_tries: 2, state, state_reason: "auth
error", try_number: 3 }));
+
+ expect(screen.getByTestId("state-reason-alert")).toHaveTextContent(
+ i18n.t(`common:taskInstance.stateReasonSummary.${titleKey}`, {
totalTries: 3, tryNumber: 3 }),
Review Comment:
Neither of these title assertions can fail on a wrong count.
i18n is never initialized under vitest. `src/i18n/config.ts` only reaches
`.init()` inside `Promise.all([resolveI18nVersion(),
resolveExtraLanguages()]).then(...)`, and `resolveI18nVersion` goes through
`VersionService.getVersion()`, which does not resolve under the test server;
`testsSetup.ts` registers no resource bundle. So
`i18n.t("common:taskInstance.stateReasonSummary.failed", {...})` hands back the
key verbatim, and `toHaveTextContent` then compares that key against the same
key rendered by the component. Interpolation values never enter it.
`DagDeactivatedBanner.test.tsx:75-76` already has the fix in this repo:
`i18n.addResourceBundle("en", "common", commonLocale, true, true)` in a
`beforeEach`.
The fixture is degenerate as well. Line 108 sets `max_tries: 2, try_number:
3`, so `totalTries` (`max_tries + 1`) and `tryNumber` are both 3, and the
assertion cannot tell the two arguments apart. I ran both mutations from a
scratch copy: reverting `Details.tsx` to interpolate a bare `max_tries`, which
is the `(3 of 2)` off-by-one from round 1, leaves all 11 tests green, and so
does swapping the two arguments.
An `up_for_retry` case with `max_tries: 3, try_number: 2` expecting
"Retrying after try 2 of 4" makes the counts distinguishable, and matches the
state where `try_number != max_tries + 1` is the norm rather than the
exception. With the resource bundle registered as well, the off-by-one mutation
fails 3 tests instead of none.
##########
airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.tsx:
##########
@@ -115,6 +122,48 @@ export const Details = () => {
return translate("common:none", { defaultValue: "None" });
};
+ const stateReasonSummary = ((): StateReasonSummary | undefined => {
+ const reason = taskInstance?.state_reason;
+
+ if (
+ reason === null ||
+ reason === undefined ||
+ taskInstance === undefined ||
+ !STATES_WITH_REASON.includes(taskInstance.state ?? "")
Review Comment:
The shared constant does not actually bind the banner, so it does not
prevent the drift it was introduced for. This `includes` check is redundant:
the two inner branches at 139 and 147 already test `state === "failed"` and
`=== "up_for_retry"` as literals, with the `return undefined` at 154 catching
everything else.
I checked by mutation rather than by reading. Four runs from a scratch copy
of `Details.tsx`:
```
widen STATES_WITH_REASON 3 failed
force the row gate true 4 failed
banner error -> warning 1 failed
force this banner gate false 0 failed (11 passed)
```
So the constant only gates the row, and the comment at 49 to 51 saying both
surfaces key off it is not true today. Add a state to the list and the row
renders it while the banner falls through to `undefined`.
One `as const` record keyed by state holding `{status, titleKey}`, driving
both the banner and the row, would make a new state fail to compile until it
has a title. It would also fold the duplicated predicate at 158 to 165 into the
same source.
--
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]