Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/airflow/sdk/api/datamodels/_generated.py: 99%

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

387 statements  

1# generated by datamodel-codegen: 

2# filename: http://0.0.0.0:8080/execution/openapi.json 

3# version: 0.82.0 

4 

5# Licensed to the Apache Software Foundation (ASF) under one 

6# or more contributor license agreements. See the NOTICE file 

7# distributed with this work for additional information 

8# regarding copyright ownership. The ASF licenses this file 

9# to you under the Apache License, Version 2.0 (the 

10# "License"); you may not use this file except in compliance 

11# with the License. You may obtain a copy of the License at 

12# 

13# http://www.apache.org/licenses/LICENSE-2.0 

14# 

15# Unless required by applicable law or agreed to in writing, 

16# software distributed under the License is distributed on an 

17# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY 

18# KIND, either express or implied. See the License for the 

19# specific language governing permissions and limitations 

20# under the License. 

21from __future__ import annotations 

22 

23from datetime import timedelta 

24from enum import Enum 

25from typing import Annotated, Any, Final, Literal 

26from uuid import UUID 

27 

28from pydantic import AwareDatetime, BaseModel, ConfigDict, Field, JsonValue, RootModel 

29 

30API_VERSION: Final[str] = "2026-10-30" 

31 

32 

33class AssetAliasReferenceAssetEventDagRun(BaseModel): 

34 """ 

35 Schema for AssetAliasModel used in AssetEventDagRunReference. 

36 """ 

37 

38 model_config = ConfigDict( 

39 extra="forbid", 

40 ) 

41 name: Annotated[str, Field(title="Name")] 

42 

43 

44class AssetProfile(BaseModel): 

45 """ 

46 Profile of an asset-like object. 

47 

48 Asset will have name, uri defined, with type set to 'Asset'. 

49 AssetNameRef will have name defined, type set to 'AssetNameRef'. 

50 AssetUriRef will have uri defined, type set to 'AssetUriRef'. 

51 AssetAlias will have name defined, type set to 'AssetAlias'. 

52 

53 Note that 'type' here is distinct from 'asset_type' the user declares on an 

54 Asset (or subclass). This field is for distinguishing between different 

55 asset-related types (Asset, AssetRef, or AssetAlias). 

56 """ 

57 

58 model_config = ConfigDict( 

59 extra="forbid", 

60 ) 

61 name: Annotated[str | None, Field(title="Name")] = None 

62 uri: Annotated[str | None, Field(title="Uri")] = None 

63 type: Annotated[str, Field(title="Type")] 

64 

65 

66class ConnectionResponse(BaseModel): 

67 """ 

68 Connection schema for responses with fields that are needed for Runtime. 

69 """ 

70 

71 conn_id: Annotated[str, Field(title="Conn Id")] 

72 conn_type: Annotated[str, Field(title="Conn Type")] 

73 host: Annotated[str | None, Field(title="Host")] 

74 schema_: Annotated[str | None, Field(alias="schema", title="Schema")] 

75 login: Annotated[str | None, Field(title="Login")] 

76 password: Annotated[str | None, Field(title="Password")] 

77 port: Annotated[int | None, Field(title="Port")] 

78 extra: Annotated[str | None, Field(title="Extra")] 

79 

80 

81class ConnectionTestConnectionResponse(BaseModel): 

82 """ 

83 Connection data returned to workers from a test request. 

84 """ 

85 

86 conn_id: Annotated[str, Field(title="Conn Id")] 

87 conn_type: Annotated[str, Field(title="Conn Type")] 

88 host: Annotated[str | None, Field(title="Host")] = None 

89 login: Annotated[str | None, Field(title="Login")] = None 

90 password: Annotated[str | None, Field(title="Password")] = None 

91 schema_: Annotated[str | None, Field(alias="schema", title="Schema")] = None 

92 port: Annotated[int | None, Field(title="Port")] = None 

93 extra: Annotated[str | None, Field(title="Extra")] = None 

94 

95 

96class ResultMessage(RootModel[str]): 

97 root: Annotated[str, Field(max_length=2000, title="Result Message")] 

98 

99 

100class ConnectionTestState(str, Enum): 

101 """ 

102 All possible states of a connection test. 

103 """ 

104 

105 PENDING = "pending" 

106 QUEUED = "queued" 

