Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/routes/dag_runs.py: 54%

93 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 

20import logging 

21from typing import Annotated 

22 

23from cadwyn import VersionedAPIRouter 

24from fastapi import HTTPException, Query, status 

25from sqlalchemy import func, select 

26from sqlalchemy.exc import NoResultFound 

27 

28from airflow.api.common.trigger_dag import trigger_dag 

29from airflow.api_fastapi.common.dagbag import DagBagDep, get_dag_for_run, resolve_run_on_latest_version 

30from airflow.api_fastapi.common.db.common import SessionDep 

31from airflow.api_fastapi.common.types import UtcDateTime 

32from airflow.api_fastapi.compat import HTTP_422_UNPROCESSABLE_CONTENT 

33from airflow.api_fastapi.execution_api.datamodels.dagrun import DagRunStateResponse, TriggerDAGRunPayload 

34from airflow.api_fastapi.execution_api.datamodels.taskinstance import DagRun 

35from airflow.api_fastapi.execution_api.datamodels.token import TIToken 

36from airflow.api_fastapi.execution_api.security import CurrentTIToken 

37from airflow.exceptions import DagNotPartitionedError, DagRunAlreadyExists, InvalidPartitionKeyError 

38from airflow.models.dag import DagModel 

39from airflow.models.dagrun import DagRun as DagRunModel 

40from airflow.models.taskinstance import TaskInstance 

41from airflow.utils.state import DagRunState 

42from airflow.utils.types import DagRunTriggeredByType, DagRunType 

43 

44router = VersionedAPIRouter() 

45 

46log = logging.getLogger(__name__) 

47 

48 

49@router.only_exists_in_older_versions 

50@router.get("/{dag_id}/previous") 

51def get_previous_dagrun_compat( 

52 dag_id: str, 

53 logical_date: UtcDateTime, 

54 session: SessionDep, 

55 state: DagRunState | None = None, 

56): 

57 """ 

58 Redirect old previous dag run request to the new endpoint. 

59 

60 This endpoint must be put before ``get_dag_run`` so not to be shadowed. 

61 Newer client versions would not see this endpoint, and be routed to 

62 ``get_dag_run`` below instead. 

63 """ 

64 return get_previous_dagrun(dag_id, logical_date, session, state) 

65 

66 

67@router.get( 

68 "/{dag_id}/{run_id}", 

69 responses={status.HTTP_404_NOT_FOUND: {"description": "Dag run not found"}}, 

70) 

71def get_dag_run(dag_id: str, run_id: str, session: SessionDep) -> DagRun: 

72 """Get detail of a Dag run.""" 

73 dr = session.scalar(select(DagRunModel).where(DagRunModel.dag_id == dag_id, DagRunModel.run_id == run_id)) 

74 if dr is None: 

75 raise HTTPException( 

76 status.HTTP_404_NOT_FOUND, 

77 detail={ 

78 "reason": "not_found", 

79 "message": f"Dag run with dag_id '{dag_id}' and run_id '{run_id}' was not found", 

80 }, 

81 ) 

82 return DagRun.model_validate(dr) 

83 

84 

85@router.post( 

86 "/{dag_id}/{run_id}", 

87 status_code=status.HTTP_204_NO_CONTENT, 

88 responses={ 

89 status.HTTP_400_BAD_REQUEST: {"description": "Dag has import errors or the run type is not allowed"}, 

90 status.HTTP_404_NOT_FOUND: {"description": "Dag not found for the given dag_id"}, 

91 status.HTTP_409_CONFLICT: {"description": "Dag run already exists for the given dag_id"}, 

92 HTTP_422_UNPROCESSABLE_CONTENT: {"description": "Invalid payload"}, 

93 }, 

94) 

95def trigger_dag_run( 

96 dag_id: str, 

97 run_id: str, 

98 payload: TriggerDAGRunPayload, 

99 session: SessionDep, 

100 token: TIToken = CurrentTIToken, 

101) -> None: 

102 """Trigger a Dag run.""" 

103 dm = session.scalar(select(DagModel).where(~DagModel.is_stale, DagModel.dag_id == dag_id).limit(1)) 

