Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/airflow/sdk/__init__.py: 21%

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

61 statements  

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. 

17from __future__ import annotations 

18 

19from typing import TYPE_CHECKING 

20 

21__all__ = [ 

22 "__version__", 

23 "AllowedKeyMapper", 

24 "Asset", 

25 "AssetAccessControl", 

26 "AssetAlias", 

27 "AssetAll", 

28 "AssetAny", 

29 "AssetOrTimeSchedule", 

30 "AssetWatcher", 

31 "AsyncCallback", 

32 "BaseAsyncOperator", 

33 "BaseBranchOperator", 

34 "BaseDeadlineReference", 

35 "BaseHook", 

36 "BaseNotifier", 

37 "BaseOperator", 

38 "BaseOperatorLink", 

39 "BaseSensorOperator", 

40 "BaseXCom", 

41 "BranchMixIn", 

42 "ChainMapper", 

43 "ChainRetryPolicy", 

44 "Connection", 

45 "Context", 

46 "CronDataIntervalTimetable", 

47 "CronTriggerTimetable", 

48 "CronPartitionTimetable", 

49 "DAG", 

50 "DagRunState", 

51 "DayWindow", 

52 "DeadlineAlert", 

53 "DeadlineReference", 

54 "DeltaDataIntervalTimetable", 

55 "DeltaTriggerTimetable", 

56 "EdgeModifier", 

57 "EventsTimetable", 

58 "ExceptionRetryPolicy", 

59 "FanOutMapper", 

60 "FixedKeyMapper", 

61 "HourWindow", 

62 "IdentityMapper", 

63 "Label", 

64 "Metadata", 

65 "MinimumCount", 

66 "MonthWindow", 

67 "MultipleCronTriggerTimetable", 

68 "NEVER_EXPIRE", 

69 "ObjectStoragePath", 

70 "Param", 

71 "ParamsDict", 

72 "PartitionedAtRuntime", 

73 "PartitionedAssetTimetable", 

74 "PartitionMapper", 

75 "PokeReturnValue", 

76 "ProductMapper", 

77 "QuarterWindow", 

78 "ResumableJobMixin", 

79 "RetryAction", 

80 "RetryDecision", 

81 "RetryPolicy", 

82 "RetryRule", 

83 "RollupMapper", 

84 "SegmentWindow", 

85 "SkipMixin", 

86 "SyncCallback", 

87 "StartOfDayMapper", 

88 "StartOfHourMapper", 

89 "StartOfMonthMapper", 

90 "StartOfQuarterMapper", 

91 "StartOfWeekMapper", 

92 "StartOfYearMapper", 

93 "TaskGroup", 

94 "TaskInstance", 

95 "TaskInstanceState", 

96 "TriggerRule", 

97 "Variable", 

98 "WaitForAll", 

99 "WeekWindow", 

100 "WeightRule", 

101 "Window", 

102 "XComArg", 

103 "YearWindow", 

104 "asset", 

105 "chain", 

106 "chain_linear", 

107 "conf", 

108 "cross_downstream", 

109 "dag", 

110 "deadline_reference", 

111 "get_current_context", 

112 "get_parsing_context", 

113 "literal", 

114 "lineage", 

115 "macros", 

116 "result", 

117 "setup", 

118 "task", 

119 "task_group", 

120 "teardown", 

121] 

122 

123__version__ = "1.4.0" 

124 

125if TYPE_CHECKING: 

126 from airflow.sdk.api.datamodels._generated import DagRunState, TaskInstanceState, TriggerRule, WeightRule 

127 from airflow.sdk.bases.branch import BaseBranchOperator, BranchMixIn 

128 from airflow.sdk.bases.hook import BaseHook 

129 from airflow.sdk.bases.notifier import BaseNotifier 

130 from airflow.sdk.bases.operator import ( 

131 BaseAsyncOperator, 

132 BaseOperator, 

133 chain, 

134 chain_linear, 

135 cross_downstream, 

136 ) 

137 from airflow.sdk.bases.operatorlink import BaseOperatorLink 

138 from airflow.sdk.bases.resumablejobmixin import ResumableJobMixin 

139 from airflow.sdk.bases.sensor import BaseSensorOperator, PokeReturnValue 

140 from airflow.sdk.bases.skipmixin import SkipMixin 

141 from airflow.sdk.bases.xcom import BaseXCom 

142 from airflow.sdk.configuration import AirflowSDKConfigParser 

143 from airflow.sdk.definitions.asset import ( 

144 Asset, 

145 AssetAccessControl, 

146 AssetAlias, 

147 AssetAll, 

148 AssetAny, 

149 AssetWatcher, 

150 ) 

151 from airflow.sdk.definitions.asset.decorators import asset 

152 from airflow.sdk.definitions.asset.metadata import Metadata 

153 from airflow.sdk.definitions.callback import AsyncCallback, SyncCallback 

154 from airflow.sdk.definitions.connection import Connection 

155 from airflow.sdk.definitions.context import Context, get_current_context, get_parsing_context 

