Coverage for async_durable_execution/_core/models.py: 97%

336 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-08-30 23:43 +0000

1"""Core durable execution data models.""" 

2 

3from __future__ import annotations 

4 

5import datetime 

6from collections.abc import Mapping, MutableMapping 

7from dataclasses import MISSING, dataclass, field, fields 

8from enum import Enum 

9from typing import ( 

10 Any, 

11 Protocol, 

12 TypeAlias, 

13 TypeVar, 

14 cast, 

15 get_args, 

16 get_origin, 

17 get_type_hints, 

18) 

19 

20# Replace with `type` it when dropping support to Python 3.11 

21ReplayChildren: TypeAlias = bool 

22OperationPayload: TypeAlias = str 

23TimeoutSeconds: TypeAlias = int 

24ModelT = TypeVar("ModelT") 

25 

26 

27class LambdaContext(Protocol): 

28 """Minimal AWS Lambda context surface used by the SDK.""" 

29 

30 aws_request_id: str 

31 log_group_name: str | None = None 

32 log_stream_name: str | None = None 

33 function_name: str | None = None 

34 memory_limit_in_mb: str | None = None 

35 function_version: str | None = None 

36 invoked_function_arn: str | None = None 

37 tenant_id: str | None = None 

38 client_context: Any | None = None 

39 identity: Any | None = None 

40 

41 def get_remaining_time_in_millis(self) -> int: ... 

42 def log(self, msg) -> None: ... 

43 

44 

45def _model_field_metadata( 

46 *, 

47 alias: str, 

48 serializer: Any = None, 

49 deserializer: Any = None, 

50 omit_if_none: bool = True, 

51 omit_if_falsey: bool = False, 

52 is_timestamp: bool = False, 

53) -> dict[str, Any]: 

54 """Build dataclass field metadata for the shared mapping engine.""" 

55 return { 

56 "alias": alias, 

57 "serializer": serializer, 

58 "deserializer": deserializer, 

59 "omit_if_none": omit_if_none, 

60 "omit_if_falsey": omit_if_falsey, 

61 "is_timestamp": is_timestamp, 

62 } 

63 

64 

65class MappingModel: 

66 """Dataclass serialized through the shared Python mapping engine.""" 

67 

68 @classmethod 

69 def from_dict( 

70 cls: type[ModelT], 

71 data: Mapping[str, Any], 

72 ) -> ModelT: 

73 """Construct a model from its serialized mapping.""" 

74 return _model_from_mapping(cls, data) 

75 

76 def to_dict(self) -> MutableMapping[str, Any]: 

77 """Convert the model to its serialized mapping.""" 

78 return _model_to_mapping(self) 

79 

80 

81class AwsApiModel: 

82 """Dataclass serialized as an AWS API mapping with native Python values.""" 

83 

84 @classmethod 

85 def from_dict(cls: type[ModelT], data: Mapping[str, Any]) -> ModelT: 

86 """Create a model from an AWS API mapping with native Python values.""" 

87 return _model_from_mapping(cls, data) 

88 

89 def to_dict(self) -> MutableMapping[str, Any]: 

90 """Convert the model to an AWS API mapping with native Python values.""" 

91 return _model_to_mapping(self) 

92 

93 

94# Backward-compatible internal alias retained for extensions importing the old name. 

95BotoSerializableModel = AwsApiModel 

96 

97 

98class JsonSerializableModel: 

99 @classmethod 

100 def from_dict(cls: type[ModelT], data: Mapping[str, Any]) -> ModelT: 

101 """Create a model from a JSON-compatible mapping.""" 

102 return _model_from_mapping(cls, data, json_mode=True) 

103 

104 def to_dict(self) -> MutableMapping[str, Any]: 

105 """Convert the model to a JSON-compatible mapping.""" 

106 return _model_to_mapping(self, json_mode=True) 

107 

108 

109_MAPPING_MODEL_TYPES = ( 

110 MappingModel, 

111 AwsApiModel, 

112 JsonSerializableModel, 

113) 

114 

115 

116def _enum_type(annotation: Any) -> type[Enum] | None: 

117 if isinstance(annotation, type) and issubclass(annotation, Enum): 

118 return annotation 

119 

120 origin = get_origin(annotation) 

121 if origin is None: 

122 return None 

123 

124 for arg in get_args(annotation): 

125 enum_cls = _enum_type(arg) 

126 if enum_cls is not None: 

127 return enum_cls 

128 

129 return None 

130 

131 

