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

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 

20import inspect 

21from collections.abc import Iterable, Mapping 

22from datetime import datetime, timedelta 

23from functools import cache 

24from typing import TYPE_CHECKING, Any 

25 

26from itsdangerous import URLSafeSerializer 

27from pendulum.tz.timezone import FixedTimezone, Timezone 

28from pydantic import ( 

29 AliasGenerator, 

30 ConfigDict, 

31 computed_field, 

32 field_serializer, 

33 field_validator, 

34) 

35 

36from airflow._shared.module_loading import qualname 

37from airflow.api_fastapi.core_api.base import BaseModel, StrictBaseModel, make_partial_model 

38from airflow.api_fastapi.core_api.datamodels.dag_tags import DagTagResponse 

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

40from airflow.configuration import conf 

41from airflow.models.dag_version import DagVersion 

42from airflow.utils.types import DagRunType 

43 

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

45 from airflow.serialization.definitions.param import SerializedParamsDict 

46 

47 

48def _is_response_safe_pod_override(value: Any) -> bool: 

49 """Whether a pod_override value is already safe to preserve in the response.""" 

50 return value is None or isinstance(value, str | int | float | Mapping | list) 

51 

52 

53@cache 

54def _get_file_token_serializer() -> URLSafeSerializer: 

55 """ 

56 Return a cached URLSafeSerializer instance. 

57 

58 Uses @cache for lazy initialization - the serializer is created on first 

59 call rather than at module import time. This avoids issues if the module 

60 is imported before configuration is fully loaded. 

61 """ 

62 return URLSafeSerializer(conf.get_mandatory_value("api", "secret_key")) 

63 

64 

65DAG_ALIAS_MAPPING: dict[str, str] = { 

66 # The keys are the names in the response, the values are the original names in the model 

67 # This is used to map the names in the response to the names in the model 

68 # See: https://github.com/apache/airflow/issues/46732 

69 "next_dagrun_logical_date": "next_dagrun", 

70 "next_dagrun_run_after": "next_dagrun_create_after", 

71} 

72 

73 

74class DAGResponse(BaseModel): 

75 """Dag serializer for responses.""" 

76 

77 model_config = ConfigDict( 

78 alias_generator=AliasGenerator( 

79 validation_alias=lambda field_name: DAG_ALIAS_MAPPING.get(field_name, field_name), 

80 ), 

81 ) 

82 

83 dag_id: str 

84 dag_display_name: str 

85 is_paused: bool 

86 is_stale: bool 

87 last_parsed_time: datetime | None 

88 last_parse_duration: float | None 

89 last_expired: datetime | None 

90 bundle_name: str | None 

91 bundle_version: str | None 

92 relative_fileloc: str | None 

93 fileloc: str 

94 description: str | None 

95 timetable_summary: str | None 

96 timetable_description: str | None 

97 timetable_partitioned: bool 

98 timetable_periodic: bool 

99 tags: list[DagTagResponse] 

100 max_active_tasks: int 

101 max_active_runs: int | None 

102 max_consecutive_failed_dag_runs: int 

103 has_task_concurrency_limits: bool 

104 has_import_errors: bool 

105 next_dagrun_logical_date: datetime | None 

106 next_dagrun_data_interval_start: datetime | None 

107 next_dagrun_data_interval_end: datetime | None 

108 next_dagrun_run_after: datetime | None 

109 allowed_run_types: list[DagRunType] | None 

110 owners: list[str] 

111 

112 @field_serializer("tags") 

113 def serialize_tags(self, tags: list[DagTagResponse]) -> list[DagTagResponse]: 

114 """Sort tags alphabetically by name.""" 

115 return sorted(tags, key=lambda tag: tag.name) 

116 

117 @field_validator("owners", mode="before") 

118 @classmethod 

119 def get_owners(cls, v: Any) -> list[str] | None: 

120 """Convert owners attribute to Dag representation.""" 

121 if not (v is None or isinstance(v, str)): 121 ↛ 122line 121 didn't jump to line 122 because the condition on line 121 was never true

122 return v 

123 

124 if v is None: 124 ↛ 125line 124 didn't jump to line 125 because the condition on line 124 was never true

125 return [] 

126 if isinstance(v, str): 126 ↛ 128line 126 didn't jump to line 128 because the condition on line 126 was always true

127 return [x.strip() for x in v.split(",")] 

128 return v 

129 

130 @field_validator("timetable_summary", mode="before") 

131 @classmethod 

132 def get_timetable_summary(cls, tts: str | None) -> str | None: 

133 """Validate the string representation of timetable_summary.""" 

134 if tts is None or tts == "None": 

135 return None 

136 return str(tts) 

137 

138 # Mypy issue https://github.com/python/mypy/issues/1362 

139 @computed_field # type: ignore[prop-decorator] 

140 @property 

141 def is_backfillable(self) -> bool: 

142 """Whether this Dag's schedule supports backfilling.""" 

143 if not self.timetable_periodic: 

144 return False 

145 if self.allowed_run_types is not None and DagRunType.BACKFILL_JOB not in self.allowed_run_types: 145 ↛ 146line 145 didn't jump to line 146 because the condition on line 145 was never true

146 return False 

147 return True 

148 

149 # Mypy issue https://github.com/python/mypy/issues/1362 

150 @computed_field # type: ignore[prop-decorator] 

151 @property 

152 def file_token(self) -> str: 

153 """Return file token.""" 

