This is an automated email from the ASF dual-hosted git repository.
henry3260 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 1ab53cd8e38 Keep the Dag asset schedule the same shape for every
caller (#73805)
1ab53cd8e38 is described below
commit 1ab53cd8e38c07c7b114384d27f76267934bbae3
Author: Henry Chen <[email protected]>
AuthorDate: Tue Oct 6 01:28:41 2026 +0800
Keep the Dag asset schedule the same shape for every caller (#73805)
next_run_assets returns only the assets the caller may read, so a UI
deriving the schedule from that list renders a different Dag depending on
who is looking at it: a smaller total, the single-asset layout instead of
the popover, or no asset schedule at all when none of the assets are
readable.
---
.../api_fastapi/core_api/datamodels/ui/assets.py | 6 +
.../api_fastapi/core_api/openapi/_private_ui.yaml | 4 +
.../api_fastapi/core_api/routes/ui/assets.py | 18 ++-
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 5 +
.../airflow/ui/openapi-gen/requests/types.gen.ts | 1 +
.../ui/src/pages/DagsList/AssetSchedule.test.tsx | 125 +++++++++++++++++++++
.../ui/src/pages/DagsList/AssetSchedule.tsx | 12 +-
.../api_fastapi/core_api/routes/ui/test_assets.py | 34 +++++-
8 files changed, 198 insertions(+), 7 deletions(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/assets.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/assets.py
index 54847eae467..7296257205b 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/assets.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/assets.py
@@ -52,4 +52,10 @@ class NextRunAssetsResponse(BaseModel):
asset_expression: MaybeAssetExpression = None
events: list[NextRunAssetEventResponse]
+ scheduling_asset_count: int = 0
+ """How many assets the Dag is scheduled on, before filtering ``events``
down to
+ the ones the caller may read. ``events`` is caller-scoped, so a UI that
derives
+ the schedule's shape from ``len(events)`` changes what it renders with the
+ caller's permissions; this count does not. It reveals nothing new — the
redacted
+ ``asset_expression`` already carries one slot per asset."""
pending_partition_count: int | None = None
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
index 170393f4cde..de12f997600 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
+++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/_private_ui.yaml
@@ -4453,6 +4453,10 @@ components:
$ref: '#/components/schemas/NextRunAssetEventResponse'
type: array
title: Events
+ scheduling_asset_count:
+ type: integer
+ title: Scheduling Asset Count
+ default: 0
pending_partition_count:
anyOf:
- type: integer
diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/assets.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/assets.py
index b23ea597716..5a8dbefb917 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/assets.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/assets.py
@@ -230,6 +230,16 @@ def next_run_assets(
)
query = readable_assets_filter.to_orm(query)
+ # Counted before the readable filter narrows ``query``: the UI derives the
+ # schedule's shape (total, and whether to render the multi-asset popover)
from
+ # this, so that shape stays the same for every caller regardless of which
+ # assets they may read.
+ scheduling_asset_count = session.scalar(
+ select(func.count())
+ .select_from(DagScheduleAssetReference)
+ .where(DagScheduleAssetReference.dag_id == dag_id)
+ )
+
if not is_partitioned:
query = query.join(
AssetDagRunQueue,
@@ -253,7 +263,11 @@ def next_run_assets(
)
for row in raw_rows
]
- model_data: dict[str, Any] = {"asset_expression": asset_expression,
"events": events}
+ model_data: dict[str, Any] = {
+ "asset_expression": asset_expression,
+ "events": events,
+ "scheduling_asset_count": scheduling_asset_count,
+ }
return NextRunAssetsResponse.model_validate(model_data)
# Partitioned Dags: enrich with per-asset received/required counts and
rollup flag.
@@ -290,6 +304,7 @@ def next_run_assets(
model_data = {
"asset_expression": asset_expression,
"events": events,
+ "scheduling_asset_count": scheduling_asset_count,
"pending_partition_count": pending_partition_count,
}
return NextRunAssetsResponse.model_validate(model_data)
@@ -364,6 +379,7 @@ def next_run_assets(
model_data = {
"asset_expression": asset_expression,
"events": events,
+ "scheduling_asset_count": scheduling_asset_count,
"pending_partition_count": pending_partition_count,
}
return NextRunAssetsResponse.model_validate(model_data)
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index 7efd0bc2abd..8479226b0c8 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -11689,6 +11689,11 @@ export const $NextRunAssetsResponse = {
type: 'array',
title: 'Events'
},
+ scheduling_asset_count: {
+ type: 'integer',
+ title: 'Scheduling Asset Count',
+ default: 0
+ },
pending_partition_count: {
anyOf: [
{
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index 154f4a802a8..ddadbd19f28 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -3041,6 +3041,7 @@ export type NextRunAssetEventResponse = {
export type NextRunAssetsResponse = {
asset_expression?: AssetExpressionAsset | AssetExpressionAlias |
AssetExpressionRef | AssetExpressionAny | AssetExpressionAll | null;
events: Array<NextRunAssetEventResponse>;
+ scheduling_asset_count?: number;
pending_partition_count?: number | null;
};
diff --git
a/airflow-core/src/airflow/ui/src/pages/DagsList/AssetSchedule.test.tsx
b/airflow-core/src/airflow/ui/src/pages/DagsList/AssetSchedule.test.tsx
new file mode 100644
index 00000000000..aa56fc65e4b
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/DagsList/AssetSchedule.test.tsx
@@ -0,0 +1,125 @@
+/*!
+ * 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 { render, screen } from "@testing-library/react";
+import { beforeEach, describe, expect, it, vi } from "vitest";
+
+import type * as OpenapiQueries from "openapi/queries";
+import type { NextRunAssetEventResponse, NextRunAssetsResponse } from
"openapi/requests/types.gen";
+
+import { Wrapper } from "src/utils/Wrapper";
+
+import { AssetSchedule } from "./AssetSchedule";
+
+vi.mock("react-i18next", () => ({
+ useTranslation: () => ({
+ i18n: { language: "en" },
+ // eslint-disable-next-line id-length
+ t: (key: string, options?: { count?: number; total?: number }) =>
+ options?.count === undefined ? key : `${key}:${options.count} of
${options.total}`,
+ }),
+}));
+
+vi.mock("openapi/queries", async (importOriginal) => {
+ const actual = await importOriginal<typeof OpenapiQueries>();
+
+ return {
+ ...actual,
+ useAssetServiceGetDagAssetQueuedEvents: vi.fn(),
+ useAssetServiceNextRunAssets: vi.fn(),
+ };
+});
+
+const { useAssetServiceGetDagAssetQueuedEvents, useAssetServiceNextRunAssets }
=
+ await import("openapi/queries");
+
+const makeEvent = (id: number, name: string): NextRunAssetEventResponse => ({
+ asset_inactive: false,
+ id,
+ is_rollup: false,
+ last_update: null,
+ mapper_error: false,
+ name,
+ received_count: 0,
+ received_keys: [],
+ required_count: 1,
+ required_keys: [],
+ uri: `s3://bucket/${name}`,
+});
+
+const nextRunResponse = (nextRun: NextRunAssetsResponse) =>
+ ({ data: nextRun, error: null, isFetching: false, isLoading: false }) as
ReturnType<
+ typeof useAssetServiceNextRunAssets
+ >;
+
+const queuedEventsResponse = () =>
+ ({
+ data: { queued_events: [], total_entries: 0 },
+ error: null,
+ isFetching: false,
+ isLoading: false,
+ }) as ReturnType<typeof useAssetServiceGetDagAssetQueuedEvents>;
+
+const renderSchedule = (nextRun: NextRunAssetsResponse) => {
+
vi.mocked(useAssetServiceNextRunAssets).mockReturnValue(nextRunResponse(nextRun));
+
vi.mocked(useAssetServiceGetDagAssetQueuedEvents).mockReturnValue(queuedEventsResponse());
+
+ render(
+ <AssetSchedule dagId="dag_id" timetablePartitioned={false}
timetableSummary="Every day at midnight" />,
+ { wrapper: Wrapper },
+ );
+};
+
+describe("AssetSchedule", () => {
+ beforeEach(() => {
+ vi.clearAllMocks();
+ });
+
+ it("shows the Dag's full asset total when the caller may read only some of
them", () => {
+ renderSchedule({
+ events: [makeEvent(1, "visible_asset")],
+ scheduling_asset_count: 3,
+ });
+
+ expect(screen.getByRole("button")).toHaveTextContent("assetSchedule:0 of
3");
+ });
+
+ it("still renders an asset schedule when the caller may read none of the
assets", () => {
+ renderSchedule({ events: [], scheduling_asset_count: 3 });
+
+ expect(screen.getByRole("button")).toHaveTextContent("assetSchedule:0 of
3");
+ expect(screen.queryByText("Every day at
midnight")).not.toBeInTheDocument();
+ });
+
+ it("falls back to the timetable summary when the Dag has no scheduling
assets", () => {
+ renderSchedule({ events: [], scheduling_asset_count: 0 });
+
+ expect(screen.getByText("Every day at midnight")).toBeInTheDocument();
+ });
+
+ it("renders the single-asset view only when the Dag is scheduled on one
asset", () => {
+ renderSchedule({
+ events: [makeEvent(1, "only_asset")],
+ scheduling_asset_count: 1,
+ });
+
+ expect(screen.getByRole("link", { name: "only_asset"
})).toBeInTheDocument();
+ expect(screen.queryByRole("button")).not.toBeInTheDocument();
+ });
+});
diff --git a/airflow-core/src/airflow/ui/src/pages/DagsList/AssetSchedule.tsx
b/airflow-core/src/airflow/ui/src/pages/DagsList/AssetSchedule.tsx
index 4b82cabefb5..65ec10a6477 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagsList/AssetSchedule.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagsList/AssetSchedule.tsx
@@ -118,13 +118,17 @@ export const AssetSchedule = ({ assetExpression, dagId,
timetablePartitioned, ti
0,
)
: pendingEvents.length;
+ // `events` only carries the assets the caller may read, so the total and the
+ // layout branches below come from `scheduling_asset_count` instead — the
schedule
+ // a Dag shows should not change with who is looking at it.
+ const schedulingAssetCount = nextRun?.scheduling_asset_count ?? 0;
const scheduledTotal = timetablePartitioned
? nextRunEvents.reduce((sum, event) => sum + (event.required_count ?? 1),
0)
- : nextRunEvents.length;
+ : schedulingAssetCount;
const isLoading = isNextRunLoading || (!timetablePartitioned &&
isQueuedEventsLoading);
- if (!nextRunEvents.length) {
+ if (!schedulingAssetCount) {
return (
<HStack>
<FiDatabase style={{ display: "inline", flexShrink: 0 }} />
@@ -160,7 +164,7 @@ export const AssetSchedule = ({ assetExpression, dagId,
timetablePartitioned, ti
// pendingCount === 1: render single-asset view with inactive warning.
const [partitionedAsset] = nextRunEvents;
- if (nextRunEvents.length === 1 && partitionedAsset !== undefined) {
+ if (schedulingAssetCount === 1 && partitionedAsset !== undefined) {
const requiredCount = partitionedAsset.required_count ?? 1;
const receivedCount = partitionedAsset.received_count ?? 0;
const requiredKeys = partitionedAsset.required_keys ?? [];
@@ -214,7 +218,7 @@ export const AssetSchedule = ({ assetExpression, dagId,
timetablePartitioned, ti
const [asset] = nextRunEvents;
- if (nextRunEvents.length === 1 && asset !== undefined) {
+ if (schedulingAssetCount === 1 && asset !== undefined) {
const requiredCount = asset.required_count ?? 1;
const receivedCount = asset.received_count ?? 0;
const requiredKeys = asset.required_keys ?? [];
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_assets.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_assets.py
index 2f1079f2f9c..e26dc173651 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_assets.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_assets.py
@@ -84,8 +84,9 @@ class TestNextRunAssets:
dag_maker.create_dagrun()
dag_maker.sync_dagbag_to_db()
- # 4 queries for the endpoint plus 1 to resolve the assets the caller
may read.
- with assert_queries_count(5):
+ # 4 queries for the endpoint, 1 to resolve the assets the caller may
read and
+ # 1 for the caller-independent scheduling asset count.
+ with assert_queries_count(6):
response = test_client.get("/next_run_assets/upstream")
assert response.status_code == 200
@@ -118,6 +119,7 @@ class TestNextRunAssets:
"asset_inactive": False,
}
],
+ "scheduling_asset_count": 1,
"pending_partition_count": None,
}
@@ -186,6 +188,32 @@ class TestNextRunAssets:
assert response.status_code == 200
assert [event["name"] for event in response.json()["events"]] ==
["visible_asset"]
+ assert response.json()["scheduling_asset_count"] == 2
+
+ @mock.patch(
+
"airflow.api_fastapi.auth.managers.base_auth_manager.BaseAuthManager.get_authorized_assets",
+ autospec=True,
+ )
+ def test_scheduling_asset_count_is_unaffected_when_no_asset_is_readable(
+ self, mock_get_authorized_assets, test_client, dag_maker
+ ):
+ with dag_maker(
+ dag_id="all_hidden_upstream",
+ schedule=[
+ Asset(uri="s3://bucket/hidden1", name="hidden_asset1"),
+ Asset(uri="s3://bucket/hidden2", name="hidden_asset2"),
+ ],
+ serialized=True,
+ ):
+ EmptyOperator(task_id="task1")
+ dag_maker.sync_dagbag_to_db()
+ mock_get_authorized_assets.return_value = set()
+
+ response = test_client.get("/next_run_assets/all_hidden_upstream")
+
+ assert response.status_code == 200
+ assert response.json()["events"] == []
+ assert response.json()["scheduling_asset_count"] == 2
@mock.patch(
"airflow.api_fastapi.auth.managers.base_auth_manager.BaseAuthManager.get_authorized_assets",
@@ -240,6 +268,7 @@ class TestNextRunAssets:
assert [event["name"] for event in events] == ["part_visible"]
assert events[0]["received_keys"] == ["2024-01-01"]
assert events[0]["required_keys"] == ["2024-01-01"]
+ assert response.json()["scheduling_asset_count"] == 2
def test_should_respond_401(self, unauthenticated_test_client):
response = unauthenticated_test_client.get("/next_run_assets/upstream")
@@ -336,6 +365,7 @@ class TestNextRunAssets:
"asset_inactive": False,
},
],
+ "scheduling_asset_count": 2,
"pending_partition_count": None,
}