From 6bbb4e80b8926bc4212efb13efba2cd3ded0dbff Mon Sep 17 00:00:00 2001 From: imrichardwu Date: Thu, 30 Jul 2026 08:21:46 -0700 Subject: [PATCH 1/3] 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. --- .../api_fastapi/core_api/routes/ui/deadlines.py | 1 + .../core_api/routes/ui/test_deadlines.py | 14 ++++++++++++++ 2 files changed, 15 insertions(+) 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 06eda42ed8957..a5ebb85e05272 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 @@ -177,6 +177,7 @@ def get_dag_deadline_alerts( select(SerializedDagModel) .where(SerializedDagModel.dag_id == dag_id) .order_by(SerializedDagModel.id.desc()) + .limit(1) ) if not serialized_dag: 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 acadab3a6b16f..f8fd5690ad481 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 @@ -17,10 +17,14 @@ from __future__ import annotations +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.deadline import Deadline from airflow.models.deadline_alert import DeadlineAlert from airflow.models.serialized_dag import SerializedDagModel @@ -473,6 +477,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") From 8ce7f99b962a1b951a837f1e458cbb568b8ab21f Mon Sep 17 00:00:00 2001 From: imrichardwu Date: Thu, 30 Jul 2026 20:11:08 -0700 Subject: [PATCH 2/3] Sort imports in deadline alerts UI route tests --- .../tests/unit/api_fastapi/core_api/routes/ui/test_deadlines.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 929857098fa0f..bebf57a44ae24 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 @@ -17,8 +17,8 @@ from __future__ import annotations -from unittest.mock import Mock from datetime import timedelta +from unittest.mock import Mock import pytest from fastapi import HTTPException From fff634e4af7dbfaada96892f67bcafc5f98db717 Mon Sep 17 00:00:00 2001 From: imrichardwu Date: Fri, 31 Jul 2026 18:28:55 -0700 Subject: [PATCH 3/3] 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. --- .../core_api/openapi/_private_ui.yaml | 8 +++ .../core_api/routes/ui/deadlines.py | 23 +++++--- .../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 | 54 +++++++++++++++++++ 10 files changed, 100 insertions(+), 17 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 a6d31b9956bf8..58d097827eaf4 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 @@ -1172,6 +1172,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 511d4fa09ff65..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,23 +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()) - .limit(1) ) + 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 0b1fda69bba10..f11d2d808d1c1 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts @@ -1002,12 +1002,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 091573243dfdd..637e37a62f400 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts @@ -1928,18 +1928,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 2550f81601131..d7d3fecb52c44 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts @@ -1928,18 +1928,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 9da0d1e5ab10a..1b993afa9b01c 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts @@ -1928,18 +1928,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 2b694fe3b5a75..4d2dc42adcaba 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts @@ -1928,18 +1928,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 c31f14f957b33..88ba2d33d68a8 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 @@ -4915,6 +4915,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` @@ -4929,6 +4930,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 2c6a870e0e323..f19669c72fef8 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 @@ -4672,6 +4672,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 bebf57a44ae24..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 @@ -26,6 +26,7 @@ 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 @@ -522,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