154 payload = { 

155 "bundle_name": self.bundle_name, 

156 "relative_fileloc": self.relative_fileloc, 

157 } 

158 return _get_file_token_serializer().dumps(payload) 

159 

160 

161class DAGPatchBody(StrictBaseModel): 

162 """Dag Serializer for updatable bodies.""" 

163 

164 is_paused: bool 

165 

166 

167DAGPatchBodyPartial = make_partial_model(DAGPatchBody) 

168 

169 

170class DAGCollectionResponse(BaseModel): 

171 """Dag Collection serializer for responses.""" 

172 

173 dags: Iterable[DAGResponse] 

174 total_entries: int 

175 

176 

177class DAGDetailsResponse(DAGResponse): 

178 """Specific serializer for Dag Details responses.""" 

179 

180 model_config = ConfigDict( 

181 from_attributes=True, 

182 alias_generator=AliasGenerator( 

183 validation_alias=lambda field_name: { 

184 "dag_run_timeout": "dagrun_timeout", 

185 "last_parsed": "last_loaded", 

186 "template_search_path": "template_searchpath", 

187 **DAG_ALIAS_MAPPING, 

188 }.get(field_name, field_name), 

189 ), 

190 ) 

191 

192 catchup: bool 

193 dag_run_timeout: timedelta | None 

194 asset_expression: dict | None 

195 doc_md: str | None 

196 start_date: datetime | None 

197 end_date: datetime | None 

198 is_paused_upon_creation: bool | None 

199 params: Mapping | None 

200 render_template_as_native_obj: bool 

201 template_search_path: list[str] | None 

202 timezone: str | None 

203 last_parsed: datetime | None 

204 default_args: Mapping | None 

205 rerun_with_latest_version: bool | None = None 

206 owner_links: dict[str, str] | None = None 

207 is_favorite: bool = False 

208 active_runs_count: int = 0 

209 

210 @field_validator("timezone", mode="before") 

211 @classmethod 

212 def get_timezone(cls, tz: Timezone | FixedTimezone) -> str | None: 

213 """Convert timezone attribute to string representation.""" 

214 if tz is None: 214 ↛ 215line 214 didn't jump to line 215 because the condition on line 214 was never true

215 return None 

216 return str(tz) 

217 

218 @field_validator("doc_md", mode="before") 

219 @classmethod 

220 def get_doc_md(cls, doc_md: str | None) -> str | None: 

221 """Clean indentation in doc md.""" 

222 if doc_md is None: 

223 return None 

224 return inspect.cleandoc(doc_md) 

225 

226 @field_validator("default_args", mode="before") 

227 @classmethod 

228 def get_default_args(cls, default_args: Mapping | None) -> Mapping | None: 

229 """ 

230 Sanitize default_args for the API response. 

231 

232 Targets the common case where ``executor_config["pod_override"]`` is a 

233 Kubernetes ``V1Pod``: when the value is not a JSON primitive 

234 (``None``/``str``/``int``/``float``) or a ``Mapping``/``list``, it is 

235 rewritten to a fully-qualified type-name string so the response stays 

236 valid JSON. The container check is shallow — a ``Mapping`` or ``list`` 

237 whose contents are themselves non-serializable (e.g. nested ``V1Pod``) 

238 will still raise during response serialization, as will any other 

239 non-JSON values elsewhere in ``default_args``. 

240 """ 

241 if default_args is None: 241 ↛ 242line 241 didn't jump to line 242 because the condition on line 241 was never true

242 return None 

243 executor_config = default_args.get("executor_config") 

244 if not (isinstance(executor_config, Mapping) and "pod_override" in executor_config): 244 ↛ 247line 244 didn't jump to line 247 because the condition on line 244 was always true

245 return default_args 

246 

247 pod_override = executor_config["pod_override"] 

248 if _is_response_safe_pod_override(pod_override): 

249 return default_args 

250 

251 sanitized_executor_config = dict(executor_config) 

252 sanitized_executor_config["pod_override"] = qualname(pod_override) 

253 result = dict(default_args) 

254 result["executor_config"] = sanitized_executor_config 

255 return result 

256 

257 @field_validator("params", mode="before") 

258 @classmethod 

259 def get_params(cls, params: SerializedParamsDict | None) -> dict | None: 

260 """Convert params attribute to dict representation.""" 

261 if params is None: 261 ↛ 262line 261 didn't jump to line 262 because the condition on line 261 was never true

262 return None 

263 return {k: v.dump() for k, v in params.items()} 

264 

265 # Mypy issue https://github.com/python/mypy/issues/1362 

266 @computed_field(deprecated=True) # type: ignore[prop-decorator] 

267 @property 

268 def concurrency(self) -> int: 

269 """ 

270 Return max_active_tasks as concurrency. 

271 

272 Deprecated: Use max_active_tasks instead. 

273 """ 

274 return self.max_active_tasks 

275 

276 # Mypy issue https://github.com/python/mypy/issues/1362 

277 @computed_field # type: ignore[prop-decorator] 

278 @property 

279 def latest_dag_version(self) -> DagVersionResponse | None: 

280 """Return the latest DagVersion.""" 

281 latest_dag_version = DagVersion.get_latest_version( 

282 self.dag_id, load_dag_model=True, load_bundle_model=True 

283 ) 

284 if latest_dag_version is None: 284 ↛ 285line 284 didn't jump to line 285 because the condition on line 284 was never true

285 return latest_dag_version 

286 return DagVersionResponse.model_validate(latest_dag_version)