From 6d9c87db6df9deb682f0c12b4b0978dd98394512 Mon Sep 17 00:00:00 2001 From: Alex Mazzeo Date: Mon, 5 Oct 2026 10:18:04 -0700 Subject: [PATCH] First draft of nexus handler callbacks --- .github/workflows/ci.yml | 2 +- CHANGELOG.md | 8 + pyproject.toml | 4 + scripts/gen_payload_visitor.py | 6 +- temporalio/bridge/_visitor.py | 15 + temporalio/client/_impl.py | 5 + temporalio/client/_interceptor.py | 2 + temporalio/client/_nexus.py | 122 ++- temporalio/nexus/__init__.py | 6 + temporalio/nexus/_notification.py | 364 +++++++ tests/__init__.py | 2 +- tests/nexus/test_nexus_handler_callbacks.py | 1045 +++++++++++++++++++ tests/nexus/test_nexus_type_errors.py | 477 +++++++++ uv.lock | 132 ++- 14 files changed, 2115 insertions(+), 75 deletions(-) create mode 100644 temporalio/nexus/_notification.py create mode 100644 tests/nexus/test_nexus_handler_callbacks.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 3105c764b..ee3d3613d 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -232,7 +232,7 @@ jobs: - run: uv add --dev --python 3.10 "googleapis-common-protos==1.70.0" - run: uv add --python 3.10 "protobuf<4" - run: uv sync --all-extras - - run: cargo install --locked nexgen --version 0.2.4 --features advanced --force + - run: cargo install --locked nexgen --version 0.2.7 --features advanced --force - run: poe build-develop - run: poe gen-protos - name: Check generation unchanged diff --git a/CHANGELOG.md b/CHANGELOG.md index b5f0a94ec..33d0b6beb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -24,6 +24,14 @@ to include examples, links to docs, or any other relevant information. instead of from an instance (i.e. its first parameter is an unbound `self`). - `workflow.new_random()` accepts an optional `name` that is mixed into the seed, so differently named generators, and `workflow.random()`, produce different sequences. +- **Experimental**: Standalone Nexus operations can report their outcome to a Nexus + service with completion callbacks. Pass `completion_callbacks` to + `NexusClient.start_operation` or `NexusClient.execute_operation`. Create the + callbacks with `temporalio.nexus.create_completion_callback`. When the + operation finishes, the server calls an on-complete operation on a worker for the + task queue of the callback. The call carries the result or the failure, and a source + context that you supply. The new `temporalio.nexus.notifications` module has the + input and output types of on-complete operations. ### Changed diff --git a/pyproject.toml b/pyproject.toml index 845e5f980..147f2a5b2 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -292,3 +292,7 @@ exclude = ["temporalio/bridge/target/**/*", "temporalio/bridge/sdk-core/.git"] # Prevent uv commands from building the package by default package = false exclude-newer = "2 weeks" + +[tool.uv.sources] +# TODO: Remove once a nexus-rpc release records the service of each operation +nexus-rpc = { git = "https://github.com/nexus-rpc/sdk-python", rev = "e1d70abbb9378a8398fd472bd466c5dfa9f1aca7" } diff --git a/scripts/gen_payload_visitor.py b/scripts/gen_payload_visitor.py index 7fab85dcd..ba5ddc112 100644 --- a/scripts/gen_payload_visitor.py +++ b/scripts/gen_payload_visitor.py @@ -12,6 +12,7 @@ sys.path.insert(0, str(base_dir)) from temporalio.api.common.v1.message_pb2 import Payload, Payloads, SearchAttributes +from temporalio.api.notificationservice.v1 import OnCompleteRequest, OnCompleteResponse from temporalio.bridge.proto.nexus import NexusTaskCompletion from temporalio.bridge.proto.workflow_activation.workflow_activation_pb2 import ( WorkflowActivation, @@ -440,11 +441,14 @@ def write_bridge_visitors() -> None: out_path = base_dir / "temporalio" / "bridge" / "_visitor.py" # Build root descriptors: WorkflowActivation, WorkflowActivationCompletion, - # NexusTaskCompletion, and the system Nexus operation roots. + # NexusTaskCompletion, the system Nexus operation roots, and the notification + # service messages delivered to Nexus handler callbacks. roots: list[Descriptor] = [ WorkflowActivation.DESCRIPTOR, WorkflowActivationCompletion.DESCRIPTOR, NexusTaskCompletion.DESCRIPTOR, + OnCompleteRequest.DESCRIPTOR, + OnCompleteResponse.DESCRIPTOR, ] + discover_system_nexus_roots() code = VisitorGenerator().generate(roots) diff --git a/temporalio/bridge/_visitor.py b/temporalio/bridge/_visitor.py index d8fa05edf..ead09a986 100644 --- a/temporalio/bridge/_visitor.py +++ b/temporalio/bridge/_visitor.py @@ -624,6 +624,21 @@ async def _visit_coresdk_nexus_NexusTaskCompletion( elif o.HasField("failure"): await self._visit_temporal_api_failure_v1_Failure(fs, o.failure) + async def _visit_temporal_api_notificationservice_v1_OnCompleteRequest( + self, fs: VisitorFunctions, o: Any + ): + if o.HasField("source_context"): + await self._visit_temporal_api_common_v1_Payload(fs, o.source_context) + if o.HasField("success"): + await self._visit_temporal_api_common_v1_Payload(fs, o.success) + elif o.HasField("failure"): + await self._visit_temporal_api_failure_v1_Failure(fs, o.failure) + + async def _visit_temporal_api_notificationservice_v1_OnCompleteResponse( + self, fs: VisitorFunctions, o: Any + ) -> None: + pass + async def _visit_temporal_api_common_v1_Header(self, fs: VisitorFunctions, o: Any): for v in o.fields.values(): await self._visit_temporal_api_common_v1_Payload(fs, v) diff --git a/temporalio/client/_impl.py b/temporalio/client/_impl.py index 0c1e67b85..eb2da0a70 100644 --- a/temporalio/client/_impl.py +++ b/temporalio/client/_impl.py @@ -1580,6 +1580,11 @@ async def start_nexus_operation( "temporalio.api.enums.v1.NexusOperationIdConflictPolicy.ValueType", int(input.id_conflict_policy), ), + # Callback source contexts use the started operation's serialization + # context, like the rest of the request. + completion_callbacks=[ + await cb._to_proto(data_converter) for cb in input.completion_callbacks + ], ) if input.schedule_to_close_timeout is not None: diff --git a/temporalio/client/_interceptor.py b/temporalio/client/_interceptor.py index 68077ebc9..9de9a9c51 100644 --- a/temporalio/client/_interceptor.py +++ b/temporalio/client/_interceptor.py @@ -18,6 +18,7 @@ import temporalio.api.common.v1 import temporalio.api.workflowservice.v1 import temporalio.common +import temporalio.nexus from temporalio.converter import DataConverter if TYPE_CHECKING: @@ -605,6 +606,7 @@ class StartNexusOperationInput: headers: Mapping[str, str] rpc_metadata: Mapping[str, str | bytes] rpc_timeout: timedelta | None + completion_callbacks: Sequence[temporalio.nexus.CompletionCallback[Any]] = () @dataclass diff --git a/temporalio/client/_nexus.py b/temporalio/client/_nexus.py index 1f0a7338e..fa22c11c2 100644 --- a/temporalio/client/_nexus.py +++ b/temporalio/client/_nexus.py @@ -490,11 +490,15 @@ async def start_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[OutputT]: ... - # Overload for string operation name + # Overload for string operation name with result_type. Callbacks must accept + # the result type. @overload @abstractmethod async def start_operation( @@ -505,17 +509,46 @@ async def start_operation( id: str, id_reuse_policy: temporalio.common.NexusOperationIDReusePolicy = temporalio.common.NexusOperationIDReusePolicy.ALLOW_DUPLICATE, id_conflict_policy: temporalio.common.NexusOperationIDConflictPolicy = temporalio.common.NexusOperationIDConflictPolicy.FAIL, - result_type: type[OutputT] | None = None, + result_type: type[OutputT], schedule_to_close_timeout: timedelta | None = None, schedule_to_start_timeout: timedelta | None = None, start_to_close_timeout: timedelta | None = None, search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[OutputT]: ... + # Overload for string operation name without result_type. Callbacks of any + # output type are accepted. + @overload + @abstractmethod + async def start_operation( + self, + operation: str, + arg: Any, + *, + id: str, + id_reuse_policy: temporalio.common.NexusOperationIDReusePolicy = temporalio.common.NexusOperationIDReusePolicy.ALLOW_DUPLICATE, + id_conflict_policy: temporalio.common.NexusOperationIDConflictPolicy = temporalio.common.NexusOperationIDConflictPolicy.FAIL, + result_type: None = None, + schedule_to_close_timeout: timedelta | None = None, + schedule_to_start_timeout: timedelta | None = None, + start_to_close_timeout: timedelta | None = None, + search_attributes: temporalio.common.TypedSearchAttributes | None = None, + summary: str | None = None, + headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), + rpc_metadata: Mapping[str, str | bytes] = {}, + rpc_timeout: timedelta | None = None, + ) -> NexusOperationHandle[Any]: ... + # Overload for workflow_run_operation methods @overload @abstractmethod @@ -536,6 +569,9 @@ async def start_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[OutputT]: ... @@ -560,6 +596,9 @@ async def start_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[OutputT]: ... @@ -584,6 +623,9 @@ async def start_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[OutputT]: ... @@ -607,6 +649,9 @@ async def start_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[OutputT]: ... @@ -636,6 +681,9 @@ async def start_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[OutputT]: ... @@ -656,6 +704,7 @@ async def start_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[temporalio.nexus.CompletionCallback[Any]] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[Any]: @@ -686,6 +735,12 @@ async def start_operation( search_attributes: Search attributes for the operation. summary: Summary for the operation. headers: Headers to attach to the Nexus request. + completion_callbacks: Callbacks that report the outcome of the + operation. Create them with + :py:func:`temporalio.nexus.create_completion_callback`. Each + callback must accept the output type of the operation. For a + string operation name, type checkers check callbacks against + ``result_type`` if it is set. rpc_metadata: Headers used on the RPC call. rpc_timeout: Optional RPC deadline to set for the RPC call. @@ -711,11 +766,15 @@ async def execute_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> OutputT: ... - # Overload for string operation name + # Overload for string operation name with result_type. Callbacks must accept + # the result type. @overload @abstractmethod async def execute_operation( @@ -726,17 +785,46 @@ async def execute_operation( id: str, id_reuse_policy: temporalio.common.NexusOperationIDReusePolicy = temporalio.common.NexusOperationIDReusePolicy.ALLOW_DUPLICATE, id_conflict_policy: temporalio.common.NexusOperationIDConflictPolicy = temporalio.common.NexusOperationIDConflictPolicy.FAIL, - result_type: type[OutputT] | None = None, + result_type: type[OutputT], schedule_to_close_timeout: timedelta | None = None, schedule_to_start_timeout: timedelta | None = None, start_to_close_timeout: timedelta | None = None, search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> OutputT: ... + # Overload for string operation name without result_type. Callbacks of any + # output type are accepted. + @overload + @abstractmethod + async def execute_operation( + self, + operation: str, + arg: Any, + *, + id: str, + id_reuse_policy: temporalio.common.NexusOperationIDReusePolicy = temporalio.common.NexusOperationIDReusePolicy.ALLOW_DUPLICATE, + id_conflict_policy: temporalio.common.NexusOperationIDConflictPolicy = temporalio.common.NexusOperationIDConflictPolicy.FAIL, + result_type: None = None, + schedule_to_close_timeout: timedelta | None = None, + schedule_to_start_timeout: timedelta | None = None, + start_to_close_timeout: timedelta | None = None, + search_attributes: temporalio.common.TypedSearchAttributes | None = None, + summary: str | None = None, + headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), + rpc_metadata: Mapping[str, str | bytes] = {}, + rpc_timeout: timedelta | None = None, + ) -> Any: ... + # Overload for workflow_run_operation methods @overload @abstractmethod @@ -757,6 +845,9 @@ async def execute_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> OutputT: ... @@ -781,6 +872,9 @@ async def execute_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> OutputT: ... @@ -805,6 +899,9 @@ async def execute_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> OutputT: ... @@ -829,6 +926,9 @@ async def execute_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> OutputT: ... @@ -858,6 +958,9 @@ async def execute_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[ + temporalio.nexus.CompletionCallback[OutputT] + ] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> OutputT: ... @@ -878,6 +981,7 @@ async def execute_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[temporalio.nexus.CompletionCallback[Any]] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> Any: @@ -910,6 +1014,12 @@ async def execute_operation( search_attributes: Search attributes for the operation. summary: Summary for the operation. headers: Headers to attach to the Nexus request. + completion_callbacks: Callbacks that report the outcome of the + operation. Create them with + :py:func:`temporalio.nexus.create_completion_callback`. Each + callback must accept the output type of the operation. For a + string operation name, type checkers check callbacks against + ``result_type`` if it is set. rpc_metadata: Headers used on the RPC call. rpc_timeout: Optional RPC deadline to set for the RPC call. @@ -974,6 +1084,7 @@ async def start_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[temporalio.nexus.CompletionCallback[Any]] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> NexusOperationHandle[Any]: @@ -1005,6 +1116,7 @@ async def start_operation( headers=dict(headers) if headers else {}, rpc_metadata=rpc_metadata, rpc_timeout=rpc_timeout, + completion_callbacks=completion_callbacks, ) ) @@ -1023,6 +1135,7 @@ async def execute_operation( search_attributes: temporalio.common.TypedSearchAttributes | None = None, summary: str | None = None, headers: Mapping[str, str] | None = None, + completion_callbacks: Sequence[temporalio.nexus.CompletionCallback[Any]] = (), rpc_metadata: Mapping[str, str | bytes] = {}, rpc_timeout: timedelta | None = None, ) -> Any: @@ -1046,6 +1159,7 @@ async def execute_operation( headers=headers, rpc_metadata=rpc_metadata, rpc_timeout=rpc_timeout, + completion_callbacks=completion_callbacks, ) return await handle.result() diff --git a/temporalio/nexus/__init__.py b/temporalio/nexus/__init__.py index 3abc9b0f2..bb7fea18a 100644 --- a/temporalio/nexus/__init__.py +++ b/temporalio/nexus/__init__.py @@ -8,6 +8,10 @@ temporal_operation, workflow_run_operation, ) +from ._notification import ( + CompletionCallback, + create_completion_callback, +) from ._operation_context import ( Info, LoggerAdapter, @@ -38,6 +42,7 @@ "CancelActivityOptions", "CancelWorkflowRunOptions", "CancelUpdateWorkflowOptions", + "CompletionCallback", "Info", "LoggerAdapter", "NexusCallback", @@ -45,6 +50,7 @@ "TemporalCancelOperationContext", "TemporalStartOperationContext", "client", + "create_completion_callback", "in_operation", "info", "is_worker_shutdown", diff --git a/temporalio/nexus/_notification.py b/temporalio/nexus/_notification.py new file mode 100644 index 000000000..31d954d8d --- /dev/null +++ b/temporalio/nexus/_notification.py @@ -0,0 +1,364 @@ +from __future__ import annotations + +import typing +from abc import ABC, abstractmethod +from collections.abc import Awaitable, Callable +from dataclasses import dataclass +from typing import Any, Generic, TypeVar, overload + +import nexusrpc + +import temporalio.api.common.v1 +import temporalio.common +import temporalio.converter + +from . import notifications +from ._operation_context import ( + TemporalStartOperationContext, + WorkflowRunOperationContext, +) +from ._temporal_client import TemporalNexusClient, TemporalOperationResult +from ._token import WorkflowHandle +from ._util import get_operation_factory + +OutputT = TypeVar("OutputT", contravariant=True) +HandlerOutputT = TypeVar("HandlerOutputT") +HandlerSourceContextT = TypeVar("HandlerSourceContextT") + + +class CompletionCallback(ABC, Generic[OutputT]): + """Callback that reports the outcome of a Nexus operation. + + The type parameter is the output type that the callback accepts. Create + callbacks with :py:func:`create_completion_callback`. Do not subclass this class. + + .. warning:: + This API is experimental and unstable. + """ + + @abstractmethod + async def _to_proto( + self, data_converter: temporalio.converter.DataConverter + ) -> temporalio.api.common.v1.Callback: + """Convert to proto representation.""" + ... + + +@dataclass(frozen=True, kw_only=True) +class _OnCompleteCallback(CompletionCallback[OutputT]): + """Callback that calls an on-complete operation.""" + + service: str + """Service name.""" + + operation: str + """Operation name.""" + + task_queue: str + """Task queue of the worker that runs the operation.""" + + source_context: Any + """Value that the operation receives with the outcome. + + The client encodes it with the serialization context of the source operation. + """ + + async def _to_proto( + self, data_converter: temporalio.converter.DataConverter + ) -> temporalio.api.common.v1.Callback: + """Convert to proto representation.""" + [payload] = await data_converter.encode([self.source_context]) + return temporalio.api.common.v1.Callback( + nexus_handler=temporalio.api.common.v1.Callback.NexusHandler( + service=self.service, + operation=self.operation, + task_queue_name=self.task_queue, + source_context=payload, + ) + ) + + +def _describe(obj: Any) -> str: + if isinstance(obj, nexusrpc.Operation): + return f"operation {obj.name!r}" + name = getattr(obj, "__qualname__", None) + return name if isinstance(name, str) else repr(obj) + + +def _check_on_complete_operation( + service: nexusrpc.ServiceDefinition, + operation: nexusrpc.OperationDefinition[Any, Any], +) -> None: + input_type = operation.input_type + if ( + not ( + input_type is notifications.OnCompleteRequest + or typing.get_origin(input_type) is notifications.OnCompleteRequest + ) + or operation.output_type is not notifications.OnCompleteResponse + ): + raise ValueError( + f"Operation {operation.name!r} of service {service.name!r} is not an " + "on-complete operation: expected input OnCompleteRequest and output " + f"OnCompleteResponse, got input {operation.input_type!r} and output " + f"{operation.output_type!r}" + ) + + +# Service and operation names. Without a handler signature the output type is +# unknown, so the callback accepts the output of any operation. +@overload +def create_completion_callback( + *, + service: str, + operation: str, + task_queue: str, + source_context: Any, +) -> CompletionCallback[Any]: ... + + +# Service definition operations. A handler that receives RawValue can accept the +# output of any operation. +@overload +def create_completion_callback( + *, + operation: nexusrpc.Operation[ + notifications.OnCompleteRequest[ + temporalio.common.RawValue, HandlerSourceContextT + ], + notifications.OnCompleteResponse, + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[Any]: ... + + +@overload +def create_completion_callback( + *, + operation: nexusrpc.Operation[ + notifications.OnCompleteRequest[HandlerOutputT, HandlerSourceContextT], + notifications.OnCompleteResponse, + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[HandlerOutputT]: ... + + +# @sync_operation handler methods defined with async def +@overload +def create_completion_callback( + *, + operation: Callable[ + [ + Any, + nexusrpc.handler.StartOperationContext, + notifications.OnCompleteRequest[ + temporalio.common.RawValue, HandlerSourceContextT + ], + ], + Awaitable[notifications.OnCompleteResponse], + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[Any]: ... + + +@overload +def create_completion_callback( + *, + operation: Callable[ + [ + Any, + nexusrpc.handler.StartOperationContext, + notifications.OnCompleteRequest[HandlerOutputT, HandlerSourceContextT], + ], + Awaitable[notifications.OnCompleteResponse], + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[HandlerOutputT]: ... + + +# @sync_operation handler methods defined with def +@overload +def create_completion_callback( + *, + operation: Callable[ + [ + Any, + nexusrpc.handler.StartOperationContext, + notifications.OnCompleteRequest[ + temporalio.common.RawValue, HandlerSourceContextT + ], + ], + notifications.OnCompleteResponse, + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[Any]: ... + + +@overload +def create_completion_callback( + *, + operation: Callable[ + [ + Any, + nexusrpc.handler.StartOperationContext, + notifications.OnCompleteRequest[HandlerOutputT, HandlerSourceContextT], + ], + notifications.OnCompleteResponse, + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[HandlerOutputT]: ... + + +# @workflow_run_operation handler methods +@overload +def create_completion_callback( + *, + operation: Callable[ + [ + Any, + WorkflowRunOperationContext, + notifications.OnCompleteRequest[ + temporalio.common.RawValue, HandlerSourceContextT + ], + ], + Awaitable[WorkflowHandle[notifications.OnCompleteResponse]], + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[Any]: ... + + +@overload +def create_completion_callback( + *, + operation: Callable[ + [ + Any, + WorkflowRunOperationContext, + notifications.OnCompleteRequest[HandlerOutputT, HandlerSourceContextT], + ], + Awaitable[WorkflowHandle[notifications.OnCompleteResponse]], + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[HandlerOutputT]: ... + + +# @temporal_operation handler methods +@overload +def create_completion_callback( + *, + operation: Callable[ + [ + Any, + TemporalStartOperationContext, + TemporalNexusClient, + notifications.OnCompleteRequest[ + temporalio.common.RawValue, HandlerSourceContextT + ], + ], + Awaitable[TemporalOperationResult[notifications.OnCompleteResponse]], + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[Any]: ... + + +@overload +def create_completion_callback( + *, + operation: Callable[ + [ + Any, + TemporalStartOperationContext, + TemporalNexusClient, + notifications.OnCompleteRequest[HandlerOutputT, HandlerSourceContextT], + ], + Awaitable[TemporalOperationResult[notifications.OnCompleteResponse]], + ], + task_queue: str, + source_context: HandlerSourceContextT, +) -> CompletionCallback[HandlerOutputT]: ... + + +def create_completion_callback( + *, + service: str | None = None, + operation: str | nexusrpc.Operation[Any, Any] | Callable[..., Any], + task_queue: str, + source_context: Any, +) -> CompletionCallback[Any]: + """Create a callback that reports the outcome of a standalone Nexus operation. + + When the operation finishes, the server calls an on-complete operation on a worker + for ``task_queue``. The call carries the result or the failure, and + ``source_context``. + + An on-complete operation takes + :py:class:`temporalio.nexus.notifications.OnCompleteRequest` and returns + :py:class:`temporalio.nexus.notifications.OnCompleteResponse`. The server can + deliver a callback more than once. Make the operation idempotent. + + .. warning:: + This API is experimental and unstable. + + Args: + service: Service name. Required if ``operation`` is a name. Otherwise, leave + it unset. + operation: Handler method, service definition operation, or operation name. + An operation name, or an operation that receives + :py:class:`temporalio.common.RawValue`, accepts any output. + task_queue: Task queue of the worker that runs the operation. + source_context: Value that the operation receives with the outcome. + + Returns: + Callback for the ``completion_callbacks`` argument of + :py:meth:`temporalio.client.NexusClient.start_operation`. + + Raises: + ValueError: If ``operation`` is not an on-complete operation of a Nexus + service, if ``service`` and an operation name are not given together, or + if either name is empty. + """ + if service is not None or isinstance(operation, str): + if not isinstance(service, str) or not isinstance(operation, str): + raise ValueError( + "Give a service name with an operation name. A handler method or a " + "service definition operation identifies its own service" + ) + if not service or not operation: + raise ValueError("Service and operation names must not be empty") + return _OnCompleteCallback( + service=service, + operation=operation, + task_queue=task_queue, + source_context=source_context, + ) + + op = ( + operation + if isinstance(operation, nexusrpc.Operation) + else get_operation_factory(operation)[1] + ) + if op is None: + raise ValueError(f"{_describe(operation)} is not a Nexus operation") + if op.service is None: + raise ValueError( + f"{_describe(operation)} is not an operation of a Nexus service. Use an " + "operation of a class decorated with @nexusrpc.service or " + "@nexusrpc.handler.service_handler" + ) + definition = op.service.operation_definitions[op.name] + _check_on_complete_operation(op.service, definition) + return _OnCompleteCallback( + service=op.service.name, + operation=definition.name, + task_queue=task_queue, + source_context=source_context, + ) diff --git a/tests/__init__.py b/tests/__init__.py index eff2c8adb..84f4d7882 100644 --- a/tests/__init__.py +++ b/tests/__init__.py @@ -1 +1 @@ -DEV_SERVER_DOWNLOAD_VERSION = "v1.8.3-server-1.32.0-162.0" +DEV_SERVER_DOWNLOAD_VERSION = "v1.9.2-nexus-handler-callbacks" diff --git a/tests/nexus/test_nexus_handler_callbacks.py b/tests/nexus/test_nexus_handler_callbacks.py new file mode 100644 index 000000000..5363d2e6d --- /dev/null +++ b/tests/nexus/test_nexus_handler_callbacks.py @@ -0,0 +1,1045 @@ +"""Tests for Nexus handler completion callbacks.""" + +from __future__ import annotations + +import asyncio +import dataclasses +import uuid +from collections.abc import AsyncIterator, Callable, Sequence +from contextlib import asynccontextmanager +from dataclasses import dataclass +from datetime import timedelta +from typing import Any + +import nexusrpc +import nexusrpc.handler +import pytest +from nexusrpc.handler import service_handler + +import temporalio.api.common.v1 +import temporalio.converter +import temporalio.nexus.notifications as notifications +from temporalio import activity, nexus, workflow +from temporalio.client import ( + Client, + NexusOperationFailureError, + WorkflowUpdateStage, +) +from temporalio.common import RawValue +from temporalio.converter import PayloadCodec +from temporalio.exceptions import ApplicationError, CancelledError, TerminatedError +from temporalio.nexus import WorkflowRunOperationContext, workflow_run_operation +from temporalio.nexus._notification import _OnCompleteCallback +from temporalio.nexus.notifications.models import ( + OnCompleteRequestResultFailure, + OnCompleteRequestResultSuccess, +) +from temporalio.testing import WorkflowEnvironment +from temporalio.worker import Worker +from tests.helpers import assert_eventually +from tests.helpers.nexus import make_nexus_endpoint_name +from tests.nexus.test_standalone_operations import ( + BlockingHandlerWorkflow, + EchoHandlerWorkflow, + EchoInput, + EchoOutput, + StandaloneTestService, + StandaloneTestServiceHandler, + _RecordingInterceptor, +) + + +@dataclasses.dataclass +class NotificationValue: + message: str + + +class PrefixCodec(PayloadCodec): + """Reversible codec that rejects payloads it did not encode.""" + + async def encode( + self, payloads: Sequence[temporalio.api.common.v1.Payload] + ) -> list[temporalio.api.common.v1.Payload]: + return [ + temporalio.api.common.v1.Payload( + metadata={"encoding": b"test/prefix"}, + data=b"prefix:" + payload.SerializeToString(), + ) + for payload in payloads + ] + + async def decode( + self, payloads: Sequence[temporalio.api.common.v1.Payload] + ) -> list[temporalio.api.common.v1.Payload]: + decoded: list[temporalio.api.common.v1.Payload] = [] + for payload in payloads: + if payload.metadata.get("encoding") != b"test/prefix": + raise RuntimeError( + f"unexpected payload passed to codec: {dict(payload.metadata)}" + ) + decoded.append( + temporalio.api.common.v1.Payload.FromString( + payload.data.removeprefix(b"prefix:") + ) + ) + return decoded + + +# --------------------------------------------------------------------------- +# Creating completion callbacks +# --------------------------------------------------------------------------- + + +@nexusrpc.handler.service_handler(name="custom.notification.service") +class CustomNotificationHandler: + @nexusrpc.handler.sync_operation(name="CustomOnComplete") + async def on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: notifications.OnCompleteRequest[NotificationValue, NotificationValue], + ) -> notifications.OnCompleteResponse: + return notifications.OnCompleteResponse() + + +@nexusrpc.handler.service_handler(name="sub.notification.service") +class SubNotificationHandler(CustomNotificationHandler): + pass + + +@nexusrpc.service(name="renamed.notification.service") +class RenamedNotificationService: + on_complete: nexusrpc.Operation[ + notifications.OnCompleteRequest[NotificationValue, NotificationValue], + notifications.OnCompleteResponse, + ] = nexusrpc.Operation(name="RenamedOnComplete") + + +@nexusrpc.service(name="sub.renamed.notification.service") +class SubRenamedNotificationService(RenamedNotificationService): + pass + + +@nexusrpc.handler.service_handler(service=RenamedNotificationService) +class RenamedNotificationHandler: + @nexusrpc.handler.sync_operation + async def on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: notifications.OnCompleteRequest[NotificationValue, NotificationValue], + ) -> notifications.OnCompleteResponse: + return notifications.OnCompleteResponse() + + +@nexusrpc.handler.service_handler(name="shape.notification.service") +class NotificationShapeHandler: + @nexusrpc.handler.sync_operation(name="BareOnComplete") + async def bare_on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: notifications.OnCompleteRequest, + ) -> notifications.OnCompleteResponse: + return notifications.OnCompleteResponse() + + @nexusrpc.handler.sync_operation(name="WrongOutput") + async def wrong_output( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: notifications.OnCompleteRequest[NotificationValue, NotificationValue], + ) -> str: + return "" + + +@nexusrpc.handler.service_handler(name="kinds.notification.service") +class OperationKindsNotificationHandler: + @nexusrpc.handler.sync_operation(name="DefOnComplete") + def def_on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: notifications.OnCompleteRequest[NotificationValue, NotificationValue], + ) -> notifications.OnCompleteResponse: + return notifications.OnCompleteResponse() + + @workflow_run_operation(name="WorkflowRunOnComplete") + async def workflow_run_on_complete( + self, + _ctx: WorkflowRunOperationContext, + _input: notifications.OnCompleteRequest[NotificationValue, NotificationValue], + ) -> nexus.WorkflowHandle[notifications.OnCompleteResponse]: + raise NotImplementedError + + @nexus.temporal_operation(name="TemporalOnComplete") + async def temporal_on_complete( + self, + _ctx: nexus.TemporalStartOperationContext, + _client: nexus.TemporalNexusClient, + _input: notifications.OnCompleteRequest[NotificationValue, NotificationValue], + ) -> nexus.TemporalOperationResult[notifications.OnCompleteResponse]: + raise NotImplementedError + + +@nexusrpc.handler.service_handler(name="not-a-notification-service") +class OtherOperationHandler: + @nexusrpc.handler.sync_operation(name="OtherOperation") + async def other_operation( + self, + _ctx: nexusrpc.handler.StartOperationContext, + input: str, + ) -> str: + return input + + +def test_completion_callback_handler_method() -> None: + """A handler method resolves to its service and registered operation name.""" + callback = nexus.create_completion_callback( + operation=CustomNotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + assert callback == _OnCompleteCallback( + task_queue="notifications", + source_context=NotificationValue("context"), + service="custom.notification.service", + operation="CustomOnComplete", + ) + + +def test_completion_callback_operation_kinds() -> None: + """def sync_operation, workflow_run_operation, and temporal_operation methods are supported.""" + for operation, name in ( + (OperationKindsNotificationHandler.def_on_complete, "DefOnComplete"), + ( + OperationKindsNotificationHandler.workflow_run_on_complete, + "WorkflowRunOnComplete", + ), + (OperationKindsNotificationHandler.temporal_on_complete, "TemporalOnComplete"), + ): + callback = nexus.create_completion_callback( + operation=operation, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + assert callback == _OnCompleteCallback( + task_queue="notifications", + source_context=NotificationValue("context"), + service="kinds.notification.service", + operation=name, + ) + + +def test_completion_callback_renamed_by_service_definition() -> None: + """Definition operations and handler methods use the definition's operation name.""" + for operation in ( + RenamedNotificationService.on_complete, + RenamedNotificationHandler.on_complete, + ): + callback = nexus.create_completion_callback( + operation=operation, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + assert callback == _OnCompleteCallback( + task_queue="notifications", + source_context=NotificationValue("context"), + service="renamed.notification.service", + operation="RenamedOnComplete", + ) + + +def test_completion_callback_inherited_operations() -> None: + """Inherited operations resolve to the service they are accessed on.""" + handler_callback = nexus.create_completion_callback( + operation=SubNotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + assert handler_callback == _OnCompleteCallback( + task_queue="notifications", + source_context=NotificationValue("context"), + service="sub.notification.service", + operation="CustomOnComplete", + ) + definition_callback = nexus.create_completion_callback( + operation=SubRenamedNotificationService.on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + assert definition_callback == _OnCompleteCallback( + task_queue="notifications", + source_context=NotificationValue("context"), + service="sub.renamed.notification.service", + operation="RenamedOnComplete", + ) + + +def test_completion_callback_local_handler_class() -> None: + """Handler classes defined in a function are supported.""" + + @nexusrpc.handler.service_handler(name="local.notification.service") + class LocalNotificationHandler: + @nexusrpc.handler.sync_operation(name="OnComplete") + async def on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: notifications.OnCompleteRequest[ + NotificationValue, NotificationValue + ], + ) -> notifications.OnCompleteResponse: + return notifications.OnCompleteResponse() + + callback = nexus.create_completion_callback( + operation=LocalNotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + assert callback == _OnCompleteCallback( + task_queue="notifications", + source_context=NotificationValue("context"), + service="local.notification.service", + operation="OnComplete", + ) + + +def test_completion_callback_names() -> None: + """A service name is given with an operation name, and only with names.""" + assert nexus.create_completion_callback( + service="external.notification.service", + operation="ExternalOnComplete", + task_queue="notifications", + source_context=NotificationValue("context"), + ) == _OnCompleteCallback( + task_queue="notifications", + source_context=NotificationValue("context"), + service="external.notification.service", + operation="ExternalOnComplete", + ) + invalid_names = "^Give a service name with an operation name" + with pytest.raises(ValueError, match=invalid_names): + nexus.create_completion_callback( # type: ignore[call-overload] + operation="ExternalOnComplete", # pyright: ignore[reportArgumentType] + task_queue="notifications", + source_context=NotificationValue("context"), + ) + with pytest.raises(ValueError, match=invalid_names): + nexus.create_completion_callback( # type: ignore[call-overload] + service="custom.notification.service", + operation=CustomNotificationHandler.on_complete, # pyright: ignore[reportArgumentType] + task_queue="notifications", + source_context=NotificationValue("context"), + ) + for service, operation in (("", "ExternalOnComplete"), ("external", "")): + with pytest.raises( + ValueError, match="^Service and operation names must not be empty$" + ): + nexus.create_completion_callback( + service=service, + operation=operation, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + + +def test_completion_callback_rejects_invalid_operations() -> None: + """Operations that are not on-complete operations of a service are rejected.""" + + class UndecoratedHandler: + async def on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: notifications.OnCompleteRequest[ + NotificationValue, NotificationValue + ], + ) -> notifications.OnCompleteResponse: + return notifications.OnCompleteResponse() + + class UndecoratedOperationHandler: + @nexusrpc.handler.sync_operation(name="OnComplete") + async def on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: notifications.OnCompleteRequest[ + NotificationValue, NotificationValue + ], + ) -> notifications.OnCompleteResponse: + return notifications.OnCompleteResponse() + + with pytest.raises( + ValueError, match="^.*UndecoratedHandler.on_complete is not a Nexus operation$" + ): + nexus.create_completion_callback( + operation=UndecoratedHandler.on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + with pytest.raises( + ValueError, + match=( + "^.*UndecoratedOperationHandler.on_complete is not an operation of a " + "Nexus service" + ), + ): + nexus.create_completion_callback( + operation=UndecoratedOperationHandler.on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + with pytest.raises(ValueError, match="^operation 'Free' is not an operation"): + nexus.create_completion_callback( + operation=nexusrpc.Operation[ + notifications.OnCompleteRequest[NotificationValue, NotificationValue], + notifications.OnCompleteResponse, + ](name="Free"), + task_queue="notifications", + source_context=NotificationValue("context"), + ) + with pytest.raises( + ValueError, + match=( + "^Operation 'OtherOperation' of service 'not-a-notification-service' is " + "not an on-complete operation: expected input OnCompleteRequest and " + "output OnCompleteResponse, got input and output " + "$" + ), + ): + nexus.create_completion_callback( + operation=OtherOperationHandler.other_operation, # type: ignore[arg-type] + task_queue="notifications", + source_context=NotificationValue("context"), + ) + + +def test_completion_callback_checks_operation_shape() -> None: + """Operations need an OnCompleteRequest input and an OnCompleteResponse output.""" + assert nexus.create_completion_callback( + operation=NotificationShapeHandler.bare_on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) == _OnCompleteCallback( + task_queue="notifications", + source_context=NotificationValue("context"), + service="shape.notification.service", + operation="BareOnComplete", + ) + with pytest.raises( + ValueError, + match="'WrongOutput' .* is not an on-complete operation: .* output $", + ): + nexus.create_completion_callback( + operation=NotificationShapeHandler.wrong_output, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + + +async def test_on_complete_callback_to_proto() -> None: + """The callback proto targets the handler operation with an encoded context.""" + callback = nexus.create_completion_callback( + operation=CustomNotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + proto = await callback._to_proto(temporalio.converter.DataConverter.default) + + assert proto.nexus_handler.service == "custom.notification.service" + assert proto.nexus_handler.operation == "CustomOnComplete" + assert proto.nexus_handler.task_queue_name == "notifications" + assert await temporalio.converter.DataConverter.default.decode( + [proto.nexus_handler.source_context], [NotificationValue] + ) == [NotificationValue("context")] + + +async def test_on_complete_callback_to_proto_applies_codec() -> None: + """The source context is encoded with the full data converter, including codecs.""" + data_converter = dataclasses.replace( + temporalio.converter.default(), payload_codec=PrefixCodec() + ) + + callback = nexus.create_completion_callback( + operation=CustomNotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationValue("context"), + ) + proto = await callback._to_proto(data_converter) + + source_context = proto.nexus_handler.source_context + assert source_context.metadata["encoding"] == b"test/prefix" + assert await data_converter.decode([source_context], [NotificationValue]) == [ + NotificationValue("context") + ] + + +# --------------------------------------------------------------------------- +# Delivering completion callbacks +# --------------------------------------------------------------------------- + + +@dataclass(frozen=True) +class CallbackContext: + label: str + + +@workflow.defn +class FailingHandlerWorkflow: + @workflow.run + async def run(self, input: EchoInput) -> EchoOutput: + raise ApplicationError(input.value, "detail", non_retryable=True) + + +@nexusrpc.service +class CallbackTestService: + fail_async: nexusrpc.Operation[EchoInput, EchoOutput] + + +@service_handler(service=CallbackTestService) +class CallbackTestServiceHandler: + @workflow_run_operation + async def fail_async( + self, ctx: WorkflowRunOperationContext, input: EchoInput + ) -> nexus.WorkflowHandle[EchoOutput]: + return await ctx.start_workflow( + FailingHandlerWorkflow.run, + input, + id=str(uuid.uuid4()), + ) + + +@nexusrpc.handler.service_handler(name="callback.test.NotificationService") +class EchoOnCompleteHandler: + def __init__(self) -> None: + # Keyed by label so redelivered callbacks don't count twice + self.requests: dict[ + str, notifications.OnCompleteRequest[EchoOutput, CallbackContext] + ] = {} + + @nexusrpc.handler.sync_operation(name="OnComplete") + async def on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + input: notifications.OnCompleteRequest[EchoOutput, CallbackContext], + ) -> notifications.OnCompleteResponse: + self.requests[input.source_context.label] = input + return notifications.OnCompleteResponse() + + +@nexusrpc.handler.service_handler(name="callback.test.RawNotificationService") +class RawOnCompleteHandler: + def __init__(self) -> None: + self.requests: dict[ + str, notifications.OnCompleteRequest[RawValue, CallbackContext] + ] = {} + + @nexusrpc.handler.sync_operation(name="OnComplete") + async def on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + input: notifications.OnCompleteRequest[RawValue, CallbackContext], + ) -> notifications.OnCompleteResponse: + self.requests[input.source_context.label] = input + return notifications.OnCompleteResponse() + + +@dataclass(frozen=True) +class RecordCompletionInput: + label: str + output: EchoOutput + + +@nexusrpc.handler.service_handler(name="callback.test.ActivityNotificationService") +class ActivityOnCompleteHandler: + def __init__(self) -> None: + self.outputs: dict[str, EchoOutput] = {} + + @activity.defn + async def record_completion( + self, input: RecordCompletionInput + ) -> notifications.OnCompleteResponse: + self.outputs[input.label] = input.output + return notifications.OnCompleteResponse() + + @nexus.temporal_operation + async def on_complete( + self, + _ctx: nexus.TemporalStartOperationContext, + client: nexus.TemporalNexusClient, + input: notifications.OnCompleteRequest[EchoOutput, CallbackContext], + ) -> nexus.TemporalOperationResult[notifications.OnCompleteResponse]: + assert isinstance(input.result, OnCompleteRequestResultSuccess) + return await client.start_activity( + self.record_completion, + RecordCompletionInput( + label=input.source_context.label, output=input.result.value + ), + id=f"on-complete-{input.source_context.label}", + start_to_close_timeout=timedelta(seconds=10), + ) + + +@asynccontextmanager +async def _callback_workers( + client: Client, + env: WorkflowEnvironment, + completion_handler: object, + operation_handler: StandaloneTestServiceHandler | None = None, + completion_activities: Sequence[Callable[..., Any]] = (), +) -> AsyncIterator[tuple[str, str]]: + """Run operation and completion workers, yielding (endpoint, completion queue).""" + operation_task_queue = str(uuid.uuid4()) + completion_task_queue = str(uuid.uuid4()) + endpoint_name = make_nexus_endpoint_name(operation_task_queue) + async with ( + Worker( + client, + task_queue=operation_task_queue, + nexus_service_handlers=[ + operation_handler or StandaloneTestServiceHandler(), + CallbackTestServiceHandler(), + ], + workflows=[ + EchoHandlerWorkflow, + BlockingHandlerWorkflow, + FailingHandlerWorkflow, + ], + ), + Worker( + client, + task_queue=completion_task_queue, + nexus_service_handlers=[completion_handler], + activities=completion_activities, + ), + ): + endpoint = await env.create_nexus_endpoint(endpoint_name, operation_task_queue) + try: + yield endpoint_name, completion_task_queue + finally: + await env.delete_nexus_endpoint(endpoint) + + +@pytest.mark.requires_local_server +async def test_completion_callback_delivers_success( + client: Client, env: WorkflowEnvironment +): + """A successful operation delivers its result and source context to the handler.""" + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + completion_handler = EchoOnCompleteHandler() + async with _callback_workers(client, env, completion_handler) as ( + endpoint_name, + completion_task_queue, + ): + handle = await client.create_nexus_client( + service=StandaloneTestService, endpoint=endpoint_name + ).start_operation( + StandaloneTestService.echo_async, + EchoInput(value="hello"), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=[ + nexus.create_completion_callback( + operation=EchoOnCompleteHandler.on_complete, + task_queue=completion_task_queue, + source_context=CallbackContext(label="success"), + ) + ], + ) + + assert await handle.result() == EchoOutput(value="hello") + + async def received_callbacks() -> None: + assert set(completion_handler.requests) == {"success"} + + await assert_eventually(received_callbacks) + request = completion_handler.requests["success"] + assert request.source_context == CallbackContext(label="success") + assert request.result == OnCompleteRequestResultSuccess(EchoOutput("hello")) + + +@pytest.mark.requires_local_server +async def test_completion_callback_delivers_failure( + client: Client, env: WorkflowEnvironment +): + """A failed operation delivers the handler workflow's failure to the handler.""" + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + completion_handler = EchoOnCompleteHandler() + async with _callback_workers(client, env, completion_handler) as ( + endpoint_name, + completion_task_queue, + ): + handle = await client.create_nexus_client( + service=CallbackTestService, endpoint=endpoint_name + ).start_operation( + CallbackTestService.fail_async, + EchoInput(value="callback test failure"), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=[ + nexus.create_completion_callback( + operation=EchoOnCompleteHandler.on_complete, + task_queue=completion_task_queue, + source_context=CallbackContext(label="failure"), + ) + ], + ) + + with pytest.raises(NexusOperationFailureError): + await handle.result() + + async def received_callbacks() -> None: + assert set(completion_handler.requests) == {"failure"} + + await assert_eventually(received_callbacks) + request = completion_handler.requests["failure"] + assert request.source_context == CallbackContext(label="failure") + assert isinstance(request.result, OnCompleteRequestResultFailure) + error = request.result.value + assert isinstance(error, ApplicationError) + assert error.message == "callback test failure" + assert error.non_retryable + + +@pytest.mark.requires_local_server +@pytest.mark.parametrize("outcome", ["cancel", "terminate"]) +async def test_completion_callback_delivers_cancellation_and_termination( + client: Client, env: WorkflowEnvironment, outcome: str +): + """Canceled and terminated operations deliver their failure to the handler.""" + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + operation_handler = StandaloneTestServiceHandler() + completion_handler = EchoOnCompleteHandler() + blocking_input = f"{outcome}-{uuid.uuid4()}" + async with _callback_workers( + client, env, completion_handler, operation_handler + ) as (endpoint_name, completion_task_queue): + handle = await client.create_nexus_client( + service=StandaloneTestService, endpoint=endpoint_name + ).start_operation( + StandaloneTestService.blocking_async, + EchoInput(value=blocking_input), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=[ + nexus.create_completion_callback( + operation=EchoOnCompleteHandler.on_complete, + task_queue=completion_task_queue, + source_context=CallbackContext(label=outcome), + ) + ], + ) + await asyncio.wait_for(operation_handler.started_blocking.wait(), timeout=10) + + if outcome == "cancel": + await handle.cancel(reason="callback test cancel") + else: + await handle.terminate(reason="callback test termination") + # Terminating the operation leaves its handler workflow running + await client.get_workflow_handle( + f"blocking_async-{blocking_input}" + ).terminate() + + with pytest.raises(NexusOperationFailureError): + await handle.result() + + async def received_callbacks() -> None: + assert set(completion_handler.requests) == {outcome} + + await assert_eventually(received_callbacks) + request = completion_handler.requests[outcome] + assert request.source_context == CallbackContext(label=outcome) + assert isinstance(request.result, OnCompleteRequestResultFailure) + if outcome == "cancel": + assert isinstance(request.result.value, CancelledError) + else: + assert isinstance(request.result.value, TerminatedError) + assert request.result.value.message == "callback test termination" + + +@pytest.mark.requires_local_server +async def test_multiple_completion_callbacks_on_operation( + client: Client, env: WorkflowEnvironment +): + """Every callback attached to an operation is passed to interceptors and delivered.""" + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + interceptor = _RecordingInterceptor() + config = client.config() + config["interceptors"] = [interceptor] + intercepted_client = Client(**config) + completion_handler = EchoOnCompleteHandler() + async with _callback_workers(client, env, completion_handler) as ( + endpoint_name, + completion_task_queue, + ): + completion_callbacks = [ + nexus.create_completion_callback( + operation=EchoOnCompleteHandler.on_complete, + task_queue=completion_task_queue, + source_context=CallbackContext(label=label), + ) + for label in ("first", "second") + ] + handle = await intercepted_client.create_nexus_client( + service=StandaloneTestService, endpoint=endpoint_name + ).start_operation( + StandaloneTestService.echo_async, + EchoInput(value="hello"), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=completion_callbacks, + ) + + [start_input] = interceptor.start_calls + assert list(start_input.completion_callbacks) == completion_callbacks + assert await handle.result() == EchoOutput(value="hello") + + async def received_callbacks() -> None: + assert set(completion_handler.requests) == {"first", "second"} + + await assert_eventually(received_callbacks) + for request in completion_handler.requests.values(): + assert request.result == OnCompleteRequestResultSuccess(EchoOutput("hello")) + + +@pytest.mark.requires_local_server +async def test_raw_value_completion_handler_for_multiple_operations( + client: Client, env: WorkflowEnvironment +): + """A RawValue handler receives completions from operations with different results.""" + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + operation_handler = StandaloneTestServiceHandler() + completion_handler = RawOnCompleteHandler() + blocking_input = f"multiple-operations-{uuid.uuid4()}" + async with _callback_workers( + client, env, completion_handler, operation_handler + ) as (endpoint_name, completion_task_queue): + nexus_client = client.create_nexus_client( + service=StandaloneTestService, endpoint=endpoint_name + ) + echo_handle = await nexus_client.start_operation( + StandaloneTestService.echo_async, + EchoInput(value="echo"), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=[ + nexus.create_completion_callback( + operation=RawOnCompleteHandler.on_complete, + task_queue=completion_task_queue, + source_context=CallbackContext(label="echo_async"), + ) + ], + ) + blocking_handle = await nexus_client.start_operation( + StandaloneTestService.blocking_async, + EchoInput(value=blocking_input), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=[ + nexus.create_completion_callback( + operation=RawOnCompleteHandler.on_complete, + task_queue=completion_task_queue, + source_context=CallbackContext(label="blocking_async"), + ) + ], + ) + + await asyncio.wait_for(operation_handler.started_blocking.wait(), timeout=10) + await client.get_workflow_handle( + f"blocking_async-{blocking_input}" + ).start_update( + BlockingHandlerWorkflow.unblock, + wait_for_stage=WorkflowUpdateStage.COMPLETED, + ) + + assert await echo_handle.result() == EchoOutput(value="echo") + assert await blocking_handle.result() == EchoOutput(value=blocking_input) + expected_results = { + "echo_async": EchoOutput(value="echo"), + "blocking_async": EchoOutput(value=blocking_input), + } + + async def received_callbacks() -> None: + assert set(completion_handler.requests) == set(expected_results) + + await assert_eventually(received_callbacks) + for label, request in completion_handler.requests.items(): + assert isinstance(request.result, OnCompleteRequestResultSuccess) + # A RawValue handler receives the payload, so convert it here + result = client.data_converter.payload_converter.from_payload( + request.result.value.payload, EchoOutput + ) + assert result == expected_results[label] + + +@pytest.mark.requires_local_server +async def test_temporal_operation_completion_handler_starts_activity( + client: Client, env: WorkflowEnvironment +): + """A temporal_operation completion handler can start a standalone activity.""" + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + completion_handler = ActivityOnCompleteHandler() + label = str(uuid.uuid4()) + async with _callback_workers( + client, + env, + completion_handler, + completion_activities=[completion_handler.record_completion], + ) as (endpoint_name, completion_task_queue): + result = await client.create_nexus_client( + service=StandaloneTestService, endpoint=endpoint_name + ).execute_operation( + StandaloneTestService.echo_async, + EchoInput(value="hello"), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=[ + nexus.create_completion_callback( + operation=ActivityOnCompleteHandler.on_complete, + task_queue=completion_task_queue, + source_context=CallbackContext(label=label), + ) + ], + ) + + assert result == EchoOutput(value="hello") + + async def activity_recorded_output() -> None: + assert completion_handler.outputs == {label: EchoOutput(value="hello")} + + await assert_eventually(activity_recorded_output) + activity_handle = client.get_activity_handle( + f"on-complete-{label}", result_type=notifications.OnCompleteResponse + ) + assert await activity_handle.result() == notifications.OnCompleteResponse() + + +@pytest.mark.requires_local_server +@pytest.mark.parametrize("success", [True, False]) +async def test_completion_callback_by_name_with_codec( + client: Client, env: WorkflowEnvironment, success: bool +): + """Name-based callbacks from execute_operation round-trip through a payload codec.""" + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + config = client.config() + config["data_converter"] = dataclasses.replace( + client.data_converter, payload_codec=PrefixCodec() + ) + codec_client = Client(**config) + completion_handler = RawOnCompleteHandler() + async with _callback_workers(codec_client, env, completion_handler) as ( + endpoint_name, + completion_task_queue, + ): + completion_callbacks = [ + nexus.create_completion_callback( + service="callback.test.RawNotificationService", + operation="OnComplete", + task_queue=completion_task_queue, + source_context=CallbackContext(label="codec"), + ) + ] + if success: + result = await codec_client.create_nexus_client( + service=StandaloneTestService, endpoint=endpoint_name + ).execute_operation( + StandaloneTestService.echo_async, + EchoInput(value="encoded"), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=completion_callbacks, + ) + assert result == EchoOutput(value="encoded") + else: + with pytest.raises(NexusOperationFailureError): + await codec_client.create_nexus_client( + service=CallbackTestService, endpoint=endpoint_name + ).execute_operation( + CallbackTestService.fail_async, + EchoInput(value="encoded failure"), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=completion_callbacks, + ) + + async def received_callbacks() -> None: + assert set(completion_handler.requests) == {"codec"} + + await assert_eventually(received_callbacks) + request = completion_handler.requests["codec"] + assert request.source_context == CallbackContext(label="codec") + if success: + assert isinstance(request.result, OnCompleteRequestResultSuccess) + # The worker decoded the codec before handing the handler the raw value + raw = request.result.value.payload + assert raw.metadata["encoding"] != b"test/prefix" + assert codec_client.data_converter.payload_converter.from_payload( + raw, EchoOutput + ) == EchoOutput(value="encoded") + else: + assert isinstance(request.result, OnCompleteRequestResultFailure) + error = request.result.value + assert isinstance(error, ApplicationError) + assert error.message == "encoded failure" + # Failure details can only be read if the worker decoded the codec + assert error.details == ("detail",) + + +@pytest.mark.requires_local_server +async def test_completion_callback_typed_handler_with_codec( + client: Client, env: WorkflowEnvironment +): + """A typed handler receives the result decoded through a payload codec.""" + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + config = client.config() + config["data_converter"] = dataclasses.replace( + client.data_converter, payload_codec=PrefixCodec() + ) + codec_client = Client(**config) + completion_handler = EchoOnCompleteHandler() + async with _callback_workers(codec_client, env, completion_handler) as ( + endpoint_name, + completion_task_queue, + ): + result = await codec_client.create_nexus_client( + service=StandaloneTestService, endpoint=endpoint_name + ).execute_operation( + StandaloneTestService.echo_async, + EchoInput(value="encoded"), + id=str(uuid.uuid4()), + schedule_to_close_timeout=timedelta(seconds=30), + completion_callbacks=[ + nexus.create_completion_callback( + operation=EchoOnCompleteHandler.on_complete, + task_queue=completion_task_queue, + source_context=CallbackContext(label="typed-codec"), + ) + ], + ) + + assert result == EchoOutput(value="encoded") + + async def received_callbacks() -> None: + assert set(completion_handler.requests) == {"typed-codec"} + + await assert_eventually(received_callbacks) + request = completion_handler.requests["typed-codec"] + assert request.source_context == CallbackContext(label="typed-codec") + assert request.result == OnCompleteRequestResultSuccess(EchoOutput("encoded")) diff --git a/tests/nexus/test_nexus_type_errors.py b/tests/nexus/test_nexus_type_errors.py index 946e34035..727299d1f 100644 --- a/tests/nexus/test_nexus_type_errors.py +++ b/tests/nexus/test_nexus_type_errors.py @@ -10,7 +10,9 @@ import nexusrpc +import temporalio.common import temporalio.nexus +import temporalio.nexus.notifications from temporalio import activity, workflow from temporalio.client import Client, NexusOperationHandle from temporalio.nexus import TemporalOperationStartHandlerFunc @@ -27,6 +29,297 @@ class MyOutput: pass +@dataclass +class NotificationSourceContext: + value: str + + +@nexusrpc.handler.service_handler( + name="temporal.notificationservice.v1.NotificationService" +) +class NotificationHandler: + @nexusrpc.handler.sync_operation(name="OnComplete") + async def on_complete( + self, + _context: nexusrpc.handler.StartOperationContext, + _input: temporalio.nexus.notifications.OnCompleteRequest[ + MyOutput, NotificationSourceContext + ], + ) -> temporalio.nexus.notifications.OnCompleteResponse: + return temporalio.nexus.notifications.OnCompleteResponse() + + +@nexusrpc.handler.service_handler(name="raw-value.notification.service") +class RawValueNotificationHandler: + @nexusrpc.handler.sync_operation(name="OnRawComplete") + async def on_raw_complete( + self, + _context: nexusrpc.handler.StartOperationContext, + _input: temporalio.nexus.notifications.OnCompleteRequest[ + temporalio.common.RawValue, NotificationSourceContext + ], + ) -> temporalio.nexus.notifications.OnCompleteResponse: + return temporalio.nexus.notifications.OnCompleteResponse() + + +@nexusrpc.handler.service_handler(name="kinds.notification.service") +class OperationKindsNotificationHandler: + @nexusrpc.handler.sync_operation + def def_on_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: temporalio.nexus.notifications.OnCompleteRequest[ + MyOutput, NotificationSourceContext + ], + ) -> temporalio.nexus.notifications.OnCompleteResponse: + return temporalio.nexus.notifications.OnCompleteResponse() + + @nexusrpc.handler.sync_operation + def def_on_raw_complete( + self, + _ctx: nexusrpc.handler.StartOperationContext, + _input: temporalio.nexus.notifications.OnCompleteRequest[ + temporalio.common.RawValue, NotificationSourceContext + ], + ) -> temporalio.nexus.notifications.OnCompleteResponse: + return temporalio.nexus.notifications.OnCompleteResponse() + + @temporalio.nexus.workflow_run_operation + async def workflow_run_on_complete( + self, + _ctx: temporalio.nexus.WorkflowRunOperationContext, + _input: temporalio.nexus.notifications.OnCompleteRequest[ + MyOutput, NotificationSourceContext + ], + ) -> temporalio.nexus.WorkflowHandle[ + temporalio.nexus.notifications.OnCompleteResponse + ]: + raise NotImplementedError + + @temporalio.nexus.workflow_run_operation + async def workflow_run_on_raw_complete( + self, + _ctx: temporalio.nexus.WorkflowRunOperationContext, + _input: temporalio.nexus.notifications.OnCompleteRequest[ + temporalio.common.RawValue, NotificationSourceContext + ], + ) -> temporalio.nexus.WorkflowHandle[ + temporalio.nexus.notifications.OnCompleteResponse + ]: + raise NotImplementedError + + @temporalio.nexus.temporal_operation + async def temporal_on_complete( + self, + _ctx: temporalio.nexus.TemporalStartOperationContext, + _client: temporalio.nexus.TemporalNexusClient, + _input: temporalio.nexus.notifications.OnCompleteRequest[ + MyOutput, NotificationSourceContext + ], + ) -> temporalio.nexus.TemporalOperationResult[ + temporalio.nexus.notifications.OnCompleteResponse + ]: + raise NotImplementedError + + @temporalio.nexus.temporal_operation + async def temporal_on_raw_complete( + self, + _ctx: temporalio.nexus.TemporalStartOperationContext, + _client: temporalio.nexus.TemporalNexusClient, + _input: temporalio.nexus.notifications.OnCompleteRequest[ + temporalio.common.RawValue, NotificationSourceContext + ], + ) -> temporalio.nexus.TemporalOperationResult[ + temporalio.nexus.notifications.OnCompleteResponse + ]: + raise NotImplementedError + + +@nexusrpc.service(name="definition.notification.service") +class NotificationDefinition: + on_complete: nexusrpc.Operation[ + temporalio.nexus.notifications.OnCompleteRequest[ + MyOutput, NotificationSourceContext + ], + temporalio.nexus.notifications.OnCompleteResponse, + ] + on_raw_complete: nexusrpc.Operation[ + temporalio.nexus.notifications.OnCompleteRequest[ + temporalio.common.RawValue, NotificationSourceContext + ], + temporalio.nexus.notifications.OnCompleteResponse, + ] + + +# The declared types are wrong, so each error message shows the inferred type. +def completion_callback_type_inference() -> None: + # overloads infer the handler's output type + # assert-type-error-pyright: 'Type "CompletionCallback\[MyOutput\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # assert-type-error-pyright: 'Type "CompletionCallback\[MyOutput\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=NotificationDefinition.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # assert-type-error-pyright: 'Type "CompletionCallback\[Any\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=RawValueNotificationHandler.on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # assert-type-error-pyright: 'Type "CompletionCallback\[Any\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + service="named.notification.service", + operation="OnComplete", + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # def sync_operation, workflow_run_operation, and temporal_operation handler + # methods are supported + # assert-type-error-pyright: 'Type "CompletionCallback\[MyOutput\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=OperationKindsNotificationHandler.def_on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # assert-type-error-pyright: 'Type "CompletionCallback\[Any\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=OperationKindsNotificationHandler.def_on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # assert-type-error-pyright: 'Type "CompletionCallback\[MyOutput\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=OperationKindsNotificationHandler.workflow_run_on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # assert-type-error-pyright: 'Type "CompletionCallback\[Any\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=OperationKindsNotificationHandler.workflow_run_on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # assert-type-error-pyright: 'Type "CompletionCallback\[MyOutput\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=OperationKindsNotificationHandler.temporal_on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + # assert-type-error-pyright: 'Type "CompletionCallback\[Any\]" is not assignable to declared type "None"' + _: None = temporalio.nexus.create_completion_callback( # type: ignore + operation=OperationKindsNotificationHandler.temporal_on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + + +# These calls fail at runtime, so they are only type checked. +def invalid_completion_callbacks() -> None: + # source_context must match the handler's source context type + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=NotificationHandler.on_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=NotificationDefinition.on_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=NotificationDefinition.on_raw_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=RawValueNotificationHandler.on_raw_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=OperationKindsNotificationHandler.def_on_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=OperationKindsNotificationHandler.def_on_raw_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=OperationKindsNotificationHandler.workflow_run_on_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=OperationKindsNotificationHandler.workflow_run_on_raw_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=OperationKindsNotificationHandler.temporal_on_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( + operation=OperationKindsNotificationHandler.temporal_on_raw_complete, # type: ignore[arg-type] + task_queue="notifications", + # assert-type-error-pyright: 'Argument of type "MyInput" cannot be assigned to parameter "source_context"' + source_context=MyInput(), # type: ignore + ) + + # operation names require a service name + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( # type: ignore[call-overload] + # assert-type-error-pyright: 'Argument of type "Literal\['OnComplete'\]" cannot be assigned to parameter "operation"' + operation="OnComplete", # type: ignore + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + + # operations other than names already identify their service + temporalio.nexus.create_completion_callback( + service="named.notification.service", + # assert-type-error-pyright: 'cannot be assigned to parameter "operation" of type "str"' + operation=NotificationHandler.on_complete, # type: ignore + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + + # all arguments are keyword arguments + # assert-type-error-pyright: 'No overloads for "create_completion_callback" match' + temporalio.nexus.create_completion_callback( # type: ignore[call-overload] + NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ) + + @workflow.defn class MyNoArgProcWorkflow: @workflow.run @@ -788,6 +1081,190 @@ async def standalone_operation_type_tests(): ) _defn_handle_output: MyOutput = await _defn_handle.result() + # completion callback output types must match the started operation + _callback_handle: NexusOperationHandle[ + MyOutput + ] = await nexus_client.start_operation( + MyService.my_sync_operation, + MyInput(), + id="op-with-callback", + completion_callbacks=[ + temporalio.nexus.create_completion_callback( + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + _mixed_callback_handle: NexusOperationHandle[ + MyOutput + ] = await nexus_client.start_operation( + MyService.my_sync_operation, + MyInput(), + id="op-with-mixed-callbacks", + completion_callbacks=[ + temporalio.nexus.create_completion_callback( + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + # RawValue handlers accept the output of any operation + temporalio.nexus.create_completion_callback( + operation=RawValueNotificationHandler.on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + # RawValue and by-name callbacks are compatible with any output type + await nexus_client.start_operation( + MyService.my_temporal_operation, + 0, + id="op-with-any-output-callbacks", + completion_callbacks=[ + temporalio.nexus.create_completion_callback( + operation=RawValueNotificationHandler.on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + # operation names have no handler signature, so they accept any output type + temporalio.nexus.create_completion_callback( + service="named.notification.service", + operation="OnComplete", + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + # assert-type-error-pyright: 'No overloads for "start_operation" match' + await nexus_client.start_operation( # type: ignore + MyService.my_temporal_operation, + 0, + id="op-with-wrong-callback", + completion_callbacks=[ # type: ignore[arg-type] + # assert-type-error-pyright: '"list\[CompletionCallback\[MyOutput\]\]" cannot be assigned to parameter "completion_callbacks"' + temporalio.nexus.create_completion_callback( # type: ignore + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + _callback_output: MyOutput = await nexus_client.execute_operation( + MyService.my_sync_operation, + MyInput(), + id="execute-with-callback", + completion_callbacks=[ + temporalio.nexus.create_completion_callback( + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + temporalio.nexus.create_completion_callback( + operation=RawValueNotificationHandler.on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + # assert-type-error-pyright: 'No overloads for "execute_operation" match' + await nexus_client.execute_operation( # type: ignore + MyService.my_temporal_operation, + 0, + id="execute-with-wrong-callback", + completion_callbacks=[ # type: ignore[arg-type] + # assert-type-error-pyright: '"list\[CompletionCallback\[MyOutput\]\]" cannot be assigned to parameter "completion_callbacks"' + temporalio.nexus.create_completion_callback( # type: ignore + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + + # with a string operation name, callbacks must accept result_type + _str_callback_handle: NexusOperationHandle[ + MyOutput + ] = await nexus_client.start_operation( + "my_sync_operation", + MyInput(), + id="str-op-with-callbacks", + result_type=MyOutput, + completion_callbacks=[ + temporalio.nexus.create_completion_callback( + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + temporalio.nexus.create_completion_callback( + operation=RawValueNotificationHandler.on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + _str_callback_output: MyOutput = await nexus_client.execute_operation( + "my_sync_operation", + MyInput(), + id="str-execute-with-callbacks", + result_type=MyOutput, + completion_callbacks=[ + temporalio.nexus.create_completion_callback( + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + temporalio.nexus.create_completion_callback( + operation=RawValueNotificationHandler.on_raw_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + # without result_type, any callback is accepted + await nexus_client.start_operation( + "my_sync_operation", + MyInput(), + id="str-op-without-result-type", + completion_callbacks=[ + temporalio.nexus.create_completion_callback( + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + # assert-type-error-pyright: 'No overloads for "start_operation" match' + await nexus_client.start_operation( # type: ignore + "my_sync_operation", + MyInput(), + id="str-op-with-wrong-callback", + result_type=str, + completion_callbacks=[ # type: ignore[arg-type] + # assert-type-error-pyright: '"list\[CompletionCallback\[MyOutput\]\]" cannot be assigned to parameter "completion_callbacks"' + temporalio.nexus.create_completion_callback( # type: ignore + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + # assert-type-error-pyright: 'No overloads for "execute_operation" match' + await nexus_client.execute_operation( # type: ignore + "my_sync_operation", + MyInput(), + id="str-execute-with-wrong-callback", + result_type=str, + completion_callbacks=[ # type: ignore[arg-type] + # assert-type-error-pyright: '"list\[CompletionCallback\[MyOutput\]\]" cannot be assigned to parameter "completion_callbacks"' + temporalio.nexus.create_completion_callback( # type: ignore + operation=NotificationHandler.on_complete, + task_queue="notifications", + source_context=NotificationSourceContext(value="ok"), + ), + ], + ) + # result_type is not allowed when an operation is provided await nexus_client.start_operation( # assert-type-error-pyright: 'cannot be assigned to parameter "operation" of type "str"' diff --git a/uv.lock b/uv.lock index 6795afdf4..084bf04f0 100644 --- a/uv.lock +++ b/uv.lock @@ -10,7 +10,7 @@ resolution-markers = [ ] [options] -exclude-newer = "0001-01-01T00:00:00Z" # This has no effect and is included for backwards compatibility when using relative exclude-newer values. +exclude-newer = "2026-09-18T18:50:56.899109Z" exclude-newer-span = "P2W" [[package]] @@ -258,13 +258,13 @@ name = "anthropic" version = "1.3.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "anyio" }, - { name = "docstring-parser" }, - { name = "httpx2" }, - { name = "jiter" }, - { name = "pydantic" }, - { name = "sniffio" }, - { name = "typing-extensions" }, + { name = "anyio", marker = "python_full_version >= '3.11'" }, + { name = "docstring-parser", marker = "python_full_version >= '3.11'" }, + { name = "httpx2", marker = "python_full_version >= '3.11'" }, + { name = "jiter", marker = "python_full_version >= '3.11'" }, + { name = "pydantic", marker = "python_full_version >= '3.11'" }, + { name = "sniffio", marker = "python_full_version >= '3.11'" }, + { name = "typing-extensions", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/b4/50/463166f02179ab279edb61de1589a6f69cb3838d6a2fb6f2c92a3f8042f1/anthropic-1.3.0.tar.gz", hash = "sha256:6873492a77ede8849a161ab1bc78bc9a1e492a006d0b5bb4c57ac77845df838a", size = 1148177, upload-time = "2026-09-01T17:37:10.392Z" } wheels = [ @@ -942,13 +942,13 @@ name = "deepagents" version = "0.7.12" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "langchain" }, - { name = "langchain-anthropic" }, - { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" } }, - { name = "langchain-google-genai" }, - { name = "langsmith" }, - { name = "packaging" }, - { name = "wcmatch" }, + { name = "langchain", marker = "python_full_version >= '3.11'" }, + { name = "langchain-anthropic", marker = "python_full_version >= '3.11'" }, + { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, + { name = "langchain-google-genai", marker = "python_full_version >= '3.11'" }, + { name = "langsmith", marker = "python_full_version >= '3.11'" }, + { name = "packaging", marker = "python_full_version >= '3.11'" }, + { name = "wcmatch", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/b9/af/35d8adf1181a27f2d98f535d81246d73328297a5daf99c4130a5c667494e/deepagents-0.7.12.tar.gz", hash = "sha256:e5968af37f505d79bea2bf7217dad0bdaff6689950e33a1a0f82c25f7ac9fd7e", size = 297304, upload-time = "2026-09-01T18:50:18.6Z" } wheels = [ @@ -1023,7 +1023,7 @@ name = "exceptiongroup" version = "1.3.1" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "typing-extensions" }, + { name = "typing-extensions", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" } wheels = [ @@ -1597,8 +1597,8 @@ name = "httpcore2" version = "2.12.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "h11" }, - { name = "truststore" }, + { name = "h11", marker = "python_full_version != '3.12.*' or sys_platform != 'emscripten'" }, + { name = "truststore", marker = "python_full_version != '3.12.*' or sys_platform != 'emscripten'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/be/ad/f4f0e57345f1870f3e8cb624e058d7eca6e5a27d33bcc3311d9b618734cd/httpcore2-2.12.0.tar.gz", hash = "sha256:9293522bba0aa7c4c8e9e3f040c16575bd8868e155a77fa30c7a9085a5eae648", size = 67548, upload-time = "2026-08-18T13:22:08.211Z" } wheels = [ @@ -2012,9 +2012,9 @@ name = "langchain" version = "1.3.18" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" } }, - { name = "langgraph", version = "1.2.11", source = { registry = "https://pypi.org/simple" } }, - { name = "pydantic" }, + { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, + { name = "langgraph", version = "1.2.11", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, + { name = "pydantic", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/31/e9/0e5425e522b624d7c5c0911957b467b82890b36b0099bf727951f2ee3059/langchain-1.3.18.tar.gz", hash = "sha256:74ce99294f6f2c82ee64c3df39daa6a22085ee1449a718065b3604a88d78bf4d", size = 637052, upload-time = "2026-08-27T17:33:12.702Z" } wheels = [ @@ -2026,9 +2026,9 @@ name = "langchain-anthropic" version = "1.7.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "anthropic" }, - { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" } }, - { name = "pydantic" }, + { name = "anthropic", marker = "python_full_version >= '3.11'" }, + { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, + { name = "pydantic", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/56/fc/52f6d1d6069bafb08626e204c89c49c8dd4a536eedbb94f0b7e78668594d/langchain_anthropic-1.7.0.tar.gz", hash = "sha256:d48e3c118ff8d3eea83f17b50234a2d2ff491a2375d565f212eb990e7e3856cb", size = 750068, upload-time = "2026-08-27T15:23:59.261Z" } wheels = [ @@ -2043,15 +2043,15 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "jsonpatch" }, - { name = "langchain-protocol" }, - { name = "langsmith" }, - { name = "packaging" }, - { name = "pydantic" }, - { name = "pyyaml" }, - { name = "tenacity" }, - { name = "typing-extensions" }, - { name = "uuid-utils" }, + { name = "jsonpatch", marker = "python_full_version < '3.11'" }, + { name = "langchain-protocol", marker = "python_full_version < '3.11'" }, + { name = "langsmith", marker = "python_full_version < '3.11'" }, + { name = "packaging", marker = "python_full_version < '3.11'" }, + { name = "pydantic", marker = "python_full_version < '3.11'" }, + { name = "pyyaml", marker = "python_full_version < '3.11'" }, + { name = "tenacity", marker = "python_full_version < '3.11'" }, + { name = "typing-extensions", marker = "python_full_version < '3.11'" }, + { name = "uuid-utils", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/2a/b9/e937d0a90b26540bff07e7a7c64349f3b29c2dcc36257cd1cd3fdce17f2a/langchain_core-1.4.9.tar.gz", hash = "sha256:f8078901145bed0466755277500a5a22822a7b628808c4c0a28d4fc88895fcf2", size = 967294, upload-time = "2026-07-08T20:06:54.191Z" } wheels = [ @@ -2069,16 +2069,16 @@ resolution-markers = [ "(python_full_version >= '3.11' and python_full_version < '3.13' and sys_platform != 'emscripten') or (python_full_version == '3.11.*' and sys_platform == 'emscripten')", ] dependencies = [ - { name = "httpx" }, - { name = "jsonpatch" }, - { name = "langchain-protocol" }, - { name = "langsmith" }, - { name = "packaging" }, - { name = "pydantic" }, - { name = "pyyaml" }, - { name = "tenacity" }, - { name = "typing-extensions" }, - { name = "uuid-utils" }, + { name = "httpx", marker = "python_full_version >= '3.11'" }, + { name = "jsonpatch", marker = "python_full_version >= '3.11'" }, + { name = "langchain-protocol", marker = "python_full_version >= '3.11'" }, + { name = "langsmith", marker = "python_full_version >= '3.11'" }, + { name = "packaging", marker = "python_full_version >= '3.11'" }, + { name = "pydantic", marker = "python_full_version >= '3.11'" }, + { name = "pyyaml", marker = "python_full_version >= '3.11'" }, + { name = "tenacity", marker = "python_full_version >= '3.11'" }, + { name = "typing-extensions", marker = "python_full_version >= '3.11'" }, + { name = "uuid-utils", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/90/12/aff76ca89c219ebe6f9dd3c5dbc4e3b1cf5450e9fc7037dccad23d45cd7a/langchain_core-1.6.1.tar.gz", hash = "sha256:1b156cb395aac4f009a8a1b38a574c7d948fe2d5f74c96e0d8a5017b4149e04f", size = 1003359, upload-time = "2026-08-27T19:31:14.956Z" } wheels = [ @@ -2090,10 +2090,10 @@ name = "langchain-google-genai" version = "4.4.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "filetype" }, - { name = "google-genai" }, - { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" } }, - { name = "pydantic" }, + { name = "filetype", marker = "python_full_version >= '3.11'" }, + { name = "google-genai", marker = "python_full_version >= '3.11'" }, + { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, + { name = "pydantic", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/cf/4e/41798c80b574d958d189e049f13d64eb9246a66623465290f2cbfd641759/langchain_google_genai-4.4.0.tar.gz", hash = "sha256:7871beec56ac07b719f77c46997845db6ff2267b817bffb0f97877053b0895d7", size = 378415, upload-time = "2026-09-01T20:15:45.816Z" } wheels = [ @@ -2120,12 +2120,12 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "langchain-core", version = "1.4.9", source = { registry = "https://pypi.org/simple" } }, - { name = "langgraph-checkpoint" }, - { name = "langgraph-prebuilt" }, - { name = "langgraph-sdk" }, - { name = "pydantic" }, - { name = "xxhash" }, + { name = "langchain-core", version = "1.4.9", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "langgraph-checkpoint", marker = "python_full_version < '3.11'" }, + { name = "langgraph-prebuilt", marker = "python_full_version < '3.11'" }, + { name = "langgraph-sdk", marker = "python_full_version < '3.11'" }, + { name = "pydantic", marker = "python_full_version < '3.11'" }, + { name = "xxhash", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/41/4b/0d1130e26b41a99dcc88353bbe7162a1f255c4db746bd94024268e6af27b/langgraph-1.2.9.tar.gz", hash = "sha256:385f87bc1802c35af7e0aa479278ecba8582d103515eb48256cb2ddcd42d0bd4", size = 722869, upload-time = "2026-07-10T01:30:14.985Z" } wheels = [ @@ -2143,12 +2143,12 @@ resolution-markers = [ "(python_full_version >= '3.11' and python_full_version < '3.13' and sys_platform != 'emscripten') or (python_full_version == '3.11.*' and sys_platform == 'emscripten')", ] dependencies = [ - { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" } }, - { name = "langgraph-checkpoint" }, - { name = "langgraph-prebuilt" }, - { name = "langgraph-sdk" }, - { name = "pydantic" }, - { name = "xxhash" }, + { name = "langchain-core", version = "1.6.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, + { name = "langgraph-checkpoint", marker = "python_full_version >= '3.11'" }, + { name = "langgraph-prebuilt", marker = "python_full_version >= '3.11'" }, + { name = "langgraph-sdk", marker = "python_full_version >= '3.11'" }, + { name = "pydantic", marker = "python_full_version >= '3.11'" }, + { name = "xxhash", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/56/0d/c8e7ee98896659e1b6555db0ab115a9ca899844744645d5d894032bab1d7/langgraph-1.2.11.tar.gz", hash = "sha256:9ecfe11e50d338b34b15cf4d8a442642de103e8ae6971320efba84e4542eb363", size = 725753, upload-time = "2026-08-11T14:00:36.945Z" } wheels = [ @@ -2836,14 +2836,10 @@ wheels = [ [[package]] name = "nexus-rpc" version = "1.4.0" -source = { registry = "https://pypi.org/simple" } +source = { git = "https://github.com/nexus-rpc/sdk-python?rev=e1d70abbb9378a8398fd472bd466c5dfa9f1aca7#e1d70abbb9378a8398fd472bd466c5dfa9f1aca7" } dependencies = [ { name = "typing-extensions" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/35/d5/cd1ffb202b76ebc1b33c1332a3416e55a39929006982adc2b1eb069aaa9b/nexus_rpc-1.4.0.tar.gz", hash = "sha256:3b8b373d4865671789cc43623e3dc0bcbf192562e40e13727e17f1c149050fba", size = 82367, upload-time = "2026-02-25T22:01:34.053Z" } -wheels = [ - { url = "https://files.pythonhosted.org/packages/11/52/6327a5f4fda01207205038a106a99848a41c83e933cd23ea2cab3d2ebc6c/nexus_rpc-1.4.0-py3-none-any.whl", hash = "sha256:14c953d3519113f8ccec533a9efdb6b10c28afef75d11cdd6d422640c40b3a49", size = 29645, upload-time = "2026-02-25T22:01:33.122Z" }, -] [[package]] name = "nh3" @@ -4614,8 +4610,8 @@ name = "secretstorage" version = "3.5.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "cryptography" }, - { name = "jeepney" }, + { name = "cryptography", marker = "python_full_version != '3.12.*' or sys_platform != 'emscripten'" }, + { name = "jeepney", marker = "python_full_version != '3.12.*' or sys_platform != 'emscripten'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/1c/03/e834bcd866f2f8a49a85eaff47340affa3bfa391ee9912a952a1faa68c7b/secretstorage-3.5.0.tar.gz", hash = "sha256:f04b8e4689cbce351744d5537bf6b1329c6fc68f91fa666f60a380edddcd11be", size = 19884, upload-time = "2025-11-23T19:02:53.191Z" } wheels = [ @@ -4903,7 +4899,7 @@ requires-dist = [ { name = "langgraph", marker = "extra == 'langgraph'", specifier = ">=1.1.0" }, { name = "langsmith", marker = "extra == 'langsmith'", specifier = ">=0.7.34,<0.13" }, { name = "mcp", marker = "extra == 'google-adk'", specifier = ">=1.24,<2" }, - { name = "nexus-rpc", specifier = "==1.4.0" }, + { name = "nexus-rpc", git = "https://github.com/nexus-rpc/sdk-python?rev=e1d70abbb9378a8398fd472bd466c5dfa9f1aca7" }, { name = "opentelemetry-api", marker = "extra == 'cloud-run-worker-otel'", specifier = ">=1.26,<2" }, { name = "opentelemetry-api", marker = "extra == 'lambda-worker-otel'", specifier = ">=1.26,<2" }, { name = "opentelemetry-api", marker = "extra == 'opentelemetry'", specifier = ">=1.26,<2" }, @@ -5483,7 +5479,7 @@ name = "wcmatch" version = "11.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "bracex" }, + { name = "bracex", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/16/25/1da725838132221e33568973da484ff43813662ccc06ebf7f6e3abddfcd5/wcmatch-11.0.tar.gz", hash = "sha256:55d95c2447789712774b198ceec72939e88b5618f1f8f0a9b605bf7740b63b96", size = 141360, upload-time = "2026-07-10T05:50:24.183Z" } wheels = [