kaxil commented on code in PR #74352: URL: https://github.com/apache/airflow/pull/74352#discussion_r4222339996
########## airflow-core/src/airflow/ui/src/hooks/useTaskInstanceView.ts: ########## @@ -0,0 +1,67 @@ +/*! + * 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 { useParams, useSearchParams } from "react-router-dom"; + +import { + useTaskInstanceServiceGetMappedTaskInstance, + useTaskInstanceServiceGetTaskInstanceTryDetails, +} from "openapi/queries"; + +import { SearchParamsKeys } from "src/constants/searchParams"; +import { useTaskInstanceCoordinates } from "src/hooks/useTaskInstanceCoordinates"; +import { isStatePending, useAutoRefresh } from "src/utils"; + +export const isExactTryView = (searchParams: URLSearchParams) => + searchParams.has(SearchParamsKeys.TRY_NUMBER) && + searchParams.has(SearchParamsKeys.REGION_ID) && + searchParams.has(SearchParamsKeys.REGION_INDEX); + +export const useTaskInstanceView = () => { + const { dagId = "", mapIndex = "-1", runId = "", taskId = "" } = useParams(); + const [searchParams] = useSearchParams(); + const coordinates = useTaskInstanceCoordinates(); + const tryParameter = searchParams.get(SearchParamsKeys.TRY_NUMBER); + const exactTry = isExactTryView(searchParams); + const refetchInterval = useAutoRefresh({ dagId }); + const params = { ...coordinates, dagId, dagRunId: runId, mapIndex: Number(mapIndex), taskId }; + const live = useTaskInstanceServiceGetMappedTaskInstance(params, undefined, { + enabled: !Number.isNaN(params.mapIndex), + refetchInterval: (query) => isStatePending(query.state.data?.state) && refetchInterval, + ...(exactTry ? { retry: false } : {}), + staleTime: 0, + }); + const history = useTaskInstanceServiceGetTaskInstanceTryDetails( + { ...params, taskTryNumber: Number(tryParameter) }, + undefined, + { enabled: exactTry }, + ); + const historical = + exactTry && + history.data !== undefined && + (live.data?.id !== history.data.id || live.data.try_number !== history.data.try_number); Review Comment: While `live` is still loading, `live.data?.id !== history.data.id` is already true, so `historical` turns on as soon as the try-details query wins the race. Every regional link from `getTaskInstanceLink` carries `try_number`, so a fresh load of a live loop or mapped TI can hit the `Navigate replace` to `/logs` in `TaskInstance.tsx` and switch off the required-actions and HITL queries before the live row arrives. Could this count as historical only once `live` has settled (`live.isSuccess || live.isError`), with a hook test where history resolves first? ########## airflow-core/src/airflow/ui/src/pages/TaskInstance/TaskInstance.tsx: ########## @@ -97,66 +114,75 @@ export const TaskInstance = () => { const refetchInterval = useAutoRefresh({ dagId }); const parsedMapIndex = parseInt(mapIndex, 10); - const { - data: taskInstance, - error, - isLoading, - } = useTaskInstanceServiceGetMappedTaskInstance( - { - dagId, - dagRunId: runId, - mapIndex: parsedMapIndex, - taskId, - }, - undefined, - { - enabled: !isNaN(parsedMapIndex), - refetchInterval: (query) => isStatePending(query.state.data?.state) && refetchInterval, - staleTime: 0, - }, - ); - - const { summariesByRunId } = useGridTiSummariesStream({ dagId, runIds: runId ? [runId] : [] }); + const { summariesByRunId } = useGridTiSummariesStream({ + dagId, + runIds: !historical && runId ? [runId] : [], + }); const gridTISummaries = summariesByRunId.get(runId); const taskInstanceSummary = gridTISummaries?.task_instances.find((ti) => ti.task_id === taskId); const taskCount = Object.entries(taskInstanceSummary?.child_states ?? {}) .map(([_state, count]) => count) .reduce((sum, val) => sum + val, 0); + const scopedTabs = tabs + .filter((tab) => !historical || ["details", logsTabValue].includes(tab.value)) + .map((tab) => ({ + ...tab, + search: tab.search ?? (coordinateSearch.toString() || undefined), + })); const newTabs = - taskInstance && taskInstance.map_index > -1 + taskInstance && taskInstance.map_index > -1 && !isRegional && !historical ? [ - ...tabs.slice(0, 1), + ...scopedTabs.slice(0, 1), { icon: <MdOutlineTask />, label: translate("tabs.mappedTaskInstances_other", { count: Number(taskCount), }), value: "task_instances", }, - ...tabs.slice(1), + ...scopedTabs.slice(1), ] - : tabs; + : scopedTabs; - const { tabs: requiredActionTabs } = useRequiredActionTabs({ dagId, dagRunId: runId, taskId }, newTabs, { - autoRedirect: true, - refetchInterval: isStatePending(taskInstance?.state) && refetchInterval, - }); + const { tabs: requiredActionTabs } = useRequiredActionTabs( + { ...coordinates, dagId, dagRunId: runId, mapIndex: parsedMapIndex, taskId }, + newTabs, + { + autoRedirect: true, + enabled: !historical, + refetchInterval: isStatePending(taskInstance?.state) && refetchInterval, + }, + ); const { tabs: displayTabs } = useHITLReviewTabs({ dagId, dagRunId: runId, taskId }, requiredActionTabs, { + ...coordinates, + enabled: !historical, mapIndex: parsedMapIndex, refetchInterval: isStatePending(taskInstance?.state) && refetchInterval, }); + const taskPath = getTaskInstanceLink({ dagId, dagRunId: runId, mapIndex: parsedMapIndex, taskId }); + + if (historical && ![`${taskPath}/details`, `${taskPath}/logs`, taskPath].includes(location.pathname)) { Review Comment: For a regional TI the Required Actions tab carries the region params, so picking an earlier try in `HITLResponse` sets `try_number` and makes the view historical. `/required_actions` isn't in this allow-list, so the user gets bounced to Logs instead of seeing that try's response. Any HITL operator with retries inside a loop or an `.expand()` hits this; before this change the same click loaded that try's HITL detail. ########## airflow-core/src/airflow/ui/src/components/ActionAccordion/columns.tsx: ########## @@ -27,10 +27,10 @@ import type { DataTableFeatures } from "src/components/DataTable/features"; import type { MetaColumn } from "src/components/DataTable/types"; import { StateBadge } from "src/components/StateBadge"; -// Stable per-row key; dag_run_id keeps the same task distinct across runs (past/future -// expansion), and map_index disambiguates the mapped instances of one task. export const taskInstanceKey = (ti: TaskInstanceResponse): string => - `${ti.dag_run_id}:${ti.task_id}:${ti.map_index}`; + ti.region_id === "00000000-0000-0000-0000-000000000000" + ? `${ti.dag_run_id}:${ti.task_id}:${ti.map_index}` + : ti.id; Review Comment: For regional rows this key is the per-try UUID, and `prepare_db_for_next_try` gives the successor row a new `uuid7()`. The clear dry run keeps polling while anything is pending, so if an excluded TI retries while the dialog is open it comes back as kept under its new id, the stale key keeps `hasExclusions` true, and the new id goes out in `task_instance_ids`, which clears the TI the user unticked. Keying on run, task, map index, region id and region index would survive the retry. ########## airflow-core/src/airflow/ui/src/queries/usePatchTaskInstance.ts: ########## @@ -61,6 +62,7 @@ export const usePatchTaskInstance = ({ const onSuccessFn = async () => { const queryKeys = [ + [useDagRunServiceGetExecutionKey, { dagId, dagRunId }], UseTaskInstanceServiceGetTaskInstanceKeyFn({ dagId, dagRunId, taskId }), Review Comment: `UseTaskInstanceServiceGetTaskInstanceKeyFn` builds a filter with `regionId: undefined` and `regionIndex: undefined`, and TanStack's `partialMatchKey` compares those as values, so it no longer matches the `DagBreadcrumb` query that now carries the coordinates. Marking a loop TI or editing its note leaves the breadcrumb state stale. A literal prefix like the one on line 70 avoids that. ########## airflow-core/src/airflow/ui/src/queries/useClearTaskInstances.ts: ########## @@ -113,6 +115,8 @@ export const useClearTaskInstances = ({ ]; const queryKeys = [ + [useDagRunServiceGetExecutionKey, { dagId, dagRunId }], + [useTaskInstanceServiceGetMappedTaskInstanceKey, { dagId, dagRunId }], Review Comment: These run-scoped keys use the hook's own `dagRunId`, but the `hasExclusions` path in `ClearTaskInstanceDialog` fires one clear per run with `task_instance_ids` and no `task_ids`, so the other runs' Execution, Gantt and mapped-TI caches stay stale (finished runs don't poll). `variables.requestBody.dag_run_id ?? dagRunId` would cover them. With this prefix in place the per-task `taskInstanceKeys` block above still adds nothing and could go. ########## airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearExecutionDialog.tsx: ########## @@ -0,0 +1,169 @@ +/*! + * 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 { useEffect, useState } from "react"; + +import { Button, Stack, Text, Textarea } from "@chakra-ui/react"; +import { useTranslation } from "react-i18next"; + +import type { ClearTaskInstancesBody, ExecutionTaskResponse } from "openapi/requests/types.gen"; + +import { Checkbox, Modal } from "src/system-components"; + +import { ErrorAlert } from "src/components/ErrorAlert"; + +import { + useClearKeepTaskStateDefault, + useClearPreventRunningTaskDefault, + useClearTaskInstanceDefaultOptions, +} from "src/hooks/useUserSettings"; +import { useClearTaskInstances } from "src/queries/useClearTaskInstances"; +import { useClearTaskInstancesDryRun } from "src/queries/useClearTaskInstancesDryRun"; + +type SelectedExecution = Pick< + ExecutionTaskResponse, + "id" | "map_index" | "note" | "region_id" | "region_index" | "task_display_name" | "task_id" +>; + +export const ClearExecutionDialog = ({ + dagId, + executions, + onClose, + open, + runId, +}: { + readonly dagId: string; + readonly executions: Array<SelectedExecution>; + readonly onClose: () => void; + readonly open: boolean; + readonly runId: string; +}) => { + const { t: translate } = useTranslation("dag"); + const [defaultOptions] = useClearTaskInstanceDefaultOptions(); + const [preventRunningDefault] = useClearPreventRunningTaskDefault(); + const [keepTaskStateDefault] = useClearKeepTaskStateDefault(); + const initialNote = (executions.length === 1 ? executions[0]?.note : undefined) ?? ""; + const defaultDownstream = defaultOptions.includes("downstream"); + const [downstream, setDownstream] = useState(defaultDownstream); + const [later, setLater] = useState(true); + const [whole, setWhole] = useState(false); + const [preventRunning, setPreventRunning] = useState(preventRunningDefault); + const [keepTaskState, setKeepTaskState] = useState(keepTaskStateDefault); + const [note, setNote] = useState(initialNote); + + useEffect(() => { + if (!open) { + setDownstream(defaultDownstream); + setLater(true); + setWhole(false); + setPreventRunning(preventRunningDefault); + setKeepTaskState(keepTaskStateDefault); + setNote(initialNote); + } + }, [open, defaultDownstream, preventRunningDefault, keepTaskStateDefault, initialNote]); + const mappedIds = executions + .filter( + (ti) => + ti.map_index >= 0 || + (ti.region_id !== "00000000-0000-0000-0000-000000000000" && ti.region_index === -1), + ) + .map((ti) => ti.id); + const requestBody: ClearTaskInstancesBody = { + dag_run_id: runId, + include_downstream: downstream, + include_later_loop_iterations: later, + keep_task_state: keepTaskState, + only_failed: false, + prevent_running_task: preventRunning, + task_instance_ids: executions.map((ti) => ti.id), + whole_expansion_ids: whole ? mappedIds : [], + }; + const preview = useClearTaskInstancesDryRun({ + dagId, + options: { enabled: open, retry: false }, + requestBody, + }); + const clear = useClearTaskInstances({ dagId, dagRunId: runId, onSuccessConfirm: onClose }); + + return ( + <Modal + footerActions={ + <Button + disabled={preview.isPending || preview.isError} + loading={clear.isPending} + onClick={() => + clear.mutate({ + dagId, + requestBody: { + ...requestBody, + dry_run: false, + note: note === initialNote ? undefined : note || undefined, + }, + }) + } + > + {translate("execution.clearSelected")} + </Button> + } + onOpenChange={(details) => { + if (!details.open) { + onClose(); + } + }} + open={open} + title={translate("execution.clearSelected")} + > + <Stack gap={4}> + <ErrorAlert error={preview.error} /> + <Checkbox checked={downstream} onCheckedChange={(details) => setDownstream(details.checked === true)}> + {translate("execution.clearDownstream")} + </Checkbox> + <Checkbox checked={later} onCheckedChange={(details) => setLater(details.checked === true)}> + {translate("execution.clearLater")} + </Checkbox> + {mappedIds.length > 0 ? ( + <Checkbox checked={whole} onCheckedChange={(details) => setWhole(details.checked === true)}> + {translate("execution.clearWhole")} + </Checkbox> + ) : undefined} + <Checkbox + checked={preventRunning} + onCheckedChange={(details) => setPreventRunning(details.checked === true)} + > + {translate("dags:runAndTaskActions.options.preventRunningTasks")} + </Checkbox> + <Checkbox + checked={keepTaskState} + onCheckedChange={(details) => setKeepTaskState(details.checked === true)} + > + {translate("dags:runAndTaskActions.options.keepTaskState")} + </Checkbox> + {preview.data === undefined ? undefined : ( + <Text>{translate("execution.clearAffected", { count: preview.data.total_entries })}</Text> Review Comment: With "Clear later loop iterations" on by default, a bare count makes it hard to tell what is about to be cleared. `preview.data` is already what `ActionAccordion` takes in the other clear dialog, so could this show the affected list too? ########## airflow-core/src/airflow/ui/src/pages/Run/Execution.tsx: ########## @@ -0,0 +1,292 @@ +/*! + * 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 { useEffect, useState } from "react"; + +import { Box, Button, Heading, HStack, Link, Stack, Text } from "@chakra-ui/react"; +import { useTranslation } from "react-i18next"; +import { Link as RouterLink, useParams, useSearchParams } from "react-router-dom"; + +import { useDagRunServiceGetDagRun, useDagRunServiceGetExecution } from "openapi/queries"; +import type { ExecutionRegionResponse, ExecutionTaskResponse } from "openapi/requests/types.gen"; + +import { Checkbox, Pagination, ProgressBar } from "src/system-components"; + +import { ClearExecutionDialog } from "src/components/Clear/TaskInstance/ClearExecutionDialog"; +import { ErrorAlert } from "src/components/ErrorAlert"; +import { StateBadge } from "src/components/StateBadge"; + +import { SearchParamsKeys } from "src/constants/searchParams"; +import { isStatePending, useAutoRefresh } from "src/utils"; +import { getTaskInstanceLink } from "src/utils/links"; + +const PAGE_SIZE = 100; + +type Group = { + index?: number; + nodeId?: string; + tasks: Array<ExecutionTaskResponse>; +}; + +const groupExecutions = (tasks: Array<ExecutionTaskResponse>, regions: Array<ExecutionRegionResponse>) => { + const byRegion = new Map(regions.map((region) => [region.id, region])); + const groups = new Map<string, Group>(); + + for (const task of tasks) { + const region = byRegion.get(task.region_id); + const mapped = region?.node_id === task.task_id; + const parentId = region?.parent_region_id; + const parent = parentId === undefined || parentId === null ? undefined : byRegion.get(parentId); + const nodeId = mapped ? parent?.node_id : region?.node_id; + const index = + nodeId === undefined + ? undefined + : mapped + ? (region.parent_region_index ?? undefined) + : task.region_index; + const key = JSON.stringify([nodeId, index]); + const group = groups.get(key) ?? { index, nodeId, tasks: [] }; + + group.tasks.push(task); + groups.set(key, group); + } + + return [...groups.entries()].sort( + ([, first], [, second]) => + (first.nodeId ?? "").localeCompare(second.nodeId ?? "") || (first.index ?? -1) - (second.index ?? -1), + ); +}; + +const executionLink = (task: ExecutionTaskResponse) => { + const path = getTaskInstanceLink( + { dagId: task.dag_id, dagRunId: task.dag_run_id, mapIndex: task.map_index, taskId: task.task_id }, + "logs", + ); + const query = new URLSearchParams({ + region_id: task.region_id, + region_index: String(task.region_index), + try_number: String(task.try_number), + }); + + return `${path}?${query}`; +}; + +const TaskRow = ({ + onSelect, + selected, + selectLabel, + task, +}: { + readonly onSelect: () => void; + readonly selected: boolean; + readonly selectLabel: string; + readonly task: ExecutionTaskResponse; +}) => { + const { t: translate } = useTranslation(); + + return ( + <HStack justify="space-between" py={1}> + <Checkbox aria-label={selectLabel} checked={selected} onCheckedChange={onSelect} /> + <Link asChild> + <RouterLink to={executionLink(task)}> + {task.task_display_name} + {task.map_index >= 0 ? ` [${task.map_index}]` : ""} + </RouterLink> + </Link> + <StateBadge state={task.state}>{translate(`common:states.${task.state ?? "none"}`)}</StateBadge> + </HStack> + ); +}; + +const ExecutionView = () => { + const { dagId = "", runId = "" } = useParams(); + const { t: translate } = useTranslation("dag"); + const [searchParams, setSearchParams] = useSearchParams(); + const parsedOffset = Number(searchParams.get(SearchParamsKeys.EXECUTION_OFFSET) ?? 0); + const offset = Number.isInteger(parsedOffset) && parsedOffset >= 0 ? parsedOffset : 0; + const refresh = useAutoRefresh({ dagId }); + const { data: dagRun } = useDagRunServiceGetDagRun({ dagId, dagRunId: runId }, undefined, { + refetchInterval: (query) => isStatePending(query.state.data?.state) && refresh, + }); + const { data, error, isLoading } = useDagRunServiceGetExecution( + { dagId, dagRunId: runId, limit: PAGE_SIZE, offset }, + undefined, + { refetchInterval: isStatePending(dagRun?.state) && refresh }, Review Comment: Gating on the run's state stops polling as soon as the run finishes, but the run and execution queries are on separate timers, so the last fetch can still show the final task as running next to a successful run. Could this also keep polling while any row in `data` is pending? ########## airflow-core/src/airflow/ui/src/pages/Events/Events.tsx: ########## @@ -256,18 +273,20 @@ export const Events = () => { ...dagIdArg, ...eventArg, limit: pagination.pageSize, - mapIndex: mapIndexNumber, + mapIndex: regional ? undefined : mapIndexNumber, offset: pagination.pageIndex * pagination.pageSize, orderBy, + taskInstanceId: regional ? selectedTask?.id : undefined, Review Comment: Regional TIs (so every mapped TI too) are still filtered by the live try's `taskInstanceId` with `tryNumber` dropped, so earlier tries and the mark, clear and note action logs (null `task_instance_id`) are hidden and the try filter does nothing. `AssetEvents.tsx:107` has the same shape: `sourceTaskInstanceId: taskInstance?.id` loses asset events from earlier tries after a clear, while a sentinel TI still lists every try through `sourceMapIndex`. ########## airflow-core/src/airflow/ui/src/layouts/Details/Gantt/utils.ts: ########## @@ -325,13 +312,19 @@ export const getGanttSegmentTo = ({ // Clone the pre-parsed params so mutations don't leak across segments. const searchParams = new URLSearchParams(baseSearchParams); - const maxTryForTask = maxTryByTaskId.get(taskId) ?? 1; - const isOlderTry = tryNumber !== undefined && tryNumber < maxTryForTask; - if (isOlderTry) { - searchParams.set(SearchParamsKeys.TRY_NUMBER, tryNumber.toString()); + if (item.regionId !== undefined && item.regionIndex !== undefined) { + searchParams.set("region_id", item.regionId); + searchParams.set("region_index", item.regionIndex.toString()); } else { + searchParams.delete("region_id"); + searchParams.delete("region_index"); + } + + if (tryNumber === undefined) { searchParams.delete(SearchParamsKeys.TRY_NUMBER); + } else { + searchParams.set(SearchParamsKeys.TRY_NUMBER, tryNumber.toString()); Review Comment: This now pins `try_number` and the region params on every Gantt segment, including the latest try and tasks in the sentinel region. Clicking a running non-loop task gives `?region_id=00000000-...®ion_index=-1&try_number=1`, which `isExactTryView` treats as an exact-try view, so the page stays on try 1 once the task retries. Before this change `try_number` was only set for tries older than the task's latest one. Could that come back (keyed per task and region), and skip the region params for the sentinel? ########## airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearTaskInstanceDialog.tsx: ########## @@ -399,4 +400,27 @@ const ClearTaskInstanceDialog = (props: Props) => { ); }; -export default ClearTaskInstanceDialog; +const ScopedClearTaskInstanceDialog = (props: Props) => { + const { allMapped } = props; + + if (allMapped) { + return <ClearTaskInstanceDialog {...props} />; + } + const { onClose, open, taskInstance } = props; + + if (taskInstance.region_id !== "00000000-0000-0000-0000-000000000000") { Review Comment: `expand_mapped_task` still calls `DynamicRegion.get_or_create` for every expansion in the sentinel region, so a plain `.expand()` TI with no loop has a real region and still lands in `ClearExecutionDialog`: no upstream, past/future, only-failed or run-on-latest options, no affected list, and a "Clear later loop iterations" checkbox that means nothing there. The same predicate is unchanged in `MarkTaskInstanceAsDialog.tsx:49`, `BulkMarkTaskInstancesAsButton.tsx:56` and `TaskInstance.tsx:58`, which still hides the Mapped Task Instances tab. That changes every Dag using dynamic mapping, so I'd like this keyed on loop membership before it goes in. ########## airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearExecutionDialog.tsx: ########## @@ -0,0 +1,169 @@ +/*! + * 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 { useEffect, useState } from "react"; + +import { Button, Stack, Text, Textarea } from "@chakra-ui/react"; +import { useTranslation } from "react-i18next"; + +import type { ClearTaskInstancesBody, ExecutionTaskResponse } from "openapi/requests/types.gen"; + +import { Checkbox, Modal } from "src/system-components"; + +import { ErrorAlert } from "src/components/ErrorAlert"; + +import { + useClearKeepTaskStateDefault, + useClearPreventRunningTaskDefault, + useClearTaskInstanceDefaultOptions, +} from "src/hooks/useUserSettings"; +import { useClearTaskInstances } from "src/queries/useClearTaskInstances"; +import { useClearTaskInstancesDryRun } from "src/queries/useClearTaskInstancesDryRun"; + +type SelectedExecution = Pick< + ExecutionTaskResponse, + "id" | "map_index" | "note" | "region_id" | "region_index" | "task_display_name" | "task_id" +>; + +export const ClearExecutionDialog = ({ + dagId, + executions, + onClose, + open, + runId, +}: { + readonly dagId: string; + readonly executions: Array<SelectedExecution>; + readonly onClose: () => void; + readonly open: boolean; + readonly runId: string; +}) => { + const { t: translate } = useTranslation("dag"); + const [defaultOptions] = useClearTaskInstanceDefaultOptions(); + const [preventRunningDefault] = useClearPreventRunningTaskDefault(); + const [keepTaskStateDefault] = useClearKeepTaskStateDefault(); + const initialNote = (executions.length === 1 ? executions[0]?.note : undefined) ?? ""; + const defaultDownstream = defaultOptions.includes("downstream"); + const [downstream, setDownstream] = useState(defaultDownstream); + const [later, setLater] = useState(true); + const [whole, setWhole] = useState(false); + const [preventRunning, setPreventRunning] = useState(preventRunningDefault); + const [keepTaskState, setKeepTaskState] = useState(keepTaskStateDefault); + const [note, setNote] = useState(initialNote); + + useEffect(() => { + if (!open) { + setDownstream(defaultDownstream); + setLater(true); + setWhole(false); + setPreventRunning(preventRunningDefault); + setKeepTaskState(keepTaskStateDefault); + setNote(initialNote); + } + }, [open, defaultDownstream, preventRunningDefault, keepTaskStateDefault, initialNote]); + const mappedIds = executions + .filter( + (ti) => + ti.map_index >= 0 || + (ti.region_id !== "00000000-0000-0000-0000-000000000000" && ti.region_index === -1), + ) + .map((ti) => ti.id); + const requestBody: ClearTaskInstancesBody = { + dag_run_id: runId, + include_downstream: downstream, + include_later_loop_iterations: later, Review Comment: There's no `run_on_latest_version` in this request, so the backend falls back to the Dag's `rerun_with_latest_version` and then the config, which defaults to `False`. `ClearTaskInstanceDialog` auto-checks run on latest version when the versions differ, and every non-sentinel TI now routes here, so clearing a mapped or loop TI after deploying a fix reruns it on the old bundle with no way to choose. Could this reuse `getRunOnLatestVersionState` and `useRerunWithLatestVersion` the way the other dialog does? ########## airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearExecutionDialog.test.tsx: ########## @@ -0,0 +1,280 @@ +/*! + * 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/vitest"; +import { fireEvent, render, screen, waitFor } from "@testing-library/react"; +import { afterEach, beforeAll, expect, it, vi } from "vitest"; + +import { TaskInstanceService } from "openapi/requests"; + +import { + CLEAR_KEEP_TASK_STATE_KEY, + CLEAR_PREVENT_RUNNING_TASK_KEY, + CLEAR_TASK_INSTANCE_DEFAULT_OPTIONS_KEY, +} from "src/constants/localStorage"; +import i18n from "src/i18n/config"; +import { Wrapper } from "src/utils/Wrapper"; + +import dagTranslations from "../../../../public/i18n/locales/en/dag.json"; +import dagsTranslations from "../../../../public/i18n/locales/en/dags.json"; +import { ClearExecutionDialog } from "./ClearExecutionDialog"; + +beforeAll(() => { + i18n.addResourceBundle("en", "dag", dagTranslations, true, true); + i18n.addResourceBundle("en", "dags", dagsTranslations, true, true); +}); +afterEach(() => { + vi.restoreAllMocks(); + localStorage.clear(); +}); + +const execution = { + id: "selected-uuid", + map_index: -1, + note: "existing note", + region_id: "11111111-1111-1111-1111-111111111111", + region_index: 1, + task_display_name: "work", + task_id: "body.work", +}; + +const findClearRequest = (clear: { mock: { calls: Array<Array<unknown>> } }) => + (clear.mock.calls as Array<[{ requestBody: { dry_run?: boolean; note?: string | null } }]>) + .map(([request]) => request) + .find((request) => request.requestBody.dry_run === false); + +it("keeps stale UUID scope visible when the clear preview rejects an archived execution", async () => { + const clear = vi + .spyOn(TaskInstanceService, "postClearTaskInstances") + .mockRejectedValue({ body: { detail: "Selected execution was archived" }, status: 409 }); + const onClose = vi.fn(); + + render( + <ClearExecutionDialog + dagId="dag" + executions={[ + { + id: "archived-uuid", + map_index: -1, + region_id: "11111111-1111-1111-1111-111111111111", + region_index: 3, + task_display_name: "work", + task_id: "body.work", + }, + ]} + onClose={onClose} + open + runId="run" + />, + { wrapper: Wrapper }, + ); + expect(await screen.findByText("Selected execution was archived")).toBeVisible(); + expect(screen.getByRole("button", { name: "Clear selected executions" })).toBeDisabled(); + expect(clear).toHaveBeenCalledTimes(1); + expect(clear.mock.lastCall?.[0].requestBody).toMatchObject({ + dry_run: true, + task_instance_ids: ["archived-uuid"], + }); + expect(onClose).not.toHaveBeenCalled(); +}); + +it("previews and clears UUID seeds with separate whole-expansion and later-loop intent", async () => { + const clear = vi + .spyOn(TaskInstanceService, "postClearTaskInstances") + .mockResolvedValue({ task_instances: [], total_entries: 0 }); + + render( + <ClearExecutionDialog + dagId="dag" + executions={[ + { + id: "selected-uuid", + map_index: 1, + region_id: "11111111-1111-1111-1111-111111111111", + region_index: 1, + task_display_name: "mapped", + task_id: "body.mapped", + }, + ]} + onClose={vi.fn()} + open + runId="run" + />, + { wrapper: Wrapper }, + ); + + await waitFor(() => + expect(clear.mock.lastCall?.[0]).toMatchObject({ + dagId: "dag", + requestBody: { + dag_run_id: "run", + dry_run: true, + include_downstream: true, + include_later_loop_iterations: true, + only_failed: false, + task_instance_ids: ["selected-uuid"], + whole_expansion_ids: [], + }, + }), + ); + fireEvent.click(screen.getByLabelText("Clear whole mapped expansions")); + fireEvent.click(screen.getByLabelText("Clear later loop iterations")); + await waitFor(() => + expect(clear.mock.lastCall?.[0]).toMatchObject({ + dagId: "dag", + requestBody: { + dry_run: true, + include_later_loop_iterations: false, + whole_expansion_ids: ["selected-uuid"], + }, + }), + ); + await waitFor(() => + expect(screen.getByRole("button", { name: "Clear selected executions" })).toBeEnabled(), + ); + fireEvent.click(screen.getByRole("button", { name: "Clear selected executions" })); + await waitFor(() => + expect( + clear.mock.calls.map(([request]) => request).find((request) => request.requestBody.dry_run === false), + ).toMatchObject({ + dagId: "dag", + requestBody: { + dry_run: false, + include_downstream: true, + include_later_loop_iterations: false, + task_instance_ids: ["selected-uuid"], + whole_expansion_ids: ["selected-uuid"], + }, + }), + ); +}); + +it("sends the saved clear defaults in the dry run and the real request", async () => { + localStorage.setItem(CLEAR_PREVENT_RUNNING_TASK_KEY, JSON.stringify(true)); + localStorage.setItem(CLEAR_KEEP_TASK_STATE_KEY, JSON.stringify(true)); + localStorage.setItem(CLEAR_TASK_INSTANCE_DEFAULT_OPTIONS_KEY, JSON.stringify([])); + const clear = vi + .spyOn(TaskInstanceService, "postClearTaskInstances") + .mockResolvedValue({ task_instances: [], total_entries: 1 }); + + render(<ClearExecutionDialog dagId="dag" executions={[execution]} onClose={vi.fn()} open runId="run" />, { + wrapper: Wrapper, + }); + + await waitFor(() => + expect(clear.mock.lastCall?.[0].requestBody).toMatchObject({ + dry_run: true, + include_downstream: false, + keep_task_state: true, + prevent_running_task: true, + }), + ); + expect(await screen.findByText("1 execution will be cleared.")).toBeVisible(); + expect(screen.getByLabelText("Prevent rerun if task is running")).toBeChecked(); + expect(screen.getByLabelText("Keep task state and resume")).toBeChecked(); + fireEvent.click(screen.getByLabelText("Keep task state and resume")); + fireEvent.click(screen.getByLabelText("Prevent rerun if task is running")); + await waitFor(() => + expect(clear.mock.lastCall?.[0].requestBody).toMatchObject({ + dry_run: true, + keep_task_state: false, + prevent_running_task: false, + }), + ); + await waitFor(() => + expect(screen.getByRole("button", { name: "Clear selected executions" })).toBeEnabled(), + ); + fireEvent.click(screen.getByRole("button", { name: "Clear selected executions" })); + await waitFor(() => + expect(findClearRequest(clear)?.requestBody).toMatchObject({ + keep_task_state: false, + prevent_running_task: false, + }), + ); +}); + +it("hides the affected count while the dry run is pending", async () => { + vi.spyOn(TaskInstanceService, "postClearTaskInstances").mockReturnValue( + new Promise(() => undefined) as ReturnType<typeof TaskInstanceService.postClearTaskInstances>, + ); + + render(<ClearExecutionDialog dagId="dag" executions={[execution]} onClose={vi.fn()} open runId="run" />, { + wrapper: Wrapper, + }); + + expect(await screen.findByLabelText("Clear downstream tasks")).toBeVisible(); + expect(screen.queryByText(/executions? will be cleared/u)).not.toBeInTheDocument(); +}); + +it("seeds the note from the execution and only sends it when edited", async () => { + const clear = vi + .spyOn(TaskInstanceService, "postClearTaskInstances") + .mockResolvedValue({ task_instances: [], total_entries: 1 }); + + render(<ClearExecutionDialog dagId="dag" executions={[execution]} onClose={vi.fn()} open runId="run" />, { + wrapper: Wrapper, + }); + + const note = await screen.findByLabelText("Reason for clearing (optional)"); + + expect(note).toHaveValue("existing note"); + await waitFor(() => + expect(screen.getByRole("button", { name: "Clear selected executions" })).toBeEnabled(), + ); + fireEvent.click(screen.getByRole("button", { name: "Clear selected executions" })); + await waitFor(() => expect(findClearRequest(clear)).toBeDefined()); + expect(findClearRequest(clear)?.requestBody.note).toBeUndefined(); Review Comment: Neither note test types a new note and submits it, so a change that never sends `note` would keep both green. A case that edits the textarea and asserts the clear request carries the new text would pin it. ########## airflow-core/src/airflow/ui/src/hooks/useTaskInstanceCoordinates.ts: ########## @@ -0,0 +1,31 @@ +/*! + * 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 { useSearchParams } from "react-router-dom"; + +import { SearchParamsKeys } from "src/constants/searchParams"; + +export const useTaskInstanceCoordinates = () => { + const [searchParams] = useSearchParams(); + const regionIndex = searchParams.get(SearchParamsKeys.REGION_INDEX); + + return { Review Comment: The breadcrumb, plugin and version lookups pick up the coordinates now, but the row delete still doesn't: `DeleteTaskInstanceButton` and `useDeleteTaskInstance` send only `mapIndex`, though `DeleteTaskInstanceData` accepts `regionId` and `regionIndex`. Deleting a loop pass from the Task Instances list goes out with no region, and the scope resolver rejects it with "A loop producer requires an explicit scope...". ########## airflow-core/src/airflow/ui/src/utils/links.ts: ########## @@ -39,9 +39,20 @@ export const getTaskInstanceLink = ( const tabPath = tab === undefined ? "" : `/${tab}`; if ("dag_id" in tiOrParams) { - return `/dags/${tiOrParams.dag_id}/runs/${tiOrParams.dag_run_id}/tasks/${tiOrParams.task_id}${ + const path = `/dags/${tiOrParams.dag_id}/runs/${tiOrParams.dag_run_id}/tasks/${tiOrParams.task_id}${ tiOrParams.map_index >= 0 ? `/mapped/${tiOrParams.map_index}` : "" }${tabPath}`; + + if (tiOrParams.region_id === "00000000-0000-0000-0000-000000000000") { + return path; + } + const query = new URLSearchParams({ + region_id: tiOrParams.region_id, + region_index: String(tiOrParams.region_index), + try_number: String(tiOrParams.try_number), Review Comment: `getTaskInstanceLink` still puts `try_number` on every regional link, so a running mapped TI opened from the Task Instances list is in exact-try mode and flips to the historical view once it retries or gets cleared. After the flip the header shows the try as it was when the page loaded (often still running), since the try-details query doesn't poll, and `TaskTrySelect` stops polling too, so newer tries never show up in the picker. `executionLink` in `Execution.tsx:83` pins `try_number` the same way, sentinel tasks included. ########## airflow-core/src/airflow/ui/src/pages/Run/Execution.tsx: ########## @@ -0,0 +1,292 @@ +/*! + * 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 { useEffect, useState } from "react"; + +import { Box, Button, Heading, HStack, Link, Stack, Text } from "@chakra-ui/react"; +import { useTranslation } from "react-i18next"; +import { Link as RouterLink, useParams, useSearchParams } from "react-router-dom"; + +import { useDagRunServiceGetDagRun, useDagRunServiceGetExecution } from "openapi/queries"; +import type { ExecutionRegionResponse, ExecutionTaskResponse } from "openapi/requests/types.gen"; + +import { Checkbox, Pagination, ProgressBar } from "src/system-components"; + +import { ClearExecutionDialog } from "src/components/Clear/TaskInstance/ClearExecutionDialog"; +import { ErrorAlert } from "src/components/ErrorAlert"; +import { StateBadge } from "src/components/StateBadge"; + +import { SearchParamsKeys } from "src/constants/searchParams"; +import { isStatePending, useAutoRefresh } from "src/utils"; +import { getTaskInstanceLink } from "src/utils/links"; + +const PAGE_SIZE = 100; + +type Group = { + index?: number; + nodeId?: string; + tasks: Array<ExecutionTaskResponse>; +}; + +const groupExecutions = (tasks: Array<ExecutionTaskResponse>, regions: Array<ExecutionRegionResponse>) => { + const byRegion = new Map(regions.map((region) => [region.id, region])); + const groups = new Map<string, Group>(); + + for (const task of tasks) { + const region = byRegion.get(task.region_id); + const mapped = region?.node_id === task.task_id; + const parentId = region?.parent_region_id; + const parent = parentId === undefined || parentId === null ? undefined : byRegion.get(parentId); + const nodeId = mapped ? parent?.node_id : region?.node_id; + const index = + nodeId === undefined + ? undefined + : mapped + ? (region.parent_region_index ?? undefined) + : task.region_index; + const key = JSON.stringify([nodeId, index]); + const group = groups.get(key) ?? { index, nodeId, tasks: [] }; + + group.tasks.push(task); + groups.set(key, group); + } + + return [...groups.entries()].sort( + ([, first], [, second]) => + (first.nodeId ?? "").localeCompare(second.nodeId ?? "") || (first.index ?? -1) - (second.index ?? -1), + ); +}; + +const executionLink = (task: ExecutionTaskResponse) => { + const path = getTaskInstanceLink( + { dagId: task.dag_id, dagRunId: task.dag_run_id, mapIndex: task.map_index, taskId: task.task_id }, + "logs", + ); + const query = new URLSearchParams({ + region_id: task.region_id, + region_index: String(task.region_index), + try_number: String(task.try_number), + }); + + return `${path}?${query}`; +}; + +const TaskRow = ({ + onSelect, + selected, + selectLabel, + task, +}: { + readonly onSelect: () => void; + readonly selected: boolean; + readonly selectLabel: string; + readonly task: ExecutionTaskResponse; +}) => { + const { t: translate } = useTranslation(); + + return ( + <HStack justify="space-between" py={1}> + <Checkbox aria-label={selectLabel} checked={selected} onCheckedChange={onSelect} /> + <Link asChild> + <RouterLink to={executionLink(task)}> + {task.task_display_name} + {task.map_index >= 0 ? ` [${task.map_index}]` : ""} + </RouterLink> + </Link> + <StateBadge state={task.state}>{translate(`common:states.${task.state ?? "none"}`)}</StateBadge> + </HStack> + ); +}; + +const ExecutionView = () => { + const { dagId = "", runId = "" } = useParams(); + const { t: translate } = useTranslation("dag"); + const [searchParams, setSearchParams] = useSearchParams(); + const parsedOffset = Number(searchParams.get(SearchParamsKeys.EXECUTION_OFFSET) ?? 0); + const offset = Number.isInteger(parsedOffset) && parsedOffset >= 0 ? parsedOffset : 0; + const refresh = useAutoRefresh({ dagId }); + const { data: dagRun } = useDagRunServiceGetDagRun({ dagId, dagRunId: runId }, undefined, { + refetchInterval: (query) => isStatePending(query.state.data?.state) && refresh, + }); + const { data, error, isLoading } = useDagRunServiceGetExecution( + { dagId, dagRunId: runId, limit: PAGE_SIZE, offset }, + undefined, + { refetchInterval: isStatePending(dagRun?.state) && refresh }, + ); + const totalEntries = data?.total_entries ?? 0; + + useEffect(() => { + if (data !== undefined && offset > 0 && offset >= totalEntries) { + setSearchParams( + (previous) => { + const updated = new URLSearchParams(previous); + const lastOffset = Math.max(0, Math.ceil(totalEntries / PAGE_SIZE) - 1) * PAGE_SIZE; + + if (lastOffset === 0) { + updated.delete(SearchParamsKeys.EXECUTION_OFFSET); + } else { + updated.set(SearchParamsKeys.EXECUTION_OFFSET, String(lastOffset)); + } + + return updated; + }, + { replace: true }, + ); + } + }, [data, offset, setSearchParams, totalEntries]); + const [expanded, setExpanded] = useState(new Set<string>()); + const [selected, setSelected] = useState(new Map<string, ExecutionTaskResponse>()); + const [clearing, setClearing] = useState(false); + const row = (task: ExecutionTaskResponse) => ( + <TaskRow + key={task.id} + onSelect={() => + setSelected((previous) => { + const updated = new Map(previous); + + if (updated.has(task.id)) { + updated.delete(task.id); + } else { + updated.set(task.id, task); + } + + return updated; + }) + } + selected={selected.has(task.id)} + selectLabel={translate("execution.select", { task: task.task_display_name })} + task={task} + /> + ); + const groups = groupExecutions(data?.task_instances ?? [], data?.regions ?? []); + const page = (next: number) => + setSearchParams((previous) => { + const updated = new URLSearchParams(previous); + + updated.set(SearchParamsKeys.EXECUTION_OFFSET, String((next - 1) * PAGE_SIZE)); + + return updated; + }); + + return ( + <Stack gap={4} p={4}> + <Heading size="lg">{translate("execution.title")}</Heading> + <Button alignSelf="start" disabled={selected.size === 0} onClick={() => setClearing(true)}> + {translate("execution.clearSelected")} + </Button> + {clearing ? ( + <ClearExecutionDialog + dagId={dagId} + executions={[...selected.values()]} + onClose={() => { + setClearing(false); + setSelected(new Map()); Review Comment: Cancel still throws the selection away, since `onClose` resets `selected` whether or not anything was cleared. The entries are also never checked against the latest page, so a selected task that retried before the user hits Clear still sends its archived UUID and the dry run rejects it. Resetting only after a successful clear and dropping ids that aren't on the current page would handle both. -- 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]