132def _accepts_plain_string(annotation: Any) -> bool: 

133 if annotation is str: 

134 return True 

135 

136 origin = get_origin(annotation) 

137 if origin is None: 

138 return False 

139 

140 return any(_accepts_plain_string(arg) for arg in get_args(annotation)) 

141 

142 

143def _model_type(annotation: Any) -> type[Any] | None: 

144 if isinstance(annotation, type) and issubclass(annotation, _MAPPING_MODEL_TYPES): 

145 return annotation 

146 

147 origin = get_origin(annotation) 

148 if isinstance(origin, type) and issubclass(origin, _MAPPING_MODEL_TYPES): 148 ↛ 149line 148 didn't jump to line 149 because the condition on line 148 was never true

149 return origin 

150 if origin is None: 

151 return None 

152 

153 for arg in get_args(annotation): 

154 model_cls = _model_type(arg) 

155 if model_cls is not None: 

156 return model_cls 

157 

158 return None 

159 

160 

161def _deserialize_value( 

162 value: Any, 

163 annotation: Any, 

164 metadata: Mapping[str, Any], 

165 *, 

166 json_mode: bool = False, 

167) -> Any: 

168 if value is None: 

169 return None 

170 

171 custom_deserializer = metadata.get("deserializer") 

172 if custom_deserializer is not None: 172 ↛ 173line 172 didn't jump to line 173 because the condition on line 172 was never true

173 return custom_deserializer(value) 

174 

175 if json_mode and metadata.get("is_timestamp", False): 

176 return TimestampConverter.from_unix_millis(value) 

177 

178 model_cls = _model_type(annotation) 

179 if model_cls is not None and isinstance(value, Mapping): 

180 return _model_from_mapping(model_cls, value, json_mode=json_mode) 

181 

182 enum_cls = _enum_type(annotation) 

183 if enum_cls is not None: 

184 try: 

185 return enum_cls(value) 

186 except ValueError: 

187 if isinstance(value, str) and _accepts_plain_string(annotation): 187 ↛ 189line 187 didn't jump to line 189 because the condition on line 187 was always true

188 return value 

189 raise 

190 

191 origin = get_origin(annotation) 

192 if origin is list: 

193 args = get_args(annotation) 

194 item_annotation = args[0] if args else Any 

195 return [ 

196 _deserialize_value( 

197 item, 

198 item_annotation, 

199 {}, 

200 json_mode=json_mode, 

201 ) 

202 for item in value 

203 ] 

204 

205 return value 

206 

207 

208def _serialize_value( 

209 value: Any, 

210 metadata: Mapping[str, Any], 

211 *, 

212 json_mode: bool = False, 

213) -> Any: 

214 if value is None: 

215 return None 

216 

217 custom_serializer = metadata.get("serializer") 

218 if custom_serializer is not None: 218 ↛ 219line 218 didn't jump to line 219 because the condition on line 218 was never true

219 return custom_serializer(value) 

220 

221 if json_mode and metadata.get("is_timestamp", False): 

222 return TimestampConverter.to_unix_millis(value) 

223 

224 if isinstance(value, Enum): 

225 return value.value 

226 

227 if isinstance(value, _MAPPING_MODEL_TYPES): 

228 return _model_to_mapping(value, json_mode=json_mode) 

229 

230 if isinstance(value, list): 

231 return [_serialize_value(item, {}, json_mode=json_mode) for item in value] 

232 

233 return value 

234 

235 

236def _model_from_mapping( 

237 model_cls: type[ModelT], 

238 data: Mapping[str, Any], 

239 *, 

240 json_mode: bool = False, 

241) -> ModelT: 

242 kwargs: dict[str, Any] = {} 

243 type_hints = get_type_hints(model_cls) 

244 

245 for model_field in fields(cast("Any", model_cls)): 

246 alias = model_field.metadata.get("alias", model_field.name) 

247 annotation = type_hints.get(model_field.name, model_field.type) 

248 

249 if alias in data: 

250 kwargs[model_field.name] = _deserialize_value( 

251 data[alias], 

252 annotation, 

253 model_field.metadata, 

254 json_mode=json_mode, 

255 ) 

256 continue 

257 

258 if ( 

259 model_field.default is not MISSING 

260 or model_field.default_factory is not MISSING 

261 ): 

262 continue 

263 

264 raise KeyError(alias) 

265 

266 return model_cls(**kwargs) 

267 

268 

