|
21 | 21 | from typing import ( |
22 | 22 | Any, |
23 | 23 | ClassVar, |
| 24 | + Generic, |
24 | 25 | Literal, |
25 | 26 | NewType, |
26 | 27 | TypeVar, |
|
51 | 52 | ) |
52 | 53 |
|
53 | 54 | _sym_db = google.protobuf.symbol_database.Default() |
| 55 | +ValueT = TypeVar("ValueT") |
| 56 | +TransferTypeT = TypeVar("TransferTypeT") |
| 57 | +_TRANSFER_TYPE_CONVERTER_ATTR = "__temporal_transfer_type_converter" |
| 58 | + |
| 59 | + |
| 60 | +class TransferTypeConverter(Generic[ValueT, TransferTypeT], ABC): |
| 61 | + """Converter between a user-facing value and a transfer type value. |
| 62 | +
|
| 63 | + .. warning:: |
| 64 | + This API is experimental and subject to change. |
| 65 | + """ |
| 66 | + |
| 67 | + transfer_type: type[TransferTypeT] | None = None |
| 68 | + """Optional type hint for the transfer type to use when decoding payloads. |
| 69 | +
|
| 70 | + .. warning:: |
| 71 | + This API is experimental and subject to change. |
| 72 | + """ |
| 73 | + |
| 74 | + @abstractmethod |
| 75 | + def to_transfer_type(self, value: ValueT) -> TransferTypeT: |
| 76 | + """Convert a user-facing value to its transfer type value. |
| 77 | +
|
| 78 | + .. warning:: |
| 79 | + This API is experimental and subject to change. |
| 80 | + """ |
| 81 | + raise NotImplementedError |
| 82 | + |
| 83 | + @abstractmethod |
| 84 | + def from_transfer_type(self, value: TransferTypeT) -> ValueT: |
| 85 | + """Convert a transfer type value to its user-facing value. |
| 86 | +
|
| 87 | + .. warning:: |
| 88 | + This API is experimental and subject to change. |
| 89 | + """ |
| 90 | + raise NotImplementedError |
| 91 | + |
| 92 | + |
| 93 | +class _TransferTypeConvertibleDecorator(Generic[ValueT, TransferTypeT]): |
| 94 | + def __init__( |
| 95 | + self, converter_type: type[TransferTypeConverter[ValueT, TransferTypeT]] |
| 96 | + ) -> None: |
| 97 | + self._converter_type = converter_type |
| 98 | + |
| 99 | + def __call__(self, cls: type[ValueT]) -> type[ValueT]: |
| 100 | + if hasattr(cls, _TRANSFER_TYPE_CONVERTER_ATTR): |
| 101 | + raise TypeError("class already has a transfer type converter") |
| 102 | + setattr(cls, _TRANSFER_TYPE_CONVERTER_ATTR, self._converter_type()) |
| 103 | + return cls |
| 104 | + |
| 105 | + |
| 106 | +def transfer_type_convertible( |
| 107 | + converter_type: type[TransferTypeConverter[ValueT, TransferTypeT]], |
| 108 | +) -> _TransferTypeConvertibleDecorator[ValueT, TransferTypeT]: |
| 109 | + """Decorate a class with a transfer type converter class. |
| 110 | +
|
| 111 | + .. warning:: |
| 112 | + This API is experimental and subject to change. |
| 113 | + """ |
| 114 | + return _TransferTypeConvertibleDecorator(converter_type) |
| 115 | + |
| 116 | + |
| 117 | +def _get_transfer_type_converter( |
| 118 | + value_type: object, |
| 119 | +) -> TransferTypeConverter[Any, Any] | None: |
| 120 | + converter = getattr(value_type, _TRANSFER_TYPE_CONVERTER_ATTR, None) |
| 121 | + if isinstance(converter, TransferTypeConverter): |
| 122 | + return converter |
| 123 | + return None |
54 | 124 |
|
55 | 125 |
|
56 | 126 | class PayloadConverter(ABC): |
@@ -514,6 +584,74 @@ def from_payload( |
514 | 584 | raise RuntimeError("Failed parsing") from err |
515 | 585 |
|
516 | 586 |
|
| 587 | +class _TemporalTransferTypePayloadConverter(PayloadConverter, WithSerializationContext): |
| 588 | + """Payload converter wrapper for registered Temporal transfer type converters. |
| 589 | +
|
| 590 | + Values with a registered transfer type converter are first converted to their |
| 591 | + transfer type value, then encoded by the wrapped payload converter. When |
| 592 | + decoding to a type with a registered transfer type converter, the wrapped |
| 593 | + converter first decodes the payload to the transfer type value and this wrapper |
| 594 | + constructs the requested user-facing type from it. |
| 595 | + """ |
| 596 | + |
| 597 | + _inner_payload_converter: PayloadConverter |
| 598 | + |
| 599 | + def __init__(self, inner_payload_converter: PayloadConverter) -> None: |
| 600 | + """Create a Temporal transfer type payload converter.""" |
| 601 | + self._inner_payload_converter = inner_payload_converter |
| 602 | + |
| 603 | + @staticmethod |
| 604 | + def wrap(payload_converter: PayloadConverter) -> PayloadConverter: |
| 605 | + """Wrap a payload converter unless it is already wrapped.""" |
| 606 | + if isinstance(payload_converter, _TemporalTransferTypePayloadConverter): |
| 607 | + return payload_converter |
| 608 | + return _TemporalTransferTypePayloadConverter(payload_converter) |
| 609 | + |
| 610 | + def to_payloads( |
| 611 | + self, values: Sequence[Any] |
| 612 | + ) -> list[temporalio.api.common.v1.Payload]: |
| 613 | + """See base class.""" |
| 614 | + transfer_type_values: list[Any] = [] |
| 615 | + for value in values: |
| 616 | + converter = _get_transfer_type_converter(type(value)) |
| 617 | + if converter is not None: |
| 618 | + value = converter.to_transfer_type(value) |
| 619 | + transfer_type_values.append(value) |
| 620 | + return self._inner_payload_converter.to_payloads(transfer_type_values) |
| 621 | + |
| 622 | + def from_payloads( |
| 623 | + self, |
| 624 | + payloads: Sequence[temporalio.api.common.v1.Payload], |
| 625 | + type_hints: list[type] | None = None, |
| 626 | + ) -> list[Any]: |
| 627 | + """See base class.""" |
| 628 | + if type_hints is None: |
| 629 | + return self._inner_payload_converter.from_payloads(payloads, None) |
| 630 | + converters = [ |
| 631 | + _get_transfer_type_converter(type_hint) for type_hint in type_hints |
| 632 | + ] |
| 633 | + inner_type_hints = [ |
| 634 | + converter.transfer_type if converter is not None else type_hint |
| 635 | + for converter, type_hint in zip(converters, type_hints) |
| 636 | + ] |
| 637 | + values = self._inner_payload_converter.from_payloads( |
| 638 | + payloads, typing.cast("list[type]", inner_type_hints) |
| 639 | + ) |
| 640 | + return [ |
| 641 | + converter.from_transfer_type(value) if converter is not None else value |
| 642 | + for value, converter in zip(values, converters) |
| 643 | + ] |
| 644 | + |
| 645 | + def with_context(self, context: SerializationContext) -> Self: |
| 646 | + """Return a new instance with context set on the inner converter.""" |
| 647 | + if not isinstance(self._inner_payload_converter, WithSerializationContext): |
| 648 | + return self |
| 649 | + inner_payload_converter = self._inner_payload_converter.with_context(context) |
| 650 | + if inner_payload_converter is self._inner_payload_converter: |
| 651 | + return self |
| 652 | + return type(self)(inner_payload_converter) |
| 653 | + |
| 654 | + |
517 | 655 | class AdvancedJSONEncoder(json.JSONEncoder): |
518 | 656 | """Advanced JSON encoder. |
519 | 657 |
|
|
0 commit comments