Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/datamodels/dag_run.py: 99%

156 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 collections.abc import Iterable 

21from datetime import datetime 

22from enum import Enum 

23from typing import TYPE_CHECKING, Any 

24 

25from pydantic import AliasPath, AwareDatetime, Field, NonNegativeInt, model_validator 

26 

27from airflow._shared.timezones import timezone 

28from airflow.api_fastapi.core_api.base import BaseModel, StrictBaseModel 

29from airflow.api_fastapi.core_api.datamodels.dag_versions import DagVersionResponse 

30from airflow.timetables.base import DataInterval 

31from airflow.utils.state import DagRunState 

32from airflow.utils.types import DagRunTriggeredByType, DagRunType 

33 

34if TYPE_CHECKING: 34 ↛ 35line 34 didn't jump to line 35 because the condition on line 34 was never true

35 from airflow.serialization.definitions.dag import SerializedDAG 

36 

37 

38class DagRunMutableStates(str, Enum): 

39 """Dag Run states from which the run may be mutated (patched, deleted).""" 

40 

41 QUEUED = DagRunState.QUEUED 

42 SUCCESS = DagRunState.SUCCESS 

43 FAILED = DagRunState.FAILED 

44 

45 

46class DAGRunPatchBody(StrictBaseModel): 

47 """Dag Run Serializer for PATCH requests.""" 

48 

49 state: DagRunMutableStates | None = None 

50 note: str | None = Field(None, max_length=1000) 

51 

52 

53class BulkDAGRunBody(StrictBaseModel): 

54 """Request body for bulk operations on Dag Runs.""" 

55 

56 dag_run_id: str 

57 dag_id: str | None = None 

58 state: DagRunMutableStates | None = None 

59 note: str | None = Field(None, max_length=1000) 

60 

61 

62class PartitionSelectorMixin(StrictBaseModel): 

63 """Partition filter fields shared by bulk-clear and clearPartitions bodies.""" 

64 

65 partition_key: str | None = Field( 

66 default=None, 

67 description="Select runs by exact partition key match. Mutually exclusive with the other partition selectors.", 

68 ) 

69 partition_date_start: datetime | None = Field( 

70 default=None, 

71 description=( 

72 "Inclusive start of the partition date window. " 

73 "The value is interpreted in the Dag's timetable timezone. " 

74 "Mutually exclusive with the other partition selectors." 

75 ), 

76 ) 

77 partition_date_end: datetime | None = Field( 

78 default=None, 

79 description=( 

80 "Inclusive end of the partition date window. " 

81 "The value is interpreted in the Dag's timetable timezone. " 

82 "Mutually exclusive with the other partition selectors." 

83 ), 

84 ) 

85 

86 @property 

87 def has_partition_selectors(self) -> bool: 

88 return ( 

89 self.partition_key is not None 

90 or self.partition_date_start is not None 

91 or self.partition_date_end is not None 

92 ) 

93 

94 def _validate_partition_date_window_order(self) -> None: 

95 if ( 

96 self.partition_date_start is not None 

97 and self.partition_date_end is not None 

98 and self.partition_date_start > self.partition_date_end 

99 ): 

100 raise ValueError("partition_date_start must be on or before partition_date_end.") 

101 

102 def _check_exactly_one_selection_mode( 

103 self, *, extra_selector_active: bool, extra_selector_name: str 

104 ) -> None: 

105 has_partition_key = self.partition_key is not None 

106 has_partition_date_window = ( 

107 self.partition_date_start is not None or self.partition_date_end is not None 

108 ) 

109 modes_active = sum([extra_selector_active, has_partition_key, has_partition_date_window]) 

110 if modes_active != 1: 

111 raise ValueError( 

112 f"Exactly one of {extra_selector_name}, partition_key, or a partition date window " 

113 "(partition_date_start / partition_date_end) must be provided." 

114 ) 

115 self._validate_partition_date_window_order() 

116 

117 

118class BaseDAGRunClear(StrictBaseModel): 

119 """Shared options for the single-run and bulk Dag Run clear endpoints.""" 

120 

121 dry_run: bool = True 

122 only_failed: bool = False 

123 only_new: bool = Field( 

124 default=False, 

125 description="Only queue newly added tasks in the latest Dag version without clearing existing tasks.", 

126 ) 