104 if not dm: 104 ↛ 105line 104 didn't jump to line 105 because the condition on line 104 was never true

105 raise HTTPException( 

106 status.HTTP_404_NOT_FOUND, 

107 detail={"reason": "not_found", "message": f"Dag with dag_id: '{dag_id}' not found"}, 

108 ) 

109 

110 if dm.has_import_errors: 110 ↛ 111line 110 didn't jump to line 111 because the condition on line 110 was never true

111 raise HTTPException( 

112 status.HTTP_400_BAD_REQUEST, 

113 detail={ 

114 "reason": "import_errors", 

115 "message": f"Dag with dag_id '{dag_id}' has import errors and cannot be triggered", 

116 }, 

117 ) 

118 

119 if dm.allowed_run_types is not None and DagRunType.OPERATOR_TRIGGERED not in dm.allowed_run_types: 119 ↛ 120line 119 didn't jump to line 120 because the condition on line 119 was never true

120 raise HTTPException( 

121 status.HTTP_400_BAD_REQUEST, 

122 detail={ 

123 "reason": "denied_run_type", 

124 "message": f"Dag with dag_id '{dag_id}' does not allow operator-triggered runs", 

125 }, 

126 ) 

127 

128 # Inherit triggering_user_name from the calling task's DagRun so chains of 

129 # TriggerDagRunOperator preserve the original human user across child runs. 

130 parent_ti = session.get(TaskInstance, token.id) 

131 triggering_user_name = parent_ti.dag_run.triggering_user_name if parent_ti else None 

132 

133 try: 

134 trigger_dag( 

135 dag_id=dag_id, 

136 run_id=run_id, 

137 run_type=DagRunType.OPERATOR_TRIGGERED, 

138 conf=payload.conf, 

139 logical_date=payload.logical_date, 

140 run_after=payload.run_after, 

141 triggered_by=DagRunTriggeredByType.OPERATOR, 

142 triggering_user_name=triggering_user_name, 

143 replace_microseconds=False, 

144 partition_key=payload.partition_key, 

145 note=payload.note, 

146 session=session, 

147 ) 

148 except DagRunAlreadyExists: 

149 raise HTTPException( 

150 status.HTTP_409_CONFLICT, 

151 detail={ 

152 "reason": "already_exists", 

153 "message": f"A run already exists for Dag '{dag_id}' with run_id '{run_id}'", 

154 }, 

155 ) 

156 except DagNotPartitionedError as e: 

157 raise HTTPException( 

158 status.HTTP_400_BAD_REQUEST, 

159 detail={"reason": "not_partitioned", "message": str(e)}, 

160 ) 

161 except InvalidPartitionKeyError as e: 

162 raise HTTPException( 

163 status.HTTP_400_BAD_REQUEST, 

164 detail={"reason": "invalid_partition_key", "message": str(e)}, 

165 ) from e 

166 except ValueError as e: 

167 raise HTTPException( 

168 status.HTTP_400_BAD_REQUEST, 

169 detail={ 

170 "reason": "value_error", 

171 "message": str(e), 

172 }, 

173 ) from e 

174 

175 

176@router.post( 

177 "/{dag_id}/{run_id}/clear", 

178 status_code=status.HTTP_204_NO_CONTENT, 

179 responses={ 

180 status.HTTP_400_BAD_REQUEST: {"description": "Dag has import errors and cannot be triggered"}, 

181 status.HTTP_404_NOT_FOUND: {"description": "Dag not found for the given dag_id"}, 

182 HTTP_422_UNPROCESSABLE_CONTENT: {"description": "Invalid payload"}, 

183 }, 

184) 

185def clear_dag_run( 

186 dag_id: str, 

187 run_id: str, 

188 session: SessionDep, 

189 dag_bag: DagBagDep, 

190) -> None: 

191 """Clear a Dag run.""" 

192 dm = session.scalar(select(DagModel).where(~DagModel.is_stale, DagModel.dag_id == dag_id).limit(1)) 

193 if not dm: 193 ↛ 194line 193 didn't jump to line 194 because the condition on line 193 was never true