269def _model_to_mapping( 

270 model: Any, *, json_mode: bool = False 

271) -> MutableMapping[str, Any]: 

272 result: MutableMapping[str, Any] = {} 

273 

274 for model_field in fields(model): 

275 alias = model_field.metadata.get("alias", model_field.name) 

276 value = getattr(model, model_field.name) 

277 

278 if value is None and model_field.metadata.get("omit_if_none", True): 

279 continue 

280 if not value and model_field.metadata.get("omit_if_falsey", False): 280 ↛ 281line 280 didn't jump to line 281 because the condition on line 280 was never true

281 continue 

282 

283 result[alias] = _serialize_value( 

284 value, 

285 model_field.metadata, 

286 json_mode=json_mode, 

287 ) 

288 

289 return result 

290 

291 

292class OperationAction(Enum): 

293 """State transition requested when checkpointing an operation.""" 

294 

295 START = "START" 

296 SUCCEED = "SUCCEED" 

297 FAIL = "FAIL" 

298 RETRY = "RETRY" 

299 CANCEL = "CANCEL" 

300 

301 

302class OperationStatus(Enum): 

303 """Persisted lifecycle status of an operation in execution history.""" 

304 

305 STARTED = "STARTED" 

306 PENDING = "PENDING" 

307 READY = "READY" 

308 SUCCEEDED = "SUCCEEDED" 

309 FAILED = "FAILED" 

310 CANCELLED = "CANCELLED" 

311 TIMED_OUT = "TIMED_OUT" 

312 STOPPED = "STOPPED" 

313 

314 

315class CallbackTimeoutType(Enum): 

316 """Timeout categories surfaced for callback failures.""" 

317 

318 TIMEOUT = "Callback.Timeout" 

319 HEARTBEAT = "Callback.Heartbeat" 

320 

321 

322class OperationSubType(Enum): 

323 """Fine-grained operation kind used in execution history.""" 

324 

325 STEP = "Step" 

326 WAIT = "Wait" 

327 CALLBACK = "Callback" 

328 RUN_IN_CHILD_CONTEXT = "RunInChildContext" 

329 MAP = "Map" 

330 MAP_ITERATION = "MapIteration" 

331 PARALLEL = "Parallel" 

332 PARALLEL_BRANCH = "ParallelBranch" 

333 WAIT_FOR_CALLBACK = "WaitForCallback" 

334 WAIT_FOR_CONDITION = "WaitForCondition" 

335 CHAINED_INVOKE = "ChainedInvoke" 

336 EXECUTION = "Execution" 

337 

338 

339OperationSubTypeValue: TypeAlias = OperationSubType | str 

340 

341 

342class OperationType(Enum): 

343 """Top-level operation categories persisted by the durable backend.""" 

344 

345 EXECUTION = "EXECUTION" 

346 CONTEXT = "CONTEXT" 

347 STEP = "STEP" 

348 WAIT = "WAIT" 

349 CALLBACK = "CALLBACK" 

350 CHAINED_INVOKE = "CHAINED_INVOKE" 

351 

352 

353@dataclass(frozen=True) 

354class OperationIdentifier: 

355 """Container for operation id, parent id, and name.""" 

356 

357 operation_id: str | None 

358 sub_type: OperationSubTypeValue 

359 parent_id: str | None = None 

360 name: str | None = None 

361 operation_type: OperationType | None = field(default=None, kw_only=True) 

362 

363 def require_operation_id(self) -> str: 

364 """Return the operation id for non-root operations.""" 

365 if self.operation_id is None: 

366 msg = "operation_id is required for non-execution operations" 

367 raise ValueError(msg) 

368 return self.operation_id 

369 

370 @classmethod 

371 def create_execution_op(cls) -> OperationIdentifier: 

372 return cls(None, OperationSubType.EXECUTION, None, None) 

373 

374 

375class InvocationStatus(Enum): 

376 """Overall result of a single durable Lambda invocation.""" 

377 

378 SUCCEEDED = "SUCCEEDED" 

379 FAILED = "FAILED" 

380 PENDING = "PENDING" 

381 

382 # Used internally only: the invocation failed and the backend will retry 

383 RETRY = "RETRY" 

384 

385 

386@dataclass(frozen=True) 

387class ErrorObject(AwsApiModel): 

388 """Serializable representation of an exception captured by the SDK.""" 

389 

390 message: str | None = field( 

391 default=None, metadata=_model_field_metadata(alias="ErrorMessage") 

392 ) 