107 RUNNING = "running" 

108 SUCCESS = "success" 

109 FAILED = "failed" 

110 

111 

112class DagResponse(BaseModel): 

113 """ 

114 Schema for DAG response. 

115 """ 

116 

117 dag_id: Annotated[str, Field(title="Dag Id")] 

118 is_paused: Annotated[bool, Field(title="Is Paused")] 

119 bundle_name: Annotated[str | None, Field(title="Bundle Name")] 

120 bundle_version: Annotated[str | None, Field(title="Bundle Version")] 

121 relative_fileloc: Annotated[str | None, Field(title="Relative Fileloc")] 

122 owners: Annotated[str | None, Field(title="Owners")] 

123 tags: Annotated[list[str], Field(title="Tags")] 

124 next_dagrun: Annotated[AwareDatetime | None, Field(title="Next Dagrun")] 

125 

126 

127class DagRunAssetReference(BaseModel): 

128 """ 

129 DagRun serializer for asset responses. 

130 """ 

131 

132 model_config = ConfigDict( 

133 extra="forbid", 

134 ) 

135 run_id: Annotated[str, Field(title="Run Id")] 

136 dag_id: Annotated[str, Field(title="Dag Id")] 

137 logical_date: Annotated[AwareDatetime | None, Field(title="Logical Date")] 

138 start_date: Annotated[AwareDatetime | None, Field(title="Start Date")] 

139 end_date: Annotated[AwareDatetime | None, Field(title="End Date")] 

140 state: Annotated[str, Field(title="State")] 

141 data_interval_start: Annotated[AwareDatetime | None, Field(title="Data Interval Start")] 

142 data_interval_end: Annotated[AwareDatetime | None, Field(title="Data Interval End")] 

143 partition_key: Annotated[str | None, Field(title="Partition Key")] 

144 

145 

146class DagRunState(str, Enum): 

147 """ 

148 All possible states that a DagRun can be in. 

149 

150 These are "shared" with TaskInstanceState in some parts of the code, 

151 so please ensure that their values always match the ones with the 

152 same name in TaskInstanceState. 

153 """ 

154 

155 QUEUED = "queued" 

156 RUNNING = "running" 

157 SUCCESS = "success" 

158 FAILED = "failed" 

159 

160 

161class DagRunStateResponse(BaseModel): 

162 """ 

163 Schema for DAG Run State response. 

164 """ 

165 

166 state: DagRunState 

167 

168 

169class DagRunType(str, Enum): 

170 """ 

171 Class with DagRun types. 

172 """ 

173 

174 BACKFILL = "backfill" 

175 SCHEDULED = "scheduled" 

176 MANUAL = "manual" 

177 OPERATOR_TRIGGERED = "operator_triggered" 

178 ASSET_TRIGGERED = "asset_triggered" 

179 ASSET_MATERIALIZATION = "asset_materialization" 

180 

181 

182class HITLUser(BaseModel): 

183 """ 

184 Schema for a Human-in-the-loop users. 

185 """ 

186 

187 id: Annotated[str, Field(title="Id")] 

188 name: Annotated[str, Field(title="Name")] 

189 

190 

191class HTTPExceptionResponse(BaseModel): 

192 """ 

193 HTTPException Model used for error response. 

194 """ 

195 

196 detail: Annotated[str | dict[str, Any], Field(title="Detail")] 

197 

198 

199class InactiveAssetsResponse(BaseModel): 

200 """ 

201 Response for inactive assets. 

202 """ 

203 

204 inactive_assets: Annotated[list[AssetProfile] | None, Field(title="Inactive Assets")] = None 

205 

206 

207class IntermediateTIState(str, Enum): 

208 """ 

209 States that a Task Instance can be in that indicate it is not yet in a terminal or running state. 

210 """ 

211 

212 SCHEDULED = "scheduled" 

213 QUEUED = "queued" 

214 RESTARTING = "restarting" 

215 UP_FOR_RETRY = "up_for_retry" 

216 UP_FOR_RESCHEDULE = "up_for_reschedule" 

217 DEFERRED = "deferred" 

218 AWAITING_INPUT = "awaiting_input" 

219 

220 

221class PrevSuccessfulDagRunResponse(BaseModel): 

222 """ 

223 Schema for response with previous successful DagRun information for Task Template Context. 

224 """ 