127 run_on_latest_version: bool | None = Field( 

128 default=None, 

129 description="(Experimental) Run on the latest bundle version of the Dag after clearing. " 

130 "If not specified, falls back to the DAG-level ``rerun_with_latest_version`` parameter, " 

131 "then the ``[core] rerun_with_latest_version`` config option, " 

132 "and finally ``False``.", 

133 ) 

134 note: str | None = Field(default=None, max_length=1000) 

135 

136 @model_validator(mode="before") 

137 @classmethod 

138 def validate_only_new_only_failed_mutually_exclusive(cls, data: Any) -> Any: 

139 if data.get("only_new") and data.get("only_failed"): 

140 raise ValueError("only_new and only_failed are mutually exclusive") 

141 return data 

142 

143 

144class DAGRunClearBody(BaseDAGRunClear): 

145 """Dag Run serializer for clear endpoint body.""" 

146 

147 

148class BulkDAGRunClearBody(BaseDAGRunClear, PartitionSelectorMixin): 

149 """Request body for the bulk clear Dag Runs endpoint.""" 

150 

151 dag_runs: list[BulkDAGRunBody] = Field(default_factory=list) 

152 

153 @model_validator(mode="after") 

154 def validate_exactly_one_selection_mode(self) -> BulkDAGRunClearBody: 

155 self._check_exactly_one_selection_mode( 

156 extra_selector_active=bool(self.dag_runs), 

157 extra_selector_name="dag_runs (non-empty)", 

158 ) 

159 return self 

160 

161 

162class DAGRunResponse(BaseModel): 

163 """Dag Run serializer for responses.""" 

164 

165 dag_run_id: str = Field(validation_alias="run_id") 

166 dag_id: str 

167 logical_date: datetime | None 

168 queued_at: datetime | None 

169 start_date: datetime | None 

170 end_date: datetime | None 

171 duration: float | None 

172 data_interval_start: datetime | None 

173 data_interval_end: datetime | None 

174 run_after: datetime 

175 last_scheduling_decision: datetime | None 

176 run_type: DagRunType 

177 state: DagRunState 

178 triggered_by: DagRunTriggeredByType | None 

179 triggering_user_name: str | None 

180 conf: dict | None 

181 note: str | None 

182 dag_versions: list[DagVersionResponse] 

183 bundle_version: str | None 

184 dag_display_name: str = Field(validation_alias=AliasPath("dag_model", "dag_display_name")) 

185 partition_key: str | None 

186 partition_date: datetime | None 

187 

188 

189class DAGRunCollectionResponse(BaseModel): 

190 """ 

191 Dag Run collection response supporting both offset and cursor pagination. 

192 

193 A single flat model is used instead of a discriminated union 

194 (``Annotated[Offset | Cursor, Field(discriminator=...)]``) because 

195 the OpenAPI ``oneOf`` + ``discriminator`` construct is not handled 

196 correctly by ``@hey-api/openapi-ts`` / ``@7nohe/openapi-react-query-codegen``: 

197 return types degrade to ``unknown`` in JSDoc and can produce 

198 incorrect TypeScript types (see hey-api/openapi-ts#1613, #3270). 

199 """ 

200 

201 dag_runs: Iterable[DAGRunResponse] 

202 total_entries: int | None = Field( 

203 default=None, 

204 description="Number of matching items. For offset pagination this is the exact total. " 

205 "For cursor pagination it is capped at ``total_entries_limit``; a value equal to that " 

206 "limit means at least that many items match.", 

207 ) 

208 total_entries_limit: int | None = Field( 

209 default=None, 

210 description="Cap applied to ``total_entries`` under cursor pagination. ``null`` for offset " 

211 "pagination, where ``total_entries`` is exact.", 

212 ) 

213 next_cursor: str | None = Field( 

214 default=None, 

215 description="Token pointing to the next page. Populated for cursor pagination, " 

216 "``null`` when using offset pagination or when there is no next page.", 

217 ) 

218 previous_cursor: str | None = Field( 

219 default=None, 

220 description="Token pointing to the previous page. Populated for cursor pagination, " 

221 "``null`` when using offset pagination or when on the first page.", 

222 ) 

223 

224 

225class TriggerDAGRunPostBody(StrictBaseModel): 

226 """Trigger Dag Run Serializer for POST body.""" 

227 

228 dag_run_id: str | None = None 

229 data_interval_start: AwareDatetime | None = None 

230 data_interval_end: AwareDatetime | None = None 

231 logical_date: AwareDatetime | None 

232 run_after: datetime | None = Field(default_factory=timezone.utcnow) 

