dheerajturaga commented on code in PR #74368: URL: https://github.com/apache/airflow/pull/74368#discussion_r4209428001
########## airflow-core/src/airflow/ui/src/utils/assetEvents.ts: ########## @@ -0,0 +1,29 @@ +/*! + * 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 { DAGRunResponse } from "openapi/requests/types.gen"; + +/** + * Whether a run can have consumed asset events. + * + * Only an asset-triggered run records the events that caused it; every other run type, including + * a materialization, is created without any. The tab and the page it opens share this so one + * cannot offer a tab the other never fills. + */ +export const canHaveUpstreamAssetEvents = (dagRun: DAGRunResponse | undefined): boolean => + dagRun === undefined || dagRun.run_type === "asset_triggered"; Review Comment: Scheduled runs can also have consumed asset events now. With `AssetAndTimeSchedule` (#58543), `_create_dagruns_for_dags` in the scheduler creates the run with `run_type=SCHEDULED` and then calls `created_run.consumed_asset_events.extend(gated_asset_events)`. `/upstreamAssetEvents` returns those events whatever the run type. With this helper, those runs lose the Asset Events tab entirely. The page's `enabled` check already had this gap on `main`. But hiding the tab, plus the docstring ("every other run type … is created without any") and the `"scheduled"` case in `useDagRunTabs.test.tsx`, would lock it in. Could the backend answer this instead, e.g. a `has_consumed_asset_events` flag on the Dag run response? Then the UI doesn't have to track which run types the scheduler attaches events to. Minor, related: returning `true` for an `undefined` run mixes "unknown" with "can". That's why `pages/Run/AssetEvents.tsx` has to add `dagRun !== undefined &&`. Keeping the "unknown, so keep the tab" default inside `useDagRunTabs` would make the helper safer to reuse. --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting ########## airflow-core/src/airflow/ui/src/hooks/useTaskInstanceTabs.ts: ########## @@ -0,0 +1,74 @@ +/*! + * 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 { useLocation } from "react-router-dom"; + +import { useTaskServiceGetTask } from "openapi/queries"; + +import { TaskInstanceTab } from "src/constants/tab"; +import type { TabItem } from "src/hooks/useRequiredActionTabs"; + +type Params = { + dagId: string; + /** The version this task instance ran under; omit to describe the Dag as it is now. */ + dagVersionId?: string; + /** False until the task instance is known, so the task is never fetched unpinned first. */ + isVersionKnown: boolean; + taskId: string; +}; + +/** + * Drops task-instance tabs the Dag definition rules out. + * + * Both rules are answered from the task alone, so they hold before the task has run and cost no + * extra request: a task with no outlets can never emit an asset event, and one with no template + * fields never renders anything. + * + * Tabs whose emptiness is only knowable from observed data are left alone. An empty tab is a + * smaller cost than a tab that disappears once its query comes back. + * + * A tab the user is currently on is always kept, so deep links keep working and the tab bar + * never renders with nothing selected. + */ +export const useTaskInstanceTabs = <T extends TabItem>( + { dagId, dagVersionId, isVersionKnown, taskId }: Params, + tabs: Array<T>, +): { tabs: Array<T> } => { + const { pathname } = useLocation(); + const lastSegment = pathname.split("/").pop() ?? ""; + + // Pinned to the version this instance ran under, so an older instance is judged by the + // definition it actually ran, not by a later edit. A run can span versions, so this comes + // from the task instance rather than the run. + // + // Held until the instance is known: fetching unpinned first would answer from the latest + // version, then answer again from the pinned one, which is exactly the flicker this avoids. + const { data: task } = useTaskServiceGetTask({ dagId, dagVersionId, taskId }, undefined, { + enabled: Boolean(dagId) && Boolean(taskId) && isVersionKnown, + }); + + // Unknown until the task loads, so everything stays put rather than flickering out and back. + // `has_outlets` is optional, so a server that omits it leaves the tab alone too: only a + // definite "no outlets" hides it. + const canFill: Record<string, boolean> = { + [TaskInstanceTab.AssetEvents]: task?.has_outlets !== false, + [TaskInstanceTab.RenderedTemplates]: task === undefined || (task.template_fields ?? []).length > 0, + }; + + return { tabs: tabs.filter((tab) => (canFill[tab.value] ?? true) || tab.value === lastSegment) }; Review Comment: This interacts with the user's default task instance tab setting. "Rendered Templates" and "Asset Events" are both options there. `DefaultTab` redirects to the chosen tab, `lastSegment` matches it, and the tab is kept. As soon as the user clicks another tab, it disappears from the bar. Should `DefaultTab` fall back to Logs when the chosen tab is ruled out for this task? --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting ########## airflow-core/src/airflow/ui/src/hooks/useTaskInstanceTabs.ts: ########## @@ -0,0 +1,74 @@ +/*! + * 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 { useLocation } from "react-router-dom"; + +import { useTaskServiceGetTask } from "openapi/queries"; + +import { TaskInstanceTab } from "src/constants/tab"; +import type { TabItem } from "src/hooks/useRequiredActionTabs"; + +type Params = { + dagId: string; + /** The version this task instance ran under; omit to describe the Dag as it is now. */ + dagVersionId?: string; + /** False until the task instance is known, so the task is never fetched unpinned first. */ + isVersionKnown: boolean; + taskId: string; +}; + +/** + * Drops task-instance tabs the Dag definition rules out. + * + * Both rules are answered from the task alone, so they hold before the task has run and cost no + * extra request: a task with no outlets can never emit an asset event, and one with no template + * fields never renders anything. + * + * Tabs whose emptiness is only knowable from observed data are left alone. An empty tab is a + * smaller cost than a tab that disappears once its query comes back. + * + * A tab the user is currently on is always kept, so deep links keep working and the tab bar + * never renders with nothing selected. + */ +export const useTaskInstanceTabs = <T extends TabItem>( + { dagId, dagVersionId, isVersionKnown, taskId }: Params, + tabs: Array<T>, +): { tabs: Array<T> } => { + const { pathname } = useLocation(); + const lastSegment = pathname.split("/").pop() ?? ""; + + // Pinned to the version this instance ran under, so an older instance is judged by the + // definition it actually ran, not by a later edit. A run can span versions, so this comes + // from the task instance rather than the run. + // + // Held until the instance is known: fetching unpinned first would answer from the latest + // version, then answer again from the pinned one, which is exactly the flicker this avoids. + const { data: task } = useTaskServiceGetTask({ dagId, dagVersionId, taskId }, undefined, { + enabled: Boolean(dagId) && Boolean(taskId) && isVersionKnown, + }); + + // Unknown until the task loads, so everything stays put rather than flickering out and back. + // `has_outlets` is optional, so a server that omits it leaves the tab alone too: only a + // definite "no outlets" hides it. + const canFill: Record<string, boolean> = { + [TaskInstanceTab.AssetEvents]: task?.has_outlets !== false, Review Comment: Edge case: the task instance Asset Events page lists events by source Dag, run, task and map index, so it includes events from every try. This check looks only at the task instance's *current* `dag_version_id`. Clearing with "run on latest version" moves that to the latest version (`taskinstance.py`, `clear_task_instances`). If the latest version dropped the outlets, events from earlier tries are still there but the tab is gone. Probably acceptable, but worth noting in the docstring if so. --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting ########## airflow-core/src/airflow/api_fastapi/core_api/datamodels/tasks.py: ########## @@ -81,11 +81,22 @@ class TaskResponse(BaseModel): params: abc.MutableMapping | None class_ref: dict | None is_mapped: bool | None + has_outlets: bool = Field( Review Comment: `validate_model` always sets this, so could it be a required `has_outlets: bool`, like `is_mapped`? With `default=False`, the OpenAPI schema marks it optional. The UI type becomes `has_outlets?: boolean`, which is where the `!== false` workaround in the hook comes from. And airflow-ctl's generated model becomes `bool | None = False`, so an older server that doesn't send the field would read as "no outlets". --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting ########## airflow-core/src/airflow/api_fastapi/core_api/datamodels/tasks.py: ########## @@ -81,11 +81,22 @@ class TaskResponse(BaseModel): params: abc.MutableMapping | None class_ref: dict | None is_mapped: bool | None + has_outlets: bool = Field( + default=False, + description="Whether the task declares any ``outlets``, i.e. whether running it can " + "produce asset events.", + ) @model_validator(mode="before") @classmethod def validate_model(cls, task: Any) -> Any: - task.__dict__.update({"class_ref": _get_class_ref(task), "is_mapped": task.is_mapped}) + task.__dict__.update( + { + "class_ref": _get_class_ref(task), + "is_mapped": task.is_mapped, + "has_outlets": bool(task.outlets), Review Comment: Question: `bool(task.outlets)` only sees statically declared outlets. `outlets` passes `.expand()` validation because it's a `BaseOperator` init argument. `SerializedMappedOperator.outlets` reads only `partial_kwargs`, so `.partial(...).expand(outlets=[...])` reports `has_outlets=False` while each mapped task instance emits events. The same goes for outlets added to `self.outlets` inside `execute()`, which the task runner reads after execution. If either pattern is supported, the tab would be hidden for tasks that do have events. --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting ########## airflow-core/src/airflow/api_fastapi/core_api/routes/public/tasks.py: ########## @@ -102,9 +103,31 @@ def get_tasks( ), dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK))], ) -def get_task(dag_id: str, task_id, session: SessionDep, dag_bag: DagBagDep) -> TaskResponse: - """Get simplified representation of a task.""" - dag = get_latest_version_of_dag(dag_bag, dag_id, session) +def get_task( + dag_id: str, + task_id, + session: SessionDep, + dag_bag: DagBagDep, + dag_version_id: UUID | None = None, +) -> TaskResponse: + """ + Get simplified representation of a task. + + ``dag_version_id`` pins the lookup to one Dag version, so a caller describing a past task + instance sees the task as it was defined then rather than as it is now. A Dag run can span + several versions, so the version is taken per task instance rather than per run. + """ + if dag_version_id is None: + dag = get_latest_version_of_dag(dag_bag, dag_id, session) + else: + versioned_dag = dag_bag.get_dag(dag_version_id, session=session) Review Comment: nit: `dag_bag.get_dag()` deserializes and caches the Dag for any version ID *before* the `dag_id` check. It's not an access problem, since the response is a 404. But a request on one Dag can make the API server deserialize another Dag's version into the shared cache, and pinned lookups of old versions push other entries out of it. A cheap `select(DagVersion.dag_id).where(DagVersion.id == dag_version_id)` first would avoid both. --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting ########## airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_tasks.py: ########## @@ -278,6 +283,69 @@ def test_should_respond_200_serialized(self, test_client, testing_dag_bundle): assert response.status_code == 200 assert response.json() == expected + def test_dag_version_id_describes_the_task_as_it_was(self, test_client, dag_maker, session): + """A later edit must not change how an earlier version's task is reported.""" + dag_id = "test_versioned_task_dag" + + # A new DagVersion is cut only when the serialized Dag changes under a new bundle version, + # so each version is written the way the shared multi-version fixture does it. + with dag_maker(dag_id, session=session, bundle_version="commit-one"): + EmptyOperator(task_id=self.task_id, retries=4) + session.commit() + first_version_id = str(DagVersion.get_version(dag_id, session=session).id) + + with dag_maker(dag_id, session=session, bundle_version="commit-two"): + EmptyOperator(task_id=self.task_id, retries=0) + session.commit() + + assert str(DagVersion.get_version(dag_id, session=session).id) != first_version_id + + test_client.app.dependency_overrides[dag_bag_from_app] = DBDagBag + url = f"{self.api_prefix}/{dag_id}/tasks/{self.task_id}" + + latest = test_client.get(url) + assert latest.status_code == 200 + assert latest.json()["retries"] == 0 + + pinned = test_client.get(url, params={"dag_version_id": first_version_id}) + assert pinned.status_code == 200 + assert pinned.json()["retries"] == 4 + + def test_dag_version_id_of_another_dag_is_not_found(self, test_client, testing_dag_bundle): Review Comment: nit: this covers the `dag_id` mismatch. Could you also add a case with a random UUID, to cover the `versioned_dag is None` branch? --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting ########## airflow-core/src/airflow/api_fastapi/core_api/routes/public/tasks.py: ########## @@ -102,9 +103,31 @@ def get_tasks( ), dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK))], ) -def get_task(dag_id: str, task_id, session: SessionDep, dag_bag: DagBagDep) -> TaskResponse: - """Get simplified representation of a task.""" - dag = get_latest_version_of_dag(dag_bag, dag_id, session) +def get_task( + dag_id: str, + task_id, + session: SessionDep, + dag_bag: DagBagDep, + dag_version_id: UUID | None = None, +) -> TaskResponse: + """ + Get simplified representation of a task. + + ``dag_version_id`` pins the lookup to one Dag version, so a caller describing a past task + instance sees the task as it was defined then rather than as it is now. A Dag run can span + several versions, so the version is taken per task instance rather than per run. + """ + if dag_version_id is None: + dag = get_latest_version_of_dag(dag_bag, dag_id, session) + else: + versioned_dag = dag_bag.get_dag(dag_version_id, session=session) + # Looked up by version id alone, so the Dag in the path still has to be checked. + if versioned_dag is None or versioned_dag.dag_id != dag_id: + raise HTTPException( + status.HTTP_404_NOT_FOUND, + f"Dag version {dag_version_id} was not found for dag {dag_id}", Review Comment: nit: we write "Dag" in prose: ```suggestion f"Dag version {dag_version_id} was not found for Dag {dag_id}", ``` --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting ########## airflow-core/src/airflow/ui/src/hooks/useDagRunTabs.ts: ########## @@ -0,0 +1,50 @@ +/*! + * 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 { useLocation } from "react-router-dom"; + +import type { DAGRunResponse } from "openapi/requests/types.gen"; + +import { DagRunTab } from "src/constants/tab"; +import type { TabItem } from "src/hooks/useRequiredActionTabs"; +import { canHaveUpstreamAssetEvents } from "src/utils/assetEvents"; + +/** + * Drops Dag run tabs the run itself rules out. + * + * Reading how the run was triggered rather than the Dag's current schedule keeps the answer correct + * after the Dag is edited: a run triggered back when the Dag was asset-scheduled still shows its + * events, and a new cron run of a formerly asset-scheduled Dag does not. + * + * A tab the user is currently on is always kept, so deep links keep working and the tab bar never + * renders with nothing selected. + */ +export const useDagRunTabs = <T extends TabItem>( + dagRun: DAGRunResponse | undefined, + tabs: Array<T>, +): { tabs: Array<T> } => { + const { pathname } = useLocation(); + const lastSegment = pathname.split("/").pop() ?? ""; + + // Unknown until the run loads, so everything stays put rather than flickering out and back. + const canFill: Record<string, boolean> = { + [DagRunTab.AssetEvents]: canHaveUpstreamAssetEvents(dagRun), + }; + + return { tabs: tabs.filter((tab) => (canFill[tab.value] ?? true) || tab.value === lastSegment) }; Review Comment: nit: this `lastSegment` + `filter` pair, and the doc paragraph above it, are repeated in `useTaskInstanceTabs.ts`. `lastSegment` also appears in `layouts/Details/NavTabs.tsx`. A small shared helper, e.g. `filterTabsByRules(tabs, canFill)`, would keep the three in sync. --- Drafted-by: Claude Code (Opus 5.5); reviewed by @dheerajturaga before posting -- 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]
