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 "Connection", 

44 "Context", 

45 "CronDataIntervalTimetable", 

46 "CronTriggerTimetable", 

47 "CronPartitionTimetable", 

48 "DAG", 

49 "DagRunState", 

50 "DayWindow", 

51 "DeadlineAlert", 

52 "DeadlineReference", 

53 "DeltaDataIntervalTimetable", 

54 "DeltaTriggerTimetable", 

55 "EdgeModifier", 

56 "EventsTimetable", 

57 "ExceptionRetryPolicy", 

58 "FanOutMapper", 

59 "FixedKeyMapper", 

60 "HourWindow", 

61 "IdentityMapper", 

62 "Label", 

63 "Metadata", 

64 "MinimumCount", 

65 "MonthWindow", 

66 "MultipleCronTriggerTimetable", 

67 "NEVER_EXPIRE", 

68 "ObjectStoragePath", 

69 "Param", 

70 "ParamsDict", 

71 "PartitionedAtRuntime", 

72 "PartitionedAssetTimetable", 

73 "PartitionMapper", 

74 "PokeReturnValue", 

75 "ProductMapper", 

76 "QuarterWindow", 

77 "ResumableJobMixin", 

78 "RetryAction", 

79 "RetryDecision", 

80 "RetryPolicy", 

81 "RetryRule", 

82 "RollupMapper", 

83 "SegmentWindow", 

84 "SkipMixin", 

85 "SyncCallback", 

86 "StartOfDayMapper", 

87 "StartOfHourMapper", 

88 "StartOfMonthMapper", 

89 "StartOfQuarterMapper", 

90 "StartOfWeekMapper", 

91 "StartOfYearMapper", 

92 "TaskGroup", 

93 "TaskInstance", 

94 "TaskInstanceState", 

95 "TriggerRule", 

96 "Variable", 

97 "WaitForAll", 

98 "WeekWindow", 

99 "WeightRule", 

100 "Window", 

101 "XComArg", 

102 "YearWindow", 

103 "asset", 

104 "chain", 

105 "chain_linear", 

106 "conf", 

107 "cross_downstream", 

108 "dag", 

109 "deadline_reference", 

110 "get_current_context", 

111 "get_parsing_context", 

112 "literal", 

113 "lineage", 

114 "macros", 

115 "result", 

116 "setup", 

117 "task", 

118 "task_group", 

119 "teardown", 

120] 

121 

122__version__ = "1.4.0" 

123 

124if TYPE_CHECKING: 

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

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

127 from airflow.sdk.bases.hook import BaseHook 

128 from airflow.sdk.bases.notifier import BaseNotifier 

129 from airflow.sdk.bases.operator import ( 

130 BaseAsyncOperator, 

131 BaseOperator, 

132 chain, 

133 chain_linear, 

134 cross_downstream, 

135 ) 

136 from airflow.sdk.bases.operatorlink import BaseOperatorLink 

137 from airflow.sdk.bases.resumablejobmixin import ResumableJobMixin 

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

139 from airflow.sdk.bases.skipmixin import SkipMixin 

140 from airflow.sdk.bases.xcom import BaseXCom 

141 from airflow.sdk.configuration import AirflowSDKConfigParser 

142 from airflow.sdk.definitions.asset import ( 

143 Asset, 

144 AssetAccessControl, 

145 AssetAlias, 

146 AssetAll, 

147 AssetAny, 

148 AssetWatcher, 

149 ) 

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

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

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

153 from airflow.sdk.definitions.connection import Connection 

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

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

156 from airflow.sdk.definitions.deadline import ( 

157 BaseDeadlineReference, 

158 DeadlineAlert, 

159 DeadlineReference, 

160 deadline_reference, 

161 ) 

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

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

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

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

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

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

168 PartitionMapper, 

169 RollupMapper, 

170 ) 

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

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

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

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

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

176 FanOutMapper, 

177 StartOfDayMapper, 

178 StartOfHourMapper, 

179 StartOfMonthMapper, 

180 StartOfQuarterMapper, 

181 StartOfWeekMapper, 

182 StartOfYearMapper, 

183 ) 

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

185 MinimumCount, 

186 WaitForAll, 

187 ) 

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

189 DayWindow, 

190 HourWindow, 

191 MonthWindow, 

192 QuarterWindow, 

193 SegmentWindow, 

194 WeekWindow, 

195 Window, 

196 YearWindow, 

197 ) 

198 from airflow.sdk.definitions.retry_policy import ( 

199 ExceptionRetryPolicy, 

200 RetryAction, 

201 RetryDecision, 

202 RetryPolicy, 

203 RetryRule, 

204 ) 