225 

226 data_interval_start: Annotated[AwareDatetime | None, Field(title="Data Interval Start")] = None 

227 data_interval_end: Annotated[AwareDatetime | None, Field(title="Data Interval End")] = None 

228 start_date: Annotated[AwareDatetime | None, Field(title="Start Date")] = None 

229 end_date: Annotated[AwareDatetime | None, Field(title="End Date")] = None 

230 

231 

232class PreviousTIResponse(BaseModel): 

233 """ 

234 Schema for response with previous TaskInstance information. 

235 """ 

236 

237 task_id: Annotated[str, Field(title="Task Id")] 

238 dag_id: Annotated[str, Field(title="Dag Id")] 

239 run_id: Annotated[str, Field(title="Run Id")] 

240 logical_date: Annotated[AwareDatetime | None, Field(title="Logical Date")] = None 

241 start_date: Annotated[AwareDatetime | None, Field(title="Start Date")] = None 

242 end_date: Annotated[AwareDatetime | None, Field(title="End Date")] = None 

243 state: Annotated[str | None, Field(title="State")] = None 

244 try_number: Annotated[int, Field(title="Try Number")] 

245 map_index: Annotated[int | None, Field(title="Map Index")] = -1 

246 duration: Annotated[float | None, Field(title="Duration")] = None 

247 

248 

249class TIAwaitingInputStatePayload(BaseModel): 

250 """ 

251 Schema for parking a TaskInstance in an awaiting_input state (Human-in-the-loop, no trigger). 

252 """ 

253 

254 model_config = ConfigDict( 

255 extra="forbid", 

256 ) 

257 state: Annotated[Literal["awaiting_input"], Field(title="State")] = "awaiting_input" 

258 timeout: Annotated[timedelta | None, Field(title="Timeout")] = None 

259 next_method: Annotated[str, Field(title="Next Method")] 

260 next_kwargs: Annotated[dict[str, JsonValue | None] | None, Field(title="Next Kwargs")] = None 

261 rendered_map_index: Annotated[str | None, Field(title="Rendered Map Index")] = None 

262 

263 

264class TIDeferredStatePayload(BaseModel): 

265 """ 

266 Schema for updating TaskInstance to a deferred state. 

267 """ 

268 

269 model_config = ConfigDict( 

270 extra="forbid", 

271 ) 

272 state: Annotated[Literal["deferred"], Field(title="State")] = "deferred" 

273 classpath: Annotated[str, Field(title="Classpath")] 

274 trigger_kwargs: Annotated[dict[str, JsonValue | None] | str | None, Field(title="Trigger Kwargs")] = None 

275 trigger_timeout: Annotated[timedelta | None, Field(title="Trigger Timeout")] = None 

276 queue: Annotated[str | None, Field(title="Queue")] = None 

277 next_method: Annotated[str, Field(title="Next Method")] 

278 next_kwargs: Annotated[dict[str, JsonValue | None] | None, Field(title="Next Kwargs")] = None 

279 rendered_map_index: Annotated[str | None, Field(title="Rendered Map Index")] = None 

280 

281 

282class TIEnterRunningPayload(BaseModel): 

283 """ 

284 Schema for updating TaskInstance to 'RUNNING' state with minimal required fields. 

285 """ 

286 

287 model_config = ConfigDict( 

288 extra="forbid", 

289 ) 

290 state: Annotated[Literal["running"], Field(title="State")] = "running" 

291 hostname: Annotated[str, Field(title="Hostname")] 

292 unixname: Annotated[str, Field(title="Unixname")] 

293 pid: Annotated[int, Field(title="Pid")] 

294 start_date: Annotated[AwareDatetime, Field(title="Start Date")] 

295 

296 

297class TIHeartbeatInfo(BaseModel): 

298 """ 

299 Schema for TaskInstance heartbeat endpoint. 

300 """ 

301 

302 model_config = ConfigDict( 

303 extra="forbid", 

304 ) 

305 hostname: Annotated[str, Field(title="Hostname")] 

306 pid: Annotated[int, Field(title="Pid")] 

307 

308 

309class TIRescheduleStatePayload(BaseModel): 

310 """ 

311 Schema for updating TaskInstance to a up_for_reschedule state. 

312 """ 

313 