393 type: str | None = field( 

394 default=None, metadata=_model_field_metadata(alias="ErrorType") 

395 ) 

396 data: str | None = field( 

397 default=None, metadata=_model_field_metadata(alias="ErrorData") 

398 ) 

399 stack_trace: list[str] | None = field( 

400 default=None, metadata=_model_field_metadata(alias="StackTrace") 

401 ) 

402 

403 @classmethod 

404 def from_exception(cls, exception: Exception) -> ErrorObject: 

405 return cls( 

406 message=str(exception), 

407 type=type(exception).__name__, 

408 data=None, 

409 stack_trace=None, 

410 ) 

411 

412 @classmethod 

413 def from_message(cls, message: str) -> ErrorObject: 

414 return cls( 

415 message=message, 

416 type=None, 

417 data=None, 

418 stack_trace=None, 

419 ) 

420 

421 

422@dataclass(frozen=True) 

423class DurableExecutionInvocationOutput(AwsApiModel): 

424 """Representation the DurableExecutionInvocationOutput. This is what the Durable lambda handler returns. 

425 

426 If the execution has been already completed via an update to the EXECUTION operation via CheckpointDurableExecution, 

427 payload must be empty for SUCCEEDED/FAILED status. 

428 """ 

429 

430 status: InvocationStatus = field(metadata=_model_field_metadata(alias="Status")) 

431 result: str | None = field( 

432 default=None, metadata=_model_field_metadata(alias="Result") 

433 ) 

434 error: ErrorObject | None = field( 

435 default=None, metadata=_model_field_metadata(alias="Error") 

436 ) 

437 

438 @classmethod 

439 def create_succeeded(cls, result: str) -> DurableExecutionInvocationOutput: 

440 return cls(status=InvocationStatus.SUCCEEDED, result=result) 

441 

442 

443@dataclass(frozen=True) 

444class ExecutionDetails(AwsApiModel): 

445 """Extra fields stored on the root execution operation.""" 

446 

447 input_payload: str | None = field( 

448 default=None, 

449 metadata=_model_field_metadata(alias="InputPayload", omit_if_none=False), 

450 ) 

451 

452 

453@dataclass(frozen=True) 

454class ContextDetails(AwsApiModel): 

455 """Checkpoint payload stored for child-context style operations.""" 

456 

457 replay_children: ReplayChildren = field( 

458 default=False, metadata=_model_field_metadata(alias="ReplayChildren") 

459 ) 

460 result: OperationPayload | None = field( 

461 default=None, metadata=_model_field_metadata(alias="Result") 

462 ) 

463 error: ErrorObject | None = field( 

464 default=None, metadata=_model_field_metadata(alias="Error") 

465 ) 

466 

467 

468@dataclass(frozen=True) 

469class StepDetails(AwsApiModel): 

470 """Checkpoint payload stored for durable steps and polling checks.""" 

471 

472 attempt: int = field(default=0, metadata=_model_field_metadata(alias="Attempt")) 

473 next_attempt_timestamp: datetime.datetime | None = field( 

474 default=None, 

475 metadata=_model_field_metadata(alias="NextAttemptTimestamp", is_timestamp=True), 

476 ) 

477 result: OperationPayload | None = field( 

478 default=None, 

479 metadata=_model_field_metadata(alias="Result"), 

480 ) 

481 error: ErrorObject | None = field( 

482 default=None, 

483 metadata=_model_field_metadata(alias="Error"), 

484 ) 

485 

486 

487@dataclass(frozen=True) 

488class WaitDetails(AwsApiModel): 

489 """Checkpoint payload stored for durable waits.""" 

490 

491 scheduled_end_timestamp: datetime.datetime | None = field( 

492 default=None, 

493 metadata=_model_field_metadata( 

494 alias="ScheduledEndTimestamp", is_timestamp=True 

495 ), 

496 ) 

497 

498 

499@dataclass(frozen=True) 

500class CallbackDetails(AwsApiModel): 

501 """Checkpoint payload stored for callbacks and callback results.""" 

502 

503 callback_id: str = field(metadata=_model_field_metadata(alias="CallbackId")) 

504 result: str | None = field( 

505 default=None, metadata=_model_field_metadata(alias="Result") 

506 ) 

507 error: ErrorObject | None = field( 

508 default=None, metadata=_model_field_metadata(alias="Error") 

509 ) 

510 

511 

512@dataclass(frozen=True) 

513class ChainedInvokeDetails(AwsApiModel): 

514 """Checkpoint payload stored for durable invokes.""" 

