Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down
5 changes: 3 additions & 2 deletions airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1002,12 +1002,13 @@ export const UseDeadlinesServiceGetDeadlinesKeyFn = ({ dagId, dagRunId, deadline
export type DeadlinesServiceGetDagDeadlineAlertsDefaultResponse = Awaited<ReturnType<typeof DeadlinesService.getDagDeadlineAlerts>>;
export type DeadlinesServiceGetDagDeadlineAlertsQueryResult<TData = DeadlinesServiceGetDagDeadlineAlertsDefaultResponse, TError = unknown> = UseQueryResult<TData, TError>;
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<unknown>) => [useDeadlinesServiceGetDagDeadlineAlertsKey, ...(queryKey ?? [{ dagId, limit, offset, orderBy }])];
versionNumber?: number;
}, queryKey?: Array<unknown>) => [useDeadlinesServiceGetDagDeadlineAlertsKey, ...(queryKey ?? [{ dagId, limit, offset, orderBy, versionNumber }])];
export type StructureServiceStructureDataDefaultResponse = Awaited<ReturnType<typeof StructureService.structureData>>;
export type StructureServiceStructureDataQueryResult<TData = StructureServiceStructureDataDefaultResponse, TError = unknown> = UseQueryResult<TData, TError>;
export const useStructureServiceStructureDataKey = "StructureServiceStructureData";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
6 changes: 4 additions & 2 deletions airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
6 changes: 4 additions & 2 deletions airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1928,18 +1928,20 @@ export const useDeadlinesServiceGetDeadlines = <TData = Common.DeadlinesServiceG
* 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 useDeadlinesServiceGetDagDeadlineAlerts = <TData = Common.DeadlinesServiceGetDagDeadlineAlertsDefaultResponse, TError = unknown, TQueryKey extends Array<unknown> = unknown[]>({ dagId, limit, offset, orderBy }: {
export const useDeadlinesServiceGetDagDeadlineAlerts = <TData = Common.DeadlinesServiceGetDagDeadlineAlertsDefaultResponse, TError = unknown, TQueryKey extends Array<unknown> = unknown[]>({ dagId, limit, offset, orderBy, versionNumber }: {
dagId: string;
limit?: number;
offset?: number;
orderBy?: string[];
}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, "queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy }, queryKey), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy }) as TData, ...options });
versionNumber?: number;
}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, "queryKey" | "queryFn">) => useQuery<TData, TError>({ 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.
Expand Down
6 changes: 4 additions & 2 deletions airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1928,18 +1928,20 @@ export const useDeadlinesServiceGetDeadlinesSuspense = <TData = Common.Deadlines
* 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 useDeadlinesServiceGetDagDeadlineAlertsSuspense = <TData = Common.DeadlinesServiceGetDagDeadlineAlertsDefaultResponse, TError = unknown, TQueryKey extends Array<unknown> = unknown[]>({ dagId, limit, offset, orderBy }: {
export const useDeadlinesServiceGetDagDeadlineAlertsSuspense = <TData = Common.DeadlinesServiceGetDagDeadlineAlertsDefaultResponse, TError = unknown, TQueryKey extends Array<unknown> = unknown[]>({ dagId, limit, offset, orderBy, versionNumber }: {
dagId: string;
limit?: number;
offset?: number;
orderBy?: string[];
}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, "queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey: Common.UseDeadlinesServiceGetDagDeadlineAlertsKeyFn({ dagId, limit, offset, orderBy }, queryKey), queryFn: () => DeadlinesService.getDagDeadlineAlerts({ dagId, limit, offset, orderBy }) as TData, ...options });
versionNumber?: number;
}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>, "queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ 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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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
Expand Down
Loading