From 14a92687196cc0b102998c7f38715f96bb6ce3e7 Mon Sep 17 00:00:00 2001 From: Richard Date: Sun, 2 Aug 2026 22:37:04 -0700 Subject: [PATCH] [v3-3-test] Limit deadline alerts endpoint to the latest serialized Dag row (#70804) * The deadline alerts endpoint fetches the latest serialized Dag by ordering all rows for a dag_id descending and reading the first one. Without a LIMIT, the database materializes every serialized version of the Dag when only the newest is needed. Cap the lookup at a single row. * Sort imports in deadline alerts UI route tests * UI: Allow reading deadline alerts of a specific Dag version A backfill run executes the Dag version that was current when it was queued, whose deadline alerts may differ from the ones deployed today. The endpoint could only ever resolve the newest serialized Dag, so the UI reported present-day alerts for historical runs and misrepresented what those runs were actually held to. Resolving the version through DagVersion.version_number also removes a reliance on the serialized_dag primary key sorting chronologically, which only held because uuid7 happens to be time-ordered and would give the wrong row if an older version were serialized after a newer one. Only the primary key is read, so the serialized Dag blob no longer crosses the wire for a lookup that discards it. (cherry picked from commit 1b3121d03b6967796cbaf370f7ecf9c66664cae5) Co-authored-by: Richard --- .../core_api/openapi/_private_ui.yaml | 8 +++ .../core_api/routes/ui/deadlines.py | 22 ++++-- .../airflow/ui/openapi-gen/queries/common.ts | 5 +- .../ui/openapi-gen/queries/ensureQueryData.ts | 6 +- .../ui/openapi-gen/queries/prefetch.ts | 6 +- .../airflow/ui/openapi-gen/queries/queries.ts | 6 +- .../ui/openapi-gen/queries/suspense.ts | 6 +- .../ui/openapi-gen/requests/services.gen.ts | 2 + .../ui/openapi-gen/requests/types.gen.ts | 1 + .../core_api/routes/ui/test_deadlines.py | 67 +++++++++++++++++++ 10 files changed, 113 insertions(+), 16 deletions(-) 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 b6b45a2bacb01..2caff1f719fb6 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 @@ -853,6 +853,14 @@ paths: schema: type: string title: Dag Id + - name: version_number + in: query + required: false + schema: + anyOf: + - type: integer + - type: 'null' + title: Version Number - name: limit in: query required: false diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/deadlines.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/deadlines.py index 2d080609c3569..4127e98104e80 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/deadlines.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/deadlines.py @@ -41,6 +41,7 @@ ) from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc from airflow.api_fastapi.core_api.security import ReadableDagRunsFilterDep, requires_access_dag +from airflow.models.dag_version import DagVersion from airflow.models.dagrun import DagRun from airflow.models.deadline import Deadline from airflow.models.deadline_alert import DeadlineAlert @@ -171,22 +172,31 @@ def get_dag_deadline_alerts( ).dynamic_depends(default="created_at") ), ], + version_number: int | None = None, ) -> DeadlineAlertCollectionResponse: """Get all deadline alerts defined on a Dag.""" - serialized_dag = session.scalar( - select(SerializedDagModel) + serialized_dag_select = ( + select(SerializedDagModel.id) + .join(DagVersion, SerializedDagModel.dag_version_id == DagVersion.id) .where(SerializedDagModel.dag_id == dag_id) - .order_by(SerializedDagModel.id.desc()) ) + if version_number is None: + serialized_dag_select = serialized_dag_select.order_by(DagVersion.version_number.desc()).limit(1) + not_found_detail = f"Dag with id {dag_id} was not found" + else: + serialized_dag_select = serialized_dag_select.where(DagVersion.version_number == version_number) + not_found_detail = f"Dag with id {dag_id} and version number {version_number} was not found" - if not serialized_dag: + serialized_dag_id = session.scalar(serialized_dag_select) + + if not serialized_dag_id: raise HTTPException( status.HTTP_404_NOT_FOUND, - f"Dag with id {dag_id} was not found", + not_found_detail, ) query = select(DeadlineAlert).where( - DeadlineAlert.serialized_dag_id == serialized_dag.id, + DeadlineAlert.serialized_dag_id == serialized_dag_id, ) alerts_select, total_entries = paginated_select( diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts index eb45a01065f5a..a6a00c9d84d5c 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts @@ -954,12 +954,13 @@ export const UseDeadlinesServiceGetDeadlinesKeyFn = ({ dagId, dagRunId, deadline export type DeadlinesServiceGetDagDeadlineAlertsDefaultResponse = Awaited>; export type DeadlinesServiceGetDagDeadlineAlertsQueryResult = UseQueryResult; export const useDeadlinesServiceGetDagDeadlineAlertsKey = "DeadlinesServiceGetDagDeadlineAlerts"; -export const UseDeadlinesServiceGetDagDeadlineAlertsKeyFn = ({ dagId, limit, offset, orderBy }: { +export const UseDeadlinesServiceGetDagDeadlineAlertsKeyFn = ({ dagId, limit, offset, orderBy, versionNumber }: { dagId: string; limit?: number; offset?: number; orderBy?: string[]; -}, queryKey?: Array) => [useDeadlinesServiceGetDagDeadlineAlertsKey, ...(queryKey ?? [{ dagId, limit, offset, orderBy }])]; + versionNumber?: number; +}, queryKey?: Array) => [useDeadlinesServiceGetDagDeadlineAlertsKey, ...(queryKey ?? [{ dagId, limit, offset, orderBy, versionNumber }])]; export type StructureServiceStructureDataDefaultResponse = Awaited>; export type StructureServiceStructureDataQueryResult = UseQueryResult; export const useStructureServiceStructureDataKey = "StructureServiceStructureData"; diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts index 230bc8340fa28..79e656190cd6c 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts @@ -1835,18 +1835,20 @@ export const ensureUseDeadlinesServiceGetDeadlinesData = (queryClient: QueryClie * Get all deadline alerts defined on a Dag. * @param data The data for the request. * @param data.dagId +* @param data.versionNumber * @param data.limit * @param data.offset * @param data.orderBy Attributes to order by, multi criteria sort is supported. Prefix with `-` for descending order. Supported attributes: `id, created_at, name` * @returns DeadlineAlertCollectionResponse Successful Response * @throws ApiError */ -export const ensureUseDeadlinesServiceGetDagDeadlineAlertsData = (queryClient: QueryClient, { dagId, limit, offset, orderBy }: { +export const ensureUseDeadlinesServiceGetDagDeadlineAlertsData = (queryClient: QueryClient, { dagId, limit, offset, orderBy, versionNumber }: { dagId: string; limit?: number; offset?: number; orderBy?: string[]; -}) => queryClient.ensureQueryData({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy }), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy }) }); + versionNumber?: number; +}) => queryClient.ensureQueryData({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy, versionNumber }), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy, versionNumber }) }); /** * Structure Data * Get Structure Data. diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts index 38d37e4a5cc5c..f0c3e04621759 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts @@ -1835,18 +1835,20 @@ export const prefetchUseDeadlinesServiceGetDeadlines = (queryClient: QueryClient * Get all deadline alerts defined on a Dag. * @param data The data for the request. * @param data.dagId +* @param data.versionNumber * @param data.limit * @param data.offset * @param data.orderBy Attributes to order by, multi criteria sort is supported. Prefix with `-` for descending order. Supported attributes: `id, created_at, name` * @returns DeadlineAlertCollectionResponse Successful Response * @throws ApiError */ -export const prefetchUseDeadlinesServiceGetDagDeadlineAlerts = (queryClient: QueryClient, { dagId, limit, offset, orderBy }: { +export const prefetchUseDeadlinesServiceGetDagDeadlineAlerts = (queryClient: QueryClient, { dagId, limit, offset, orderBy, versionNumber }: { dagId: string; limit?: number; offset?: number; orderBy?: string[]; -}) => queryClient.prefetchQuery({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy }), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy }) }); + versionNumber?: number; +}) => queryClient.prefetchQuery({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy, versionNumber }), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy, versionNumber }) }); /** * Structure Data * Get Structure Data. diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts index 336bb8516eeee..250cf1af4e299 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts @@ -1835,18 +1835,20 @@ export const useDeadlinesServiceGetDeadlines = = unknown[]>({ dagId, limit, offset, orderBy }: { +export const useDeadlinesServiceGetDagDeadlineAlerts = = unknown[]>({ dagId, limit, offset, orderBy, versionNumber }: { dagId: string; limit?: number; offset?: number; orderBy?: string[]; -}, queryKey?: TQueryKey, options?: Omit, "queryKey" | "queryFn">) => useQuery({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy }, queryKey), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy }) as TData, ...options }); + versionNumber?: number; +}, queryKey?: TQueryKey, options?: Omit, "queryKey" | "queryFn">) => useQuery({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy, versionNumber }, queryKey), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy, versionNumber }) as TData, ...options }); /** * Structure Data * Get Structure Data. diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts index 00e04931bebae..f39911e7ca37b 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts @@ -1835,18 +1835,20 @@ export const useDeadlinesServiceGetDeadlinesSuspense = = unknown[]>({ dagId, limit, offset, orderBy }: { +export const useDeadlinesServiceGetDagDeadlineAlertsSuspense = = unknown[]>({ dagId, limit, offset, orderBy, versionNumber }: { dagId: string; limit?: number; offset?: number; orderBy?: string[]; -}, queryKey?: TQueryKey, options?: Omit, "queryKey" | "queryFn">) => useSuspenseQuery({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy }, queryKey), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy }) as TData, ...options }); + versionNumber?: number; +}, queryKey?: TQueryKey, options?: Omit, "queryKey" | "queryFn">) => useSuspenseQuery({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy, versionNumber }, queryKey), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy, versionNumber }) as TData, ...options }); /** * Structure Data * Get Structure Data. diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts index bbb5cfb65ef01..cdb0bdacf1c5b 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts @@ -4792,6 +4792,7 @@ export class DeadlinesService { * Get all deadline alerts defined on a Dag. * @param data The data for the request. * @param data.dagId + * @param data.versionNumber * @param data.limit * @param data.offset * @param data.orderBy Attributes to order by, multi criteria sort is supported. Prefix with `-` for descending order. Supported attributes: `id, created_at, name` @@ -4806,6 +4807,7 @@ export class DeadlinesService { dag_id: data.dagId }, query: { + version_number: data.versionNumber, limit: data.limit, offset: data.offset, order_by: data.orderBy 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 be10ca6ccc17c..41a49101b8324 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 @@ -4482,6 +4482,7 @@ export type GetDagDeadlineAlertsData = { * Attributes to order by, multi criteria sort is supported. Prefix with `-` for descending order. Supported attributes: `id, created_at, name` */ orderBy?: Array<(string)>; + versionNumber?: number | null; }; export type GetDagDeadlineAlertsResponse = DeadlineAlertCollectionResponse; diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_deadlines.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_deadlines.py index 6c83044e17c3f..369729118c6c2 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_deadlines.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_deadlines.py @@ -18,11 +18,15 @@ from __future__ import annotations from datetime import timedelta +from unittest.mock import Mock import pytest +from fastapi import HTTPException from sqlalchemy import select from airflow._shared.timezones import timezone +from airflow.api_fastapi.core_api.routes.ui.deadlines import get_dag_deadline_alerts +from airflow.models.dag_version import DagVersion from airflow.models.deadline import Deadline from airflow.models.deadline_alert import DeadlineAlert from airflow.models.serialized_dag import SerializedDagModel @@ -475,6 +479,16 @@ def test_should_response_403(self, unauthorized_test_client): class TestGetDagDeadlineAlerts: """Tests for GET /dags/{dag_id}/deadlineAlerts.""" + def test_limits_serialized_dag_lookup(self): + session = Mock() + session.scalar.return_value = None + + with pytest.raises(HTTPException, match=f"Dag with id {DAG_ID} was not found"): + get_dag_deadline_alerts(DAG_ID, session, limit=100, offset=0, order_by=Mock()) + + statement = session.scalar.call_args.args[0] + assert "LIMIT 1" in str(statement.compile(compile_kwargs={"literal_binds": True})) + def test_returns_deadline_alerts_for_dag(self, test_client): """Returns all deadline alerts defined on the DAG.""" response = test_client.get(f"/dags/{DAG_ID}/deadlineAlerts") @@ -509,6 +523,59 @@ def test_should_response_404_for_nonexistent_dag(self, test_client): response = test_client.get("/dags/nonexistent_dag/deadlineAlerts") assert response.status_code == 404 + def _make_two_versions_with_alerts(self, dag_maker, session): + """Seed a Dag with two serialized versions, each carrying a distinct alert.""" + dag_id = "dag_versioned_alerts" + with dag_maker(dag_id, serialized=True, session=session): + EmptyOperator(task_id="task") + dag_maker.sync_dagbag_to_db() + # A structural change forces a new serialized version. + with dag_maker(dag_id, serialized=True, session=session): + EmptyOperator(task_id="task") + EmptyOperator(task_id="task2") + dag_maker.sync_dagbag_to_db() + session.commit() + + rows = session.execute( + select(SerializedDagModel, DagVersion.version_number) + .join(DagVersion, SerializedDagModel.dag_version_id == DagVersion.id) + .where(SerializedDagModel.dag_id == dag_id) + ).all() + serdag_by_version = {version_number: serdag for serdag, version_number in rows} + for version_number, serdag in serdag_by_version.items(): + session.add( + DeadlineAlert( + serialized_dag_id=serdag.id, + name=f"v{version_number}_alert", + reference=DeadlineReference.DAGRUN_QUEUED_AT.serialize_reference(), + interval=3600.0, + callback_def={"path": _CALLBACK_PATH}, + ) + ) + session.commit() + return dag_id + + @pytest.mark.parametrize( + ("params", "expected_alert"), + [ + pytest.param({}, "v2_alert", id="default_returns_latest_version"), + pytest.param({"version_number": 1}, "v1_alert", id="older_version"), + pytest.param({"version_number": 2}, "v2_alert", id="latest_version"), + ], + ) + def test_version_number_scopes_alerts(self, test_client, dag_maker, session, params, expected_alert): + dag_id = self._make_two_versions_with_alerts(dag_maker, session) + response = test_client.get(f"/dags/{dag_id}/deadlineAlerts", params=params) + assert response.status_code == 200 + data = response.json() + assert [alert["name"] for alert in data["deadline_alerts"]] == [expected_alert] + + def test_unknown_version_number_returns_404(self, test_client, dag_maker, session): + dag_id = self._make_two_versions_with_alerts(dag_maker, session) + response = test_client.get(f"/dags/{dag_id}/deadlineAlerts", params={"version_number": 999}) + assert response.status_code == 404 + assert response.json()["detail"] == f"Dag with id {dag_id} and version number 999 was not found" + @pytest.mark.parametrize("order_by", ["interval", "-interval"]) def test_order_by_interval_is_rejected(self, test_client, order_by): """``interval`` is a serialized-JSON column (timedelta/VariableInterval dict), so DB-level