515 

516 result: str | None = field( 

517 default=None, metadata=_model_field_metadata(alias="Result") 

518 ) 

519 error: ErrorObject | None = field( 

520 default=None, metadata=_model_field_metadata(alias="Error") 

521 ) 

522 

523 

524@dataclass(frozen=True) 

525class StepOptions(AwsApiModel): 

526 """Additional options recorded on step retries.""" 

527 

528 next_attempt_delay_seconds: int = field( 

529 default=0, 

530 metadata=_model_field_metadata(alias="NextAttemptDelaySeconds"), 

531 ) 

532 

533 

534@dataclass(frozen=True) 

535class WaitOptions(AwsApiModel): 

536 """ 

537 Wait Options provides details regarding suspension. 

538 

539 As of 2025/10/27: 

540 

541 - `wait_seconds` accepts values between 1, and 31622400 

542 - When wait_second seconds does not exist,then we default to 1 

543 

544 """ 

545 

546 wait_seconds: int = field( 

547 default=1, metadata=_model_field_metadata(alias="WaitSeconds") 

548 ) 

549 

550 

551@dataclass(frozen=True) 

552class CallbackOptions(AwsApiModel): 

553 """ 

554 Callback options provides details about the callback, wrt timeout 

555 and heartbeat checks. 

556 

557 As of 2025/10/27: 

558 - When timeout_seconds == 0, then the callback has no timeout 

559 - When heartbeat_timeout_seconds == 0, then the callback has no timeout 

560 

561 - When timeout_seconds is not present, then default is 0 

562 - When heartbeat_timeout_seconds, then default is 0 

563 

564 """ 

565 

566 timeout_seconds: TimeoutSeconds = field( 

567 default=0, metadata=_model_field_metadata(alias="TimeoutSeconds") 

568 ) 

569 heartbeat_timeout_seconds: int = field( 

570 default=0, 

571 metadata=_model_field_metadata(alias="HeartbeatTimeoutSeconds"), 

572 ) 

573 

574 

575@dataclass(frozen=True) 

576class ChainedInvokeOptions(AwsApiModel): 

577 """ 

578 As of 2025/10/27: 

579 - Chained invoke options only contains a function name 

580 """ 

581 

582 function_name: str = field(metadata=_model_field_metadata(alias="FunctionName")) 

583 tenant_id: str | None = field( 

584 default=None, metadata=_model_field_metadata(alias="TenantId") 

585 ) 

586 

587 

588@dataclass(frozen=True) 

589class ContextOptions(AwsApiModel): 

590 """Extra flags recorded for child-context operations.""" 

591 

592 replay_children: ReplayChildren = field( 

593 default=False, 

594 metadata=_model_field_metadata(alias="ReplayChildren"), 

595 ) 

596 

597 

598@dataclass(frozen=True) 

599class OperationUpdate(AwsApiModel): 

600 """Update an Operation. Use this to create a checkpoint. 

601 

602 See the various create_ factory class methods to instantiate me. 

603 """ 

604 

605 operation_id: str = field(metadata=_model_field_metadata(alias="Id")) 

606 operation_type: OperationType = field(metadata=_model_field_metadata(alias="Type")) 

607 action: OperationAction = field(metadata=_model_field_metadata(alias="Action")) 

608 parent_id: str | None = field( 

609 default=None, 

610 metadata=_model_field_metadata(alias="ParentId", omit_if_falsey=True), 

611 ) 

612 name: str | None = field( 

613 default=None, 

614 metadata=_model_field_metadata(alias="Name", omit_if_falsey=True), 

615 ) 

616 sub_type: OperationSubTypeValue | None = field( 

617 default=None, 

618 metadata=_model_field_metadata(alias="SubType"), 

619 ) 

620 payload: str | None = field( 

621 default=None, 

622 metadata=_model_field_metadata(alias="Payload"), 

623 ) 

624 error: ErrorObject | None = field( 

625 default=None, metadata=_model_field_metadata(alias="Error") 

626 ) 

627 context_options: ContextOptions | None = field( 

628 default=None, 

629 metadata=_model_field_metadata(alias="ContextOptions"), 

630 ) 

631 step_options: StepOptions | None = field( 

632 default=None, 

633 metadata=_model_field_metadata(alias="StepOptions"), 

634 ) 

635 wait_options: WaitOptions | None = field( 

636 default=None, 

637 metadata=_model_field_metadata(alias="WaitOptions"), 

638 ) 

