kaxil commented on code in PR #73030:
URL: https://github.com/apache/airflow/pull/73030#discussion_r4097928747


##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py:
##########
@@ -89,6 +90,17 @@ 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:
   Not reopening the clearing question, you scoped that to its own PR and I 
agree it's wider than this change. This is about the new field's contract 
rather than the column's lifecycle: the UI now gates both surfaces on 
`failed`/`up_for_retry`, but the API field carries no such qualification and 
the OpenAPI entry is bare `title: State Reason`. `retry_reason` is cleared in 
exactly one place, 
[`ti_run`](https://github.com/apache/airflow/blob/74e7573a5606e4fa7bf83616ab1aef5106a1bb62/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py#L247-L256),
 so a TI sitting in `scheduled` or `queued` behind a paused DAG or a full pool 
returns a non-null `state_reason` describing an attempt that's over. 
`dag_run.py:420` serves this model for the dry-run clear endpoint, where that 
reads oddly.
   
   A `description=` on both `Field(...)`s saying the value may describe a 
previous attempt would fix the part this PR owns, and it's permanent public 
surface either way.



##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py:
##########
@@ -89,6 +90,17 @@ 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")
+
+    @field_validator("state_reason", mode="after")
+    @classmethod
+    def redact_state_reason(cls, v: str | None) -> str | None:
+        # A retry policy composes this from the exception text, and a policy 
may opt out of the
+        # worker-side redaction, so the same string that would be masked in a 
task log can reach
+        # here unmasked.
+        if v is None:
+            return None
+        return str(redact(v))

Review Comment:
   This runs in the wrong process to mask what the worker wrote. `redact(v)` 
with no `name` does only pattern substitution against secrets registered by 
`mask_secret()` **in this process** 
([`_redact`](https://github.com/apache/airflow/blob/74e7573a5606e4fa7bf83616ab1aef5106a1bb62/airflow-core/src/airflow/_shared/secrets_masker/secrets_masker.py#L415-L421)
 reaches only `if self.replacer:` for a bare `str`; there is no name-based 
fallback the way 
[`connections.py:48`](https://github.com/apache/airflow/blob/74e7573a5606e4fa7bf83616ab1aef5106a1bb62/airflow-core/src/airflow/api_fastapi/core_api/datamodels/connections.py#L47-L52)
 gets by passing `field_info.field_name`). The masker is a per-process `@cache` 
singleton and `add_mask` is its only writer; nothing on the TI-details request 
path seeds it, so on a fresh api-server process the pattern set holds only the 
metadata-DB password, Fernet key, broker URL and JWT secret. I measured it: 
`redact("auth: the token hunter2 expired")` comes ba
 ck unchanged, while a string carrying the `sql_alchemy_conn` password masks.
   
   Your standalone check masked because one process served both the task's 
connection fetch (`Connection`'s `@reconstructor` calls 
`mask_secret(self.password)`) and the UI read. That stops holding on an 
api-server restart between the failure and the page view, with more than one 
replica (which [config.yml 
recommends](https://github.com/apache/airflow/blob/74e7573a5606e4fa7bf83616ab1aef5106a1bb62/airflow-core/src/airflow/config_templates/config.yml#L1804-L1806)
 over raising `[api] workers`), or for a secret that never passed through the 
api-server at all, such as a Variable whose key isn't in 
`DEFAULT_SENSITIVE_FIELDS`.
   
   The precedent for this exact worker to API to UI flow is 400 lines up 
`task_runner.py` from where the reason is composed: [rendered template fields 
are redacted in the task 
process](https://github.com/apache/airflow/blob/74e7573a5606e4fa7bf83616ab1aef5106a1bb62/task-sdk/src/airflow/sdk/execution_time/task_runner.py#L1355-L1365),
 with the comment "This ensures that the secrets those are registered via 
mask_secret() on workers / dag processor are properly masked on the UI." 
Redacting alongside the existing `decision.reason[:500]` truncation would bind; 
keep these validators as defence in depth. If you'd rather not touch #73027's 
code here, saying in the docs that masking is best-effort would at least stop 
the control reading as guaranteed.
   
   Same applies to the copy on `TaskInstanceHistoryResponse`. Minor, both 
models: `str(...)` around `redact` is dead, every path reachable from a `str` 
input returns a `str` including the fail-closed `"<redaction-failed>"`; 
`cast("str", redact(...))` is what `policies/retry.py:218` uses.



##########
airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py:
##########
@@ -244,8 +245,30 @@ def test_should_respond_200(self, test_client, session):
             "trigger": None,
             "triggerer_job": None,
             "team_name": None,
+            "state_reason": None,
         }
 
+    def test_should_include_state_reason(self, test_client, session):
+        self.create_task_instances(session, task_instances=[{"retry_reason": 
"auth error, do not retry"}])
+        response = test_client.get(
+            
"/dags/example_python_operator/dagRuns/TEST_DAG_RUN_ID/taskInstances/print_the_context"
+        )
+        assert response.status_code == 200
+        assert response.json()["state_reason"] == "auth error, do not retry"
+
+    @pytest.mark.enable_redact
+    def test_should_redact_secrets_in_state_reason(self, test_client, session):
+        """A policy may compose the reason from an unredacted exception, so 
mask on the way out."""
+        mask_secret("hunter2")
+        self.create_task_instances(
+            session, task_instances=[{"retry_reason": "auth: the token hunter2 
expired"}]
+        )
+        response = test_client.get(
+            
"/dags/example_python_operator/dagRuns/TEST_DAG_RUN_ID/taskInstances/print_the_context"
+        )
+        assert response.status_code == 200
+        assert "hunter2" not in response.json()["state_reason"]

Review Comment:
   A `not in` assertion is satisfied by destroying the whole value. I mutated 
both validators to collapse any secret-bearing string to `"***"` and all four 
state-reason tests stayed green, so nothing here notices if the reason text 
disappears entirely. The real output is `'auth: the token *** expired'`; 
asserting that exact string instead pins both halves. Same for the history copy 
at 2859.
   
   Separately, `mask_secret("hunter2")` writes to the cached process-global 
masker with no teardown, and there's no `reset_secrets_masker` fixture in the 
`api_fastapi` conftests or `tests_common`, so the pattern outlives the test for 
the rest of the xdist worker's session.



##########
airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.test.tsx:
##########
@@ -0,0 +1,189 @@
+/*!
+ * 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 { beforeEach, 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 commonLocale from "../../../public/i18n/locales/en/common.json";
+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", () => {
+  // Without the bundle, i18n.t() echoes the key back and every title 
assertion compares a key
+  // against itself, so interpolated counts are never checked.
+  beforeEach(() => {
+    i18n.addResourceBundle("en", "common", commonLocale, true, true);
+  });
+
+  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([
+    { maxTries: 2, state: "failed", titleKey: "failed", totalTries: 3, 
tryNumber: 3 },
+    // try_number != max_tries + 1 is the norm mid-retry, and the differing 
numbers are what make
+    // a swapped or off-by-one interpolation visible.
+    { maxTries: 3, state: "up_for_retry", titleKey: "upForRetry", totalTries: 
4, tryNumber: 2 },
+  ] as const)(
+    "titles the banner for a $state task",
+    ({ maxTries, state, titleKey, totalTries, tryNumber }) => {
+      renderDetails(
+        buildTaskInstance({
+          max_tries: maxTries,
+          state,
+          state_reason: "auth error",
+          try_number: tryNumber,
+        }),
+      );
+
+      expect(screen.getByTestId("state-reason-alert")).toHaveTextContent(
+        i18n.t(`common:taskInstance.stateReasonSummary.${titleKey}`, { 
totalTries, tryNumber }),
+      );
+    },
+  );
+
+  // Chakra encodes `status` in a generated class rather than a DOM attribute, 
so the error/warning
+  // distinction can only be pinned as "the two states do not render 
identically".
+  it("styles a failed banner differently from an up_for_retry one", () => {
+    const { unmount } = renderDetails(buildTaskInstance({ state: "failed", 
state_reason: "auth error" }));
+    const failedClass = screen.getByTestId("state-reason-alert").className;
+
+    unmount();
+    renderDetails(buildTaskInstance({ state: "up_for_retry", state_reason: 
"auth error" }));
+
+    
expect(screen.getByTestId("state-reason-alert").className).not.toBe(failedClass);
+  });
+
+  // The reason is only cleared once the task next reaches RUNNING, so a 
cleared task keeps a
+  // reason describing the previous attempt. Gating on state is what stops it 
being shown, and it
+  // has to cover the row as well as the banner or the stale text just moves 
down the page.
+  it.each(["queued", "running", "success", null] as const)(
+    "renders neither the banner nor the row for a %s task that still carries a 
reason",
+    (state) => {
+      renderDetails(buildTaskInstance({ state, state_reason: "auth error, do 
not retry" }));
+
+      
expect(screen.queryByTestId("state-reason-alert")).not.toBeInTheDocument();
+      
expect(screen.queryByText(i18n.t("common:taskInstance.stateReason"))).not.toBeInTheDocument();
+      expect(screen.queryByText("auth error, do not 
retry")).not.toBeInTheDocument();
+    },
+  );
+
+  it("keeps an earlier failed try's reason while the task is running again", 
() => {
+    renderDetails(buildTaskInstance({ state: "running", state_reason: null }), 
{
+      state: "failed",
+      state_reason: "try 1: auth error",
+    });
+
+    expect(screen.queryByTestId("state-reason-alert")).not.toBeInTheDocument();
+    expect(screen.getByText("try 1: auth error")).toBeInTheDocument();
+  });
+
+  it("titles the banner with the try counts so it is distinct from the per-try 
row", () => {

Review Comment:
   This one has no detection power the `it.each` case above doesn't already 
have: same `failed` state, same counts, only the reason string differs, and the 
reason isn't asserted. Because `tryNumber === totalTries === 3` here it can't 
even catch a swapped interpolation. I ran five production mutations with and 
without it and the kill counts were identical. Either drop it or have it assert 
the banner body contains `"rate limit"`, which would make it about the 
banner/row split its name claims.



##########
providers/common/ai/docs/retry_policies.rst:
##########
@@ -133,7 +133,9 @@ When a task fails, either policy:
    ``retry_reason``, on a FAIL as well as a RETRY: ``<category>: <reasoning>``
    from ``LLMRetryPolicy``, or one line such as
    ``category=network confidence=0.91 threshold=0.60 action=retry delay=10s``
-   from ``ClassifierRetryPolicy``.
+   from ``ClassifierRetryPolicy``. The REST API exposes it as ``state_reason``

Review Comment:
   Both halves of this sentence are unconditional. The REST field and the page 
land in 3.4.0, but the note at the top of the page still says `Requires Airflow 
>= 3.3.0` and the provider floors at `apache-airflow>=3.0.0`; the paragraph 
four lines down already qualifies its sibling claim with "Recording on a FAIL 
requires Airflow 3.4.0", so the same qualifier fits here. And the page only 
shows the reason while the task is `failed` or `up_for_retry`, which is worth 
the half-clause given you deliberately gated it.



-- 
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]

Reply via email to