314 model_config = ConfigDict( 

315 extra="forbid", 

316 ) 

317 state: Annotated[Literal["up_for_reschedule"], Field(title="State")] = "up_for_reschedule" 

318 reschedule_date: Annotated[AwareDatetime, Field(title="Reschedule Date")] 

319 end_date: Annotated[AwareDatetime, Field(title="End Date")] 

320 

321 

322class TIRetryStatePayload(BaseModel): 

323 """ 

324 Schema for updating TaskInstance to up_for_retry. 

325 """ 

326 

327 model_config = ConfigDict( 

328 extra="forbid", 

329 ) 

330 state: Annotated[Literal["up_for_retry"], Field(title="State")] = "up_for_retry" 

331 end_date: Annotated[AwareDatetime, Field(title="End Date")] 

332 rendered_map_index: Annotated[str | None, Field(title="Rendered Map Index")] = None 

333 retry_delay_seconds: Annotated[float | None, Field(title="Retry Delay Seconds")] = None 

334 retry_reason: Annotated[str | None, Field(title="Retry Reason")] = None 

335 

336 

337class TISkippedDownstreamTasksStatePayload(BaseModel): 

338 """ 

339 Schema for updating downstream tasks to a skipped state. 

340 """ 

341 

342 model_config = ConfigDict( 

343 extra="forbid", 

344 ) 

345 tasks: Annotated[list[str | tuple[str, int]], Field(title="Tasks")] 

346 

347 

348class TISuccessStatePayload(BaseModel): 

349 """ 

350 Schema for updating TaskInstance to success state. 

351 """ 

352 

353 model_config = ConfigDict( 

354 extra="forbid", 

355 ) 

356 state: Annotated[Literal["success"], Field(title="State")] = "success" 

357 end_date: Annotated[AwareDatetime, Field(title="End Date")] 

358 task_outlets: Annotated[list[AssetProfile] | None, Field(title="Task Outlets")] = None 

359 outlet_events: Annotated[list[dict[str, Any]] | None, Field(title="Outlet Events")] = None 

360 rendered_map_index: Annotated[str | None, Field(title="Rendered Map Index")] = None 

361 

362 

363class TITargetStatePayload(BaseModel): 

364 """ 

365 Schema for updating TaskInstance to a target state, excluding terminal and running states. 

366 """ 

367 

368 model_config = ConfigDict( 

369 extra="forbid", 

370 ) 

371 state: IntermediateTIState 

372 

373 

374class TaskBreadcrumbsResponse(BaseModel): 

375 """ 

376 Response for task breadcrumbs. 

377 """ 

378 

379 breadcrumbs: Annotated[list[dict[str, Any]], Field(title="Breadcrumbs")] 

380 

381 

382class TaskInstanceState(str, Enum): 

383 """ 

384 All possible states that a Task Instance can be in. 

385 

386 Note that None is also allowed, so always use this in a type hint with Optional. 

387 """ 

388 

389 REMOVED = "removed" 

390 SCHEDULED = "scheduled" 

391 QUEUED = "queued" 

392 RUNNING = "running" 

393 SUCCESS = "success" 

394 RESTARTING = "restarting" 

395 FAILED = "failed" 

396 UP_FOR_RETRY = "up_for_retry" 

397 UP_FOR_RESCHEDULE = "up_for_reschedule" 

398 UPSTREAM_FAILED = "upstream_failed" 

399 SKIPPED = "skipped" 

400 DEFERRED = "deferred" 

401 AWAITING_INPUT = "awaiting_input" 

402 

403 

404class TaskStateStorePutBody(BaseModel): 

405 """ 

406 Request body for setting a task state store value. 

407 """ 

408 

409 model_config = ConfigDict( 

410 extra="forbid", 

411 ) 

412 value: JsonValue | None 

413 expires_at: Annotated[AwareDatetime | None, Field(title="Expires At")] = None 

414 

415 

416class TaskStateStoreResponse(BaseModel): 

417 """ 

418 Task state store value returned to a worker. 

419 """ 

420 

421 model_config = ConfigDict( 

422 extra="forbid", 

423 ) 

424 value: JsonValue | None 

425 

426 

427class TaskStatesResponse(BaseModel): 

428 """ 

429 Response for task states with run_id, task and state. 

430 """ 

431 

432 task_states: Annotated[dict[str, Any], Field(title="Task States")] 

