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
« 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 collections.abc import Iterable
21from datetime import datetime
22from enum import Enum
23from typing import TYPE_CHECKING, Any
25from pydantic import AliasPath, AwareDatetime, Field, NonNegativeInt, model_validator
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
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
38class DagRunMutableStates(str, Enum):
39 """Dag Run states from which the run may be mutated (patched, deleted)."""
41 QUEUED = DagRunState.QUEUED
42 SUCCESS = DagRunState.SUCCESS
43 FAILED = DagRunState.FAILED
46class DAGRunPatchBody(StrictBaseModel):
47 """Dag Run Serializer for PATCH requests."""
49 state: DagRunMutableStates | None = None
50 note: str | None = Field(None, max_length=1000)
53class BulkDAGRunBody(StrictBaseModel):
54 """Request body for bulk operations on Dag Runs."""
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)
62class PartitionSelectorMixin(StrictBaseModel):
63 """Partition filter fields shared by bulk-clear and clearPartitions bodies."""
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 )
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 )
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.")
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()
118class BaseDAGRunClear(StrictBaseModel):
119 """Shared options for the single-run and bulk Dag Run clear endpoints."""
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)
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
144class DAGRunClearBody(BaseDAGRunClear):
145 """Dag Run serializer for clear endpoint body."""
148class BulkDAGRunClearBody(BaseDAGRunClear, PartitionSelectorMixin):
149 """Request body for the bulk clear Dag Runs endpoint."""
151 dag_runs: list[BulkDAGRunBody] = Field(default_factory=list)
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
162class DAGRunResponse(BaseModel):
163 """Dag Run serializer for responses."""
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
189class DAGRunCollectionResponse(BaseModel):
190 """
191 Dag Run collection response supporting both offset and cursor pagination.
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 """
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 )
225class TriggerDAGRunPostBody(StrictBaseModel):
226 """Trigger Dag Run Serializer for POST body."""
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)
234 conf: dict | None = Field(default_factory=dict)
235 note: str | None = None
236 partition_key: str | None = None
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
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)
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 )
266 partition_date = dag.timetable.resolve_partition_date(self.partition_key)
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 }
280class DAGRunsBatchBody(StrictBaseModel):
281 """List Dag Runs body for batch endpoint."""
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
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
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
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
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
309 duration_gte: float | None = None
310 duration_gt: float | None = None
311 duration_lte: float | None = None
312 duration_lt: float | None = None
314 conf_contains: str | None = None
317class ClearPartitionsBody(PartitionSelectorMixin):
318 """Request body for the clearPartitions endpoint (column-reset: set partition fields to None)."""
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 )
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
342class ClearPartitionsResponse(BaseModel):
343 """Response for the clearPartitions endpoint."""
345 dag_runs_cleared: int
346 task_instances_cleared: int
347 dry_run: bool