639 callback_options: CallbackOptions | None = field( 

640 default=None, 

641 metadata=_model_field_metadata(alias="CallbackOptions"), 

642 ) 

643 chained_invoke_options: ChainedInvokeOptions | None = field( 

644 default=None, 

645 metadata=_model_field_metadata(alias="ChainedInvokeOptions"), 

646 ) 

647 

648 @classmethod 

649 def create_callback( 

650 cls, identifier: OperationIdentifier, callback_options: CallbackOptions 

651 ) -> OperationUpdate: 

652 """Create an instance of OperationUpdate for type:CALLBACK, action:START""" 

653 return cls( 

654 operation_id=identifier.require_operation_id(), 

655 parent_id=identifier.parent_id, 

656 operation_type=OperationType.CALLBACK, 

657 sub_type=identifier.sub_type, 

658 action=OperationAction.START, 

659 name=identifier.name, 

660 callback_options=callback_options, 

661 ) 

662 

663 @classmethod 

664 def create_context_start( 

665 cls, identifier: OperationIdentifier, sub_type: OperationSubTypeValue 

666 ) -> OperationUpdate: 

667 """Create an instance of OperationUpdate for type: CONTEXT, action: START.""" 

668 return cls( 

669 operation_id=identifier.require_operation_id(), 

670 parent_id=identifier.parent_id, 

671 operation_type=OperationType.CONTEXT, 

672 sub_type=sub_type, 

673 action=OperationAction.START, 

674 name=identifier.name, 

675 ) 

676 

677 @classmethod 

678 def create_context_succeed( 

679 cls, 

680 identifier: OperationIdentifier, 

681 payload: str, 

682 sub_type: OperationSubTypeValue, 

683 context_options: ContextOptions | None = None, 

684 ) -> OperationUpdate: 

685 """Create an instance of OperationUpdate for type: CONTEXT, action: SUCCEED.""" 

686 return cls( 

687 operation_id=identifier.require_operation_id(), 

688 parent_id=identifier.parent_id, 

689 operation_type=OperationType.CONTEXT, 

690 sub_type=sub_type, 

691 action=OperationAction.SUCCEED, 

692 name=identifier.name, 

693 payload=payload, 

694 context_options=context_options, 

695 ) 

696 

697 @classmethod 

698 def create_context_fail( 

699 cls, 

700 identifier: OperationIdentifier, 

701 error: ErrorObject, 

702 sub_type: OperationSubTypeValue, 

703 ) -> OperationUpdate: 

704 """Create an instance of OperationUpdate for type: CONTEXT, action: FAIL.""" 

705 return cls( 

706 operation_id=identifier.require_operation_id(), 

707 parent_id=identifier.parent_id, 

708 operation_type=OperationType.CONTEXT, 

709 sub_type=sub_type, 

710 action=OperationAction.FAIL, 

711 name=identifier.name, 

712 error=error, 

713 ) 

714 

715 @classmethod 

716 def create_execution_succeed(cls, payload: str) -> OperationUpdate: 

717 """Create an instance of OperationUpdate for type: EXECUTION, action: SUCCEED.""" 

718 return cls( 

719 operation_id=f"execution-result-{int(datetime.datetime.now(tz=datetime.timezone.utc).timestamp() * 1000)}", 

720 operation_type=OperationType.EXECUTION, 

721 action=OperationAction.SUCCEED, 

722 payload=payload, 

723 ) 

724 

725 @classmethod 

726 def create_execution_fail(cls, error: ErrorObject) -> OperationUpdate: 

727 """Create an instance of OperationUpdate for type: EXECUTION, action: FAIL.""" 

728 return cls( 

729 operation_id=f"execution-result-{int(datetime.datetime.now(tz=datetime.timezone.utc).timestamp() * 1000)}", 

730 operation_type=OperationType.EXECUTION, 

731 action=OperationAction.FAIL, 

732 error=error, 

733 ) 

734 

735 @classmethod 

736 def create_step_succeed( 

737 cls, identifier: OperationIdentifier, payload: str 

738 ) -> OperationUpdate: 

739 """Create an instance of OperationUpdate for type: STEP, action: SUCCEED.""" 

740 return cls( 

741 operation_id=identifier.require_operation_id(), 

742 parent_id=identifier.parent_id, 

743 operation_type=OperationType.STEP, 

744 sub_type=identifier.sub_type, 

745 action=OperationAction.SUCCEED, 

746 name=identifier.name, 

747 payload=payload, 

748 ) 