433 

434 

435class TerminalStateNonSuccess(str, Enum): 

436 """ 

437 TaskInstance states that can be reported without extra information. 

438 """ 

439 

440 FAILED = "failed" 

441 SKIPPED = "skipped" 

442 REMOVED = "removed" 

443 UPSTREAM_FAILED = "upstream_failed" 

444 SERVER_TERMINATED = "server_terminated" 

445 

446 

447class TriggerDAGRunPayload(BaseModel): 

448 """ 

449 Schema for Trigger DAG Run API request. 

450 """ 

451 

452 model_config = ConfigDict( 

453 extra="forbid", 

454 ) 

455 logical_date: Annotated[AwareDatetime | None, Field(title="Logical Date")] = None 

456 run_after: Annotated[AwareDatetime | None, Field(title="Run After")] = None 

457 conf: Annotated[dict[str, Any] | None, Field(title="Conf")] = None 

458 reset_dag_run: Annotated[bool | None, Field(title="Reset Dag Run")] = False 

459 partition_key: Annotated[str | None, Field(title="Partition Key")] = None 

460 note: Annotated[str | None, Field(title="Note")] = None 

461 

462 

463class UpdateHITLDetailPayload(BaseModel): 

464 """ 

465 Schema for writing the response part of a Human-in-the-loop detail for a specific task instance. 

466 """ 

467 

468 ti_id: Annotated[UUID, Field(title="Ti Id")] 

469 chosen_options: Annotated[list[str], Field(min_length=1, title="Chosen Options")] 

470 params_input: Annotated[dict[str, Any] | None, Field(title="Params Input")] = None 

471 

472 

473class ValidationError(BaseModel): 

474 loc: Annotated[list[str | int], Field(title="Location")] 

475 msg: Annotated[str, Field(title="Message")] 

476 type: Annotated[str, Field(title="Error Type")] 

477 input: Annotated[Any | None, Field(title="Input")] = None 

478 ctx: Annotated[dict[str, Any] | None, Field(title="Context")] = None 

479 

480 

481class VariableKeysResponse(BaseModel): 

482 """ 

483 Variable keys schema for list responses. 

484 """ 

485 

486 model_config = ConfigDict( 

487 extra="forbid", 

488 ) 

489 keys: Annotated[list[str], Field(title="Keys")] 

490 total_entries: Annotated[int, Field(title="Total Entries")] 

491 

492 

493class VariablePostBody(BaseModel): 

494 """ 

495 Request body schema for creating variables. 

496 """ 

497 

498 model_config = ConfigDict( 

499 extra="forbid", 

500 ) 

501 val: Annotated[str | None, Field(title="Val")] 

502 description: Annotated[str | None, Field(title="Description")] = None 

503 

504 

505class VariableResponse(BaseModel): 

506 """ 

507 Variable schema for responses with fields that are needed for Runtime. 

508 """ 

509 

510 model_config = ConfigDict( 

511 extra="forbid", 

512 ) 

513 key: Annotated[str, Field(title="Key")] 

514 value: Annotated[str | None, Field(title="Value")] 

515 

516 

517class XComResponse(BaseModel): 

518 """ 

519 XCom schema for responses with fields that are needed for Runtime. 

520 """ 

521 

522 key: Annotated[str, Field(title="Key")] 

523 value: JsonValue | None 

524 

525 

526class XComSequenceIndexResponse(RootModel[JsonValue | None]): 

527 root: Annotated[ 

528 JsonValue | None, 

529 Field( 

530 description="XCom schema with minimal structure for index-based access.", 

531 title="XComSequenceIndexResponse", 

532 ), 

533 ] 

534 

535 

536class XComSequenceSliceResponse(RootModel[list[JsonValue | None]]): 

537 """ 

538 XCom schema with minimal structure for slice-based access. 

539 """ 

540 

541 root: Annotated[ 

542 list[JsonValue | None], 

543 Field( 

544 description="XCom schema with minimal structure for slice-based access.", 

545 title="XComSequenceSliceResponse", 

546 ), 

547 ] 

548 

549 

550class TaskInstance(BaseModel): 

551 """ 

552 Schema for TaskInstance model with minimal required fields needed for Runtime. 

553 """ 

554 

555 id: Annotated[UUID, Field(title="Id")] 

