This is an automated email from the ASF dual-hosted git repository.
amoghrajesh pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 8b3d4d315fa Surface the retry policy decision on task instances page
(#73030)
8b3d4d315fa is described below
commit 8b3d4d315faafc7c4a22ab3dc208128c67b3b868
Author: Amogh Desai <[email protected]>
AuthorDate: Tue Sep 29 13:44:13 2026 +0530
Surface the retry policy decision on task instances page (#73030)
---
.../core_api/datamodels/task_instance_history.py | 19 ++-
.../core_api/datamodels/task_instances.py | 18 ++-
.../api_fastapi/core_api/openapi/_private_ui.yaml | 9 ++
.../core_api/openapi/v2-rest-api-generated.yaml | 18 +++
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 24 ++++
.../airflow/ui/openapi-gen/requests/types.gen.ts | 8 ++
.../airflow/ui/public/i18n/locales/en/common.json | 5 +
.../ui/src/pages/TaskInstance/Details.test.tsx | 137 +++++++++++++++++++++
.../airflow/ui/src/pages/TaskInstance/Details.tsx | 16 +++
.../ui/src/pages/TaskInstance/Header.test.tsx | 54 ++++++++
.../airflow/ui/src/pages/TaskInstance/Header.tsx | 20 +++
.../ui/src/pages/TaskInstance/stateReason.ts | 32 +++++
.../core_api/routes/public/test_hitl.py | 1 +
.../core_api/routes/public/test_task_instances.py | 93 ++++++++++++++
.../src/airflowctl/api/datamodels/generated.py | 14 +++
providers/common/ai/docs/retry_policies.rst | 16 ++-
.../src/airflow/sdk/execution_time/task_runner.py | 8 +-
.../task_sdk/execution_time/test_task_runner.py | 24 ++++
18 files changed, 507 insertions(+), 9 deletions(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instance_history.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instance_history.py
index c913d91c1e6..6aba7ec12b3 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instance_history.py
+++
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instance_history.py
@@ -17,14 +17,16 @@
from __future__ import annotations
from datetime import datetime
-from typing import Annotated
+from typing import Annotated, cast
from pydantic import (
AliasPath,
BeforeValidator,
Field,
+ field_validator,
)
+from airflow._shared.secrets_masker import redact
from airflow.api_fastapi.core_api.base import BaseModel
from airflow.api_fastapi.core_api.datamodels.dag_versions import
DagVersionResponse
from airflow.utils.state import TaskInstanceState
@@ -62,6 +64,21 @@ class TaskInstanceHistoryResponse(BaseModel):
executor: str | None
executor_config: Annotated[str, BeforeValidator(str)]
dag_version: DagVersionResponse | None
+ state_reason: str | None = Field(
+ default=None,
+ validation_alias="retry_reason",
+ description=(
+ "The reason the task instance reached its current state, as
recorded by a retry policy. May describe a previous attempt: it is cleared only
when the task next starts running, so a task waiting to be retried or re-run
can still carry the reason its last attempt ended."
+ ),
+ )
+
+ @field_validator("state_reason", mode="after")
+ @classmethod
+ def redact_state_reason(cls, v: str | None) -> str | None:
+ # See TaskInstanceResponse.
+ if v is None:
+ return None
+ return cast("str", redact(v))
class TaskInstanceHistoryCollectionResponse(BaseModel):
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
index d138629e73a..a4d65063b60 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
@@ -18,7 +18,7 @@ from __future__ import annotations
from collections.abc import Iterable
from datetime import datetime
-from typing import Annotated, Any
+from typing import Annotated, Any, cast
from uuid import UUID
from pydantic import (
@@ -34,6 +34,7 @@ from pydantic import (
model_validator,
)
+from airflow._shared.secrets_masker import redact
from airflow.api_fastapi.core_api.base import BaseModel, StrictBaseModel
from airflow.api_fastapi.core_api.datamodels.dag_versions import
DagVersionResponse
from airflow.api_fastapi.core_api.datamodels.job import JobResponse
@@ -89,6 +90,21 @@ 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",
+ description=(
+ "The reason the task instance reached its current state, as
recorded by a retry policy. May describe a previous attempt: it is cleared only
when the task next starts running, so a task waiting to be retried or re-run
can still carry the reason its last attempt ended."
+ ),
+ )
+
+ @field_validator("state_reason", mode="after")
+ @classmethod
+ def redact_state_reason(cls, v: str | None) -> str | None:
+ # The worker already redacts this. Kept for rows written by an older
task-sdk.
+ if v is None:
+ return None
+ return cast("str", redact(v))
class TaskInstanceCollectionResponse(BaseModel):
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
index acb0fcf5512..1681b462b37 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
+++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
@@ -4935,6 +4935,15 @@ components:
- type: string
- type: 'null'
title: Team Name
+ state_reason:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: State Reason
+ description: 'The reason the task instance reached its current
state, as
+ recorded by a retry policy. May describe a previous attempt: it is
cleared
+ only when the task next starts running, so a task waiting to be
retried
+ or re-run can still carry the reason its last attempt ended.'
type: object
required:
- id
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index dabcc47c9b7..415b36445d7 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -16495,6 +16495,15 @@ components:
anyOf:
- $ref: '#/components/schemas/DagVersionResponse'
- type: 'null'
+ state_reason:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: State Reason
+ description: 'The reason the task instance reached its current
state, as
+ recorded by a retry policy. May describe a previous attempt: it is
cleared
+ only when the task next starts running, so a task waiting to be
retried
+ or re-run can still carry the reason its last attempt ended.'
type: object
required:
- task_id
@@ -16678,6 +16687,15 @@ components:
- type: string
- type: 'null'
title: Team Name
+ state_reason:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: State Reason
+ description: 'The reason the task instance reached its current
state, as
+ recorded by a retry policy. May describe a previous attempt: it is
cleared
+ only when the task next starts running, so a task waiting to be
retried
+ or re-run can still carry the reason its last attempt ended.'
type: object
required:
- id
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index b028696f3a8..8ec5c8c5116 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -7712,6 +7712,18 @@ export const $TaskInstanceHistoryResponse = {
type: 'null'
}
]
+ },
+ state_reason: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'State Reason',
+ description: 'The reason the task instance reached its current
state, as recorded by a retry policy. May describe a previous attempt: it is
cleared only when the task next starts running, so a task waiting to be retried
or re-run can still carry the reason its last attempt ended.'
}
},
type: 'object',
@@ -8012,6 +8024,18 @@ export const $TaskInstanceResponse = {
}
],
title: 'Team Name'
+ },
+ state_reason: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'State Reason',
+ description: 'The reason the task instance reached its current
state, as recorded by a retry policy. May describe a previous attempt: it is
cleared only when the task next starts running, so a task waiting to be retried
or re-run can still carry the reason its last attempt ended.'
}
},
type: 'object',
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index 403a6eb5785..ce784736c7a 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -2050,6 +2050,10 @@ export type TaskInstanceHistoryResponse = {
executor: string | null;
executor_config: string;
dag_version: DagVersionResponse | null;
+ /**
+ * The reason the task instance reached its current state, as recorded by
a retry policy. May describe a previous attempt: it is cleared only when the
task next starts running, so a task waiting to be retried or re-run can still
carry the reason its last attempt ended.
+ */
+ state_reason?: string | null;
};
/**
@@ -2093,6 +2097,10 @@ export type TaskInstanceResponse = {
triggerer_job: JobResponse | null;
dag_version: DagVersionResponse | null;
team_name?: string | null;
+ /**
+ * The reason the task instance reached its current state, as recorded by
a retry policy. May describe a previous attempt: it is cleared only when the
task next starts running, so a task waiting to be retried or re-run can still
carry the reason its last attempt ended.
+ */
+ state_reason?: string | null;
};
/**
diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
index fe8b8b93d0d..7a781ef45ed 100644
--- a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
+++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
@@ -454,6 +454,11 @@
"queuedWhen": "Queued At",
"renderedMapIndex": "Rendered Map Index",
"scheduledWhen": "Scheduled At",
+ "stateReason": "Reason for state",
+ "stateReasonSummary": {
+ "failed": "Stopped on try {{tryNumber}} of {{totalTries}}",
+ "upForRetry": "Retrying after try {{tryNumber}} of {{totalTries}}"
+ },
"trigger": "Trigger",
"triggerer": {
"assigned": "Assigned triggerer",
diff --git
a/airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.test.tsx
b/airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.test.tsx
new file mode 100644
index 00000000000..cc7d063bbd7
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.test.tsx
@@ -0,0 +1,137 @@
+/*!
+ * 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 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 row", () => {
+ // Without the bundle i18n.t() echoes the key, so the label assertions below
would pass blindly.
+ 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.queryByText(i18n.t("common:taskInstance.stateReason"))).not.toBeInTheDocument();
+ });
+
+ // Cleared only once the task next reaches RUNNING, so these states still
carry a stale reason.
+ 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.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.getByText("try 1: auth error")).toBeInTheDocument();
+ });
+
+ it("shows the selected try's reason in the table", () => {
+ renderDetails(buildTaskInstance({ state_reason: "latest try: rate limit"
}), {
+ state_reason: "older try: auth error",
+ });
+
+ expect(screen.getByText("older try: auth error")).toBeInTheDocument();
+ expect(screen.queryByText("latest try: rate
limit")).not.toBeInTheDocument();
+
expect(screen.getByText(i18n.t("common:taskInstance.stateReason"))).toBeInTheDocument();
+ });
+});
diff --git a/airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.tsx
b/airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.tsx
index 017e24c54c7..c35a3462dd7 100644
--- a/airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.tsx
@@ -43,6 +43,7 @@ import { isStatePending, useAutoRefresh, useDurationFormat }
from "src/utils";
import { BlockingDeps } from "./BlockingDeps";
import { ExtraLinks } from "./ExtraLinks";
import { TriggererInfo } from "./TriggererInfo";
+import { stateReasonDisplay } from "./stateReason";
export const Details = () => {
const { t: translate } = useTranslation();
@@ -115,6 +116,15 @@ export const Details = () => {
return translate("common:none", { defaultValue: "None" });
};
+ // Keyed off the selected try's own state, so an earlier failed try keeps
its reason while the
+ // current one is running again.
+ const tryStateReason =
+ tryInstance?.state_reason !== null &&
+ tryInstance?.state_reason !== undefined &&
+ stateReasonDisplay(tryInstance.state) !== undefined
+ ? tryInstance.state_reason
+ : undefined;
+
// omit kwargs from trigger
const triggerWithoutKwargs = taskInstance?.trigger
? (({ kwargs, ...rest }) => rest)(taskInstance.trigger)
@@ -162,6 +172,12 @@ export const Details = () => {
</Flex>
</Table.Cell>
</Table.Row>
+ {tryStateReason === undefined ? undefined : (
+ <Table.Row>
+ <Table.Cell>{translate("taskInstance.stateReason")}</Table.Cell>
+ <Table.Cell>{tryStateReason}</Table.Cell>
+ </Table.Row>
+ )}
<Table.Row>
<Table.Cell>{translate("taskId")}</Table.Cell>
<Table.Cell>
diff --git a/airflow-core/src/airflow/ui/src/pages/TaskInstance/Header.test.tsx
b/airflow-core/src/airflow/ui/src/pages/TaskInstance/Header.test.tsx
index 34ab9a45cb9..d805f3ed052 100644
--- a/airflow-core/src/airflow/ui/src/pages/TaskInstance/Header.test.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/TaskInstance/Header.test.tsx
@@ -25,6 +25,7 @@ import type { 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 { Header } from "./Header";
// Action buttons and note preview pull in mutation/permission wiring that is
@@ -77,3 +78,56 @@ describe("Header", () => {
expect(screen.queryByText(i18n.t("common:dagDetails.team"))).not.toBeInTheDocument();
});
});
+
+const renderHeader = (overrides: Partial<TaskInstanceResponse>) =>
+ render(<Header taskInstance={{ ...baseTaskInstance, ...overrides }} />, {
wrapper: Wrapper });
+
+describe("Header state reason banner", () => {
+ // Without the bundle i18n.t() echoes the key, so the titles below would
assert nothing.
+ beforeEach(() => {
+ i18n.addResourceBundle("en", "common", commonLocale, true, true);
+ });
+
+ it("does not render when there is no reason", () => {
+ renderHeader({ state: "failed", state_reason: null });
+
+ expect(screen.queryByTestId("state-reason-alert")).not.toBeInTheDocument();
+ });
+
+ // Cleared only once the task next reaches RUNNING, so these states still
carry a stale reason.
+ it.each(["queued", "running", "success", null] as const)(
+ "does not render for a %s task that still carries a reason",
+ (state) => {
+ renderHeader({ state, state_reason: "auth error, do not retry" });
+
+
expect(screen.queryByTestId("state-reason-alert")).not.toBeInTheDocument();
+ },
+ );
+
+ it.each([
+ { maxTries: 2, state: "failed", titleKey: "failed", totalTries: 3,
tryNumber: 3 },
+ // 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 }) => {
+ renderHeader({ 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 }),
+ );
+ expect(screen.getByTestId("state-reason-alert")).toHaveTextContent("auth
error");
+ },
+ );
+
+ // Chakra puts `status` in a generated class, so "not identical" is all that
can be asserted.
+ it("styles a failed banner differently from an up_for_retry one", () => {
+ const { unmount } = renderHeader({ state: "failed", state_reason: "auth
error" });
+ const failedClass = screen.getByTestId("state-reason-alert").className;
+
+ unmount();
+ renderHeader({ state: "up_for_retry", state_reason: "auth error" });
+
+
expect(screen.getByTestId("state-reason-alert").className).not.toBe(failedClass);
+ });
+});
diff --git a/airflow-core/src/airflow/ui/src/pages/TaskInstance/Header.tsx
b/airflow-core/src/airflow/ui/src/pages/TaskInstance/Header.tsx
index 850c3b42e30..26d79acd925 100644
--- a/airflow-core/src/airflow/ui/src/pages/TaskInstance/Header.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/TaskInstance/Header.tsx
@@ -24,6 +24,8 @@ import { MdOutlineTask } from "react-icons/md";
import type { TaskInstanceResponse } from "openapi/requests/types.gen";
+import { Alert } from "src/system-components";
+
import { ClearTaskInstanceButton } from "src/components/Clear";
import ClearTaskInstanceDialog from
"src/components/Clear/TaskInstance/ClearTaskInstanceDialog";
import { DagVersion } from "src/components/DagVersion";
@@ -37,6 +39,8 @@ import { useShowTeam } from "src/hooks/useShowTeam";
import { useTaskInstanceNote } from "src/queries/useTaskInstanceNote";
import { useDurationFormat } from "src/utils";
+import { stateReasonDisplay } from "./stateReason";
+
export const Header = ({ taskInstance }: { readonly taskInstance:
TaskInstanceResponse }) => {
const { t: translate } = useTranslation();
const { formatElapsed, renderDuration } = useDurationFormat();
@@ -80,8 +84,24 @@ export const Header = ({ taskInstance }: { readonly
taskInstance: TaskInstanceRe
// Stable dialog state at header/page level
const [clearOpen, setClearOpen] = useState(false);
+ // On the header, not the details tab, so it shows on every tab without
duplicating the row.
+ const stateReasonDisplayed = stateReasonDisplay(taskInstance.state);
+ const stateReason = taskInstance.state_reason;
+
return (
<Box display="flex" flexDirection="column" gap={3}>
+ {stateReasonDisplayed === undefined || stateReason === null ||
stateReason === undefined ? undefined : (
+ <Alert
+ data-testid="state-reason-alert"
+ status={stateReasonDisplayed.status}
+
title={translate(`taskInstance.stateReasonSummary.${stateReasonDisplayed.titleKey}`,
{
+ totalTries: taskInstance.max_tries + 1,
+ tryNumber: taskInstance.try_number,
+ })}
+ >
+ {stateReason}
+ </Alert>
+ )}
<HeaderCard
actions={
<>
diff --git a/airflow-core/src/airflow/ui/src/pages/TaskInstance/stateReason.ts
b/airflow-core/src/airflow/ui/src/pages/TaskInstance/stateReason.ts
new file mode 100644
index 00000000000..60f7051fab5
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/TaskInstance/stateReason.ts
@@ -0,0 +1,32 @@
+/*!
+ * 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 type { TaskInstanceState } from "openapi/requests/types.gen";
+
+type StateReasonDisplay = { status: "error" | "warning"; titleKey: string };
+
+// The header banner and the per-try row both look the state up here, so a
state added to one
+// surface cannot be forgotten on the other: it has to bring a title with it.
Keyed by
+// TaskInstanceState so a misspelt state is a compile error rather than a row
that never renders.
+const STATE_REASON_DISPLAY: Partial<Record<TaskInstanceState,
StateReasonDisplay>> = {
+ failed: { status: "error", titleKey: "failed" },
+ up_for_retry: { status: "warning", titleKey: "upForRetry" },
+};
+
+export const stateReasonDisplay = (state: TaskInstanceState | null |
undefined) =>
+ state === null || state === undefined ? undefined :
STATE_REASON_DISPLAY[state];
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py
index dab5bef1a2f..0df529b047a 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py
@@ -268,6 +268,7 @@ def expected_sample_hitl_detail_dict(sample_ti:
TaskInstance) -> dict[str, Any]:
"task_display_name": "sample_task_hitl",
"task_id": TASK_ID,
"team_name": None,
+ "state_reason": None,
"trigger": None,
"triggerer_job": None,
"try_number": 0,
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
index 4f8c3fa45a2..8d93e289a44 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
@@ -31,6 +31,7 @@ from fastapi.testclient import TestClient
from sqlalchemy import delete, func, select, update
from sqlalchemy.orm import joinedload
+from airflow._shared.secrets_masker import mask_secret
from airflow._shared.state import TaskScope
from airflow._shared.timezones.timezone import datetime
from airflow.api_fastapi.auth.managers.simple.user import SimpleAuthManagerUser
@@ -245,8 +246,40 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"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.fixture
+ def masked_secret(self):
+ """The masker is a cached process global, so drop the pattern again
for the next test."""
+ from airflow._shared.secrets_masker import _secrets_masker
+
+ masker = _secrets_masker()
+ patterns, replacer = set(masker.patterns), masker.replacer
+ mask_secret("hunter2")
+ yield
+ masker.patterns, masker.replacer = patterns, replacer
+
+ @pytest.mark.enable_redact
+ def test_should_redact_secrets_in_state_reason(self, test_client, session,
masked_secret):
+ """A policy may compose the reason from an unredacted exception, so
mask on the way out."""
+ 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 response.json()["state_reason"] == "auth: the token *** expired"
+
@conf_vars({("core", "multi_team"): "True"})
def test_should_include_team_name(self, test_client, session):
self.create_task_instances(session)
@@ -330,6 +363,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
"dag_version": {
"id": response_data["dag_version"]["id"],
"version_number": expected_version_number,
@@ -426,6 +460,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"unixname": getuser(),
},
"team_name": None,
+ "state_reason": None,
}
def test_should_respond_200_with_task_state_in_removed(self, test_client,
session):
@@ -480,6 +515,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
def test_should_respond_200_task_instance_with_rendered(self, test_client,
session):
@@ -537,6 +573,7 @@ class TestGetTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
def test_raises_404_for_nonexistent_task_instance(self, test_client):
@@ -658,6 +695,7 @@ class TestGetMappedTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
def test_should_respond_401(self, unauthenticated_test_client):
@@ -2806,8 +2844,44 @@ class TestGetTaskInstanceTry(TestTaskInstanceEndpoint):
"id": response_data["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
}
+ def test_should_include_state_reason_from_history(self, test_client,
session):
+ self.create_task_instances(
+ session,
+ task_instances=[{"state": State.SUCCESS, "retry_reason": "auth
error, do not retry"}],
+ with_ti_history=True,
+ )
+ response = test_client.get(
+
"/dags/example_python_operator/dagRuns/TEST_DAG_RUN_ID/taskInstances/print_the_context/tries/1"
+ )
+ assert response.status_code == 200
+ assert response.json()["state_reason"] == "auth error, do not retry"
+
+ @pytest.fixture
+ def masked_secret(self):
+ from airflow._shared.secrets_masker import _secrets_masker
+
+ masker = _secrets_masker()
+ patterns, replacer = set(masker.patterns), masker.replacer
+ mask_secret("hunter2")
+ yield
+ masker.patterns, masker.replacer = patterns, replacer
+
+ @pytest.mark.enable_redact
+ def test_should_redact_secrets_in_state_reason_from_history(self,
test_client, session, masked_secret):
+ self.create_task_instances(
+ session,
+ task_instances=[{"state": State.SUCCESS, "retry_reason": "auth:
the token hunter2 expired"}],
+ with_ti_history=True,
+ )
+ response = test_client.get(
+
"/dags/example_python_operator/dagRuns/TEST_DAG_RUN_ID/taskInstances/print_the_context/tries/1"
+ )
+ assert response.status_code == 200
+ assert response.json()["state_reason"] == "auth: the token *** expired"
+
@pytest.mark.parametrize("try_number", [1, 2])
def test_should_respond_200_with_different_try_numbers(self, test_client,
try_number, session):
self.create_task_instances(session, task_instances=[{"state":
State.SUCCESS}], with_ti_history=True)
@@ -2852,6 +2926,7 @@ class TestGetTaskInstanceTry(TestTaskInstanceEndpoint):
"id": response_data["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
}
@pytest.mark.parametrize("try_number", [1, 2])
@@ -2928,6 +3003,7 @@ class TestGetTaskInstanceTry(TestTaskInstanceEndpoint):
"id": response_data["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
}
def test_should_respond_200_with_task_state_in_deferred(self, test_client,
session):
@@ -2996,6 +3072,7 @@ class TestGetTaskInstanceTry(TestTaskInstanceEndpoint):
"id": response_data["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
}
def test_should_respond_200_with_task_state_in_removed(self, test_client,
session):
@@ -3043,6 +3120,7 @@ class TestGetTaskInstanceTry(TestTaskInstanceEndpoint):
"id": response_data["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
}
def test_should_respond_401(self, unauthenticated_test_client):
@@ -3118,6 +3196,7 @@ class TestGetTaskInstanceTry(TestTaskInstanceEndpoint):
"created_at": mock.ANY,
"dag_display_name": "dag_with_multiple_versions",
},
+ "state_reason": None,
}
def test_should_not_return_duplicate_runs(self, test_client, session):
@@ -3855,6 +3934,7 @@ class
TestPostClearTaskInstances(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
"try_number": 1,
"unixname": getuser(),
},
@@ -4414,6 +4494,7 @@ class TestGetTaskInstanceTries(TestTaskInstanceEndpoint):
"id":
response_data["task_instances"][0]["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
},
{
"dag_id": "example_python_operator",
@@ -4451,6 +4532,7 @@ class TestGetTaskInstanceTries(TestTaskInstanceEndpoint):
"id":
response_data["task_instances"][1]["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
},
],
"total_entries": 2,
@@ -4563,6 +4645,7 @@ class TestGetTaskInstanceTries(TestTaskInstanceEndpoint):
"id":
response_data["task_instances"][0]["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
},
{
"dag_id": "example_python_operator",
@@ -4600,6 +4683,7 @@ class TestGetTaskInstanceTries(TestTaskInstanceEndpoint):
"id":
response_data["task_instances"][1]["dag_version"]["id"],
"version_number": 1,
},
+ "state_reason": None,
},
],
"total_entries": 2,
@@ -4667,6 +4751,7 @@ class TestGetTaskInstanceTries(TestTaskInstanceEndpoint):
"created_at": mock.ANY,
"dag_display_name": "dag_with_multiple_versions",
},
+ "state_reason": None,
}
@@ -4792,6 +4877,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
],
"total_entries": 1,
@@ -5070,6 +5156,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
],
"total_entries": 1,
@@ -5210,6 +5297,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
],
"total_entries": 1,
@@ -5275,6 +5363,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
],
"total_entries": 1,
@@ -5372,6 +5461,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
],
"total_entries": 1,
@@ -5457,6 +5547,7 @@ class TestPatchTaskInstance(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
_check_task_instance_note(
@@ -5653,6 +5744,7 @@ class
TestPatchTaskInstanceDryRun(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
],
"total_entries": 1,
@@ -5943,6 +6035,7 @@ class
TestPatchTaskInstanceDryRun(TestTaskInstanceEndpoint):
"trigger": None,
"triggerer_job": None,
"team_name": None,
+ "state_reason": None,
}
],
"total_entries": 1,
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index 924d3420928..4331cdb6299 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -2414,6 +2414,13 @@ class TaskInstanceHistoryResponse(BaseModel):
executor: Annotated[str | None, Field(title="Executor")]
executor_config: Annotated[str, Field(title="Executor Config")]
dag_version: DagVersionResponse | None
+ state_reason: Annotated[
+ str | None,
+ Field(
+ description="The reason the task instance reached its current
state, as recorded by a retry policy. May describe a previous attempt: it is
cleared only when the task next starts running, so a task waiting to be retried
or re-run can still carry the reason its last attempt ended.",
+ title="State Reason",
+ ),
+ ] = None
class TaskInstanceResponse(BaseModel):
@@ -2456,6 +2463,13 @@ class TaskInstanceResponse(BaseModel):
triggerer_job: JobResponse | None
dag_version: DagVersionResponse | None
team_name: Annotated[str | None, Field(title="Team Name")] = None
+ state_reason: Annotated[
+ str | None,
+ Field(
+ description="The reason the task instance reached its current
state, as recorded by a retry policy. May describe a previous attempt: it is
cleared only when the task next starts running, so a task waiting to be retried
or re-run can still carry the reason its last attempt ended.",
+ title="State Reason",
+ ),
+ ] = None
class TaskResponse(BaseModel):
diff --git a/providers/common/ai/docs/retry_policies.rst
b/providers/common/ai/docs/retry_policies.rst
index a94f839056b..914bf71f6f5 100644
--- a/providers/common/ai/docs/retry_policies.rst
+++ b/providers/common/ai/docs/retry_policies.rst
@@ -133,7 +133,10 @@ 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``. From Airflow 3.4 the REST API exposes it as
+ ``state_reason`` on a task instance and on each try, and the Task Instance
+ page shows it under **Reason for state** while the task is failed or up for
+ retry.
This classification call is a separate model request, made by the policy
itself rather than by an operator -- it is not subject to an operator's
@@ -405,11 +408,12 @@ upper limit; zero or negative means no override, so the
task's own
``retry_delay`` and backoff apply.
``category`` and ``reasoning`` become the ``retry_reason`` (truncated to 500
-characters), recorded on both outcomes. On a RETRY the value is cleared once
the next attempt starts running;
-a FAIL is terminal, so there is no next attempt to clear it and the reason
stays
-on the row. Only the model's own words are stored -- attempt counts are left to
-whatever displays the reason. Recording on a FAIL requires Airflow 3.4.0; on
-earlier versions only the RETRY outcome is recorded.
+characters), recorded on both outcomes. On a RETRY the value is cleared once
+the next attempt starts running; a FAIL is terminal, so there is no next
+attempt to clear it and the reason stays on the row. Only the model's own words
+are stored -- attempt counts are left to whatever displays the reason.
+Recording on a FAIL requires Airflow 3.4.0; on earlier versions only the RETRY
+outcome is recorded.
Under ``ClassifierRetryPolicy`` it answers the category name and nothing else.
It does not
decide whether to retry, it does not choose the delay, and it does not explain
diff --git a/task-sdk/src/airflow/sdk/execution_time/task_runner.py
b/task-sdk/src/airflow/sdk/execution_time/task_runner.py
index 46f2fdae411..fba2ac56829 100644
--- a/task-sdk/src/airflow/sdk/execution_time/task_runner.py
+++ b/task-sdk/src/airflow/sdk/execution_time/task_runner.py
@@ -27,10 +27,11 @@ import sys
import time
from collections.abc import Callable, Iterable, Iterator, Mapping
from contextlib import ExitStack, contextmanager, suppress
+from dataclasses import replace
from datetime import datetime, timedelta, timezone
from itertools import product
from pathlib import Path
-from typing import TYPE_CHECKING, Annotated, Any, Literal
+from typing import TYPE_CHECKING, Annotated, Any, Literal, cast
from urllib.parse import quote
import attrs
@@ -1816,6 +1817,8 @@ def _evaluate_retry_policy(
Returns ``None`` when no policy is configured so the caller falls through
to the standard retry logic.
"""
+ from airflow.sdk._shared.secrets_masker import redact
+
policy = getattr(ti.task, "retry_policy", None)
if policy is None:
return None
@@ -1828,6 +1831,9 @@ def _evaluate_retry_policy(
context=context,
)
if decision.reason:
+ # Mask here, where mask_secret() registered the value: the API
server rendering this
+ # later has its own masker and does not know the worker's secrets.
+ decision = replace(decision, reason=cast("str",
redact(decision.reason)))
# Close the group so the retry policy decision is not hidden
inside "Post Execute".
log.info("::endgroup::")
log.info("Retry policy decision", action=decision.action.value,
reason=decision.reason)
diff --git a/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
b/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
index 198b6ec070b..b14e3e1bb69 100644
--- a/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
+++ b/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
@@ -64,6 +64,7 @@ from airflow.sdk import (
timezone,
)
from airflow.sdk._shared.observability.metrics.base_stats_logger import
StatsLogger
+from airflow.sdk._shared.secrets_masker import _secrets_masker
from airflow.sdk._shared.state import AssetScope, TaskScope
from airflow.sdk.api.datamodels._generated import (
AssetProfile,
@@ -1274,6 +1275,29 @@ def
test_retry_policy_retry_exhausted_reason_is_truncated(create_runtime_ti, moc
assert msg.retry_reason == "z" * 500
[email protected]_redact
+def test_retry_policy_reason_is_redacted_in_the_worker(create_runtime_ti,
mock_supervisor_comms):
+ """The reason is masked where mask_secret() registered the value, not in
the API server."""
+ _secrets_masker().add_mask("hunter2", None)
+
+ class _AlwaysFails(BaseOperator):
+ def execute(self, context):
+ raise RuntimeError("403 Forbidden: token hunter2 expired")
+
+ class _EchoPolicy(RetryPolicy):
+ def evaluate(self, exception, try_number, max_tries, context=None):
+ return RetryDecision(action=RetryAction.FAIL, reason=f"auth:
{exception}")
+
+ task = _AlwaysFails(task_id="redacted_reason", retries=2,
retry_policy=_EchoPolicy())
+ 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 == "auth: 403 Forbidden: token *** expired"
+
+
def test_plain_retries_exhausted_has_no_reason(create_runtime_ti,
mock_supervisor_comms):
"""Without a retry policy, exhausting the retry budget must not synthesize
a reason."""