749 

750 @classmethod 

751 def create_step_fail( 

752 cls, identifier: OperationIdentifier, error: ErrorObject 

753 ) -> OperationUpdate: 

754 """Create an instance of OperationUpdate for type: STEP, action: FAIL.""" 

755 return cls( 

756 operation_id=identifier.require_operation_id(), 

757 parent_id=identifier.parent_id, 

758 operation_type=OperationType.STEP, 

759 sub_type=identifier.sub_type, 

760 action=OperationAction.FAIL, 

761 name=identifier.name, 

762 error=error, 

763 ) 

764 

765 @classmethod 

766 def create_step_start(cls, identifier: OperationIdentifier) -> OperationUpdate: 

767 """Create an instance of OperationUpdate for type: STEP, action: START.""" 

768 return cls( 

769 operation_id=identifier.require_operation_id(), 

770 parent_id=identifier.parent_id, 

771 operation_type=OperationType.STEP, 

772 sub_type=identifier.sub_type, 

773 action=OperationAction.START, 

774 name=identifier.name, 

775 ) 

776 

777 @classmethod 

778 def create_step_retry( 

779 cls, 

780 identifier: OperationIdentifier, 

781 error: ErrorObject | None, 

782 next_attempt_delay_seconds: int, 

783 *, 

784 payload: str | None = None, 

785 ) -> OperationUpdate: 

786 """Create an instance of OperationUpdate for type: STEP, action: RETRY.""" 

787 return cls( 

788 operation_id=identifier.require_operation_id(), 

789 parent_id=identifier.parent_id, 

790 operation_type=OperationType.STEP, 

791 sub_type=identifier.sub_type, 

792 action=OperationAction.RETRY, 

793 name=identifier.name, 

794 payload=payload, 

795 error=error, 

796 step_options=StepOptions( 

797 next_attempt_delay_seconds=next_attempt_delay_seconds 

798 ), 

799 ) 

800 

801 @classmethod 

802 def create_invoke_start( 

803 cls, 

804 identifier: OperationIdentifier, 

805 payload: str, 

806 chained_invoke_options: ChainedInvokeOptions, 

807 ) -> OperationUpdate: 

808 """Create an instance of OperationUpdate for type: INVOKE, action: START.""" 

809 return cls( 

810 operation_id=identifier.require_operation_id(), 

811 parent_id=identifier.parent_id, 

812 operation_type=OperationType.CHAINED_INVOKE, 

813 sub_type=identifier.sub_type, 

814 action=OperationAction.START, 

815 name=identifier.name, 

816 payload=payload, 

817 chained_invoke_options=chained_invoke_options, 

818 ) 

819 

820 @classmethod 

821 def create_wait_start( 

822 cls, identifier: OperationIdentifier, wait_options: WaitOptions 

823 ) -> OperationUpdate: 

824 """Create an instance of OperationUpdate for type: WAIT, action: START.""" 

825 return cls( 

826 operation_id=identifier.require_operation_id(), 

827 parent_id=identifier.parent_id, 

828 operation_type=OperationType.WAIT, 

829 sub_type=identifier.sub_type, 

830 action=OperationAction.START, 

831 name=identifier.name, 

832 wait_options=wait_options, 

833 ) 

834 

835 

836class TimestampConverter: 

837 """Converter for datetime/Unix timestamp conversions.""" 

838 

839 @staticmethod 

840 def to_unix_millis(dt: datetime.datetime | None) -> int | None: 

841 """Convert datetime to Unix timestamp in milliseconds.""" 

842 return int(dt.timestamp() * 1000) if dt else None 

843 

844 @staticmethod 

845 def from_unix_millis(ms: int | None) -> datetime.datetime | None: 

846 """Convert Unix timestamp in milliseconds to datetime.""" 

847 return ( 

848 datetime.datetime.fromtimestamp(ms / 1000, tz=datetime.timezone.utc) 

849 if ms is not None 

850 else None 

851 ) 

852 

853 

854@dataclass(frozen=True) 

855class Operation(AwsApiModel): 

856 """Represent the Operation type for GetDurableExecutionState and CheckpointDurableExecution.""" 

857 

858 operation_id: str = field(metadata=_model_field_metadata(alias="Id")) 

859 operation_type: OperationType = field(metadata=_model_field_metadata(alias="Type")) 

860 status: OperationStatus = field(metadata=_model_field_metadata(alias="Status")) 

