kaxil commented on code in PR #74354: URL: https://github.com/apache/airflow/pull/74354#discussion_r4198311990
########## airflow-core/src/airflow/ui/src/hooks/useLoopOutcomeStats.tsx: ########## @@ -0,0 +1,90 @@ +/*! + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +import type { ReactNode } from "react"; + +import { HStack, Icon, Text } from "@chakra-ui/react"; +import { useTranslation } from "react-i18next"; +import { FiRepeat } from "react-icons/fi"; + +import { finalCriteria, reasonSentence } from "src/pages/GroupTaskInstance/LoopIterations/loopUtils"; + +import { useLoopSummary } from "src/queries/useLoopSummary"; + +export const useLoopOutcomeStats = ({ + dagId, + groupId, + runId, +}: { + dagId: string; + groupId: string; + runId: string; +}): Array<{ label: string; value: ReactNode | string }> => { + const { t: translate } = useTranslation("dag"); + const { data: summary } = useLoopSummary({ dagId, groupId, runId }); + + if (summary === undefined) { + return []; + } + + const comparison = finalCriteria(summary); + + return [ + { + label: translate("loop.outcome"), Review Comment: None of the new `loop.*` keys exist in any locale file, and neither does `taskInstance.loopIterations` in `common`. The diff doesn't touch `public/i18n/locales/`, so every new loop surface renders raw keys: the header shows "loop.outcome: loop.reason.ran_to_cap", the picker shows "loop.filter.label", the overview chart is titled "loop.history.title", and the grid tooltip reads "taskInstance.loopIterations: 3/5". The tests mock `t` as identity, so they can't catch it. Could you add a `loop` block to `en/dag.json` (including every `loop.reason.${status}` value the dynamic lookup can produce, and the interpolation params like `{{index}}`, `{{criteria}}`, `{{max}}`) plus `taskInstance.loopIterations` to `en/common.json`? ########## airflow-core/src/airflow/api_fastapi/core_api/services/public/task_coordinates.py: ########## @@ -49,6 +59,13 @@ class TaskCoordinateView: def __getattr__(self, name: str) -> Any: if name == "map_index": return self.resolver.public_map_index(self.value) + if name == "rendered_map_index" and self.resolver.public_map_index(self.value) < 0: + return None Review Comment: This returns `None` for every TI whose public map index is -1, not just loop rows. A plain unmapped, non-loop task that sets `map_index_template` gets its label rendered and stored by the task runner (`_render_map_index` runs whenever the template is set), and today `GET .../taskInstances/load` returns `"rendered_map_index": "eu-west"` for it. After this change that is `null` on the public list, get, tries, clear/patch and HITL responses, and the TaskInstances table no longer shows it either since the `: original.rendered_map_index` branch was removed. Could this fall back to the stored `_rendered_map_index` instead of `None`, so only the `region_index` fallback is suppressed? The SQL side of the `rendered_map_index` hybrid still falls back to `region_index` too, so `rendered_map_index_pattern=2` will match iteration-2 loop rows whose response label is `null`. ########## airflow-core/src/airflow/ui/src/pages/Task/Overview/LoopHistoryChart.tsx: ########## @@ -0,0 +1,145 @@ +/*! + * 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 { Box, HStack, Heading, Stat, useToken } from "@chakra-ui/react"; +import { BarElement, CategoryScale, Chart as ChartJS, Legend, LinearScale, Tooltip } from "chart.js"; +import annotationPlugin from "chartjs-plugin-annotation"; +import { Bar } from "react-chartjs-2"; +import { useTranslation } from "react-i18next"; +import { useNavigate, useParams } from "react-router-dom"; + +import { useLoopHistory } from "src/queries/useLoopHistory"; +import { getComputedCSSVariableValue } from "src/theme"; +import { median } from "src/utils/median"; + +ChartJS.register(CategoryScale, LinearScale, BarElement, Legend, Tooltip, annotationPlugin); + +export const LoopHistoryChart = ({ groupId }: { readonly groupId: string }) => { + const { dagId = "" } = useParams(); + const { t: translate } = useTranslation("dag"); + const navigate = useNavigate(); + const { data } = useLoopHistory({ dagId, groupId }); + const [successColor, failedColor, mutedColor] = useToken("colors", [ + "success.solid", + "failed.solid", + "fg.muted", + ]); + + const runs = data?.runs ?? []; + + if (runs.length === 0) { + return undefined; + } + + const converged = runs.filter((run) => run.reason === "criteria_met"); + const capHits = runs.filter((run) => run.reason === "not_converged" || run.reason === "cap_reached"); Review Comment: For an `until` loop this is always 0. The backend never emits `not_converged`, and it only sets `cap_reached` when `not group.has_until`. A loop that hits the cap without converging fails its final gate (status `failed`, reason `iteration_failed`), and one that converges on the last allowed pass is `ran_to_cap` with reason `None`. So all 14 bars can sit on the dashed cap line while the stat says "Cap hits 0". Counting `run.iterations_ran === run.max_iterations` would be recorded state rather than an inferred condition result, and it works for both loop kinds. ########## airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/grid.py: ########## @@ -43,3 +46,75 @@ class GridTISummaries(BaseModel): run_id: str dag_id: str task_instances: list[LightGridTaskInstanceSummary] + + +class LoopCriteriaComparison(BaseModel): + """Optional comparison details reported by a loop condition.""" + + field: str + op: str + target: Any = None + actual: Any = None + + +class LoopIterationSummary(BaseModel): + """Execution state for one existing iteration.""" + + index: int + state: TaskInstanceState | None = None + decision: Literal["stop", "continue"] | None = None + is_tail: bool = False + result: Any = None + criteria: LoopCriteriaComparison | None = None Review Comment: Nothing ever sets `decision`, `is_tail`, `result` or `criteria`, and the `criteria_met` / `not_converged` reasons are never produced either (`loop_run_summaries` only builds index/state/start/end and reason in {`cap_reached`, `iteration_failed`, None}). The commit message says summaries report recorded state rather than condition results, so I'd drop these fields, `LoopCriteriaComparison` and the two literals for now, along with the UI branches that read them (`decisionLabel`, `lastRanIteration`/`finalCriteria`, `buildResultTags` and its unused `hasComplex`, the "loop.target" row, the `converged` stat). `exit_task_id` on the summary is also unused by the UI and disagrees with `loop_metadata`, which returns `None` for fixed-count loops. The `criteria is None` / `decision is None` assertions in test_loop_summary.py can't fail as things stand. ########## airflow-core/src/airflow/ui/src/queries/useLoopSummary.ts: ########## @@ -0,0 +1,56 @@ +/*! + * 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 { useDagRunServiceGetDagRun, useGridServiceGetLoopSummary } from "openapi/queries"; + +import { isStatePending, useAutoRefresh } from "src/utils"; + +/** Runtime summary of a looped Task Group; polls while the loop is still running. */ +export const useLoopSummary = ({ + dagId, + groupId, + runId, +}: { + dagId: string; + groupId: string; + runId: string; +}) => { + const refetchInterval = useAutoRefresh({ dagId }); + const [searchParams] = useSearchParams(); + const { data: dagRun } = useDagRunServiceGetDagRun({ dagId, dagRunId: runId }, undefined, { + enabled: Boolean(dagId) && Boolean(runId), + refetchInterval: (query) => isStatePending(query.state.data?.state) && refetchInterval, + }); + + return useGridServiceGetLoopSummary( + { + dagId, + groupId, + loopRegionId: searchParams.get("loop_region_id") ?? undefined, + runId, + }, + undefined, + { + enabled: Boolean(dagId) && Boolean(groupId) && Boolean(runId), + refetchInterval: isStatePending(dagRun?.state) && refetchInterval, Review Comment: When the Dag run turns terminal this interval is cleared on the next render, with no last fetch. If the final gate finishes after the previous summary poll and the run is marked done before the next one, the header keeps "running" and the last iteration's state stays stale (default `staleTime` is 5 minutes, and nothing invalidates the summary on run completion). `useGridTISummaries` hits the same race and re-streams once on the pending to terminal transition; could this do the same (a `wasPending` ref plus `refetch()`)? The docstring on line 25 also says it polls while the loop is running, but it keys off the Dag run. ########## airflow-core/src/airflow/ui/src/pages/GroupTaskInstance/LoopIterations/IterationSelect.tsx: ########## @@ -0,0 +1,183 @@ +/*! + * 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 } from "react"; + +import { createListCollection, HStack, Text, VStack } from "@chakra-ui/react"; +import { useTranslation } from "react-i18next"; +import { useSearchParams } from "react-router-dom"; + +import type { LoopIterationSummary, LoopSummaryResponse } from "openapi/requests/types.gen"; + +import { Select } from "src/system-components"; + +import { StateBadge } from "src/components/StateBadge"; + +import { SearchParamsKeys } from "src/constants/searchParams"; +import { useDurationFormat } from "src/utils"; + +import { buildResultTags, decisionLabel } from "./loopUtils"; + +const ALL = "all"; + +type Props = { + readonly summary: LoopSummaryResponse; +}; + +export const IterationSelect = ({ summary }: Props) => { + const { t: translate } = useTranslation(["dag", "common"]); + const [searchParams, setSearchParams] = useSearchParams(); + const { formatElapsed } = useDurationFormat(); + const selected = searchParams.get(SearchParamsKeys.ITERATION) ?? ALL; + + const latest = summary.iterations.at(-1)?.index; + const selectedExists = summary.iterations.some((iteration) => String(iteration.index) === selected); + + useEffect(() => { + const regionId = searchParams.get(SearchParamsKeys.LOOP_REGION_ID); + + if (regionId !== null && regionId !== summary.loop_region_id) { + return; + } + const next = new URLSearchParams(searchParams); + + if ( + latest !== undefined && + (!searchParams.has(SearchParamsKeys.ITERATION) || (selected !== ALL && !selectedExists)) + ) { + next.set(SearchParamsKeys.ITERATION, String(latest)); + } + if (regionId === null && summary.loop_region_id !== null && summary.loop_region_id !== undefined) { + next.set(SearchParamsKeys.LOOP_REGION_ID, summary.loop_region_id); Review Comment: Grid, graph and gantt links all strip `loop_region_id`, so a normal visit has none. The first summary loads, then this effect writes `loop_region_id` (and `iteration`) into the URL, which changes the `useLoopSummary` query key. With no `placeholderData`, `data` goes back to `undefined`, `GroupTaskInstances` renders the plain `<TaskInstances />` and `Header` drops the outcome stats until the second fetch lands, so the picker and stats flash out and back on every entry, with a couple of extra TI/summary requests. Since the SDK rejects nested loops and loops in mapped groups, and forks collapse to one family, a run only has one invocation here. Could this skip writing `loop_region_id` unless the user picks one, and derive the default iteration in render (`param ?? String(latest)`) rather than via the effect? ########## airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py: ########## @@ -615,6 +613,9 @@ def get_task_instances( "Pass an empty string for the first page, then use ``next_cursor`` from the response. " "When ``cursor`` is provided, ``offset`` is ignored.", ), + loop_id: Annotated[str | None, Query()] = None, + iteration: Annotated[int | None, Query(ge=0)] = None, + loop_region_id: Annotated[UUID | None, Query()] = None, Review Comment: Do we want `loop_region_id` on the public v2 list endpoint (and list-shaped `loop_iterations` on the TI responses) before nested loops have a design? Authoring rejects nested loops and loops in mapped groups, the scheduler creates one root region per loop per run, and clears fork back into the same family, so there is at most one invocation per run and `parent_iterations` is always empty. Once this ships in a release these shapes are hard to change. Shipping `loop_id` + `iteration` only (and a single `loop_iteration` on the response) would drop the invocation picker, the family selection in `loop_run_summaries` and the `region_options` walk, and the plural form could come with nested-loop support. ########## airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py: ########## @@ -1490,15 +1485,27 @@ async def get_previous_task_instance( if state: query = query.where(TI.state == state) - row = (await session.execute(query.limit(1))).first() - if row is None: + first = session.execute(query.limit(1)).first() + if first is None: return None - ti, public_index, run_count = row - if run_count > 1: - raise HTTPException( - status.HTTP_409_CONFLICT, "Select explicit region coordinates for the previous task" + if first[0].region_id != SENTINEL_REGION_ID: + resolver = TaskCoordinateResolver(dag_bag, session) + requester = session.get(TI, token.id) + iterations = ( + resolver.loop_iterations(requester) + if requester is not None and (requester.dag_id, requester.task_id) == (dag_id, task_id) + else [] ) - + if iterations: + with contextlib.closing( + session.execute(query.limit(_MAX_PREVIOUS_TIS_SCANNED).execution_options(yield_per=50)) Review Comment: Each scanned row goes through `resolver.loop_iterations`, which for every not-yet-seen run calls `load_region_ancestry` (and possibly `get_dag`) on the same session while this `yield_per` stream still has unread rows. That is one or more queries per earlier run, up to 500 rows. I'm also not sure it's safe on MySQL: `yield_per` turns on server-side cursors, and mysqlclient refuses a new statement while an unbuffered result is pending. Since the scan is capped anyway, could this fetch the rows with `.all()`, call `resolver.prefetch_regions(...)`, and then filter? Pushing the position into SQL (`region_index == iteration` when the task sits directly in the loop region) would also fix the other edge: rows are ordered by `region_index` desc, so a loop with more than 500 passes in earlier runs can return `None` here even when a matching pass exists. The new test only scans two rows, so neither path is exercised. ########## airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py: ########## @@ -337,25 +329,60 @@ def test_execution_breadcrumbs_keep_regional_identity_separate_from_map_index(cl assert {row["region_index"] for row in breadcrumbs} == {1, 2} -def test_previous_ti_full_coordinates_override_default_public_index(client, dag_maker, session): - with dag_maker(serialized=True): - PythonOperator.partial(task_id="mapped", python_callable=str).expand(op_args=[[0], [1], [2]]) - dr = dag_maker.create_dagrun() - region = DynamicRegion(dag_id=dr.dag_id, run_id=dr.run_id, node_id="mapped") - session.add(region) - session.flush() - for ti in dr.task_instances: - ti.region_id = region.id [email protected] +def two_run_loop_tis(dag_maker, session): + @task_group + def body(): + EmptyOperator(task_id="task") + + with dag_maker(serialized=True) as dag: + create_loop(body, max_iterations=3) + runs = { + "old": dag_maker.create_dagrun(run_id="old", logical_date=timezone.datetime(2025, 1, 1)), + "current": dag_maker.create_dagrun(run_id="current", logical_date=timezone.datetime(2025, 1, 2)), + } + passes = {} + for name, dr in runs.items(): + region = DynamicRegion(dag_id=dr.dag_id, run_id=dr.run_id, node_id="body") + session.add(region) + session.flush() + first = next(ti for ti in dr.task_instances if ti.task_id == "body.task") + first.region_id, first.region_index, first.state = region.id, 0, State.SUCCESS + second = TaskInstance( + task=dag.get_task(first.task_id), run_id=dr.run_id, dag_version_id=first.dag_version_id + ) + second.region_id, second.region_index, second.state = region.id, 1, State.SUCCESS + session.add(second) + passes[name] = {0: first, 1: second} session.commit() + return passes + + [email protected]("requester_pass", [0, 1]) +def test_previous_ti_for_loop_task_uses_requester_pass(client, two_run_loop_tis, requester_pass): + requester = two_run_loop_tis["current"][requester_pass] + exec_app = client.app.routes[-1].app + exec_app.dependency_overrides[require_auth] = lambda: TIToken(id=requester.id, claims=TIClaims()) + + response = client.get( + f"/execution/task-instances/previous/{requester.dag_id}/{requester.task_id}", + params={"logical_date": "2025-01-02T00:00:00Z"}, + ) + + assert response.status_code == 200 + assert response.json()["run_id"] == "old" + assert response.json()["region_index"] == requester_pass + +def test_previous_ti_for_other_loop_task_returns_latest_pass(client, two_run_loop_tis): Review Comment: This runs with the fixture's default all-zeros token, so `requester` is `None` and the test never exercises a different loop task asking. Removing the `(requester.dag_id, requester.task_id) == (dag_id, task_id)` guard wouldn't fail it. Overriding `require_auth` with a TI of another task (like the test above does) would make it bind. ########## task-sdk/src/airflow/sdk/execution_time/request_handlers.py: ########## @@ -175,8 +175,6 @@ def handle_get_task_states(client: Client, msg: GetTaskStates) -> tuple[BaseMode def handle_get_previous_ti(client: Client, msg: GetPreviousTI) -> tuple[BaseModel | None, dict[str, bool]]: """Fetch the previous task instance.""" resp = client.task_instances.get_previous( Review Comment: `test_supervisor.py` still has `"region_id": None, "region_index": None` in the `get_previous_ti` case's `ClientMock.kwargs` (around line 3093), and `test_handle_requests` checks the call with `assert_called_once_with(**client_mock.kwargs)`, so that case will fail now that these kwargs aren't passed. The `region_id`/`region_index` keys in `expected_body` should stay, since `PreviousTIResponse` still has them. ########## airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py: ########## @@ -1442,46 +1442,41 @@ def get_task_instance_count( return count or 0 [email protected]( - "/previous/{dag_id}/{task_id}", - status_code=status.HTTP_200_OK, - responses=create_openapi_http_exception_doc( - [(status.HTTP_409_CONFLICT, "Explicit region coordinates are required to select the previous task")] - ), -) -async def get_previous_task_instance( +_MAX_PREVIOUS_TIS_SCANNED = 500 + + [email protected]("/previous/{dag_id}/{task_id}", status_code=status.HTTP_200_OK) +def get_previous_task_instance( dag_id: str, task_id: str, - session: AsyncSessionDep, + session: SessionDep, + dag_bag: DagBagDep, logical_date: Annotated[UtcDateTime | None, Query()] = None, map_index: Annotated[int, Query()] = -1, state: Annotated[TaskInstanceState | None, Query()] = None, - region_id: UUID | None = None, - region_index: int | None = None, + token: TIToken = CurrentTIToken, ) -> PreviousTIResponse | None: """ Get the previous task instance matching the given criteria. + When the requesting task instance is of ``task_id``, the previous task instance sits at the same Review Comment: This changes what `ti.get_previous_ti()` returns for loop tasks (the same loop position in an earlier run, where it used to 409 once the earlier run had more than one pass), and drops the region selectors from `GetPreviousTI`, the client, the supervisor schema and the version entry. The commit message only talks about UI, so it would be worth a sentence there and a line in loops.rst near the cross-run paragraph. The "otherwise the latest pass" branch is only reachable when another task asks about a loop task, since the SDK always sends its own `task_id`. ########## airflow-core/src/airflow/models/task_coordinates.py: ########## @@ -102,6 +103,46 @@ def adopt_dag(self, dag: SerializedDAG | None) -> None: if dag is not None and dag.dag_version_id is not None: self._dags.setdefault(dag.dag_version_id, dag) + def prefetch_regions(self, tis: Iterable[TaskCoordinate]) -> None: + wanted: dict[tuple[str, str], set[UUID]] = {} + for ti in tis: + if ti.region_id != UUID(int=0) and ti.region_id not in self._regions.get( + (ti.dag_id, ti.run_id), {} + ): + wanted.setdefault((ti.dag_id, ti.run_id), set()).add(ti.region_id) Review Comment: Since every mapped expansion has a region now, this prefetches ancestry for any mapped TI, loop or not, and does it once per `(dag_id, run_id)`. On the global Task Instances page (up to 100 rows across many runs, auto-refreshed while anything is pending) a Dag with mapping and no loops pays one or more extra queries per run, and `loop_iterations` then returns `[]` without reading them. Could this skip TIs whose pinned task has no `enclosing_loop` (the Dags are already cached in `_dags`), the way `_build_ti_summaries` only loads regions when the Dag has a loop, and load the rest in one `DynamicRegion.id.in_(...)` walk across runs? Minor: `UUID(int=0)` here is `SENTINEL_REGION_ID`. ########## airflow-core/src/airflow/models/task_coordinates.py: ########## @@ -102,6 +103,46 @@ def adopt_dag(self, dag: SerializedDAG | None) -> None: if dag is not None and dag.dag_version_id is not None: self._dags.setdefault(dag.dag_version_id, dag) + def prefetch_regions(self, tis: Iterable[TaskCoordinate]) -> None: + wanted: dict[tuple[str, str], set[UUID]] = {} + for ti in tis: + if ti.region_id != UUID(int=0) and ti.region_id not in self._regions.get( + (ti.dag_id, ti.run_id), {} + ): + wanted.setdefault((ti.dag_id, ti.run_id), set()).add(ti.region_id) + for (dag_id, run_id), region_ids in wanted.items(): + self._regions.setdefault((dag_id, run_id), {}).update( + load_region_ancestry(region_ids, dag_id=dag_id, run_id=run_id, session=self.session) + ) + + def loop_iterations(self, ti: TaskCoordinate) -> list[tuple[str, int]]: Review Comment: `loop_iterations` repeats what `loop_context` right below already does (find the enclosing loop, load ancestry, `loop_position`, same error), but swallows `TaskNotFound`/`ValueError` differently and is the only one using the `_regions` cache. Since nested loops are rejected, could it be `ctx = self.loop_context(ti); return [(ctx[0].node_id, ctx[1])] if ctx else []`, with `loop_context` reading through `_regions`? Otherwise a later change to how forks map to an iteration has to be made in two places. ########## task-sdk/src/airflow/sdk/definitions/_internal/loop.py: ########## @@ -107,6 +107,7 @@ def create_loop( gate = LoopGateOperator( task_id=get_unique_task_id(gate_name, task_group=group), until=until, + doc_md=inspect.cleandoc(until.__doc__) if until is not None and until.__doc__ else None, Review Comment: `until.__doc__` on a `functools.partial` (or a callable class instance) is the class docstring, so `until=functools.partial(converged, threshold=eps)` shows "partial(func, *args, **keywords) - new function with partial application..." as the loop's exit criteria in the header and the gate's docs. `create_loop` already expects callables without `__name__` (the `__loop_gate` fallback), so maybe only take the docstring when `inspect.isfunction(until) or inspect.ismethod(until)`, or unwrap `partial.func`. ########## airflow-core/src/airflow/ui/src/queries/useIsLoopGroup.ts: ########## @@ -0,0 +1,46 @@ +/*! + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +import type { GridNodeResponse } from "openapi/requests/types.gen"; + +import { useGridStructure } from "src/queries/useGridStructure"; + +const findNode = ( + nodes: Array<GridNodeResponse> | null | undefined, + groupId: string, +): GridNodeResponse | undefined => { + const list = nodes ?? []; + + return ( + list.find((node) => node.id === groupId) ?? + list.reduce<GridNodeResponse | undefined>( + (found, node) => found ?? findNode(node.children, groupId), + undefined, + ) + ); +}; + +/** The given Task Group's node from the cached grid structure (carries ``is_loop``, ``doc_md``). */ +export const useLoopGroupNode = (groupId: string): GridNodeResponse | undefined => { + const { data } = useGridStructure({}); Review Comment: `findNode` duplicates `getGroupTask` in `src/utils/groupTask.ts`, which `Task.tsx` already uses for this lookup. `useGridStructure({})` also has a different query key from Grid's and Task's (`{ limit: 1 }`), so this starts its own `/structure` request instead of reading the cached one the docstring mentions. `getGroupTask(useGridStructure({ limit: 1 }).data, groupId)` would reuse both. ########## airflow-core/src/airflow/api_fastapi/core_api/routes/ui/gantt.py: ########## @@ -94,11 +100,16 @@ def get_gantt_data( f"No task instances for dag_id={dag_id} run_id={run_id}", ) + resolver = TaskCoordinateResolver(dag_bag, session) task_instances = [ GanttTaskInstance( id=row.id, region_id=row.region_id, region_index=row.region_index, + loop_iterations=[ + LoopIterationResponse(loop_id=loop_id, iteration=iteration) + for loop_id, iteration in resolver.loop_iterations(row) Review Comment: This resolves loop iterations row by row without `resolver.prefetch_regions(results)` first, so a mapped task inside a loop (one region per iteration) costs about two queries per iteration on every Gantt poll, since `load_region_ancestry` re-reads the parent loop region each time. The rows already carry `dag_id`/`run_id`/`region_id`. The same per-TI path is now hit by `get_hitl_details` (hitl.py) and `dry_run_clear_dag_run` (services/public/dag_run.py), which build `TaskInstanceResponse` lists without going through `task_coordinate_responses`. ########## airflow-core/src/airflow/ui/src/pages/TaskInstances/index.ts: ########## @@ -16,4 +16,4 @@ * specific language governing permissions and limitations * under the License. */ -export { TaskInstances } from "./TaskInstances"; +export { getRowKey, taskInstanceColumns, TaskInstances } from "./TaskInstances"; Review Comment: Nothing imports `getRowKey`, `taskInstanceColumns` or `ColumnProps` from here; the same names elsewhere (DagsList, DagRuns, HITLTaskInstances) are their own local definitions. Could these go back to being module-private? ########## airflow-core/src/airflow/ui/src/layouts/Details/Gantt/GanttTimeline.tsx: ########## @@ -89,6 +89,7 @@ const toTooltipSummary = ( return { child_states: null, + loop_iterations: segment.loopIterations, Review Comment: `loop_iterations` isn't declared on `LightGridTaskInstanceSummaryWithWhen`, so this only type-checks because returned object literals skip excess-property checks, and the tooltip's `"loop_iterations" in taskInstance` then narrows the Gantt object to the REST types. Could it go on the wrapper next to `queued_when`/`scheduled_when` (`readonly loop_iterations?: GanttTaskInstance["loop_iterations"]`)? ########## airflow-core/src/airflow/api_fastapi/core_api/routes/ui/grid.py: ########## @@ -555,3 +582,76 @@ def _generate() -> Generator[str, None, None]: yield GridTISummaries.model_validate(summary).model_dump_json() + "\n" return StreamingResponse(content=_generate(), media_type="application/x-ndjson") + + +@grid_router.get( + "/loop/{dag_id}/{run_id}/{group_id}", + responses=create_openapi_http_exception_doc([404, 422]), + dependencies=[ + Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK_INSTANCE)), + Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN)), + ], +) +def get_loop_summary( + dag_id: str, + run_id: str, + group_id: str, + dag_bag: DagBagDep, + session: SessionDep, + loop_region_id: UUID | None = None, +) -> LoopSummaryResponse: + """Summarize one loop invocation from the run's live task instances.""" + run = session.scalar(select(DagRun).where(DagRun.dag_id == dag_id, DagRun.run_id == run_id)) + if run is None: + raise HTTPException(404, "DAG run not found") + return loop_run_summaries(run, group_id, session=session, dag_bag=dag_bag, loop_region_id=loop_region_id)[ + 0 + ] + + +@grid_router.get( + "/loop-history/{dag_id}/{group_id}", + responses=create_openapi_http_exception_doc([404]), + dependencies=[ + Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK_INSTANCE)), + Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN)), + ], +) +def get_loop_history( + dag_id: str, + group_id: str, + dag_bag: DagBagDep, + session: SessionDep, + limit: Annotated[int, Query(ge=1, le=100)] = 14, +) -> LoopHistoryResponse: + """Show loop invocations from recent runs, oldest first.""" + runs = session.scalars( + select(DagRun) + .where(DagRun.dag_id == dag_id) + .order_by(DagRun.run_after.desc(), DagRun.id.desc()) + .limit(limit) + ).all() + history = [] + for run in reversed(runs): + try: + summaries = loop_run_summaries(run, group_id, session=session, dag_bag=dag_bag) Review Comment: `loop_run_summaries` runs once per run here, each with a distinct-version query, the member TI query and at least one region load, so the default `limit=14` is around 50 queries and `limit=100` several hundred, and every live loop TI of every run comes back to Python to draw one bar. Could the member and region queries take `run_id.in_(...)` once and group in Python? Small thing while here: `/loop-history/` is the only kebab-case path next to `/ti_summaries/` and `/structure/`. ########## airflow-core/src/airflow/ui/src/queries/useLoopSummary.ts: ########## @@ -0,0 +1,56 @@ +/*! + * 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 { useDagRunServiceGetDagRun, useGridServiceGetLoopSummary } from "openapi/queries"; + +import { isStatePending, useAutoRefresh } from "src/utils"; + +/** Runtime summary of a looped Task Group; polls while the loop is still running. */ +export const useLoopSummary = ({ + dagId, + groupId, + runId, +}: { + dagId: string; + groupId: string; + runId: string; +}) => { + const refetchInterval = useAutoRefresh({ dagId }); + const [searchParams] = useSearchParams(); + const { data: dagRun } = useDagRunServiceGetDagRun({ dagId, dagRunId: runId }, undefined, { + enabled: Boolean(dagId) && Boolean(runId), + refetchInterval: (query) => isStatePending(query.state.data?.state) && refetchInterval, + }); + + return useGridServiceGetLoopSummary( + { + dagId, + groupId, + loopRegionId: searchParams.get("loop_region_id") ?? undefined, Review Comment: This PR adds `SearchParamsKeys.LOOP_REGION_ID`; could this and `LoopHistoryChart.tsx:105` use it instead of the literal? Similarly `TaskInstances.tsx` compares against `"all"`, which has to match the private `ALL` in `IterationSelect.tsx`, so exporting one constant would keep them tied. `useLoopHistory`'s `limit = 14` also just copies the API default and has no caller passing it. ########## airflow-core/src/airflow/ui/src/queries/useLoopSummary.ts: ########## @@ -0,0 +1,56 @@ +/*! + * 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 { useDagRunServiceGetDagRun, useGridServiceGetLoopSummary } from "openapi/queries"; + +import { isStatePending, useAutoRefresh } from "src/utils"; + +/** Runtime summary of a looped Task Group; polls while the loop is still running. */ +export const useLoopSummary = ({ + dagId, + groupId, + runId, +}: { + dagId: string; + groupId: string; + runId: string; +}) => { + const refetchInterval = useAutoRefresh({ dagId }); + const [searchParams] = useSearchParams(); + const { data: dagRun } = useDagRunServiceGetDagRun({ dagId, dagRunId: runId }, undefined, { + enabled: Boolean(dagId) && Boolean(runId), + refetchInterval: (query) => isStatePending(query.state.data?.state) && refetchInterval, + }); + + return useGridServiceGetLoopSummary( + { + dagId, + groupId, + loopRegionId: searchParams.get("loop_region_id") ?? undefined, + runId, + }, + undefined, + { + enabled: Boolean(dagId) && Boolean(groupId) && Boolean(runId), Review Comment: This fires for every task group, not just loops: `Header` and `GroupTaskInstances` both mount it on every `tasks/group/:groupId` page. For a plain `@task_group` the endpoint answers 422 "Task group is not a loop" (after a distinct-version query and a `get_dag`), and with `refetchInterval` set while the run is pending the errored query keeps getting re-issued every auto-refresh tick. Could `enabled` also require the group to be a loop, e.g. `&& useIsLoopGroup(groupId)` (which is currently unused)? ########## airflow-core/src/airflow/ui/src/hooks/useLoopOutcomeStats.tsx: ########## @@ -0,0 +1,90 @@ +/*! + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +import type { ReactNode } from "react"; + +import { HStack, Icon, Text } from "@chakra-ui/react"; +import { useTranslation } from "react-i18next"; +import { FiRepeat } from "react-icons/fi"; + +import { finalCriteria, reasonSentence } from "src/pages/GroupTaskInstance/LoopIterations/loopUtils"; Review Comment: This is the only hook under `src/hooks` importing from `src/pages`, and it has a single caller (`GroupTaskInstance/Header.tsx`). Could it live next to that page, or the shared helpers move to `src/utils`? In `loopUtils.ts`, `FAILED_STATES`, `ResultTag` and `lastRanIteration` are exported with no outside users, and `useLoopRuleStats`'s `includeCap` option is never passed. ########## airflow-core/src/airflow/ui/src/pages/Task/Overview/Overview.test.tsx: ########## @@ -47,6 +47,8 @@ const wrapperWithSearch = (search: string) => { }; vi.mock("openapi/queries", () => ({ + useGridServiceGetDagStructure: () => ({ data: undefined }), + useGridServiceGetLoopHistory: () => ({ data: undefined }), Review Comment: With the structure mock returning `undefined`, `loopNode` is always undefined, so `LoopHistoryChart` never mounts and the `useGridServiceGetLoopHistory` mock is unreachable. Nothing tests the `is_loop` gate in `Overview.tsx`; having the structure mock return one `is_loop: true` group and asserting the chart renders for it (and not for a plain task) would cover it. ########## airflow-core/src/airflow/ui/src/pages/TaskInstance/Details.loopContext.test.tsx: ########## @@ -0,0 +1,65 @@ +/*! + * 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 { render, screen, within } from "@testing-library/react"; +import { describe, expect, it, vi } from "vitest"; + +import type * as OpenapiQueries from "openapi/queries"; +import type { TaskInstanceHistoryResponse } from "openapi/requests/types.gen"; + +import { useTaskInstanceView } from "src/hooks/useTaskInstanceView"; +import { Wrapper } from "src/utils/Wrapper"; + +import { Details } from "./Details"; + +vi.mock("src/hooks/useTaskInstanceView", () => ({ useTaskInstanceView: vi.fn() })); Review Comment: `Details.test.tsx` already has `buildTaskInstance`/`renderDetails` and stubs the sibling panels that fetch their own data. This file mocks a different layer and leaves `BlockingDeps`/`TriggererInfo`/`TeamName` real. Could the loop case be one more test in `Details.test.tsx` using its harness and the real locale labels? -- 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]
