Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/ui/deadlines.py: 35%
49 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
1# Licensed to the Apache Software Foundation (ASF) under one
2# or more contributor license agreements. See the NOTICE file
3# distributed with this work for additional information
4# regarding copyright ownership. The ASF licenses this file
5# to you under the Apache License, Version 2.0 (the
6# "License"); you may not use this file except in compliance
7# with the License. You may obtain a copy of the License at
8#
9# http://www.apache.org/licenses/LICENSE-2.0
10#
11# Unless required by applicable law or agreed to in writing,
12# software distributed under the License is distributed on an
13# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14# KIND, either express or implied. See the License for the
15# specific language governing permissions and limitations
16# under the License.
18from __future__ import annotations
20from typing import Annotated
22from fastapi import Depends, HTTPException, status
23from sqlalchemy import select
24from sqlalchemy.orm import contains_eager, noload
26from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity
27from airflow.api_fastapi.common.db.common import SessionDep, paginated_select
28from airflow.api_fastapi.common.parameters import (
29 FilterParam,
30 QueryLimit,
31 QueryOffset,
32 RangeFilter,
33 SortParam,
34 datetime_range_filter_factory,
35 filter_param_factory,
36)
37from airflow.api_fastapi.common.router import AirflowRouter
38from airflow.api_fastapi.core_api.datamodels.ui.deadline import (
39 DeadlineAlertCollectionResponse,
40 DeadlineCollectionResponse,
41)
42from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
43from airflow.api_fastapi.core_api.security import ReadableDagRunsFilterDep, requires_access_dag
44from airflow.models.dag_version import DagVersion
45from airflow.models.dagrun import DagRun
46from airflow.models.deadline import Deadline
47from airflow.models.deadline_alert import DeadlineAlert
48from airflow.models.serialized_dag import SerializedDagModel
50deadlines_router = AirflowRouter(prefix="/dags/{dag_id}", tags=["Deadlines"])
53@deadlines_router.get(
54 "/dagRuns/{dag_run_id}/deadlines",
55 responses=create_openapi_http_exception_doc(
56 [
57 status.HTTP_400_BAD_REQUEST,
58 status.HTTP_404_NOT_FOUND,
59 ]
60 ),
61 dependencies=[
62 Depends(
63 requires_access_dag(
64 method="GET",
65 access_entity=DagAccessEntity.RUN,
66 )
67 ),
68 ],
69)
70def get_deadlines(
71 dag_id: str,
72 dag_run_id: str,
73 session: SessionDep,
74 limit: QueryLimit,
75 offset: QueryOffset,
76 readable_dag_runs_filter: ReadableDagRunsFilterDep,
77 order_by: Annotated[
78 SortParam,
79 Depends(
80 SortParam(
81 ["id", "deadline_time", "created_at", "last_updated_at", "missed"],
82 Deadline,
83 to_replace={
84 "dag_id": DagRun.dag_id,
85 "dag_run_id": DagRun.run_id,
86 "alert_name": DeadlineAlert.name,
87 },
88 ).dynamic_depends(default="deadline_time")
89 ),
90 ],
91 missed: Annotated[FilterParam[bool | None], Depends(filter_param_factory(Deadline.missed, bool | None))],
92 deadline_time: Annotated[RangeFilter, Depends(datetime_range_filter_factory("deadline_time", Deadline))],
93 last_updated_at: Annotated[
94 RangeFilter, Depends(datetime_range_filter_factory("last_updated_at", Deadline))
95 ],
96) -> DeadlineCollectionResponse:
97 """
98 Get deadlines for a Dag run.
100 This endpoint allows specifying `~` as the dag_id and dag_run_id to retrieve Deadlines for all
101 Dags and Dag runs.
102 """
103 query = (
104 select(Deadline)
105 .join(Deadline.dagrun)
106 .outerjoin(Deadline.deadline_alert)
107 .options(
108 contains_eager(Deadline.dagrun).options(noload(DagRun.deadlines)),
109 contains_eager(Deadline.deadline_alert),
110 noload(Deadline.callback),
111 )
112 )
114 if dag_run_id != "~":
115 if dag_id == "~":
116 raise HTTPException(
117 status.HTTP_400_BAD_REQUEST,
118 "dag_id is required when dag_run_id is specified",
119 )
120 query = query.where(DagRun.dag_id == dag_id, DagRun.run_id == dag_run_id)
121 elif dag_id != "~":
122 query = query.where(DagRun.dag_id == dag_id)
124 deadlines_select, total_entries = paginated_select(
125 statement=query,
126 filters=[readable_dag_runs_filter, missed, deadline_time, last_updated_at],
127 order_by=order_by,
128 offset=offset,
129 limit=limit,
130 session=session,
131 )
133 deadlines = session.scalars(deadlines_select)
135 if dag_run_id != "~" and total_entries == 0:
136 dag_run = session.scalar(select(DagRun).where(DagRun.dag_id == dag_id, DagRun.run_id == dag_run_id))
137 if not dag_run:
138 raise HTTPException(
139 status.HTTP_404_NOT_FOUND,
140 f"DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` was not found",
141 )
143 return DeadlineCollectionResponse(deadlines=deadlines, total_entries=total_entries)
146@deadlines_router.get(
147 "/deadlineAlerts",
148 responses=create_openapi_http_exception_doc(
149 [
150 status.HTTP_404_NOT_FOUND,
151 ]
152 ),
153 dependencies=[
154 Depends(
155 requires_access_dag(
156 method="GET",
157 )
158 ),
159 ],
160)
161def get_dag_deadline_alerts(
162 dag_id: str,
163 session: SessionDep,
164 limit: QueryLimit,
165 offset: QueryOffset,
166 order_by: Annotated[
167 SortParam,
168 Depends(
169 SortParam(
170 ["id", "created_at", "name"],
171 DeadlineAlert,
172 ).dynamic_depends(default="created_at")
173 ),
174 ],
175 version_number: int | None = None,
176) -> DeadlineAlertCollectionResponse:
177 """Get all deadline alerts defined on a Dag."""
178 serialized_dag_select = (
179 select(SerializedDagModel.id)
180 .join(DagVersion, SerializedDagModel.dag_version_id == DagVersion.id)
181 .where(SerializedDagModel.dag_id == dag_id)
182 )
183 if version_number is None:
184 serialized_dag_select = serialized_dag_select.order_by(DagVersion.version_number.desc()).limit(1)
185 not_found_detail = f"Dag with id {dag_id} was not found"
186 else:
187 serialized_dag_select = serialized_dag_select.where(DagVersion.version_number == version_number)
188 not_found_detail = f"Dag with id {dag_id} and version number {version_number} was not found"
190 serialized_dag_id = session.scalar(serialized_dag_select)
192 if not serialized_dag_id:
193 raise HTTPException(
194 status.HTTP_404_NOT_FOUND,
195 not_found_detail,
196 )
198 query = select(DeadlineAlert).where(
199 DeadlineAlert.serialized_dag_id == serialized_dag_id,
200 )
202 alerts_select, total_entries = paginated_select(
203 statement=query,
204 filters=None,
205 order_by=order_by,
206 offset=offset,
207 limit=limit,
208 session=session,
209 )
211 alerts = session.scalars(alerts_select)
213 return DeadlineAlertCollectionResponse(deadline_alerts=alerts, total_entries=total_entries)