Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/datamodels/assets.py: 98%
123 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 typing import TYPE_CHECKING, Annotated
24from pydantic import (
25 AliasPath,
26 AwareDatetime,
27 ConfigDict,
28 Field,
29 JsonValue,
30 NonNegativeInt,
31 StringConstraints,
32 field_validator,
33)
35from airflow._shared.secrets_masker import redact
36from airflow._shared.timezones import timezone
37from airflow.api_fastapi.core_api.base import BaseModel, StrictBaseModel
38from airflow.api_fastapi.core_api.datamodels.dag_run import TriggerDAGRunPostBody
39from airflow.models.base import ID_LEN
40from airflow.utils.types import DagRunType
42if TYPE_CHECKING: 42 ↛ 43line 42 didn't jump to line 43 because the condition on line 42 was never true
43 from airflow.serialization.definitions.dag import SerializedDAG
46class DagScheduleAssetReference(StrictBaseModel):
47 """Dag schedule reference serializer for assets."""
49 dag_id: str
50 created_at: datetime
51 updated_at: datetime
54class TaskInletAssetReference(StrictBaseModel):
55 """Task inlet reference serializer for assets."""
57 dag_id: str
58 task_id: str
59 created_at: datetime
60 updated_at: datetime
63class TaskOutletAssetReference(StrictBaseModel):
64 """Task outlet reference serializer for assets."""
66 dag_id: str
67 task_id: str
68 created_at: datetime
69 updated_at: datetime
72class LastAssetEventResponse(BaseModel):
73 """Last asset event response serializer."""
75 id: NonNegativeInt | None = None
76 timestamp: datetime | None = None
79class AssetWatcherResponse(BaseModel):
80 """Asset watcher serializer for responses."""
82 name: str
83 trigger_id: int
84 created_date: datetime
87class AssetResponse(BaseModel):
88 """Asset serializer for responses."""
90 id: int
91 name: str
92 uri: str
93 group: str
94 extra: dict[str, JsonValue] | None = None
95 created_at: datetime
96 updated_at: datetime
97 scheduled_dags: list[DagScheduleAssetReference]
98 producing_tasks: list[TaskOutletAssetReference]
99 consuming_tasks: list[TaskInletAssetReference]
100 aliases: list[AssetAliasResponse]
101 watchers: list[AssetWatcherResponse]
102 last_asset_event: LastAssetEventResponse | None = None
104 @field_validator("extra", mode="after")
105 @classmethod
106 def redact_extra(cls, v: dict):
107 return redact(v)
110class AssetCollectionResponse(BaseModel):
111 """Asset collection response."""
113 assets: list[AssetResponse]
114 total_entries: int
117class AssetAliasResponse(BaseModel):
118 """Asset alias serializer for responses."""
120 id: int
121 name: str
122 group: str
125class AssetAliasCollectionResponse(BaseModel):
126 """Asset alias collection response."""
128 asset_aliases: Iterable[AssetAliasResponse]
129 total_entries: int
132class DagRunAssetReference(StrictBaseModel):
133 """DagRun serializer for asset responses."""
135 run_id: str
136 dag_id: str
137 logical_date: datetime | None
138 start_date: datetime
139 end_date: datetime | None
140 state: str
141 data_interval_start: datetime | None
142 data_interval_end: datetime | None
143 partition_key: str | None
144 triggering: bool = Field(
145 description=(
146 "Whether this asset event triggered the referenced dag run. Only a run's most recent "
147 "consumed asset event triggers it; earlier consumed events are included in the run but "
148 "did not trigger it."
149 ),
150 )
153class AssetEventResponse(BaseModel):
154 """Asset event serializer for responses."""
156 id: int
157 asset_id: int
158 uri: str | None = Field(alias="uri", default=None)
159 name: str | None = Field(alias="name", default=None)
160 group: str | None = Field(alias="group", default=None)
161 extra: dict[str, JsonValue] | None = None
162 source_task_id: str | None = None
163 source_dag_id: str | None = None
164 source_run_id: str | None = None
165 source_map_index: int
166 created_dagruns: list[DagRunAssetReference]
167 timestamp: datetime
168 partition_key: str | None = None
170 @field_validator("extra", mode="after")
171 @classmethod
172 def redact_extra(cls, v: dict):
173 return redact(v)
176class AssetEventCollectionResponse(BaseModel):
177 """Asset event collection response."""
179 asset_events: Iterable[AssetEventResponse]
180 total_entries: int
183class QueuedEventResponse(BaseModel):
184 """Queued Event serializer for responses.."""
186 dag_id: str
187 asset_id: int
188 created_at: datetime
189 dag_display_name: str = Field(validation_alias=AliasPath("dag_model", "dag_display_name"))
192class QueuedEventCollectionResponse(BaseModel):
193 """Queued Event Collection serializer for responses."""
195 queued_events: list[QueuedEventResponse]
196 total_entries: int
199class AssetEventAccessControl(StrictBaseModel):
200 """Access control settings for asset event consumer team filtering."""
202 consumer_teams: list[str] | None = None
203 allow_global: bool = True
206class CreateAssetEventsBody(StrictBaseModel):
207 """Create asset events request."""
209 asset_id: int
210 # pattern (not strip_whitespace) so the value isn't mutated — must stay byte-identical to what
211 # `_validate_outlet_event_partition_keys` in the Execution API accepts for the same raw input.
212 partition_key: Annotated[str, StringConstraints(pattern=r"\S", max_length=ID_LEN)] | None = None
213 extra: dict = Field(default_factory=dict)
214 access_control: AssetEventAccessControl | None = None
216 @field_validator("extra", mode="after")
217 def set_from_rest_api(cls, v: dict) -> dict:
218 v["from_rest_api"] = True
219 return v
221 model_config = ConfigDict(extra="forbid")
224class MaterializeAssetBody(TriggerDAGRunPostBody):
225 """Materialize asset request."""
227 logical_date: AwareDatetime | None = None
229 def validate_context(self, dag: SerializedDAG) -> dict:
230 params = super().validate_context(dag)
231 if self.dag_run_id is None:
232 params["run_id"] = dag.timetable.generate_run_id(
233 run_type=DagRunType.ASSET_MATERIALIZATION,
234 run_after=timezone.coerce_datetime(params["run_after"]),
235 data_interval=params["data_interval"],
236 )
237 return params
239 model_config = ConfigDict(extra="forbid")