556 task_id: Annotated[str, Field(title="Task Id")] 

557 dag_id: Annotated[str, Field(title="Dag Id")] 

558 run_id: Annotated[str, Field(title="Run Id")] 

559 try_number: Annotated[int, Field(title="Try Number")] 

560 dag_version_id: Annotated[UUID, Field(title="Dag Version Id")] 

561 map_index: Annotated[int | None, Field(title="Map Index")] = -1 

562 hostname: Annotated[str | None, Field(title="Hostname")] = None 

563 context_carrier: Annotated[dict[str, Any] | None, Field(title="Context Carrier")] = None 

564 queue: Annotated[str | None, Field(title="Queue")] = "default" 

565 

566 

567class BundleInfo(BaseModel): 

568 """ 

569 Schema for telling task which bundle to run with. 

570 """ 

571 

572 name: Annotated[str, Field(title="Name")] 

573 version: Annotated[str | None, Field(title="Version")] = None 

574 version_data: Annotated[dict[str, Any] | None, Field(title="Version Data")] = None 

575 

576 

577class TerminalTIState(str, Enum): 

578 SUCCESS = "success" 

579 FAILED = "failed" 

580 SKIPPED = "skipped" 

581 UPSTREAM_FAILED = "upstream_failed" 

582 REMOVED = "removed" 

583 

584 

585class WeightRule(str, Enum): 

586 DOWNSTREAM = "downstream" 

587 UPSTREAM = "upstream" 

588 ABSOLUTE = "absolute" 

589 

590 

591class TriggerRule(str, Enum): 

592 ALL_SUCCESS = "all_success" 

593 ALL_FAILED = "all_failed" 

594 ALL_DONE = "all_done" 

595 ALL_DONE_MIN_ONE_SUCCESS = "all_done_min_one_success" 

596 ALL_DONE_SETUP_SUCCESS = "all_done_setup_success" 

597 ONE_SUCCESS = "one_success" 

598 ONE_FAILED = "one_failed" 

599 ONE_DONE = "one_done" 

600 NONE_FAILED = "none_failed" 

601 NONE_SKIPPED = "none_skipped" 

602 ALWAYS = "always" 

603 NONE_FAILED_MIN_ONE_SUCCESS = "none_failed_min_one_success" 

604 ALL_SKIPPED = "all_skipped" 

605 

606 

607class DagAttributeTypes(str, Enum): 

608 OP = "operator" 

609 TASK_GROUP = "taskgroup" 

610 

611 

612class ArgValueSchema(RootModel[dict[str, JsonValue | None]]): 

613 root: dict[str, JsonValue | None] 

614 

615 

616class AssetReferenceAssetEventDagRun(BaseModel): 

617 """ 

618 Schema for AssetModel used in AssetEventDagRunReference. 

619 """ 

620 

621 model_config = ConfigDict( 

622 extra="forbid", 

623 ) 

624 name: Annotated[str, Field(title="Name")] 

625 uri: Annotated[str, Field(title="Uri")] 

626 extra: Annotated[dict[str, JsonValue | None], Field(title="Extra")] 

627 

628 

629class AssetResponse(BaseModel): 

630 """ 

631 Asset schema for responses with fields that are needed for Runtime. 

632 """ 

633 

634 name: Annotated[str, Field(title="Name")] 

635 uri: Annotated[str, Field(title="Uri")] 

636 group: Annotated[str, Field(title="Group")] 

637 extra: Annotated[dict[str, JsonValue | None] | None, Field(title="Extra")] = None 

638 

639 

640class AssetStateStorePutBody(BaseModel): 

641 """ 

642 Request body for setting an asset state store value. 

643 """ 

644 

645 model_config = ConfigDict( 

646 extra="forbid", 

647 ) 

648 value: JsonValue | None 

649 

650 

651class AssetStateStoreResponse(BaseModel): 

652 """ 

653 Asset state store value returned to a worker. 

654 """ 

655 

656 model_config = ConfigDict( 

657 extra="forbid", 

658 ) 

659 value: JsonValue | None 

660 

661 

662class ConnectionTestResultBody(BaseModel): 

663 """ 

664 Result a worker reports back for a connection test. 

665 """ 

666 

667 model_config = ConfigDict( 

668 extra="forbid", 

669 ) 

