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

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. 

17 

18from __future__ import annotations 

19 

20from typing import Annotated 

21 

22from fastapi import Depends, HTTPException, status 

23from sqlalchemy import select 

24from sqlalchemy.orm import contains_eager, noload 

25 

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 

49 

50deadlines_router = AirflowRouter(prefix="/dags/{dag_id}", tags=["Deadlines"]) 

51 

52 

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. 

99 

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 ) 

113 

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) 

123 

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 ) 

132 

133 deadlines = session.scalars(deadlines_select) 

134 

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 ) 

142 

143 return DeadlineCollectionResponse(deadlines=deadlines, total_entries=total_entries) 

144 

145 

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" 

189 

190 serialized_dag_id = session.scalar(serialized_dag_select) 

191 

192 if not serialized_dag_id: 

193 raise HTTPException( 

194 status.HTTP_404_NOT_FOUND, 

195 not_found_detail, 

196 ) 

197 

198 query = select(DeadlineAlert).where( 

199 DeadlineAlert.serialized_dag_id == serialized_dag_id, 

200 ) 

201 

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 ) 

210 

211 alerts = session.scalars(alerts_select) 

212 

213 return DeadlineAlertCollectionResponse(deadline_alerts=alerts, total_entries=total_entries)