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]