205 from airflow.sdk.definitions.taskgroup import TaskGroup 

206 from airflow.sdk.definitions.template import literal 

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

208 AssetOrTimeSchedule, 

209 PartitionedAssetTimetable, 

210 PartitionedAtRuntime, 

211 ) 

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

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

214 CronDataIntervalTimetable, 

215 DeltaDataIntervalTimetable, 

216 ) 

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

218 CronPartitionTimetable, 

219 CronTriggerTimetable, 

220 DeltaTriggerTimetable, 

221 MultipleCronTriggerTimetable, 

222 ) 

223 from airflow.sdk.definitions.variable import Variable 

224 from airflow.sdk.definitions.xcom_arg import XComArg 

225 from airflow.sdk.execution_time import macros 

226 from airflow.sdk.execution_time.context import NEVER_EXPIRE 

227 from airflow.sdk.io.path import ObjectStoragePath 

228 from airflow.sdk.types import TaskInstance 

229 

230 conf: AirflowSDKConfigParser 

231 

232__lazy_imports: dict[str, str] = { 

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

234 "Asset": ".definitions.asset", 

235 "AssetAccessControl": ".definitions.asset", 

236 "AssetAlias": ".definitions.asset", 

237 "AssetAll": ".definitions.asset", 

238 "AssetAny": ".definitions.asset", 

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

240 "AssetWatcher": ".definitions.asset", 

241 "AsyncCallback": ".definitions.callback", 

242 "BaseAsyncOperator": ".bases.operator", 

243 "BaseBranchOperator": ".bases.branch", 

244 "BaseDeadlineReference": ".definitions.deadline", 

245 "BaseHook": ".bases.hook", 

246 "BaseNotifier": ".bases.notifier", 

247 "BaseOperator": ".bases.operator", 

248 "BaseOperatorLink": ".bases.operatorlink", 

249 "BaseSensorOperator": ".bases.sensor", 

250 "BaseXCom": ".bases.xcom", 

251 "BranchMixIn": ".bases.branch", 

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

253 "Connection": ".definitions.connection", 

254 "Context": ".definitions.context", 

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

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

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

258 "DAG": ".definitions.dag", 

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

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

261 "DeadlineAlert": ".definitions.deadline", 

262 "DeadlineReference": ".definitions.deadline", 

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

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

265 "EdgeModifier": ".definitions.edges", 

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

267 "ExceptionRetryPolicy": ".definitions.retry_policy", 

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

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

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

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

272 "Label": ".definitions.edges", 

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

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

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

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

277 "ObjectStoragePath": ".io.path", 

278 "Param": ".definitions.param", 

279 "ParamsDict": ".definitions.param", 

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

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

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

283 "PokeReturnValue": ".bases.sensor", 

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

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

286 "ResumableJobMixin": ".bases.resumablejobmixin", 

287 "RetryAction": ".definitions.retry_policy", 

288 "RetryDecision": ".definitions.retry_policy", 

289 "RetryPolicy": ".definitions.retry_policy", 

290 "RetryRule": ".definitions.retry_policy", 

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

292 "SecretCache": ".execution_time.cache", 

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

294 "SkipMixin": ".bases.skipmixin", 

295 "SyncCallback": ".definitions.callback", 

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

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

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

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

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

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

302 "TaskGroup": ".definitions.taskgroup", 

303 "TaskInstance": ".types", 

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

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

306 "Variable": ".definitions.variable", 

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

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

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

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

311 "XComArg": ".definitions.xcom_arg", 

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

313 "asset": ".definitions.asset.decorators", 

314 "chain": ".bases.operator", 

315 "chain_linear": ".bases.operator", 

316 "conf": ".configuration", 

317 "cross_downstream": ".bases.operator", 

318 "dag": ".definitions.dag", 

319 "deadline_reference": ".definitions.deadline", 

320 "NEVER_EXPIRE": ".execution_time.context", 

321 "get_current_context": ".definitions.context", 

322 "get_parsing_context": ".definitions.context", 

323 "literal": ".definitions.template", 

324 "lineage": ".lineage", 

325 "macros": ".execution_time", 

326 "result": ".definitions.decorators", 

327 "setup": ".definitions.decorators", 

328 "task": ".definitions.decorators", 

329 "task_group": ".definitions.decorators", 

330 "teardown": ".definitions.decorators", 

331} 

332 

333 

334def __getattr__(name: str): 

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

336 import importlib 

337 

338 mod = importlib.import_module(module_path, __name__) 

339 val = getattr(mod, name) 

340 

341 # Store for next time 

342 globals()[name] = val 

343 return val 

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