670 state: ConnectionTestState 

671 result_message: Annotated[ResultMessage | None, Field(title="Result Message")] = None 

672 

673 

674class HITLDetailRequest(BaseModel): 

675 """ 

676 Schema for the request part of a Human-in-the-loop detail for a specific task instance. 

677 """ 

678 

679 ti_id: Annotated[UUID, Field(title="Ti Id")] 

680 options: Annotated[list[str], Field(min_length=1, title="Options")] 

681 subject: Annotated[str, Field(title="Subject")] 

682 body: Annotated[str | None, Field(title="Body")] = None 

683 defaults: Annotated[list[str] | None, Field(title="Defaults")] = None 

684 multiple: Annotated[bool | None, Field(title="Multiple")] = False 

685 params: Annotated[dict[str, Any] | None, Field(title="Params")] = None 

686 assigned_users: Annotated[list[HITLUser] | None, Field(title="Assigned Users")] = None 

687 

688 

689class HITLDetailResponse(BaseModel): 

690 """ 

691 Schema for the response part of a Human-in-the-loop detail for a specific task instance. 

692 """ 

693 

694 response_received: Annotated[bool, Field(title="Response Received")] 

695 responded_by_user: HITLUser | None = None 

696 responded_at: Annotated[AwareDatetime | None, Field(title="Responded At")] 

697 chosen_options: Annotated[list[str] | None, Field(title="Chosen Options")] 

698 params_input: Annotated[dict[str, Any] | None, Field(title="Params Input")] = None 

699 

700 

701class HTTPValidationError(BaseModel): 

702 detail: Annotated[list[ValidationError] | None, Field(title="Detail")] = None 

703 

704 

705class LiteralArgBinding(BaseModel): 

706 """ 

707 One positional stub-task argument carrying an inline literal from the Dag file. 

708 """ 

709 

710 name: Annotated[str, Field(title="Name")] 

711 value_schema: ArgValueSchema | None = None 

712 kind: Annotated[Literal["literal"], Field(title="Kind")] 

713 value: JsonValue | None = None 

714 from_default: Annotated[bool | None, Field(title="From Default")] = False 

715 

716 

717class TITerminalStatePayload(BaseModel): 

718 """ 

719 Report a terminal outcome other than SUCCESS, or acknowledge server-requested termination. 

720 """ 

721 

722 model_config = ConfigDict( 

723 extra="forbid", 

724 ) 

725 state: TerminalStateNonSuccess 

726 end_date: Annotated[AwareDatetime, Field(title="End Date")] 

727 rendered_map_index: Annotated[str | None, Field(title="Rendered Map Index")] = None 

728 retry_reason: Annotated[str | None, Field(title="Retry Reason")] = None 

729 hostname: Annotated[str | None, Field(title="Hostname")] = None 

730 pid: Annotated[int | None, Field(title="Pid")] = None 

731 

732 

733class XComArgBinding(BaseModel): 

734 """ 

735 One positional stub-task argument pulled from an upstream task's XCom. 

736 """ 

737 

738 name: Annotated[str, Field(title="Name")] 

739 value_schema: ArgValueSchema | None = None 

740 kind: Annotated[Literal["xcom"], Field(title="Kind")] 

741 task_id: Annotated[str, Field(title="Task Id")] 

742 

743 

744class AssetEventDagRunReference(BaseModel): 

745 """ 

746 Schema for AssetEvent model used in DagRun. 

747 """ 

748 

749 model_config = ConfigDict( 

750 extra="forbid", 

751 ) 

752 asset: AssetReferenceAssetEventDagRun 

753 extra: Annotated[dict[str, JsonValue | None], Field(title="Extra")] 

754 source_task_id: Annotated[str | None, Field(title="Source Task Id")] 

755 source_dag_id: Annotated[str | None, Field(title="Source Dag Id")] 

756 source_run_id: Annotated[str | None, Field(title="Source Run Id")] 

757 source_map_index: Annotated[int | None, Field(title="Source Map Index")] 

758 source_aliases: Annotated[list[AssetAliasReferenceAssetEventDagRun], Field(title="Source Aliases")] 

759 timestamp: Annotated[AwareDatetime, Field(title="Timestamp")] 

760 partition_key: Annotated[str | None, Field(title="Partition Key")] = None 

761 

762 

