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

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 typing import TYPE_CHECKING, Annotated 

23 

24from pydantic import ( 

25 AliasPath, 

26 AwareDatetime, 

27 ConfigDict, 

28 Field, 

29 JsonValue, 

30 NonNegativeInt, 

31 StringConstraints, 

32 field_validator, 

33) 

34 

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 

41 

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 

44 

45 

46class DagScheduleAssetReference(StrictBaseModel): 

47 """Dag schedule reference serializer for assets.""" 

48 

49 dag_id: str 

50 created_at: datetime 

51 updated_at: datetime 

52 

53 

54class TaskInletAssetReference(StrictBaseModel): 

55 """Task inlet reference serializer for assets.""" 

56 

57 dag_id: str 

58 task_id: str 

59 created_at: datetime 

60 updated_at: datetime 

61 

62 

63class TaskOutletAssetReference(StrictBaseModel): 

64 """Task outlet reference serializer for assets.""" 

65 

66 dag_id: str 

67 task_id: str 

68 created_at: datetime 

69 updated_at: datetime 

70 

71 

72class LastAssetEventResponse(BaseModel): 

73 """Last asset event response serializer.""" 

74 

75 id: NonNegativeInt | None = None 

76 timestamp: datetime | None = None 

77 

78 

79class AssetWatcherResponse(BaseModel): 

80 """Asset watcher serializer for responses.""" 

81 

82 name: str 

83 trigger_id: int 

84 created_date: datetime 

85 

86 

87class AssetResponse(BaseModel): 

88 """Asset serializer for responses.""" 

89 

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 

103 

104 @field_validator("extra", mode="after") 

105 @classmethod 

106 def redact_extra(cls, v: dict): 

107 return redact(v) 

108 

109 

110class AssetCollectionResponse(BaseModel): 

111 """Asset collection response.""" 

112 

113 assets: list[AssetResponse] 

114 total_entries: int 

115 

116 

117class AssetAliasResponse(BaseModel): 

118 """Asset alias serializer for responses.""" 

119 

120 id: int 

121 name: str 

122 group: str 

123 

124 

125class AssetAliasCollectionResponse(BaseModel): 

126 """Asset alias collection response.""" 

127 

128 asset_aliases: Iterable[AssetAliasResponse] 

129 total_entries: int 

130 

131 

132class DagRunAssetReference(StrictBaseModel): 

133 """DagRun serializer for asset responses.""" 

134 

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 ) 

151 

152 

153class AssetEventResponse(BaseModel): 

154 """Asset event serializer for responses.""" 

155 

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 

169 

170 @field_validator("extra", mode="after") 

171 @classmethod 

172 def redact_extra(cls, v: dict): 

173 return redact(v) 

174 

175 

176class AssetEventCollectionResponse(BaseModel): 

177 """Asset event collection response.""" 

178 

179 asset_events: Iterable[AssetEventResponse] 

180 total_entries: int 

181 

182 

183class QueuedEventResponse(BaseModel): 

184 """Queued Event serializer for responses..""" 

185 

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")) 

190 

191 

192class QueuedEventCollectionResponse(BaseModel): 

193 """Queued Event Collection serializer for responses.""" 

194 

195 queued_events: list[QueuedEventResponse] 

196 total_entries: int 

197 

198 

199class AssetEventAccessControl(StrictBaseModel): 

200 """Access control settings for asset event consumer team filtering.""" 

201 

202 consumer_teams: list[str] | None = None 

203 allow_global: bool = True 

204 

205 

206class CreateAssetEventsBody(StrictBaseModel): 

207 """Create asset events request.""" 

208 

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 

215 

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 

220 

221 model_config = ConfigDict(extra="forbid") 

222 

223 

224class MaterializeAssetBody(TriggerDAGRunPostBody): 

225 """Materialize asset request.""" 

226 

227 logical_date: AwareDatetime | None = None 

228 

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 

238 

239 model_config = ConfigDict(extra="forbid")