861 parent_id: str | None = field( 

862 default=None, 

863 metadata=_model_field_metadata(alias="ParentId", omit_if_falsey=True), 

864 ) 

865 name: str | None = field( 

866 default=None, 

867 metadata=_model_field_metadata(alias="Name", omit_if_falsey=True), 

868 ) 

869 start_timestamp: datetime.datetime | None = field( 

870 default=None, 

871 metadata=_model_field_metadata(alias="StartTimestamp", is_timestamp=True), 

872 ) 

873 end_timestamp: datetime.datetime | None = field( 

874 default=None, 

875 metadata=_model_field_metadata(alias="EndTimestamp", is_timestamp=True), 

876 ) 

877 sub_type: OperationSubTypeValue | None = field( 

878 default=None, 

879 metadata=_model_field_metadata(alias="SubType"), 

880 ) 

881 execution_details: ExecutionDetails | None = field( 

882 default=None, 

883 metadata=_model_field_metadata(alias="ExecutionDetails"), 

884 ) 

885 context_details: ContextDetails | None = field( 

886 default=None, 

887 metadata=_model_field_metadata(alias="ContextDetails"), 

888 ) 

889 step_details: StepDetails | None = field( 

890 default=None, 

891 metadata=_model_field_metadata(alias="StepDetails"), 

892 ) 

893 wait_details: WaitDetails | None = field( 

894 default=None, 

895 metadata=_model_field_metadata(alias="WaitDetails"), 

896 ) 

897 callback_details: CallbackDetails | None = field( 

898 default=None, 

899 metadata=_model_field_metadata(alias="CallbackDetails"), 

900 ) 

901 chained_invoke_details: ChainedInvokeDetails | None = field( 

902 default=None, 

903 metadata=_model_field_metadata(alias="ChainedInvokeDetails"), 

904 ) 

905 

906 

907@dataclass(frozen=True) 

908class CheckpointUpdatedExecutionState(AwsApiModel): 

909 """Representation of the CheckpointUpdatedExecutionState structure of the DEX API.""" 

910 

911 operations: list[Operation] = field( 

912 default_factory=list, 

913 metadata=_model_field_metadata(alias="Operations"), 

914 ) 

915 next_marker: str | None = field( 

916 default=None, metadata=_model_field_metadata(alias="NextMarker") 

917 ) 

918 

919 

920@dataclass(frozen=True) 

921class CheckpointOutput(AwsApiModel): 

922 """Representation of the CheckpointDurableExecutionOutput structure of the DEX CheckpointDurableExecution API.""" 

923 

924 checkpoint_token: str | None = field( 

925 default=None, 

926 metadata=_model_field_metadata(alias="CheckpointToken"), 

927 ) 

928 new_execution_state: CheckpointUpdatedExecutionState = field( 

929 default_factory=CheckpointUpdatedExecutionState, 

930 metadata=_model_field_metadata(alias="NewExecutionState", omit_if_none=False), 

931 ) 

932 

933 

934@dataclass(frozen=True) 

935class StateOutput(AwsApiModel): 

936 """Representation of the GetDurableExecutionStateOutput structure of the DEX GetDurableExecutionState API.""" 

937 

938 operations: list[Operation] = field( 

939 default_factory=list, 

940 metadata=_model_field_metadata(alias="Operations"), 

941 ) 

942 next_marker: str | None = field( 

943 default=None, metadata=_model_field_metadata(alias="NextMarker") 

944 ) 

945 

946 

947__all__ = [ 

948 "AwsApiModel", 

949 "CallbackDetails", 

950 "CallbackOptions", 

951 "CallbackTimeoutType", 

952 "ChainedInvokeDetails", 

953 "ChainedInvokeOptions", 

954 "CheckpointOutput", 

955 "CheckpointUpdatedExecutionState", 

956 "ContextDetails", 

957 "ContextOptions", 

958 "DurableExecutionInvocationOutput", 

959 "ErrorObject", 

960 "ExecutionDetails", 

961 "InvocationStatus", 

962 "Operation", 

963 "OperationAction", 

964 "OperationPayload", 

965 "OperationStatus", 

966 "OperationSubType", 

967 "OperationSubTypeValue", 

968 "OperationType", 

969 "OperationUpdate", 

970 "ReplayChildren", 

971 "StateOutput", 

972 "StepDetails", 

973 "StepOptions", 

974 "TimeoutSeconds", 

975 "TimestampConverter", 

976 "WaitDetails", 

977 "WaitOptions", 

978 "OperationIdentifier", 

979]