kaxil commented on code in PR #74352: URL: https://github.com/apache/airflow/pull/74352#discussion_r4198280973
########## airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearExecutionDialog.tsx: ########## @@ -0,0 +1,121 @@ +/*! + * 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 { 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 { useClearTaskInstances } from "src/queries/useClearTaskInstances"; +import { useClearTaskInstancesDryRun } from "src/queries/useClearTaskInstancesDryRun"; + +type SelectedExecution = Pick< + ExecutionTaskResponse, + "id" | "map_index" | "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 [downstream, setDownstream] = useState(true); + const [later, setLater] = useState(true); + const [whole, setWhole] = useState(false); + const [note, setNote] = useState<string>(); + 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, + only_failed: false, + task_instance_ids: executions.map((ti) => ti.id), + whole_expansion_ids: whole ? mappedIds : [], + }; Review Comment: This body never sets `prevent_running_task` or `keep_task_state`, so the API defaults (both false) apply, and the user's saved clear defaults are ignored (`useClearPreventRunningTaskDefault` is true out of the box). Every mapped TI and every loop pass now clears through this dialog, so pressing Clear on a running mapped task restarts it instead of hitting the 409 that `ClearTaskInstanceDialog` guards against today, and its task state store is discarded too. Could we read `useClearPreventRunningTaskDefault` and `useClearKeepTaskStateDefault` here, show the two checkboxes, and send both values in the dry run and the real request? `downstream` could also start from `useClearTaskInstanceDefaultOptions` instead of a hardcoded `true`. ########## 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: Every mapped expansion gets its own region now, even with no loop around it, so `region_id !== sentinel` is true for every plain `.expand()` task instance in a new run, not only for loop passes. Clearing a mapped TI in a Dag with no loops therefore gets the reduced dialog: no upstream, past/future, only-failed or run-on-latest options and no affected-task list, plus a "Clear later loop iterations" checkbox that means nothing there. The same predicate greys out Past/Future in `MarkTaskInstanceAsDialog` (line 55) and `BulkMarkTaskInstancesAsButton` (line 56), and hides the Mapped Task Instances tab in `TaskInstance.tsx` (line 136), although the API resolves a non-loop mapped TI by `map_index` alone and accepts past/future for it. Could this key on loop membership instead (the `loop_iterations` field added later in the stack, or some other loop signal from the API) and keep the existing path for mapped TIs outside a loop? ########## airflow-core/src/airflow/ui/src/components/ActionAccordion/columns.tsx: ########## @@ -26,10 +26,10 @@ import { Checkbox } from "src/system-components"; 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: Keying regional rows by `ti.id` lets the user untick one loop pass in the accordion, but the exclusion path in `ClearTaskInstanceDialog` (`checkedTaskIds` and the per-run `ids.push(...)`) still turns the kept rows back into `task_id` or `[task_id, map_index]`. A loop pass has public `map_index` -1, so it goes out as the bare task id, and the loop-aware clear route then matches every current pass of that task. For example, clear an upstream sentinel task with downstream on, untick `body.work` iteration 0 and keep iteration 2: the request carries `task_ids: ["body.work"]` and iteration 0 is cleared anyway. Could the `hasExclusions` branch send `task_instance_ids: kept.map((ti) => ti.id)` per run instead? The clear body already accepts that with a single `dag_run_id`. ########## 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: Pinning `try_number` on the generic TI link puts the page in exact-try mode, so it stops following the live attempt. Opening a running mapped TI from the Task Instances list gives `?region_id=R®ion_index=0&try_number=1`. Once it retries, or the user clears it from that page's header, the live row has a new UUID and `useTaskInstanceView` flips `historical` to true: the page redirects to logs, hides most tabs and swaps in the read-only header. `TaskTrySelect` also treats any `try_number` plus `region_id` URL as history and stops polling for new tries. Should generic links carry only `region_id`/`region_index`, keeping `try_number` for explicit attempt links like Gantt segments? `executionLink` in `Execution.tsx` does the same, for sentinel tasks too. ########## airflow-core/src/airflow/ui/src/hooks/useTaskInstanceView.ts: ########## @@ -0,0 +1,63 @@ +/*! + * 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 useTaskInstanceView = () => { + const { dagId = "", mapIndex = "-1", runId = "", taskId = "" } = useParams(); + const [searchParams] = useSearchParams(); + const coordinates = useTaskInstanceCoordinates(); + const tryParameter = searchParams.get(SearchParamsKeys.TRY_NUMBER); + const exactTry = + tryParameter !== null && coordinates.regionId !== undefined && coordinates.regionIndex !== undefined; + 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, + retry: !exactTry && undefined, + 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); + + return { + error: exactTry ? (history.error ?? live.error) : live.error, Review Comment: `history.error` is `null` (not `undefined`) once the try lookup succeeds, so `??` falls through to `live.error`. When the live coordinate is gone, for example an attempt archived by a clear with no successor at that region and index, `getMappedTaskInstance` 404s: the page shows the retained try correctly, but `DetailsLayout` also shows the red error badge and `Logs` renders the 404 through `error ?? logError`, which hides any real log error. Maybe surface `live.error` only when `history.data` is undefined? The "opens an exact retained try without a live coordinate" test mocks `DetailsLayout` without `error`, so it can't see this. ########## airflow-core/src/airflow/ui/src/pages/Run/Execution.tsx: ########## @@ -0,0 +1,257 @@ +/*! + * 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 { 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 { useDagRunServiceGetExecution } from "openapi/queries"; +import type { ExecutionRegionResponse, ExecutionTaskResponse } from "openapi/requests/types.gen"; + +import { Checkbox, 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 { 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; +}) => ( + <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}>{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("execution_offset") ?? 0); + const offset = Number.isInteger(parsedOffset) && parsedOffset >= 0 ? parsedOffset : 0; + const refresh = useAutoRefresh({ dagId }); + const { data, error, isLoading } = useDagRunServiceGetExecution( + { dagId, dagRunId: runId, limit: PAGE_SIZE, offset }, + undefined, + { refetchInterval: refresh }, Review Comment: This polls every `auto_refresh_interval` for as long as the tab is open, finished runs included, because `useAutoRefresh` only checks whether the Dag is paused. The sibling run tabs gate on pending state (`Run/Details.tsx:50`, `Run/AssetEvents.tsx:38`). Could this do the same, for example `refetchInterval: () => isStatePending(dagRun?.state) && refresh` with the run from `useDagRunServiceGetDagRun`? ########## airflow-core/src/airflow/ui/src/queries/useLogs.tsx: ########## @@ -129,7 +131,15 @@ const parseLogs = ({ let parsedLines; const sources: Array<string> = []; - const logLink = taskInstance ? `${getTaskInstanceLink(taskInstance, "logs")}?try_number=${tryNumber}` : ""; + const logSearch = new URLSearchParams({ try_number: String(tryNumber) }); + + if (taskInstance?.region_id !== undefined) { Review Comment: `region_id` is a required string, so this check is always true and every task, sentinel ones included, gets `region_id=00000000-...®ion_index=-1` on its log line links. For an ordinary task on try 3, viewing try 1 and clicking a line number now lands on `?try_number=1®ion_id=<sentinel>®ion_index=-1#5`, which `useTaskInstanceView` treats as an exact retained try, so the page swaps to the history header and hides tabs. Before, the same click kept the normal page. Could this skip the sentinel the way `getTaskInstanceLink` does? ########## airflow-core/src/airflow/ui/src/pages/Events/Events.tsx: ########## @@ -255,18 +272,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: Filtering by the live attempt's `taskInstanceId` drops two kinds of rows: everything from earlier tries (each try has its own UUID now), and API action logs such as mark state, clear and note edits, which `api_fastapi/logging/decorators.py` writes with `task_instance=None`, so their `task_instance_id` is null. Every mapped TI opened from the Task Instances list carries a region, so marking a mapped TI failed and then opening its Audit Log tab doesn't show the mark event, while the same TI opened from the grid does. The user's `try_number` filter is also ignored in this mode. Could this keep the dag/run/task filters and only narrow further when the rows can actually be told apart? ########## airflow-core/src/airflow/ui/src/pages/TaskInstance/TaskInstance.tsx: ########## @@ -97,66 +116,77 @@ 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 (isLoading) { + return <ProgressBar size="xs" />; + } Review Comment: This early return swaps the whole `DetailsLayout` (grid, Gantt, tabs) for a bare progress bar whenever the TI query has no cached data. `TaskInstance` stays mounted across `:taskId`/`:mapIndex` changes, so clicking an uncached cell in the grid unmounts and remounts the grid and loses its scroll position, on ordinary cell-to-cell navigation and not only on first load. `DetailsLayout` already shows its own bar from `isLoading`, and that prop is now always false at line 179. Could we drop the early return? `historical` stays false until the try lookup resolves, so the `Navigate` below doesn't seem to need it. ########## airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearExecutionDialog.tsx: ########## @@ -0,0 +1,121 @@ +/*! + * 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 { 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 { useClearTaskInstances } from "src/queries/useClearTaskInstances"; +import { useClearTaskInstancesDryRun } from "src/queries/useClearTaskInstancesDryRun"; + +type SelectedExecution = Pick< + ExecutionTaskResponse, + "id" | "map_index" | "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 [downstream, setDownstream] = useState(true); + const [later, setLater] = useState(true); + const [whole, setWhole] = useState(false); + const [note, setNote] = useState<string>(); + 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, + only_failed: false, + 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 } })} Review Comment: From the TI header this dialog stays mounted and only `open` toggles (`Modal` unmounts just its content), so `note`, `downstream`, `later` and `whole` survive a close and a successful clear, and the next clear re-sends the old note. If the user types a note and erases it, `note: ""` goes out and the route writes the empty string over every cleared TI's note (`body.note is not None`). The textarea also starts empty instead of with the TI's current note like the old dialog. Maybe reset the state on close and send `note || undefined`? ########## 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 = () => { Review Comment: A few TI lookups on the same page still go by `(task_id, map_index)` only: `DagBreadcrumb` (getTaskInstance), `usePluginAppliesToContext` and `useSelectedVersion` (getMappedTaskInstance), and the row delete in `DeleteTaskInstanceButton`/`useDeleteTaskInstance`, which has no region parameter at all. For a loop pass with no region, `resolve_task_scope` hits "A loop producer requires an explicit scope..." and returns 400, so on a loop pass's page the breadcrumb shows no state, plugin `applies_to` sees no TI, and the version selector falls back to the run's latest Dag version instead of the pass's pinned one. Could those spread `useTaskInstanceCoordinates()` too? ########## airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearExecutionDialog.tsx: ########## @@ -0,0 +1,121 @@ +/*! + * 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 { 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 { useClearTaskInstances } from "src/queries/useClearTaskInstances"; +import { useClearTaskInstancesDryRun } from "src/queries/useClearTaskInstancesDryRun"; + +type SelectedExecution = Pick< + ExecutionTaskResponse, + "id" | "map_index" | "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 [downstream, setDownstream] = useState(true); + const [later, setLater] = useState(true); + const [whole, setWhole] = useState(false); + const [note, setNote] = useState<string>(); + 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, + only_failed: false, + 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 } })} + > + {translate("execution.clearSelected")} + </Button> + } + onOpenChange={(details) => { + if (!details.open) { + onClose(); + } + }} + open={open} + title={translate("execution.clearTitle")} + > + <Stack gap={4}> + <ErrorAlert error={preview.error ?? clear.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} + <Text>{translate("execution.clearAffected", { count: preview.data?.total_entries ?? 0 })}</Text> Review Comment: While the dry run is pending, or after it fails, this reads "0 executions will be cleared", which looks like a real answer on the confirm step of a destructive action. Could the count show only once `preview.data` exists? Also on line 98, `useClearTaskInstances` already toasts `clear.error`, so a failed clear shows up twice. Binding the alert to `preview.error` only would match `ClearTaskInstanceDialog`. ########## airflow-core/src/airflow/ui/src/layouts/Details/Grid/Bar.tsx: ########## @@ -44,7 +44,12 @@ export const Bar = ({ max, onClick, run, showVersionIndicatorMode }: Props) => { const [searchParams] = useSearchParams(); const isSelected = runId === run.run_id; - const search = searchParams.toString(); + const targetSearchParams = new URLSearchParams(searchParams); + + for (const key of ["try_number", "region_id", "region_index"]) { Review Comment: This strip loop is copied into `GridButton`, `GridTI` and `TaskNames` with string literals, and here it's redundant: `Bar` hands the stripped string to `GridButton`, which strips the same keys again. The sentinel UUID literal also appears in 15 non-test files, and a couple guard the required `region_id` with `Boolean(...)` or `ti.region_id &&`. A small module with `SENTINEL_REGION_ID`, an `isRegional(ti)` predicate and a `stripExecutionParams(params)` built on `SearchParamsKeys.REGION_ID`/`REGION_INDEX`/`TRY_NUMBER` would keep these from drifting, and would make a loop-membership check a one-place change. ########## airflow-core/src/airflow/ui/src/pages/Run/Execution.tsx: ########## @@ -0,0 +1,257 @@ +/*! + * 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 { 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 { useDagRunServiceGetExecution } from "openapi/queries"; +import type { ExecutionRegionResponse, ExecutionTaskResponse } from "openapi/requests/types.gen"; + +import { Checkbox, 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 { 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; +}) => ( + <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}>{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("execution_offset") ?? 0); + const offset = Number.isInteger(parsedOffset) && parsedOffset >= 0 ? parsedOffset : 0; + const refresh = useAutoRefresh({ dagId }); + const { data, error, isLoading } = useDagRunServiceGetExecution( + { dagId, dagRunId: runId, limit: PAGE_SIZE, offset }, + undefined, + { refetchInterval: refresh }, + ); + 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("execution_offset", String(next)); + + 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: Resetting `selected` in `onClose` means a plain Cancel throws the selection away. `selected` also holds attempt snapshots across refetches: if a selected task retries before the user presses Clear, its old UUID is archived and the dry run fails with "Selected task execution is not current in this DAG run", leaving the button disabled. Maybe reset only after a successful clear, and drop ids that are no longer on the latest page? ########## airflow-core/src/airflow/ui/src/pages/Run/Execution.tsx: ########## @@ -0,0 +1,257 @@ +/*! + * 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 { 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 { useDagRunServiceGetExecution } from "openapi/queries"; +import type { ExecutionRegionResponse, ExecutionTaskResponse } from "openapi/requests/types.gen"; + +import { Checkbox, 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 { 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; +}) => ( + <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}>{task.state ?? "none"}</StateBadge> Review Comment: This shows the raw state key, and an English "none", in every locale. `HeaderCard` uses ``translate(`common:states.${state}`)`` and `common:states.none` exists. Small thing in `dag.json` as well: `execution.clearTitle` and `execution.clearSelected` are the same string, so one key would do. ########## airflow-core/src/airflow/ui/src/pages/Run/Execution.tsx: ########## @@ -0,0 +1,257 @@ +/*! + * 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 { 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 { useDagRunServiceGetExecution } from "openapi/queries"; +import type { ExecutionRegionResponse, ExecutionTaskResponse } from "openapi/requests/types.gen"; + +import { Checkbox, 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 { 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; +}) => ( + <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}>{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("execution_offset") ?? 0); Review Comment: `execution_offset` isn't in `SearchParamsKeys`, and the Prev/Next buttons plus the "Showing x to y of z" text re-create what `system-components/Pagination` already does. If the offset is past `total_entries` (a shared URL after the live set shrinks) the page reads "Showing 0 to N of N" with Next disabled and nothing resets it, and an empty run shows a zero range with no empty state. ########## airflow-core/src/airflow/ui/src/components/TaskTrySelect.tsx: ########## @@ -26,18 +27,22 @@ import { Select } from "src/system-components"; import { StateBadge } from "src/components/StateBadge"; +import { SearchParamsKeys } from "src/constants/searchParams"; import { isStatePending, useAutoRefresh } from "src/utils"; import TaskInstanceTooltip from "./TaskInstanceTooltip"; type Props = { readonly onSelectTryNumber?: (tryNumber: number) => void; readonly selectedTryNumber?: number; - readonly taskInstance: TaskInstanceResponse; + readonly taskInstance: TaskInstanceHistoryResponse | TaskInstanceResponse; }; export const TaskTrySelect = ({ onSelectTryNumber, selectedTryNumber, taskInstance }: Props) => { const { t: translate } = useTranslation("components"); + const [searchParams] = useSearchParams(); + const inspectHistory = Review Comment: This re-derives exact-try mode with a different rule from `useTaskInstanceView` (no `region_index` check). Could the hook return `exactTry` and the callers pass it in, so the two can't drift? ########## airflow-core/src/airflow/ui/src/hooks/useTaskInstanceView.ts: ########## @@ -0,0 +1,63 @@ +/*! + * 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 useTaskInstanceView = () => { + const { dagId = "", mapIndex = "-1", runId = "", taskId = "" } = useParams(); + const [searchParams] = useSearchParams(); + const coordinates = useTaskInstanceCoordinates(); + const tryParameter = searchParams.get(SearchParamsKeys.TRY_NUMBER); + const exactTry = + tryParameter !== null && coordinates.regionId !== undefined && coordinates.regionIndex !== undefined; + 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, + retry: !exactTry && undefined, Review Comment: `!exactTry && undefined` evaluates to `false` or `undefined`, which reads like a typo. `retry: exactTry ? false : undefined` says what it means. ########## 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: With this run-wide prefix key, TanStack's partial matching already covers every `GetMappedTaskInstance` key in the run, so the per-task `taskInstanceKeys` block above invalidates nothing extra and could go. ########## airflow-core/src/airflow/ui/src/pages/TaskInstance/TaskInstance.test.tsx: ########## @@ -122,6 +125,124 @@ const Location = () => { }; describe("TaskInstance", () => { + it("opens an exact retained try without a live coordinate", async () => { + const queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } } }); + const regionId = "11111111-1111-4111-8111-111111111111"; + + vi.spyOn(TaskInstanceService, "getMappedTaskInstance").mockRejectedValue({ status: 404 }); + const history = vi.spyOn(TaskInstanceService, "getTaskInstanceTryDetails").mockResolvedValue({ + ...buildTaskInstance(TASK_A, "success", 2), + id: "retained-id", + region_id: regionId, + region_index: 3, + }); + + render( + <MemoryRouter + initialEntries={[ + `/dags/${DAG_ID}/runs/${DAG_RUN_ID}/tasks/${TASK_A}/logs?region_id=${regionId}®ion_index=3&try_number=2`, + ]} + > + <Routes> + <Route element={<TaskInstance />} path="/dags/:dagId/runs/:runId/tasks/:taskId"> + <Route element={<div />} path="*" /> + </Route> + </Routes> + </MemoryRouter>, + { wrapper: createWrapper(queryClient) }, + ); + expect(await screen.findByText("retained-id")).toBeVisible(); + expect(history).toHaveBeenCalledWith( + expect.objectContaining({ regionId, regionIndex: 3, taskTryNumber: 2 }), + ); + expect(screen.queryByRole("link", { name: "tabs.storage" })).not.toBeInTheDocument(); + }); + it.each(["logs", "xcom"])( + "opens a retained execution from %s without showing its current replacement or exposing live-only tabs", + async (tab) => { + const queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } } }); + const regionId = "11111111-1111-4111-8111-111111111111"; + + vi.spyOn(TaskInstanceService, "getMappedTaskInstance").mockResolvedValue({ + ...buildTaskInstance(TASK_A, "running", 3), + id: "current-id", + }); + const history = vi.spyOn(TaskInstanceService, "getTaskInstanceTryDetails").mockResolvedValue({ + ...buildTaskInstance(TASK_A, "success", 2), + id: "retained-id", + region_id: regionId, + region_index: 3, + }); + + render( + <MemoryRouter + initialEntries={[ + `/dags/${DAG_ID}/runs/${DAG_RUN_ID}/tasks/${TASK_A}/${tab}?region_id=${regionId}®ion_index=3&try_number=2`, + ]} + > + <Location /> + <Routes> + <Route element={<TaskInstance />} path="/dags/:dagId/runs/:runId/tasks/:taskId"> + <Route element={<div />} path="*" /> + </Route> + </Routes> + </MemoryRouter>, + { wrapper: createWrapper(queryClient) }, + ); + + expect(await screen.findByText("retained-id")).toBeVisible(); + expect(screen.getByTestId("location").textContent).toContain(`/tasks/${TASK_A}/logs?`); + expect(history).toHaveBeenCalledWith( + expect.objectContaining({ regionId, regionIndex: 3, taskTryNumber: 2 }), + ); + expect(screen.queryByText("current-id")).not.toBeInTheDocument(); + expect(screen.queryByRole("link", { name: "tabs.storage" })).not.toBeInTheDocument(); Review Comment: No tab is labelled `tabs.storage` (the state store tab is `tabs.taskStateStore`), so this and the same check on line 158 pass whether or not live-only tabs are hidden. `queryByText("current-id")` on line 198 can't fail either, since the mocked `Header` renders `task_id:state:try_number` rather than the id. Asserting `tabs.taskStateStore`/`tabs.xcom` are absent and `queryByTestId("task-instance-state")` is null would cover the tab filter and the header swap. ########## airflow-core/src/airflow/ui/src/components/MarkAs/TaskInstance/MarkTaskInstanceAsDialog.test.tsx: ########## @@ -83,6 +86,37 @@ const taskInstance: TaskInstanceResponse = { }; describe("MarkTaskInstanceAsDialog", () => { + it("marks an exact regional execution and disables cross-run scope", () => { + const regional = { + ...taskInstance, + map_index: -1, + region_id: "11111111-1111-4111-8111-111111111111", + region_index: 3, + }; + + render(<MarkTaskInstanceAsDialog onClose={vi.fn()} open state="success" taskInstance={regional} />, { + wrapper: Wrapper, + }); + expect(screen.getByRole("button", { name: /past/iu })).toBeDisabled(); + expect(screen.getByRole("button", { name: /future/iu })).toBeDisabled(); + expect(vi.mocked(usePatchTaskInstanceDryRun).mock.lastCall?.[0]).toMatchObject({ + requestBody: { + include_future: false, + include_past: false, Review Comment: `useMarkTaskInstanceDefaultOptions` defaults to `[]`, so `include_past`/`include_future` are false here whether or not the `!isRegional` gate exists. Seeding `MARK_TASK_INSTANCE_DEFAULT_OPTIONS_KEY` with `["past", "future"]` before rendering (and clearing it afterwards) would make these assertions bind to the gate. ########## airflow-core/src/airflow/ui/src/pages/Run/Execution.test.tsx: ########## @@ -0,0 +1,240 @@ +/*! + * 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 { Link, MemoryRouter, Route, Routes } from "react-router-dom"; +import { afterEach, beforeAll, describe, expect, it, vi } from "vitest"; + +import { DagRunService, TaskInstanceService } from "openapi/requests"; +import type { + ExecutionCollectionResponse, + ExecutionRegionResponse, + ExecutionTaskResponse, +} from "openapi/requests/types.gen"; + +import i18n from "src/i18n/config"; +import type * as Utils from "src/utils"; +import { BaseWrapper } from "src/utils/Wrapper"; + +import dagTranslations from "../../../public/i18n/locales/en/dag.json"; +import { Execution } from "./Execution"; + +vi.mock("src/router", () => ({ taskInstanceRoutes: [] })); + +beforeAll(() => { + i18n.addResourceBundle("en", "dag", dagTranslations, true, true); +}); + +vi.mock("src/utils", async (importOriginal) => ({ + ...(await importOriginal<typeof Utils>()), + useAutoRefresh: () => false, +})); + +const ROOT = "11111111-1111-4111-8111-111111111111"; +const FORK = "22222222-2222-4222-8222-222222222222"; +const MAPPING = "33333333-3333-4333-8333-333333333333"; + +const region = (id: string, changes: Partial<ExecutionRegionResponse> = {}): ExecutionRegionResponse => ({ + forked_from_region_id: null, + id, + node_id: "body", + parent_region_id: null, + parent_region_index: null, + resumes_from_index: 0, + ...changes, +}); + +const task = (id: string, changes: Partial<ExecutionTaskResponse> = {}): ExecutionTaskResponse => ({ + dag_id: "dag", + dag_run_id: "run", + dag_version_id: null, + duration: null, + end_date: null, + id, + map_index: -1, + operator: "EmptyOperator", + region_id: ROOT, + region_index: 0, + start_date: null, + state: "success", + task_display_name: id, + task_id: id, + try_number: 1, + ...changes, +}); + +const renderExecution = (search = "") => + render( + <BaseWrapper> + <MemoryRouter initialEntries={[`/dags/dag/runs/run/execution${search}`]}> + <Routes> + <Route element={<Execution />} path="/dags/:dagId/runs/:runId/execution" /> + </Routes> + </MemoryRouter> + </BaseWrapper>, + ); + +afterEach(() => vi.restoreAllMocks()); + +describe("Run Execution", () => { + it("resets selected UUIDs when navigating to another run", async () => { + vi.spyOn(DagRunService, "getExecution").mockResolvedValue({ + regions: [region(ROOT)], + task_instances: [task("work")], + total_entries: 1, + }); + render( + <BaseWrapper> + <MemoryRouter initialEntries={["/dags/dag/runs/first/execution"]}> + <Link to="/dags/dag/runs/second/execution">Second run</Link> + <Routes> + <Route element={<Execution />} path="/dags/:dagId/runs/:runId/execution" /> + </Routes> + </MemoryRouter> + </BaseWrapper>, + ); + fireEvent.click(await screen.findByRole("checkbox", { name: "Select work" })); + await waitFor(() => + expect(screen.getByRole("button", { name: "Clear selected executions" })).toBeEnabled(), + ); + fireEvent.click(screen.getByRole("link", { name: "Second run" })); + expect(screen.getByRole("button", { name: "Clear selected executions" })).toBeDisabled(); + }); + it("clears selected UUIDs separately when public task coordinates collide", async () => { + vi.spyOn(DagRunService, "getExecution").mockResolvedValue({ + regions: [region(ROOT), region(FORK)], + task_instances: [ + task("first", { task_id: "work" }), + task("second", { region_id: FORK, region_index: 1, task_id: "work" }), Review Comment: Two passes of one task share a `task_display_name`, so in practice both checkboxes would be labelled "Select work" (`execution.select` only uses the display name). Giving them different display names here hides that; maybe the label should include the iteration or map index. The test also stops at the dry-run call, so asserting the `dry_run: false` request would cover the clear its name describes. ########## airflow-core/src/airflow/ui/src/pages/TaskInstances/TaskInstances.test.tsx: ########## @@ -116,6 +116,7 @@ const mappedTaskInstance = { dag_id: "example_dag", dag_run_id: "manual__2026-06-07T00:00:00+00:00", map_index: 1, + region_id: "00000000-0000-0000-0000-000000000000", Review Comment: A real mapped TI never has the sentinel region now that each expansion gets its own region, so this fixture takes a path production doesn't: the list would render `.../mapped/1?region_id=...®ion_index=1&try_number=N` and key the row by `ti.id`. Could the fixture use a real region and assert the href users actually get? ########## airflow-core/src/airflow/ui/src/layouts/Details/Gantt/GanttTimeline.tsx: ########## @@ -341,13 +335,27 @@ export const GanttTimeline = ({ // Task groups don't have a try number const touchesNext = - tryNumber !== undefined && segments[segIndex + 1]?.tryNumber === tryNumber; + tryNumber !== undefined && + segments[segIndex + 1]?.taskInstanceId === segment.taskInstanceId && + segments[segIndex + 1]?.tryNumber === tryNumber; const touchesPrev = - tryNumber !== undefined && segments[segIndex - 1]?.tryNumber === tryNumber; + tryNumber !== undefined && + segments[segIndex - 1]?.taskInstanceId === segment.taskInstanceId && + segments[segIndex - 1]?.tryNumber === tryNumber; Review Comment: The new `taskInstanceId` check in `touchesNext`/`touchesPrev` exists for the loop case (adjacent same-try segments from different executions), but no `GanttTimeline` test renders two such segments, so reverting it stays green. Worth one test? -- 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]