156 from airflow.sdk.definitions.dag import DAG, dag 

157 from airflow.sdk.definitions.deadline import ( 

158 BaseDeadlineReference, 

159 DeadlineAlert, 

160 DeadlineReference, 

161 deadline_reference, 

162 ) 

163 from airflow.sdk.definitions.decorators import result, setup, task, teardown 

164 from airflow.sdk.definitions.decorators.task_group import task_group 

165 from airflow.sdk.definitions.edges import EdgeModifier, Label 

166 from airflow.sdk.definitions.param import Param, ParamsDict 

167 from airflow.sdk.definitions.partition_mappers.allowed_key import AllowedKeyMapper 

168 from airflow.sdk.definitions.partition_mappers.base import ( 

169 PartitionMapper, 

170 RollupMapper, 

171 ) 

172 from airflow.sdk.definitions.partition_mappers.chain import ChainMapper 

173 from airflow.sdk.definitions.partition_mappers.fixed_key import FixedKeyMapper 

174 from airflow.sdk.definitions.partition_mappers.identity import IdentityMapper 

175 from airflow.sdk.definitions.partition_mappers.product import ProductMapper 

176 from airflow.sdk.definitions.partition_mappers.temporal import ( 

177 FanOutMapper, 

178 StartOfDayMapper, 

179 StartOfHourMapper, 

180 StartOfMonthMapper, 

181 StartOfQuarterMapper, 

182 StartOfWeekMapper, 

183 StartOfYearMapper, 

184 ) 

185 from airflow.sdk.definitions.partition_mappers.wait_policy import ( 

186 MinimumCount, 

187 WaitForAll, 

188 ) 

189 from airflow.sdk.definitions.partition_mappers.window import ( 

190 DayWindow, 

191 HourWindow, 

192 MonthWindow, 

193 QuarterWindow, 

194 SegmentWindow, 

195 WeekWindow, 

196 Window, 

197 YearWindow, 

198 ) 

199 from airflow.sdk.definitions.retry_policy import ( 

200 ChainRetryPolicy, 

201 ExceptionRetryPolicy, 

202 RetryAction, 

203 RetryDecision, 

204 RetryPolicy, 

205 RetryRule, 

206 ) 

207 from airflow.sdk.definitions.taskgroup import TaskGroup 

208 from airflow.sdk.definitions.template import literal 

209 from airflow.sdk.definitions.timetables.assets import ( 

210 AssetOrTimeSchedule, 

211 PartitionedAssetTimetable, 

212 PartitionedAtRuntime, 

213 ) 

214 from airflow.sdk.definitions.timetables.events import EventsTimetable 

215 from airflow.sdk.definitions.timetables.interval import ( 

216 CronDataIntervalTimetable, 

217 DeltaDataIntervalTimetable, 

218 ) 

219 from airflow.sdk.definitions.timetables.trigger import ( 

220 CronPartitionTimetable, 

221 CronTriggerTimetable, 

222 DeltaTriggerTimetable, 

223 MultipleCronTriggerTimetable, 

224 ) 

225 from airflow.sdk.definitions.variable import Variable 

226 from airflow.sdk.definitions.xcom_arg import XComArg 

227 from airflow.sdk.execution_time import macros 

228 from airflow.sdk.execution_time.context import NEVER_EXPIRE 

229 from airflow.sdk.io.path import ObjectStoragePath 

230 from airflow.sdk.types import TaskInstance 

231 

232 conf: AirflowSDKConfigParser 

233 