763class AssetEventResponse(BaseModel): 

764 """ 

765 Asset event schema with fields that are needed for Runtime. 

766 """ 

767 

768 id: Annotated[int, Field(title="Id")] 

769 timestamp: Annotated[AwareDatetime, Field(title="Timestamp")] 

770 extra: Annotated[dict[str, JsonValue | None] | None, Field(title="Extra")] = None 

771 asset: AssetResponse 

772 created_dagruns: Annotated[list[DagRunAssetReference], Field(title="Created Dagruns")] 

773 source_task_id: Annotated[str | None, Field(title="Source Task Id")] = None 

774 source_dag_id: Annotated[str | None, Field(title="Source Dag Id")] = None 

775 source_run_id: Annotated[str | None, Field(title="Source Run Id")] = None 

776 source_map_index: Annotated[int | None, Field(title="Source Map Index")] = None 

777 partition_key: Annotated[str | None, Field(title="Partition Key")] = None 

778 

779 

780class AssetEventsResponse(BaseModel): 

781 """ 

782 Collection of AssetEventResponse. 

783 """ 

784 

785 asset_events: Annotated[list[AssetEventResponse], Field(title="Asset Events")] 

786 

787 

788class DagRun(BaseModel): 

789 """ 

790 Schema for DagRun model with minimal required fields needed for Runtime. 

791 """ 

792 

793 model_config = ConfigDict( 

794 extra="forbid", 

795 ) 

796 dag_id: Annotated[str, Field(title="Dag Id")] 

797 run_id: Annotated[str, Field(title="Run Id")] 

798 logical_date: Annotated[AwareDatetime | None, Field(title="Logical Date")] 

799 data_interval_start: Annotated[AwareDatetime | None, Field(title="Data Interval Start")] 

800 data_interval_end: Annotated[AwareDatetime | None, Field(title="Data Interval End")] 

801 run_after: Annotated[AwareDatetime, Field(title="Run After")] 

802 start_date: Annotated[AwareDatetime | None, Field(title="Start Date")] 

803 end_date: Annotated[AwareDatetime | None, Field(title="End Date")] 

804 clear_number: Annotated[int | None, Field(title="Clear Number")] = 0 

805 run_type: DagRunType 

806 state: DagRunState 

807 conf: Annotated[dict[str, Any] | None, Field(title="Conf")] = None 

808 triggering_user_name: Annotated[str | None, Field(title="Triggering User Name")] = None 

809 consumed_asset_events: Annotated[list[AssetEventDagRunReference], Field(title="Consumed Asset Events")] 

810 partition_key: Annotated[str | None, Field(title="Partition Key")] 

811 partition_date: Annotated[AwareDatetime | None, Field(title="Partition Date")] = None 

812 note: Annotated[str | None, Field(title="Note")] = None 

813 team_name: Annotated[str | None, Field(title="Team Name")] = None 

814 

815 

816class TaskArgBinding(RootModel[XComArgBinding | LiteralArgBinding]): 

817 root: Annotated[XComArgBinding | LiteralArgBinding, Field(discriminator="kind", title="TaskArgBinding")] 

818 

819 

820class TIRunContext(BaseModel): 

821 """ 

822 Response schema for TaskInstance run context. 

823 """ 

824 

825 dag_run: DagRun 

826 task_reschedule_count: Annotated[int | None, Field(title="Task Reschedule Count")] = 0 

827 max_tries: Annotated[int, Field(title="Max Tries")] 

828 variables: Annotated[list[VariableResponse] | None, Field(title="Variables")] = None 

829 connections: Annotated[list[ConnectionResponse] | None, Field(title="Connections")] = None 

830 next_method: Annotated[str | None, Field(title="Next Method")] = None 

831 next_kwargs: Annotated[dict[str, Any] | str | None, Field(title="Next Kwargs")] = None 

832 xcom_keys_to_clear: Annotated[list[str] | None, Field(title="Xcom Keys To Clear")] = None 

833 should_retry: Annotated[bool | None, Field(title="Should Retry")] = False 

834 start_date: Annotated[AwareDatetime | None, Field(title="Start Date")] = None 

835 arg_bindings: Annotated[list[TaskArgBinding] | None, Field(title="Arg Bindings")] = None 

836 multi_team: Annotated[bool | None, Field(title="Multi Team")] = False