233 

234 conf: dict | None = Field(default_factory=dict) 

235 note: str | None = None 

236 partition_key: str | None = None 

237 

238 @model_validator(mode="after") 

239 def check_data_intervals(self): 

240 if (self.data_interval_start is None) != (self.data_interval_end is None): 

241 raise ValueError( 

242 "Either both data_interval_start and data_interval_end must be provided or both must be None" 

243 ) 

244 return self 

245 

246 def validate_context(self, dag: SerializedDAG) -> dict: 

247 dag.validate_partition_key(self.partition_key) 

248 coerced_logical_date = timezone.coerce_datetime(self.logical_date) 

249 run_after = self.run_after or timezone.utcnow() 

250 data_interval = None 

251 if coerced_logical_date: 

252 if self.data_interval_start and self.data_interval_end: 

253 data_interval = DataInterval( 

254 start=timezone.coerce_datetime(self.data_interval_start), 

255 end=timezone.coerce_datetime(self.data_interval_end), 

256 ) 

257 else: 

258 data_interval = dag.timetable.infer_manual_data_interval(run_after=coerced_logical_date) 

259 

260 run_id = self.dag_run_id or dag.timetable.generate_run_id( 

261 run_type=DagRunType.MANUAL, 

262 run_after=timezone.coerce_datetime(run_after), 

263 data_interval=data_interval, 

264 ) 

265 

266 partition_date = dag.timetable.resolve_partition_date(self.partition_key) 

267 

268 return { 

269 "run_id": run_id, 

270 "logical_date": coerced_logical_date, 

271 "data_interval": data_interval, 

272 "run_after": run_after, 

273 "conf": self.conf, 

274 "note": self.note, 

275 "partition_key": self.partition_key, 

276 "partition_date": partition_date, 

277 } 

278 

279 

280class DAGRunsBatchBody(StrictBaseModel): 

281 """List Dag Runs body for batch endpoint.""" 

282 

283 order_by: str | None = None 

284 page_offset: NonNegativeInt = 0 

285 page_limit: NonNegativeInt = 100 

286 dag_ids: list[str] | None = None 

287 states: list[DagRunState | None] | None = None 

288 

289 run_after_gte: AwareDatetime | None = None 

290 run_after_gt: AwareDatetime | None = None 

291 run_after_lte: AwareDatetime | None = None 

292 run_after_lt: AwareDatetime | None = None 

293 

294 logical_date_gte: AwareDatetime | None = None 

295 logical_date_gt: AwareDatetime | None = None 

296 logical_date_lte: AwareDatetime | None = None 

297 logical_date_lt: AwareDatetime | None = None 

298 

299 start_date_gte: AwareDatetime | None = None 

300 start_date_gt: AwareDatetime | None = None 

301 start_date_lte: AwareDatetime | None = None 

302 start_date_lt: AwareDatetime | None = None 

303 

304 end_date_gte: AwareDatetime | None = None 

305 end_date_gt: AwareDatetime | None = None 

306 end_date_lte: AwareDatetime | None = None 

307 end_date_lt: AwareDatetime | None = None 

308 

309 duration_gte: float | None = None 

310 duration_gt: float | None = None 

311 duration_lte: float | None = None 

312 duration_lt: float | None = None 

313 

314 conf_contains: str | None = None 

315 

316 

317class ClearPartitionsBody(PartitionSelectorMixin): 

318 """Request body for the clearPartitions endpoint (column-reset: set partition fields to None).""" 

319 

320 run_id: str | None = Field( 

321 default=None, 

322 description="Select runs by exact run_id. Mutually exclusive with ``partition_key`` and partition date window.", 

323 ) 

324 clear_task_instances: bool = Field( 

325 default=False, 

326 description="Also clear task instances on the matched runs.", 

327 ) 

328 dry_run: bool = Field( 

329 default=True, 

330 description="If True, compute counts without writing any changes.", 

331 ) 

332 

333 @model_validator(mode="after") 

334 def validate_exactly_one_selector(self) -> ClearPartitionsBody: 

335 self._check_exactly_one_selection_mode( 

336 extra_selector_active=self.run_id is not None, 

337 extra_selector_name="run_id", 

338 ) 

339 return self 

340 

341 

342class ClearPartitionsResponse(BaseModel): 

343 """Response for the clearPartitions endpoint.""" 

344 

345 dag_runs_cleared: int 

346 task_instances_cleared: int 

347 dry_run: bool