234__lazy_imports: dict[str, str] = { 

235 "AllowedKeyMapper": ".definitions.partition_mappers.allowed_key", 

236 "Asset": ".definitions.asset", 

237 "AssetAccessControl": ".definitions.asset", 

238 "AssetAlias": ".definitions.asset", 

239 "AssetAll": ".definitions.asset", 

240 "AssetAny": ".definitions.asset", 

241 "AssetOrTimeSchedule": ".definitions.timetables.assets", 

242 "AssetWatcher": ".definitions.asset", 

243 "AsyncCallback": ".definitions.callback", 

244 "BaseAsyncOperator": ".bases.operator", 

245 "BaseBranchOperator": ".bases.branch", 

246 "BaseDeadlineReference": ".definitions.deadline", 

247 "BaseHook": ".bases.hook", 

248 "BaseNotifier": ".bases.notifier", 

249 "BaseOperator": ".bases.operator", 

250 "BaseOperatorLink": ".bases.operatorlink", 

251 "BaseSensorOperator": ".bases.sensor", 

252 "BaseXCom": ".bases.xcom", 

253 "BranchMixIn": ".bases.branch", 

254 "ChainMapper": ".definitions.partition_mappers.chain", 

255 "ChainRetryPolicy": ".definitions.retry_policy", 

256 "Connection": ".definitions.connection", 

257 "Context": ".definitions.context", 

258 "CronDataIntervalTimetable": ".definitions.timetables.interval", 

259 "CronTriggerTimetable": ".definitions.timetables.trigger", 

260 "CronPartitionTimetable": ".definitions.timetables.trigger", 

261 "DAG": ".definitions.dag", 

262 "DagRunState": ".api.datamodels._generated", 

263 "DayWindow": ".definitions.partition_mappers.window", 

264 "DeadlineAlert": ".definitions.deadline", 

265 "DeadlineReference": ".definitions.deadline", 

266 "DeltaDataIntervalTimetable": ".definitions.timetables.interval", 

267 "DeltaTriggerTimetable": ".definitions.timetables.trigger", 

268 "EdgeModifier": ".definitions.edges", 

269 "EventsTimetable": ".definitions.timetables.events", 

270 "ExceptionRetryPolicy": ".definitions.retry_policy", 

271 "FanOutMapper": ".definitions.partition_mappers.temporal", 

272 "FixedKeyMapper": ".definitions.partition_mappers.fixed_key", 

273 "HourWindow": ".definitions.partition_mappers.window", 

274 "IdentityMapper": ".definitions.partition_mappers.identity", 

275 "Label": ".definitions.edges", 

276 "Metadata": ".definitions.asset.metadata", 

277 "MinimumCount": ".definitions.partition_mappers.wait_policy", 

278 "MonthWindow": ".definitions.partition_mappers.window", 

279 "MultipleCronTriggerTimetable": ".definitions.timetables.trigger", 

280 "ObjectStoragePath": ".io.path", 

281 "Param": ".definitions.param", 

282 "ParamsDict": ".definitions.param", 

283 "PartitionedAtRuntime": ".definitions.timetables.assets", 

284 "PartitionedAssetTimetable": ".definitions.timetables.assets", 

285 "PartitionMapper": ".definitions.partition_mappers.base", 

286 "PokeReturnValue": ".bases.sensor", 

287 "ProductMapper": ".definitions.partition_mappers.product", 

288 "QuarterWindow": ".definitions.partition_mappers.window", 

289 "ResumableJobMixin": ".bases.resumablejobmixin", 

290 "RetryAction": ".definitions.retry_policy", 

291 "RetryDecision": ".definitions.retry_policy", 

292 "RetryPolicy": ".definitions.retry_policy", 

293 "RetryRule": ".definitions.retry_policy", 

294 "RollupMapper": ".definitions.partition_mappers.base", 

295 "SecretCache": ".execution_time.cache", 

296 "SegmentWindow": ".definitions.partition_mappers.window", 

297 "SkipMixin": ".bases.skipmixin", 

298 "SyncCallback": ".definitions.callback", 

299 "StartOfDayMapper": ".definitions.partition_mappers.temporal", 

300 "StartOfHourMapper": ".definitions.partition_mappers.temporal", 

301 "StartOfMonthMapper": ".definitions.partition_mappers.temporal", 

302 "StartOfQuarterMapper": ".definitions.partition_mappers.temporal", 

303 "StartOfWeekMapper": ".definitions.partition_mappers.temporal", 

304 "StartOfYearMapper": ".definitions.partition_mappers.temporal", 

305 "TaskGroup": ".definitions.taskgroup", 

306 "TaskInstance": ".types", 

307 "TaskInstanceState": ".api.datamodels._generated", 

308 "TriggerRule": ".api.datamodels._generated", 

309 "Variable": ".definitions.variable", 

310 "WaitForAll": ".definitions.partition_mappers.wait_policy", 

311 "WeekWindow": ".definitions.partition_mappers.window", 

312 "WeightRule": ".api.datamodels._generated", 

313 "Window": ".definitions.partition_mappers.window", 

314 "XComArg": ".definitions.xcom_arg", 

315 "YearWindow": ".definitions.partition_mappers.window", 

316 "asset": ".definitions.asset.decorators", 

317 "chain": ".bases.operator", 

318 "chain_linear": ".bases.operator", 

319 "conf": ".configuration", 

320 "cross_downstream": ".bases.operator", 

321 "dag": ".definitions.dag", 

322 "deadline_reference": ".definitions.deadline", 

323 "NEVER_EXPIRE": ".execution_time.context", 

324 "get_current_context": ".definitions.context", 

325 "get_parsing_context": ".definitions.context", 

326 "literal": ".definitions.template", 

327 "lineage": ".lineage", 

328 "macros": ".execution_time", 

329 "result": ".definitions.decorators", 

330 "setup": ".definitions.decorators", 

331 "task": ".definitions.decorators", 

332 "task_group": ".definitions.decorators", 

333 "teardown": ".definitions.decorators", 

334} 

335 

336 

337def __getattr__(name: str): 

338 if module_path := __lazy_imports.get(name): 

339 import importlib 

340 

341 mod = importlib.import_module(module_path, __name__) 

342 val = getattr(mod, name) 

343 

344 # Store for next time 

345 globals()[name] = val 

346 return val 

347 raise AttributeError(f"module {__name__!r} has no attribute {name!r}")