From b61dad7e03230f12444b125daddf5e6d4caf1747 Mon Sep 17 00:00:00 2001 From: Tim Conley Date: Thu, 16 Jul 2026 11:53:50 -0700 Subject: [PATCH 1/5] Wrap payload converters for temporal intermediate models --- temporalio/activity.py | 12 +++- temporalio/converter/__init__.py | 2 + temporalio/converter/_data_converter.py | 7 +- temporalio/converter/_payload_converter.py | 77 ++++++++++++++++++++++ temporalio/worker/_workflow_instance.py | 6 +- tests/test_converter.py | 54 +++++++++++++++ 6 files changed, 154 insertions(+), 4 deletions(-) diff --git a/temporalio/activity.py b/temporalio/activity.py index 4e632701e..1c80438e9 100644 --- a/temporalio/activity.py +++ b/temporalio/activity.py @@ -238,9 +238,17 @@ def payload_converter(self) -> temporalio.converter.PayloadConverter: self.payload_converter_class_or_instance, temporalio.converter.PayloadConverter, ): - self._payload_converter = self.payload_converter_class_or_instance + self._payload_converter = ( + temporalio.converter.TemporalIntermediatePayloadConverter.wrap( + self.payload_converter_class_or_instance + ) + ) else: - self._payload_converter = self.payload_converter_class_or_instance() + self._payload_converter = ( + temporalio.converter.TemporalIntermediatePayloadConverter.wrap( + self.payload_converter_class_or_instance() + ) + ) return self._payload_converter @property diff --git a/temporalio/converter/__init__.py b/temporalio/converter/__init__.py index 3821cbd68..7cdec2a76 100644 --- a/temporalio/converter/__init__.py +++ b/temporalio/converter/__init__.py @@ -33,6 +33,7 @@ JSONTypeConverter, JSONTypeConverterUnhandled, PayloadConverter, + TemporalIntermediatePayloadConverter, value_to_type, ) from temporalio.converter._payload_limits import ( @@ -83,6 +84,7 @@ "PayloadLimitsConfig", "PayloadSizeWarning", "SerializationContext", + "TemporalIntermediatePayloadConverter", "WithSerializationContext", "WorkflowSerializationContext", "decode_search_attributes", diff --git a/temporalio/converter/_data_converter.py b/temporalio/converter/_data_converter.py index 13b48e695..b0c87b926 100644 --- a/temporalio/converter/_data_converter.py +++ b/temporalio/converter/_data_converter.py @@ -29,6 +29,7 @@ ) from temporalio.converter._payload_converter import ( PayloadConverter, + TemporalIntermediatePayloadConverter, ) from temporalio.converter._payload_limits import ( PayloadLimitsConfig, @@ -103,7 +104,11 @@ class DataConverter(WithSerializationContext): """Server-reported limits for payloads.""" def __post_init__(self) -> None: # noqa: D105 - object.__setattr__(self, "payload_converter", self.payload_converter_class()) + object.__setattr__( + self, + "payload_converter", + TemporalIntermediatePayloadConverter.wrap(self.payload_converter_class()), + ) object.__setattr__(self, "failure_converter", self.failure_converter_class()) async def encode( diff --git a/temporalio/converter/_payload_converter.py b/temporalio/converter/_payload_converter.py index 8ee85ef72..261ad1033 100644 --- a/temporalio/converter/_payload_converter.py +++ b/temporalio/converter/_payload_converter.py @@ -514,6 +514,83 @@ def from_payload( raise RuntimeError("Failed parsing") from err +class TemporalIntermediatePayloadConverter(PayloadConverter, WithSerializationContext): + """Payload converter wrapper for generated Temporal intermediate hooks. + + Values with a ``_temporal_to_intermediate`` method are first converted to + their intermediate value, then encoded by the wrapped payload converter. When + decoding to a type with ``_temporal_from_intermediate``, the wrapped + converter first decodes the payload to the intermediate value and this + wrapper constructs the requested user-facing type from it. + """ + + _inner_payload_converter: PayloadConverter + + def __init__(self, inner_payload_converter: PayloadConverter) -> None: + """Create a Temporal intermediate payload converter.""" + self._inner_payload_converter = inner_payload_converter + + @staticmethod + def wrap(payload_converter: PayloadConverter) -> PayloadConverter: + """Wrap a payload converter unless it is already wrapped.""" + if isinstance(payload_converter, TemporalIntermediatePayloadConverter): + return payload_converter + return TemporalIntermediatePayloadConverter(payload_converter) + + def to_payloads( + self, values: Sequence[Any] + ) -> list[temporalio.api.common.v1.Payload]: + """See base class.""" + intermediate_values: list[Any] = [] + for value in values: + to_intermediate = getattr(value, "_temporal_to_intermediate", None) + if to_intermediate is not None: + value = to_intermediate(payload_converter=self._inner_payload_converter) + intermediate_values.append(value) + return self._inner_payload_converter.to_payloads(intermediate_values) + + def from_payloads( + self, + payloads: Sequence[temporalio.api.common.v1.Payload], + type_hints: list[type] | None = None, + ) -> list[Any]: + """See base class.""" + if type_hints is None: + return self._inner_payload_converter.from_payloads(payloads, None) + normalized_type_hints: list[type | None] = list(type_hints) + if len(normalized_type_hints) < len(payloads): + normalized_type_hints.extend([None] * (len(payloads) - len(type_hints))) + inner_type_hints = [ + None + if getattr(type_hint, "_temporal_from_intermediate", None) is not None + else type_hint + for type_hint in normalized_type_hints + ] + values = self._inner_payload_converter.from_payloads( + payloads, typing.cast("list[type]", inner_type_hints) + ) + return [ + from_intermediate(value, payload_converter=self._inner_payload_converter) + if ( + from_intermediate := getattr( + type_hint, "_temporal_from_intermediate", None + ) + ) + is not None + else value + for value, type_hint in zip(values, normalized_type_hints) + ] + + def with_context(self, context: SerializationContext) -> Self: + """Return a new instance with context set on the inner converter.""" + if not isinstance(self._inner_payload_converter, WithSerializationContext): + return self + inner_payload_converter = self._inner_payload_converter.with_context(context) + if inner_payload_converter is self._inner_payload_converter: + return self + return type(self)(inner_payload_converter) + + class AdvancedJSONEncoder(json.JSONEncoder): """Advanced JSON encoder. diff --git a/temporalio/worker/_workflow_instance.py b/temporalio/worker/_workflow_instance.py index 74edc66b7..4d329059b 100644 --- a/temporalio/worker/_workflow_instance.py +++ b/temporalio/worker/_workflow_instance.py @@ -246,7 +246,11 @@ def __init__(self, det: WorkflowInstanceDetails) -> None: self._defn = det.defn self._workflow_input: ExecuteWorkflowInput | None = None self._info = det.info - self._context_free_payload_converter = det.payload_converter_class() + self._context_free_payload_converter = ( + temporalio.converter.TemporalIntermediatePayloadConverter.wrap( + det.payload_converter_class() + ) + ) self._context_free_failure_converter = det.failure_converter_class() workflow_context = temporalio.converter.WorkflowSerializationContext( namespace=det.info.namespace, diff --git a/tests/test_converter.py b/tests/test_converter.py index 10365f9c1..2281513f2 100644 --- a/tests/test_converter.py +++ b/tests/test_converter.py @@ -44,6 +44,8 @@ JSONTypeConverter, JSONTypeConverterUnhandled, PayloadCodec, + PayloadConverter, + TemporalIntermediatePayloadConverter, decode_search_attributes, encode_search_attribute_values, value_to_type, @@ -254,6 +256,58 @@ def test_binary_proto(): assert decoded == proto +@dataclass +class TemporalIntermediateModel: + value: str + + @classmethod + def _temporal_from_intermediate( + cls, + intermediate: temporalio.api.common.v1.WorkflowExecution, + *, + payload_converter: PayloadConverter | None = None, + ) -> TemporalIntermediateModel: + assert payload_converter is not None + return cls(value=intermediate.workflow_id) + + def _temporal_to_intermediate( + self, *, payload_converter: PayloadConverter | None = None + ) -> temporalio.api.common.v1.WorkflowExecution: + assert payload_converter is not None + return temporalio.api.common.v1.WorkflowExecution( + workflow_id=self.value, + run_id="run-id", + ) + + +class CustomDefaultPayloadConverter(DefaultPayloadConverter): + pass + + +def test_temporal_intermediate_payload_converter_wraps_user_converter(): + data_converter = DataConverter( + payload_converter_class=CustomDefaultPayloadConverter + ) + converter = data_converter.payload_converter + assert isinstance(converter, TemporalIntermediatePayloadConverter) + value = TemporalIntermediateModel("workflow-id") + + payload = converter.to_payload(value) + + assert payload.metadata["encoding"] == b"json/protobuf" + assert ( + payload.metadata["messageType"] == b"temporal.api.common.v1.WorkflowExecution" + ) + assert all("temporal-wire" not in key for key in payload.metadata) + assert all(b"temporal-wire" not in value for value in payload.metadata.values()) + assert converter.from_payload(payload, TemporalIntermediateModel) == value + + plain_proto_payload = converter.to_payload( + temporalio.api.common.v1.WorkflowExecution(workflow_id="id1", run_id="id2") + ) + assert plain_proto_payload.metadata["encoding"] == b"json/protobuf" + + def test_encode_search_attribute_values(): with pytest.raises(TypeError, match="of type tuple not one of"): encode_search_attribute_values([("bad type",)]) # type: ignore[arg-type] From 585bc27d62e9dc5f299685289ef239c32bcb215a Mon Sep 17 00:00:00 2001 From: Tim Conley Date: Thu, 16 Jul 2026 13:04:40 -0700 Subject: [PATCH 2/5] Support intermediate hooks in system Nexus conversion --- temporalio/converter/_payload_converter.py | 4 +- temporalio/nexus/system/__init__.py | 74 ++++++++++++++++++++-- temporalio/worker/_workflow_instance.py | 4 +- tests/nexus/test_temporal_system_nexus.py | 12 +++- tests/test_converter.py | 8 +-- tests/worker/test_visitor.py | 4 +- 6 files changed, 87 insertions(+), 19 deletions(-) diff --git a/temporalio/converter/_payload_converter.py b/temporalio/converter/_payload_converter.py index 261ad1033..8cdcc7abe 100644 --- a/temporalio/converter/_payload_converter.py +++ b/temporalio/converter/_payload_converter.py @@ -545,7 +545,7 @@ def to_payloads( for value in values: to_intermediate = getattr(value, "_temporal_to_intermediate", None) if to_intermediate is not None: - value = to_intermediate(payload_converter=self._inner_payload_converter) + value = to_intermediate() intermediate_values.append(value) return self._inner_payload_converter.to_payloads(intermediate_values) @@ -570,7 +570,7 @@ def from_payloads( payloads, typing.cast("list[type]", inner_type_hints) ) return [ - from_intermediate(value, payload_converter=self._inner_payload_converter) + from_intermediate(value) if ( from_intermediate := getattr( type_hint, "_temporal_from_intermediate", None diff --git a/temporalio/nexus/system/__init__.py b/temporalio/nexus/system/__init__.py index 21c5a1408..e76756cd7 100644 --- a/temporalio/nexus/system/__init__.py +++ b/temporalio/nexus/system/__init__.py @@ -2,22 +2,82 @@ from __future__ import annotations +import contextlib +import contextvars +from collections.abc import Iterator, Sequence +from typing import Any + import temporalio.api.common.v1 import temporalio.converter from temporalio.bridge._visitor_functions import VisitorFunctions from temporalio.converter import BinaryProtoPayloadConverter, CompositePayloadConverter TEMPORAL_SYSTEM_ENDPOINT = "__temporal_system" +_user_payload_converter: contextvars.ContextVar[ + temporalio.converter.PayloadConverter | None +] = contextvars.ContextVar("temporal-system-nexus-user-payload-converter", default=None) -class SystemNexusPayloadConverter(CompositePayloadConverter): - """Payload converter for system Nexus outer envelopes.""" +@contextlib.contextmanager +def user_payload_converter_context( + payload_converter: temporalio.converter.PayloadConverter, +) -> Iterator[None]: + """Set the user payload converter for system Nexus model conversion.""" + token = _user_payload_converter.set(payload_converter) + try: + yield + finally: + _user_payload_converter.reset(token) + + +def current_user_payload_converter() -> temporalio.converter.PayloadConverter: + """Return the active user payload converter for system Nexus model conversion.""" + payload_converter = _user_payload_converter.get() + if payload_converter is None: + raise RuntimeError("System Nexus user payload converter context is not active") + return payload_converter + + +class _SystemNexusOuterPayloadConverter(CompositePayloadConverter): + """Payload converter for system Nexus outer proto envelopes.""" def __init__(self) -> None: """Create a payload converter for system Nexus outer envelopes.""" super().__init__(BinaryProtoPayloadConverter()) +class SystemNexusPayloadConverter(temporalio.converter.PayloadConverter): + """Payload converter for system Nexus outer envelopes.""" + + _user_payload_converter: temporalio.converter.PayloadConverter + _outer_payload_converter: temporalio.converter.PayloadConverter + + def __init__(self, user_payload_converter: temporalio.converter.PayloadConverter) -> None: + """Create a payload converter for system Nexus outer envelopes.""" + self._user_payload_converter = user_payload_converter + self._outer_payload_converter = ( + temporalio.converter.TemporalIntermediatePayloadConverter.wrap( + _SystemNexusOuterPayloadConverter() + ) + ) + + def to_payloads( + self, values: Sequence[Any] + ) -> list[temporalio.api.common.v1.Payload]: + """See base class.""" + with user_payload_converter_context(self._user_payload_converter): + return self._outer_payload_converter.to_payloads(values) + + def from_payloads( + self, + payloads: Sequence[temporalio.api.common.v1.Payload], + type_hints: list[type] | None = None, + ) -> list[Any]: + """See base class.""" + with user_payload_converter_context(self._user_payload_converter): + return self._outer_payload_converter.from_payloads(payloads, type_hints) + + def is_system_endpoint(endpoint: str) -> bool: """Return whether a Nexus endpoint is the Temporal system endpoint.""" return endpoint == TEMPORAL_SYSTEM_ENDPOINT @@ -33,7 +93,7 @@ async def maybe_visit_payload( if not is_system_endpoint(endpoint): return None - payload_converter = get_payload_converter() + payload_converter = _SystemNexusOuterPayloadConverter() value = payload_converter.from_payload(payload) from ._payload_visitor import PayloadVisitor @@ -43,15 +103,19 @@ async def maybe_visit_payload( return payload_converter.to_payload(value) -def get_payload_converter() -> temporalio.converter.PayloadConverter: +def get_payload_converter( + user_payload_converter: temporalio.converter.PayloadConverter, +) -> temporalio.converter.PayloadConverter: """Return the fixed payload converter for system Nexus outer envelopes.""" - return SystemNexusPayloadConverter() + return SystemNexusPayloadConverter(user_payload_converter) __all__ = [ "TEMPORAL_SYSTEM_ENDPOINT", + "current_user_payload_converter", "get_payload_converter", "is_system_endpoint", "maybe_visit_payload", "SystemNexusPayloadConverter", + "user_payload_converter_context", ] diff --git a/temporalio/worker/_workflow_instance.py b/temporalio/worker/_workflow_instance.py index 4d329059b..dd7f4571a 100644 --- a/temporalio/worker/_workflow_instance.py +++ b/temporalio/worker/_workflow_instance.py @@ -2093,7 +2093,9 @@ async def operation_handle_fn() -> OutputT: t.uncancel() # type: ignore[union-attr] payload_converter = ( - temporalio.nexus.system.get_payload_converter() + temporalio.nexus.system.get_payload_converter( + self._workflow_context_payload_converter + ) if temporalio.nexus.system.is_system_endpoint(input.endpoint) else self._context_free_payload_converter ) diff --git a/tests/nexus/test_temporal_system_nexus.py b/tests/nexus/test_temporal_system_nexus.py index b689ee8d9..4629d230b 100644 --- a/tests/nexus/test_temporal_system_nexus.py +++ b/tests/nexus/test_temporal_system_nexus.py @@ -186,7 +186,9 @@ def _new_system_nexus_request_payload() -> temporalio.api.common.v1.Payload: assert nested_payload is not None request = workflowservice_pb2.SignalWithStartWorkflowExecutionRequest() request.input.payloads.add().CopyFrom(nested_payload) - payload = nexus_system.get_payload_converter().to_payload(request) + payload = nexus_system.get_payload_converter( + temporalio.converter.PayloadConverter.default + ).to_payload(request) assert payload is not None return payload @@ -201,7 +203,9 @@ async def test_schedule_system_nexus_endpoint_ignores_operation_registry() -> No await PayloadVisitor().visit(visitor, completion) schedule = completion.successful.commands[0].schedule_nexus_operation - decoded = nexus_system.get_payload_converter().from_payload(schedule.input) + decoded = nexus_system.get_payload_converter( + temporalio.converter.PayloadConverter.default + ).from_payload(schedule.input) assert isinstance( decoded, workflowservice_pb2.SignalWithStartWorkflowExecutionRequest ) @@ -339,7 +343,9 @@ def _field_is_repeated(field: FieldDescriptor) -> bool: ], ) def test_system_nexus_proto_roundtrip(message_type: type[Message]) -> None: - payload_converter = nexus_system.get_payload_converter() + payload_converter = nexus_system.get_payload_converter( + temporalio.converter.PayloadConverter.default + ) proto_value = _build_proto_sample(message_type) payload = payload_converter.to_payload(proto_value) assert payload is not None diff --git a/tests/test_converter.py b/tests/test_converter.py index 2281513f2..0ae2a2464 100644 --- a/tests/test_converter.py +++ b/tests/test_converter.py @@ -264,16 +264,10 @@ class TemporalIntermediateModel: def _temporal_from_intermediate( cls, intermediate: temporalio.api.common.v1.WorkflowExecution, - *, - payload_converter: PayloadConverter | None = None, ) -> TemporalIntermediateModel: - assert payload_converter is not None return cls(value=intermediate.workflow_id) - def _temporal_to_intermediate( - self, *, payload_converter: PayloadConverter | None = None - ) -> temporalio.api.common.v1.WorkflowExecution: - assert payload_converter is not None + def _temporal_to_intermediate(self) -> temporalio.api.common.v1.WorkflowExecution: return temporalio.api.common.v1.WorkflowExecution( workflow_id=self.value, run_id="run-id", diff --git a/tests/worker/test_visitor.py b/tests/worker/test_visitor.py index bd4004625..aa9a931e1 100644 --- a/tests/worker/test_visitor.py +++ b/tests/worker/test_visitor.py @@ -357,7 +357,9 @@ async def _visit(self) -> None: finally: active_visits -= 1 - payload_converter = nexus_system.get_payload_converter() + payload_converter = nexus_system.get_payload_converter( + temporalio.converter.PayloadConverter.default + ) system_request = workflowservice_pb2.SignalWithStartWorkflowExecutionRequest( input=Payloads(payloads=[Payload(data=b"workflow-input")]), signal_input=Payloads(payloads=[Payload(data=b"signal-input")]), From 844deb00b8238383f99e30a97c613ba7c0c5baee Mon Sep 17 00:00:00 2001 From: Tim Conley Date: Tue, 21 Jul 2026 08:30:59 -0700 Subject: [PATCH 3/5] Rename intermediate hooks to data model hooks --- CHANGELOG.md | 6 +++ temporalio/activity.py | 13 +++---- temporalio/converter/__init__.py | 2 - temporalio/converter/_data_converter.py | 12 +++--- temporalio/converter/_payload_converter.py | 38 +++++++++---------- temporalio/nexus/system/__init__.py | 11 +++--- temporalio/worker/_workflow.py | 2 +- temporalio/worker/_workflow_instance.py | 8 +--- temporalio/worker/workflow_sandbox/_runner.py | 2 +- tests/test_converter.py | 15 ++++---- 10 files changed, 52 insertions(+), 57 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1a1a92a95..013583ebd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,6 +20,12 @@ to include examples, links to docs, or any other relevant information. ### Added +- Added SDK payload converter support for values and type hints that expose + `_temporal_to_data_model` and `_temporal_from_data_model` hooks. This lets + hook-aware types delegate their wire representation to the configured payload + converter, preserving SDK behavior such as serialization contexts for nested + payload fields. + ### Changed ### Deprecated diff --git a/temporalio/activity.py b/temporalio/activity.py index 1c80438e9..a51a8c007 100644 --- a/temporalio/activity.py +++ b/temporalio/activity.py @@ -29,6 +29,7 @@ import temporalio.bridge.proto.activity_task import temporalio.common import temporalio.converter +from temporalio.converter._payload_converter import _TemporalDataModelPayloadConverter from .types import CallableType @@ -238,16 +239,12 @@ def payload_converter(self) -> temporalio.converter.PayloadConverter: self.payload_converter_class_or_instance, temporalio.converter.PayloadConverter, ): - self._payload_converter = ( - temporalio.converter.TemporalIntermediatePayloadConverter.wrap( - self.payload_converter_class_or_instance - ) + self._payload_converter = _TemporalDataModelPayloadConverter.wrap( + self.payload_converter_class_or_instance ) else: - self._payload_converter = ( - temporalio.converter.TemporalIntermediatePayloadConverter.wrap( - self.payload_converter_class_or_instance() - ) + self._payload_converter = _TemporalDataModelPayloadConverter.wrap( + self.payload_converter_class_or_instance() ) return self._payload_converter diff --git a/temporalio/converter/__init__.py b/temporalio/converter/__init__.py index 7cdec2a76..3821cbd68 100644 --- a/temporalio/converter/__init__.py +++ b/temporalio/converter/__init__.py @@ -33,7 +33,6 @@ JSONTypeConverter, JSONTypeConverterUnhandled, PayloadConverter, - TemporalIntermediatePayloadConverter, value_to_type, ) from temporalio.converter._payload_limits import ( @@ -84,7 +83,6 @@ "PayloadLimitsConfig", "PayloadSizeWarning", "SerializationContext", - "TemporalIntermediatePayloadConverter", "WithSerializationContext", "WorkflowSerializationContext", "decode_search_attributes", diff --git a/temporalio/converter/_data_converter.py b/temporalio/converter/_data_converter.py index b0c87b926..0832cd182 100644 --- a/temporalio/converter/_data_converter.py +++ b/temporalio/converter/_data_converter.py @@ -29,7 +29,7 @@ ) from temporalio.converter._payload_converter import ( PayloadConverter, - TemporalIntermediatePayloadConverter, + _TemporalDataModelPayloadConverter, ) from temporalio.converter._payload_limits import ( PayloadLimitsConfig, @@ -104,13 +104,13 @@ class DataConverter(WithSerializationContext): """Server-reported limits for payloads.""" def __post_init__(self) -> None: # noqa: D105 - object.__setattr__( - self, - "payload_converter", - TemporalIntermediatePayloadConverter.wrap(self.payload_converter_class()), - ) + object.__setattr__(self, "payload_converter", self._new_payload_converter()) object.__setattr__(self, "failure_converter", self.failure_converter_class()) + def _new_payload_converter(self) -> PayloadConverter: + """Create a payload converter instance with SDK data model hooks enabled.""" + return _TemporalDataModelPayloadConverter.wrap(self.payload_converter_class()) + async def encode( self, values: Sequence[Any] ) -> list[temporalio.api.common.v1.Payload]: diff --git a/temporalio/converter/_payload_converter.py b/temporalio/converter/_payload_converter.py index 8cdcc7abe..7d29d65b8 100644 --- a/temporalio/converter/_payload_converter.py +++ b/temporalio/converter/_payload_converter.py @@ -514,40 +514,40 @@ def from_payload( raise RuntimeError("Failed parsing") from err -class TemporalIntermediatePayloadConverter(PayloadConverter, WithSerializationContext): - """Payload converter wrapper for generated Temporal intermediate hooks. +class _TemporalDataModelPayloadConverter(PayloadConverter, WithSerializationContext): + """Payload converter wrapper for generated Temporal data model hooks. - Values with a ``_temporal_to_intermediate`` method are first converted to - their intermediate value, then encoded by the wrapped payload converter. When - decoding to a type with ``_temporal_from_intermediate``, the wrapped - converter first decodes the payload to the intermediate value and this + Values with a ``_temporal_to_data_model`` method are first converted to + their data model value, then encoded by the wrapped payload converter. When + decoding to a type with ``_temporal_from_data_model``, the wrapped + converter first decodes the payload to the data model value and this wrapper constructs the requested user-facing type from it. """ _inner_payload_converter: PayloadConverter def __init__(self, inner_payload_converter: PayloadConverter) -> None: - """Create a Temporal intermediate payload converter.""" + """Create a Temporal data model payload converter.""" self._inner_payload_converter = inner_payload_converter @staticmethod def wrap(payload_converter: PayloadConverter) -> PayloadConverter: """Wrap a payload converter unless it is already wrapped.""" - if isinstance(payload_converter, TemporalIntermediatePayloadConverter): + if isinstance(payload_converter, _TemporalDataModelPayloadConverter): return payload_converter - return TemporalIntermediatePayloadConverter(payload_converter) + return _TemporalDataModelPayloadConverter(payload_converter) def to_payloads( self, values: Sequence[Any] ) -> list[temporalio.api.common.v1.Payload]: """See base class.""" - intermediate_values: list[Any] = [] + data_model_values: list[Any] = [] for value in values: - to_intermediate = getattr(value, "_temporal_to_intermediate", None) - if to_intermediate is not None: - value = to_intermediate() - intermediate_values.append(value) - return self._inner_payload_converter.to_payloads(intermediate_values) + to_data_model = getattr(value, "_temporal_to_data_model", None) + if to_data_model is not None: + value = to_data_model() + data_model_values.append(value) + return self._inner_payload_converter.to_payloads(data_model_values) def from_payloads( self, @@ -562,7 +562,7 @@ def from_payloads( normalized_type_hints.extend([None] * (len(payloads) - len(type_hints))) inner_type_hints = [ None - if getattr(type_hint, "_temporal_from_intermediate", None) is not None + if getattr(type_hint, "_temporal_from_data_model", None) is not None else type_hint for type_hint in normalized_type_hints ] @@ -570,11 +570,9 @@ def from_payloads( payloads, typing.cast("list[type]", inner_type_hints) ) return [ - from_intermediate(value) + from_data_model(value) if ( - from_intermediate := getattr( - type_hint, "_temporal_from_intermediate", None - ) + from_data_model := getattr(type_hint, "_temporal_from_data_model", None) ) is not None else value diff --git a/temporalio/nexus/system/__init__.py b/temporalio/nexus/system/__init__.py index e76756cd7..b35a1ce1a 100644 --- a/temporalio/nexus/system/__init__.py +++ b/temporalio/nexus/system/__init__.py @@ -11,6 +11,7 @@ import temporalio.converter from temporalio.bridge._visitor_functions import VisitorFunctions from temporalio.converter import BinaryProtoPayloadConverter, CompositePayloadConverter +from temporalio.converter._payload_converter import _TemporalDataModelPayloadConverter TEMPORAL_SYSTEM_ENDPOINT = "__temporal_system" _user_payload_converter: contextvars.ContextVar[ @@ -52,13 +53,13 @@ class SystemNexusPayloadConverter(temporalio.converter.PayloadConverter): _user_payload_converter: temporalio.converter.PayloadConverter _outer_payload_converter: temporalio.converter.PayloadConverter - def __init__(self, user_payload_converter: temporalio.converter.PayloadConverter) -> None: + def __init__( + self, user_payload_converter: temporalio.converter.PayloadConverter + ) -> None: """Create a payload converter for system Nexus outer envelopes.""" self._user_payload_converter = user_payload_converter - self._outer_payload_converter = ( - temporalio.converter.TemporalIntermediatePayloadConverter.wrap( - _SystemNexusOuterPayloadConverter() - ) + self._outer_payload_converter = _TemporalDataModelPayloadConverter.wrap( + _SystemNexusOuterPayloadConverter() ) def to_payloads( diff --git a/temporalio/worker/_workflow.py b/temporalio/worker/_workflow.py index 8e6ba2726..1ed22ec02 100644 --- a/temporalio/worker/_workflow.py +++ b/temporalio/worker/_workflow.py @@ -789,7 +789,7 @@ def _create_workflow_instance( # Create instance from details det = WorkflowInstanceDetails( - payload_converter_class=self._data_converter.payload_converter_class, + payload_converter_factory=self._data_converter._new_payload_converter, failure_converter_class=self._data_converter.failure_converter_class, interceptor_classes=self._interceptor_classes, defn=defn, diff --git a/temporalio/worker/_workflow_instance.py b/temporalio/worker/_workflow_instance.py index dd7f4571a..0b4fa0dd0 100644 --- a/temporalio/worker/_workflow_instance.py +++ b/temporalio/worker/_workflow_instance.py @@ -135,7 +135,7 @@ def set_worker_level_failure_exception_types( class WorkflowInstanceDetails: """Immutable details for creating a workflow instance.""" - payload_converter_class: type[temporalio.converter.PayloadConverter] + payload_converter_factory: Callable[[], temporalio.converter.PayloadConverter] failure_converter_class: type[temporalio.converter.FailureConverter] interceptor_classes: Sequence[type[WorkflowInboundInterceptor]] defn: temporalio.workflow._Definition @@ -246,11 +246,7 @@ def __init__(self, det: WorkflowInstanceDetails) -> None: self._defn = det.defn self._workflow_input: ExecuteWorkflowInput | None = None self._info = det.info - self._context_free_payload_converter = ( - temporalio.converter.TemporalIntermediatePayloadConverter.wrap( - det.payload_converter_class() - ) - ) + self._context_free_payload_converter = det.payload_converter_factory() self._context_free_failure_converter = det.failure_converter_class() workflow_context = temporalio.converter.WorkflowSerializationContext( namespace=det.info.namespace, diff --git a/temporalio/worker/workflow_sandbox/_runner.py b/temporalio/worker/workflow_sandbox/_runner.py index b11c9b8c4..3c32c7236 100644 --- a/temporalio/worker/workflow_sandbox/_runner.py +++ b/temporalio/worker/workflow_sandbox/_runner.py @@ -79,7 +79,7 @@ def prepare_workflow(self, defn: temporalio.workflow._Definition) -> None: # Just create with fake info which validates self.create_instance( WorkflowInstanceDetails( - payload_converter_class=temporalio.converter.DataConverter.default.payload_converter_class, + payload_converter_factory=temporalio.converter.DataConverter.default._new_payload_converter, failure_converter_class=temporalio.converter.DataConverter.default.failure_converter_class, interceptor_classes=[], defn=defn, diff --git a/tests/test_converter.py b/tests/test_converter.py index 0ae2a2464..6e3a31685 100644 --- a/tests/test_converter.py +++ b/tests/test_converter.py @@ -44,12 +44,11 @@ JSONTypeConverter, JSONTypeConverterUnhandled, PayloadCodec, - PayloadConverter, - TemporalIntermediatePayloadConverter, decode_search_attributes, encode_search_attribute_values, value_to_type, ) +from temporalio.converter._payload_converter import _TemporalDataModelPayloadConverter from temporalio.exceptions import ( ApplicationError, FailureError, @@ -261,13 +260,13 @@ class TemporalIntermediateModel: value: str @classmethod - def _temporal_from_intermediate( + def _temporal_from_data_model( cls, - intermediate: temporalio.api.common.v1.WorkflowExecution, + data_model: temporalio.api.common.v1.WorkflowExecution, ) -> TemporalIntermediateModel: - return cls(value=intermediate.workflow_id) + return cls(value=data_model.workflow_id) - def _temporal_to_intermediate(self) -> temporalio.api.common.v1.WorkflowExecution: + def _temporal_to_data_model(self) -> temporalio.api.common.v1.WorkflowExecution: return temporalio.api.common.v1.WorkflowExecution( workflow_id=self.value, run_id="run-id", @@ -278,12 +277,12 @@ class CustomDefaultPayloadConverter(DefaultPayloadConverter): pass -def test_temporal_intermediate_payload_converter_wraps_user_converter(): +def test_temporal_data_model_payload_converter_wraps_user_converter(): data_converter = DataConverter( payload_converter_class=CustomDefaultPayloadConverter ) converter = data_converter.payload_converter - assert isinstance(converter, TemporalIntermediatePayloadConverter) + assert isinstance(converter, _TemporalDataModelPayloadConverter) value = TemporalIntermediateModel("workflow-id") payload = converter.to_payload(value) From 317a8a08a03f79192ba3d78b168f8326f32a2725 Mon Sep 17 00:00:00 2001 From: Tim Conley Date: Wed, 22 Jul 2026 11:09:44 -0700 Subject: [PATCH 4/5] Use typed data model converter decorators --- CHANGELOG.md | 2 +- temporalio/converter/_payload_converter.py | 28 +++++++++++++--------- tests/test_converter.py | 26 +++++++------------- 3 files changed, 27 insertions(+), 29 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b6257471f..42ae7325e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,7 +21,7 @@ to include examples, links to docs, or any other relevant information. ### Added - Added experimental SDK payload converter support for values and type hints - decorated with `@data_model_convertible(...)` using a `DataModelConverter`. + decorated with `@data_model_convertible(...)` using a `DataModelConverter` class. This lets data-model-aware types delegate their wire representation to the configured payload converter, preserving SDK behavior such as serialization contexts for nested payload fields. diff --git a/temporalio/converter/_payload_converter.py b/temporalio/converter/_payload_converter.py index 420176d1d..258f51b05 100644 --- a/temporalio/converter/_payload_converter.py +++ b/temporalio/converter/_payload_converter.py @@ -90,22 +90,28 @@ def from_data_model(self, value: DataModelT) -> ValueT: raise NotImplementedError -def data_model_convertible( - converter: DataModelConverter[ValueT, DataModelT], -) -> Callable[[type[ValueT]], type[ValueT]]: - """Decorate a class with a data model converter. - - .. warning:: - This API is experimental and subject to change. - """ +class _DataModelConvertibleDecorator(Generic[ValueT, DataModelT]): + def __init__( + self, converter_type: type[DataModelConverter[ValueT, DataModelT]] + ) -> None: + self._converter_type = converter_type - def decorator(cls: type[ValueT]) -> type[ValueT]: + def __call__(self, cls: type[ValueT]) -> type[ValueT]: if hasattr(cls, _DATA_MODEL_CONVERTER_ATTR): raise TypeError("class already has a data model converter") - setattr(cls, _DATA_MODEL_CONVERTER_ATTR, converter) + setattr(cls, _DATA_MODEL_CONVERTER_ATTR, self._converter_type()) return cls - return decorator + +def data_model_convertible( + converter_type: type[DataModelConverter[ValueT, DataModelT]], +) -> _DataModelConvertibleDecorator[ValueT, DataModelT]: + """Decorate a class with a data model converter class. + + .. warning:: + This API is experimental and subject to change. + """ + return _DataModelConvertibleDecorator(converter_type) def _get_data_model_converter( diff --git a/tests/test_converter.py b/tests/test_converter.py index 41a3ff9af..4a11050a4 100644 --- a/tests/test_converter.py +++ b/tests/test_converter.py @@ -257,14 +257,9 @@ def test_binary_proto(): assert decoded == proto -@dataclass -class TemporalDataModelValue: - value: str - - class TemporalDataModelValueConverter( DataModelConverter[ - TemporalDataModelValue, + "TemporalDataModelValue", temporalio.api.common.v1.WorkflowExecution, ] ): @@ -285,17 +280,15 @@ def from_data_model( return TemporalDataModelValue(value=value.workflow_id) +@data_model_convertible(TemporalDataModelValueConverter) @dataclass -class TemporalDataModelValueWithoutHint: +class TemporalDataModelValue: value: str -data_model_convertible(TemporalDataModelValueConverter())(TemporalDataModelValue) - - class TemporalDataModelValueWithoutHintConverter( DataModelConverter[ - TemporalDataModelValueWithoutHint, + "TemporalDataModelValueWithoutHint", temporalio.api.common.v1.WorkflowExecution, ] ): @@ -314,9 +307,10 @@ def from_data_model( return TemporalDataModelValueWithoutHint(value=value.workflow_id) -data_model_convertible(TemporalDataModelValueWithoutHintConverter())( - TemporalDataModelValueWithoutHint -) +@data_model_convertible(TemporalDataModelValueWithoutHintConverter) +@dataclass +class TemporalDataModelValueWithoutHint: + value: str class CustomDefaultPayloadConverter(DefaultPayloadConverter): @@ -362,9 +356,7 @@ def test_temporal_data_model_payload_converter_without_data_model_type(): def test_data_model_convertible_rejects_existing_converter(): with pytest.raises(TypeError, match="already has a data model converter"): - data_model_convertible(TemporalDataModelValueConverter())( - TemporalDataModelValue - ) + data_model_convertible(TemporalDataModelValueConverter)(TemporalDataModelValue) def test_encode_search_attribute_values(): From d41282a51f966ffa18f3387555237f2916db4132 Mon Sep 17 00:00:00 2001 From: Tim Conley Date: Wed, 22 Jul 2026 11:51:59 -0700 Subject: [PATCH 5/5] Rename data model hooks to transfer types --- CHANGELOG.md | 6 +- temporalio/activity.py | 8 ++- temporalio/converter/__init__.py | 8 +-- temporalio/converter/_data_converter.py | 8 ++- temporalio/converter/_payload_converter.py | 84 +++++++++++----------- temporalio/nexus/system/__init__.py | 6 +- tests/test_converter.py | 74 ++++++++++--------- 7 files changed, 104 insertions(+), 90 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 42ae7325e..6230f2754 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,10 +21,10 @@ to include examples, links to docs, or any other relevant information. ### Added - Added experimental SDK payload converter support for values and type hints - decorated with `@data_model_convertible(...)` using a `DataModelConverter` class. - This lets data-model-aware types delegate their wire representation to the + decorated with `@transfer_type_convertible(...)` using a `TransferTypeConverter` class. + This lets types with transfer type converters delegate their wire representation to the configured payload converter, preserving SDK behavior such as serialization - contexts for nested payload fields. + contexts. - Added `TLSConfig.verification_server_name` to verify the server certificate against a fixed name instead of the connection's server name. Unlike `domain`, it does not change the TLS SNI or HTTP/2 authority values, which keep following the connected host, so it can be used when the diff --git a/temporalio/activity.py b/temporalio/activity.py index a51a8c007..3f69bc17f 100644 --- a/temporalio/activity.py +++ b/temporalio/activity.py @@ -29,7 +29,9 @@ import temporalio.bridge.proto.activity_task import temporalio.common import temporalio.converter -from temporalio.converter._payload_converter import _TemporalDataModelPayloadConverter +from temporalio.converter._payload_converter import ( + _TemporalTransferTypePayloadConverter, +) from .types import CallableType @@ -239,11 +241,11 @@ def payload_converter(self) -> temporalio.converter.PayloadConverter: self.payload_converter_class_or_instance, temporalio.converter.PayloadConverter, ): - self._payload_converter = _TemporalDataModelPayloadConverter.wrap( + self._payload_converter = _TemporalTransferTypePayloadConverter.wrap( self.payload_converter_class_or_instance ) else: - self._payload_converter = _TemporalDataModelPayloadConverter.wrap( + self._payload_converter = _TemporalTransferTypePayloadConverter.wrap( self.payload_converter_class_or_instance() ) return self._payload_converter diff --git a/temporalio/converter/__init__.py b/temporalio/converter/__init__.py index 7c9d98089..ebd2b8396 100644 --- a/temporalio/converter/__init__.py +++ b/temporalio/converter/__init__.py @@ -26,7 +26,6 @@ BinaryPlainPayloadConverter, BinaryProtoPayloadConverter, CompositePayloadConverter, - DataModelConverter, DefaultPayloadConverter, EncodingPayloadConverter, JSONPlainPayloadConverter, @@ -34,7 +33,8 @@ JSONTypeConverter, JSONTypeConverterUnhandled, PayloadConverter, - data_model_convertible, + TransferTypeConverter, + transfer_type_convertible, value_to_type, ) from temporalio.converter._search_attributes import ( @@ -66,7 +66,7 @@ "BinaryPlainPayloadConverter", "BinaryProtoPayloadConverter", "CompositePayloadConverter", - "DataModelConverter", + "TransferTypeConverter", "DataConverter", "DefaultFailureConverter", "DefaultFailureConverterWithEncodedAttributes", @@ -82,7 +82,7 @@ "SerializationContext", "WithSerializationContext", "WorkflowSerializationContext", - "data_model_convertible", + "transfer_type_convertible", "decode_search_attributes", "decode_typed_search_attributes", "default", diff --git a/temporalio/converter/_data_converter.py b/temporalio/converter/_data_converter.py index 4f184c13c..6425d6b61 100644 --- a/temporalio/converter/_data_converter.py +++ b/temporalio/converter/_data_converter.py @@ -28,7 +28,7 @@ ) from temporalio.converter._payload_converter import ( PayloadConverter, - _TemporalDataModelPayloadConverter, + _TemporalTransferTypePayloadConverter, ) from temporalio.converter._serialization_context import ( SerializationContext, @@ -95,8 +95,10 @@ def __post_init__(self) -> None: # noqa: D105 object.__setattr__(self, "failure_converter", self.failure_converter_class()) def _new_payload_converter(self) -> PayloadConverter: - """Create a payload converter instance with SDK data model hooks enabled.""" - return _TemporalDataModelPayloadConverter.wrap(self.payload_converter_class()) + """Create a payload converter instance with SDK transfer type hooks enabled.""" + return _TemporalTransferTypePayloadConverter.wrap( + self.payload_converter_class() + ) async def encode( self, values: Sequence[Any] diff --git a/temporalio/converter/_payload_converter.py b/temporalio/converter/_payload_converter.py index 258f51b05..f10b6a4e0 100644 --- a/temporalio/converter/_payload_converter.py +++ b/temporalio/converter/_payload_converter.py @@ -53,27 +53,27 @@ _sym_db = google.protobuf.symbol_database.Default() ValueT = TypeVar("ValueT") -DataModelT = TypeVar("DataModelT") -_DATA_MODEL_CONVERTER_ATTR = "__temporal_data_model_converter" +TransferTypeT = TypeVar("TransferTypeT") +_TRANSFER_TYPE_CONVERTER_ATTR = "__temporal_transfer_type_converter" -class DataModelConverter(Generic[ValueT, DataModelT], ABC): - """Converter between a user-facing value and a data model value. +class TransferTypeConverter(Generic[ValueT, TransferTypeT], ABC): + """Converter between a user-facing value and a transfer type value. .. warning:: This API is experimental and subject to change. """ - data_model_type: type[DataModelT] | None = None - """Optional data model type hint to use when decoding payloads. + transfer_type: type[TransferTypeT] | None = None + """Optional type hint for the transfer type to use when decoding payloads. .. warning:: This API is experimental and subject to change. """ @abstractmethod - def to_data_model(self, value: ValueT) -> DataModelT: - """Convert a user-facing value to its data model value. + def to_transfer_type(self, value: ValueT) -> TransferTypeT: + """Convert a user-facing value to its transfer type value. .. warning:: This API is experimental and subject to change. @@ -81,8 +81,8 @@ def to_data_model(self, value: ValueT) -> DataModelT: raise NotImplementedError @abstractmethod - def from_data_model(self, value: DataModelT) -> ValueT: - """Convert a data model value to its user-facing value. + def from_transfer_type(self, value: TransferTypeT) -> ValueT: + """Convert a transfer type value to its user-facing value. .. warning:: This API is experimental and subject to change. @@ -90,35 +90,35 @@ def from_data_model(self, value: DataModelT) -> ValueT: raise NotImplementedError -class _DataModelConvertibleDecorator(Generic[ValueT, DataModelT]): +class _TransferTypeConvertibleDecorator(Generic[ValueT, TransferTypeT]): def __init__( - self, converter_type: type[DataModelConverter[ValueT, DataModelT]] + self, converter_type: type[TransferTypeConverter[ValueT, TransferTypeT]] ) -> None: self._converter_type = converter_type def __call__(self, cls: type[ValueT]) -> type[ValueT]: - if hasattr(cls, _DATA_MODEL_CONVERTER_ATTR): - raise TypeError("class already has a data model converter") - setattr(cls, _DATA_MODEL_CONVERTER_ATTR, self._converter_type()) + if hasattr(cls, _TRANSFER_TYPE_CONVERTER_ATTR): + raise TypeError("class already has a transfer type converter") + setattr(cls, _TRANSFER_TYPE_CONVERTER_ATTR, self._converter_type()) return cls -def data_model_convertible( - converter_type: type[DataModelConverter[ValueT, DataModelT]], -) -> _DataModelConvertibleDecorator[ValueT, DataModelT]: - """Decorate a class with a data model converter class. +def transfer_type_convertible( + converter_type: type[TransferTypeConverter[ValueT, TransferTypeT]], +) -> _TransferTypeConvertibleDecorator[ValueT, TransferTypeT]: + """Decorate a class with a transfer type converter class. .. warning:: This API is experimental and subject to change. """ - return _DataModelConvertibleDecorator(converter_type) + return _TransferTypeConvertibleDecorator(converter_type) -def _get_data_model_converter( +def _get_transfer_type_converter( value_type: object, -) -> DataModelConverter[Any, Any] | None: - converter = getattr(value_type, _DATA_MODEL_CONVERTER_ATTR, None) - if isinstance(converter, DataModelConverter): +) -> TransferTypeConverter[Any, Any] | None: + converter = getattr(value_type, _TRANSFER_TYPE_CONVERTER_ATTR, None) + if isinstance(converter, TransferTypeConverter): return converter return None @@ -584,40 +584,40 @@ def from_payload( raise RuntimeError("Failed parsing") from err -class _TemporalDataModelPayloadConverter(PayloadConverter, WithSerializationContext): - """Payload converter wrapper for registered Temporal data model converters. +class _TemporalTransferTypePayloadConverter(PayloadConverter, WithSerializationContext): + """Payload converter wrapper for registered Temporal transfer type converters. - Values with a registered data model converter are first converted to their - data model value, then encoded by the wrapped payload converter. When - decoding to a type with a registered data model converter, the wrapped - converter first decodes the payload to the data model value and this wrapper + Values with a registered transfer type converter are first converted to their + transfer type value, then encoded by the wrapped payload converter. When + decoding to a type with a registered transfer type converter, the wrapped + converter first decodes the payload to the transfer type value and this wrapper constructs the requested user-facing type from it. """ _inner_payload_converter: PayloadConverter def __init__(self, inner_payload_converter: PayloadConverter) -> None: - """Create a Temporal data model payload converter.""" + """Create a Temporal transfer type payload converter.""" self._inner_payload_converter = inner_payload_converter @staticmethod def wrap(payload_converter: PayloadConverter) -> PayloadConverter: """Wrap a payload converter unless it is already wrapped.""" - if isinstance(payload_converter, _TemporalDataModelPayloadConverter): + if isinstance(payload_converter, _TemporalTransferTypePayloadConverter): return payload_converter - return _TemporalDataModelPayloadConverter(payload_converter) + return _TemporalTransferTypePayloadConverter(payload_converter) def to_payloads( self, values: Sequence[Any] ) -> list[temporalio.api.common.v1.Payload]: """See base class.""" - data_model_values: list[Any] = [] + transfer_type_values: list[Any] = [] for value in values: - converter = _get_data_model_converter(type(value)) + converter = _get_transfer_type_converter(type(value)) if converter is not None: - value = converter.to_data_model(value) - data_model_values.append(value) - return self._inner_payload_converter.to_payloads(data_model_values) + value = converter.to_transfer_type(value) + transfer_type_values.append(value) + return self._inner_payload_converter.to_payloads(transfer_type_values) def from_payloads( self, @@ -627,16 +627,18 @@ def from_payloads( """See base class.""" if type_hints is None: return self._inner_payload_converter.from_payloads(payloads, None) - converters = [_get_data_model_converter(type_hint) for type_hint in type_hints] + converters = [ + _get_transfer_type_converter(type_hint) for type_hint in type_hints + ] inner_type_hints = [ - converter.data_model_type if converter is not None else type_hint + converter.transfer_type if converter is not None else type_hint for converter, type_hint in zip(converters, type_hints) ] values = self._inner_payload_converter.from_payloads( payloads, typing.cast("list[type]", inner_type_hints) ) return [ - converter.from_data_model(value) if converter is not None else value + converter.from_transfer_type(value) if converter is not None else value for value, converter in zip(values, converters) ] diff --git a/temporalio/nexus/system/__init__.py b/temporalio/nexus/system/__init__.py index d8e201e16..14a43cb72 100644 --- a/temporalio/nexus/system/__init__.py +++ b/temporalio/nexus/system/__init__.py @@ -15,7 +15,9 @@ import temporalio.converter from temporalio.bridge._visitor_functions import VisitorFunctions from temporalio.converter import BinaryProtoPayloadConverter, CompositePayloadConverter -from temporalio.converter._payload_converter import _TemporalDataModelPayloadConverter +from temporalio.converter._payload_converter import ( + _TemporalTransferTypePayloadConverter, +) TEMPORAL_SYSTEM_ENDPOINT = "__temporal_system" _user_payload_converter: contextvars.ContextVar[ @@ -62,7 +64,7 @@ def __init__( ) -> None: """Create a payload converter for system Nexus outer envelopes.""" self._user_payload_converter = user_payload_converter - self._outer_payload_converter = _TemporalDataModelPayloadConverter.wrap( + self._outer_payload_converter = _TemporalTransferTypePayloadConverter.wrap( _SystemNexusOuterPayloadConverter() ) diff --git a/tests/test_converter.py b/tests/test_converter.py index 4a11050a4..b5e10c518 100644 --- a/tests/test_converter.py +++ b/tests/test_converter.py @@ -38,19 +38,21 @@ BinaryProtoPayloadConverter, CompositePayloadConverter, DataConverter, - DataModelConverter, DefaultFailureConverterWithEncodedAttributes, DefaultPayloadConverter, JSONPlainPayloadConverter, JSONTypeConverter, JSONTypeConverterUnhandled, PayloadCodec, - data_model_convertible, + TransferTypeConverter, decode_search_attributes, encode_search_attribute_values, + transfer_type_convertible, value_to_type, ) -from temporalio.converter._payload_converter import _TemporalDataModelPayloadConverter +from temporalio.converter._payload_converter import ( + _TemporalTransferTypePayloadConverter, +) from temporalio.exceptions import ( ApplicationError, FailureError, @@ -257,59 +259,59 @@ def test_binary_proto(): assert decoded == proto -class TemporalDataModelValueConverter( - DataModelConverter[ - "TemporalDataModelValue", +class TemporalTransferTypeValueConverter( + TransferTypeConverter[ + "TemporalTransferTypeValue", temporalio.api.common.v1.WorkflowExecution, ] ): - data_model_type = temporalio.api.common.v1.WorkflowExecution + transfer_type = temporalio.api.common.v1.WorkflowExecution - def to_data_model( - self, value: TemporalDataModelValue + def to_transfer_type( + self, value: TemporalTransferTypeValue ) -> temporalio.api.common.v1.WorkflowExecution: return temporalio.api.common.v1.WorkflowExecution( workflow_id=value.value, run_id="run-id", ) - def from_data_model( + def from_transfer_type( self, value: temporalio.api.common.v1.WorkflowExecution, - ) -> TemporalDataModelValue: - return TemporalDataModelValue(value=value.workflow_id) + ) -> TemporalTransferTypeValue: + return TemporalTransferTypeValue(value=value.workflow_id) -@data_model_convertible(TemporalDataModelValueConverter) +@transfer_type_convertible(TemporalTransferTypeValueConverter) @dataclass -class TemporalDataModelValue: +class TemporalTransferTypeValue: value: str -class TemporalDataModelValueWithoutHintConverter( - DataModelConverter[ - "TemporalDataModelValueWithoutHint", +class TemporalTransferTypeValueWithoutHintConverter( + TransferTypeConverter[ + "TemporalTransferTypeValueWithoutHint", temporalio.api.common.v1.WorkflowExecution, ] ): - def to_data_model( - self, value: TemporalDataModelValueWithoutHint + def to_transfer_type( + self, value: TemporalTransferTypeValueWithoutHint ) -> temporalio.api.common.v1.WorkflowExecution: return temporalio.api.common.v1.WorkflowExecution( workflow_id=value.value, run_id="run-id", ) - def from_data_model( + def from_transfer_type( self, value: temporalio.api.common.v1.WorkflowExecution, - ) -> TemporalDataModelValueWithoutHint: - return TemporalDataModelValueWithoutHint(value=value.workflow_id) + ) -> TemporalTransferTypeValueWithoutHint: + return TemporalTransferTypeValueWithoutHint(value=value.workflow_id) -@data_model_convertible(TemporalDataModelValueWithoutHintConverter) +@transfer_type_convertible(TemporalTransferTypeValueWithoutHintConverter) @dataclass -class TemporalDataModelValueWithoutHint: +class TemporalTransferTypeValueWithoutHint: value: str @@ -317,13 +319,13 @@ class CustomDefaultPayloadConverter(DefaultPayloadConverter): pass -def test_temporal_data_model_payload_converter_wraps_user_converter(): +def test_temporal_transfer_type_payload_converter_wraps_user_converter(): data_converter = DataConverter( payload_converter_class=CustomDefaultPayloadConverter ) converter = data_converter.payload_converter - assert isinstance(converter, _TemporalDataModelPayloadConverter) - value = TemporalDataModelValue("workflow-id") + assert isinstance(converter, _TemporalTransferTypePayloadConverter) + value = TemporalTransferTypeValue("workflow-id") payload = converter.to_payload(value) @@ -333,7 +335,7 @@ def test_temporal_data_model_payload_converter_wraps_user_converter(): ) assert all("temporal-wire" not in key for key in payload.metadata) assert all(b"temporal-wire" not in value for value in payload.metadata.values()) - assert converter.from_payload(payload, TemporalDataModelValue) == value + assert converter.from_payload(payload, TemporalTransferTypeValue) == value plain_proto_payload = converter.to_payload( temporalio.api.common.v1.WorkflowExecution(workflow_id="id1", run_id="id2") @@ -341,9 +343,9 @@ def test_temporal_data_model_payload_converter_wraps_user_converter(): assert plain_proto_payload.metadata["encoding"] == b"json/protobuf" -def test_temporal_data_model_payload_converter_without_data_model_type(): +def test_temporal_transfer_type_payload_converter_without_transfer_type_hint(): converter = DataConverter.default.payload_converter - value = TemporalDataModelValueWithoutHint("workflow-id") + value = TemporalTransferTypeValueWithoutHint("workflow-id") payload = converter.to_payload(value) @@ -351,12 +353,16 @@ def test_temporal_data_model_payload_converter_without_data_model_type(): assert ( payload.metadata["messageType"] == b"temporal.api.common.v1.WorkflowExecution" ) - assert converter.from_payload(payload, TemporalDataModelValueWithoutHint) == value + assert ( + converter.from_payload(payload, TemporalTransferTypeValueWithoutHint) == value + ) -def test_data_model_convertible_rejects_existing_converter(): - with pytest.raises(TypeError, match="already has a data model converter"): - data_model_convertible(TemporalDataModelValueConverter)(TemporalDataModelValue) +def test_transfer_type_convertible_rejects_existing_converter(): + with pytest.raises(TypeError, match="already has a transfer type converter"): + transfer_type_convertible(TemporalTransferTypeValueConverter)( + TemporalTransferTypeValue + ) def test_encode_search_attribute_values():