194 raise HTTPException( 

195 status.HTTP_404_NOT_FOUND, 

196 detail={"reason": "not_found", "message": f"Dag with dag_id: '{dag_id}' not found"}, 

197 ) 

198 

199 if dm.has_import_errors: 199 ↛ 200line 199 didn't jump to line 200 because the condition on line 199 was never true

200 raise HTTPException( 

201 status.HTTP_400_BAD_REQUEST, 

202 detail={ 

203 "reason": "import_errors", 

204 "message": f"Dag with dag_id '{dag_id}' has import errors and cannot be triggered", 

205 }, 

206 ) 

207 

208 dag_run = session.scalar( 

209 select(DagRunModel).where(DagRunModel.dag_id == dag_id, DagRunModel.run_id == run_id) 

210 ) 

211 if dag_run is None: 211 ↛ 216line 211 didn't jump to line 216 because the condition on line 211 was always true

212 raise HTTPException( 

213 status.HTTP_404_NOT_FOUND, 

214 detail={"reason": "not_found", "message": f"Dag run with run_id: '{run_id}' not found"}, 

215 ) 

216 dag = get_dag_for_run(dag_bag, dag_run=dag_run, session=session) 

217 

218 resolved_run_on_latest = resolve_run_on_latest_version(None, dag_id, session) 

219 dag.clear(run_id=run_id, run_on_latest_version=resolved_run_on_latest) 

220 

221 

222@router.get( 

223 "/{dag_id}/{run_id}/state", 

224 responses={status.HTTP_404_NOT_FOUND: {"description": "Dag run not found"}}, 

225) 

226def get_dagrun_state( 

227 dag_id: str, 

228 run_id: str, 

229 session: SessionDep, 

230) -> DagRunStateResponse: 

231 """Get a Dag run State.""" 

232 try: 

233 state: DagRunState = session.scalars( 

234 select(DagRunModel.state).where(DagRunModel.dag_id == dag_id, DagRunModel.run_id == run_id) 

235 ).one() 

236 except NoResultFound: 

237 raise HTTPException( 

238 status.HTTP_404_NOT_FOUND, 

239 detail={ 

240 "reason": "not_found", 

241 "message": f"Dag run with dag_id '{dag_id}' and run_id '{run_id}' was not found", 

242 }, 

243 ) 

244 return DagRunStateResponse(state=state) 

245 

246 

247@router.get("/count", status_code=status.HTTP_200_OK) 

248def get_dr_count( 

249 dag_id: str, 

250 session: SessionDep, 

251 logical_dates: Annotated[list[UtcDateTime] | None, Query()] = None, 

252 run_ids: Annotated[list[str] | None, Query()] = None, 

253 states: Annotated[list[str] | None, Query()] = None, 

254) -> int: 

255 """Get the count of Dag runs matching the given criteria.""" 

256 stmt = select(func.count()).select_from(DagRunModel).where(DagRunModel.dag_id == dag_id) 

257 if logical_dates: 

258 stmt = stmt.where(DagRunModel.logical_date.in_(logical_dates)) 

259 if run_ids: 

260 stmt = stmt.where(DagRunModel.run_id.in_(run_ids)) 

261 if states: 

262 stmt = stmt.where(DagRunModel.state.in_(states)) 

263 return session.scalar(stmt) or 0 

264 

265 

266@router.get("/previous", status_code=status.HTTP_200_OK) 

267def get_previous_dagrun( 

268 dag_id: str, 

269 logical_date: UtcDateTime, 

270 session: SessionDep, 

271 state: Annotated[DagRunState | None, Query()] = None, 

272) -> DagRun | None: 

273 """Get the previous Dag run before the given logical date, optionally filtered by state.""" 

274 stmt = ( 

275 select(DagRunModel) 

276 .where(DagRunModel.dag_id == dag_id, DagRunModel.logical_date < logical_date) 

277 .order_by(DagRunModel.logical_date.desc()) 

278 .limit(1) 

279 ) 

280 if state: 

281 stmt = stmt.where(DagRunModel.state == state) 

282 if not (dag_run := session.scalar(stmt)): 

283 return None 

284 return DagRun.model_validate(dag_run)