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
« prev ^ index » next coverage.py v7.16.0, created at 2026-08-30 23:43 +0000
1"""Core durable execution data models."""
3from __future__ import annotations
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)
20# Replace with `type` it when dropping support to Python 3.11
21ReplayChildren: TypeAlias = bool
22OperationPayload: TypeAlias = str
23TimeoutSeconds: TypeAlias = int
24ModelT = TypeVar("ModelT")
27class LambdaContext(Protocol):
28 """Minimal AWS Lambda context surface used by the SDK."""
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
41 def get_remaining_time_in_millis(self) -> int: ...
42 def log(self, msg) -> None: ...
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 }
65class MappingModel:
66 """Dataclass serialized through the shared Python mapping engine."""
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)
76 def to_dict(self) -> MutableMapping[str, Any]:
77 """Convert the model to its serialized mapping."""
78 return _model_to_mapping(self)
81class AwsApiModel:
82 """Dataclass serialized as an AWS API mapping with native Python values."""
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)
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)
94# Backward-compatible internal alias retained for extensions importing the old name.
95BotoSerializableModel = AwsApiModel
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)
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)
109_MAPPING_MODEL_TYPES = (
110 MappingModel,
111 AwsApiModel,
112 JsonSerializableModel,
113)
116def _enum_type(annotation: Any) -> type[Enum] | None:
117 if isinstance(annotation, type) and issubclass(annotation, Enum):
118 return annotation
120 origin = get_origin(annotation)
121 if origin is None:
122 return None
124 for arg in get_args(annotation):
125 enum_cls = _enum_type(arg)
126 if enum_cls is not None:
127 return enum_cls
129 return None
132def _accepts_plain_string(annotation: Any) -> bool:
133 if annotation is str:
134 return True
136 origin = get_origin(annotation)
137 if origin is None:
138 return False
140 return any(_accepts_plain_string(arg) for arg in get_args(annotation))
143def _model_type(annotation: Any) -> type[Any] | None:
144 if isinstance(annotation, type) and issubclass(annotation, _MAPPING_MODEL_TYPES):
145 return annotation
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
153 for arg in get_args(annotation):
154 model_cls = _model_type(arg)
155 if model_cls is not None:
156 return model_cls
158 return None
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
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)
175 if json_mode and metadata.get("is_timestamp", False):
176 return TimestampConverter.from_unix_millis(value)
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)
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
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 ]
205 return value
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
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)
221 if json_mode and metadata.get("is_timestamp", False):
222 return TimestampConverter.to_unix_millis(value)
224 if isinstance(value, Enum):
225 return value.value
227 if isinstance(value, _MAPPING_MODEL_TYPES):
228 return _model_to_mapping(value, json_mode=json_mode)
230 if isinstance(value, list):
231 return [_serialize_value(item, {}, json_mode=json_mode) for item in value]
233 return value
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)
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)
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
258 if (
259 model_field.default is not MISSING
260 or model_field.default_factory is not MISSING
261 ):
262 continue
264 raise KeyError(alias)
266 return model_cls(**kwargs)
269def _model_to_mapping(
270 model: Any, *, json_mode: bool = False
271) -> MutableMapping[str, Any]:
272 result: MutableMapping[str, Any] = {}
274 for model_field in fields(model):
275 alias = model_field.metadata.get("alias", model_field.name)
276 value = getattr(model, model_field.name)
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
283 result[alias] = _serialize_value(
284 value,
285 model_field.metadata,
286 json_mode=json_mode,
287 )
289 return result
292class OperationAction(Enum):
293 """State transition requested when checkpointing an operation."""
295 START = "START"
296 SUCCEED = "SUCCEED"
297 FAIL = "FAIL"
298 RETRY = "RETRY"
299 CANCEL = "CANCEL"
302class OperationStatus(Enum):
303 """Persisted lifecycle status of an operation in execution history."""
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"
315class CallbackTimeoutType(Enum):
316 """Timeout categories surfaced for callback failures."""
318 TIMEOUT = "Callback.Timeout"
319 HEARTBEAT = "Callback.Heartbeat"
322class OperationSubType(Enum):
323 """Fine-grained operation kind used in execution history."""
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"
339OperationSubTypeValue: TypeAlias = OperationSubType | str
342class OperationType(Enum):
343 """Top-level operation categories persisted by the durable backend."""
345 EXECUTION = "EXECUTION"
346 CONTEXT = "CONTEXT"
347 STEP = "STEP"
348 WAIT = "WAIT"
349 CALLBACK = "CALLBACK"
350 CHAINED_INVOKE = "CHAINED_INVOKE"
353@dataclass(frozen=True)
354class OperationIdentifier:
355 """Container for operation id, parent id, and name."""
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)
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
370 @classmethod
371 def create_execution_op(cls) -> OperationIdentifier:
372 return cls(None, OperationSubType.EXECUTION, None, None)
375class InvocationStatus(Enum):
376 """Overall result of a single durable Lambda invocation."""
378 SUCCEEDED = "SUCCEEDED"
379 FAILED = "FAILED"
380 PENDING = "PENDING"
382 # Used internally only: the invocation failed and the backend will retry
383 RETRY = "RETRY"
386@dataclass(frozen=True)
387class ErrorObject(AwsApiModel):
388 """Serializable representation of an exception captured by the SDK."""
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 )
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 )
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 )
422@dataclass(frozen=True)
423class DurableExecutionInvocationOutput(AwsApiModel):
424 """Representation the DurableExecutionInvocationOutput. This is what the Durable lambda handler returns.
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 """
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 )
438 @classmethod
439 def create_succeeded(cls, result: str) -> DurableExecutionInvocationOutput:
440 return cls(status=InvocationStatus.SUCCEEDED, result=result)
443@dataclass(frozen=True)
444class ExecutionDetails(AwsApiModel):
445 """Extra fields stored on the root execution operation."""
447 input_payload: str | None = field(
448 default=None,
449 metadata=_model_field_metadata(alias="InputPayload", omit_if_none=False),
450 )
453@dataclass(frozen=True)
454class ContextDetails(AwsApiModel):
455 """Checkpoint payload stored for child-context style operations."""
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 )
468@dataclass(frozen=True)
469class StepDetails(AwsApiModel):
470 """Checkpoint payload stored for durable steps and polling checks."""
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 )
487@dataclass(frozen=True)
488class WaitDetails(AwsApiModel):
489 """Checkpoint payload stored for durable waits."""
491 scheduled_end_timestamp: datetime.datetime | None = field(
492 default=None,
493 metadata=_model_field_metadata(
494 alias="ScheduledEndTimestamp", is_timestamp=True
495 ),
496 )
499@dataclass(frozen=True)
500class CallbackDetails(AwsApiModel):
501 """Checkpoint payload stored for callbacks and callback results."""
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 )
512@dataclass(frozen=True)
513class ChainedInvokeDetails(AwsApiModel):
514 """Checkpoint payload stored for durable invokes."""
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 )
524@dataclass(frozen=True)
525class StepOptions(AwsApiModel):
526 """Additional options recorded on step retries."""
528 next_attempt_delay_seconds: int = field(
529 default=0,
530 metadata=_model_field_metadata(alias="NextAttemptDelaySeconds"),
531 )
534@dataclass(frozen=True)
535class WaitOptions(AwsApiModel):
536 """
537 Wait Options provides details regarding suspension.
539 As of 2025/10/27:
541 - `wait_seconds` accepts values between 1, and 31622400
542 - When wait_second seconds does not exist,then we default to 1
544 """
546 wait_seconds: int = field(
547 default=1, metadata=_model_field_metadata(alias="WaitSeconds")
548 )
551@dataclass(frozen=True)
552class CallbackOptions(AwsApiModel):
553 """
554 Callback options provides details about the callback, wrt timeout
555 and heartbeat checks.
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
561 - When timeout_seconds is not present, then default is 0
562 - When heartbeat_timeout_seconds, then default is 0
564 """
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 )
575@dataclass(frozen=True)
576class ChainedInvokeOptions(AwsApiModel):
577 """
578 As of 2025/10/27:
579 - Chained invoke options only contains a function name
580 """
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 )
588@dataclass(frozen=True)
589class ContextOptions(AwsApiModel):
590 """Extra flags recorded for child-context operations."""
592 replay_children: ReplayChildren = field(
593 default=False,
594 metadata=_model_field_metadata(alias="ReplayChildren"),
595 )
598@dataclass(frozen=True)
599class OperationUpdate(AwsApiModel):
600 """Update an Operation. Use this to create a checkpoint.
602 See the various create_ factory class methods to instantiate me.
603 """
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
836class TimestampConverter:
837 """Converter for datetime/Unix timestamp conversions."""
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
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 )
854@dataclass(frozen=True)
855class Operation(AwsApiModel):
856 """Represent the Operation type for GetDurableExecutionState and CheckpointDurableExecution."""
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 )
907@dataclass(frozen=True)
908class CheckpointUpdatedExecutionState(AwsApiModel):
909 """Representation of the CheckpointUpdatedExecutionState structure of the DEX API."""
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 )
920@dataclass(frozen=True)
921class CheckpointOutput(AwsApiModel):
922 """Representation of the CheckpointDurableExecutionOutput structure of the DEX CheckpointDurableExecution API."""
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 )
934@dataclass(frozen=True)
935class StateOutput(AwsApiModel):
936 """Representation of the GetDurableExecutionStateOutput structure of the DEX GetDurableExecutionState API."""
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 )
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]