From a7c37daec0a4b633d3ab0d4c758c80010d216335 Mon Sep 17 00:00:00 2001 From: Alex Mazzeo Date: Tue, 4 Aug 2026 11:41:57 -0700 Subject: [PATCH 1/2] Only send request ID when in a backing nexus context when starting a SAA --- temporalio/nexus/_operation_context.py | 2 +- tests/nexus/test_link_propagation.py | 3 ++- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/temporalio/nexus/_operation_context.py b/temporalio/nexus/_operation_context.py index a4c983943..cae533650 100644 --- a/temporalio/nexus/_operation_context.py +++ b/temporalio/nexus/_operation_context.py @@ -814,10 +814,10 @@ def _apply_nexus_context_to_start_activity_request( # pyright: ignore[reportUnu req.on_conflict_options.attach_completion_callbacks = True req.on_conflict_options.attach_links = True - req.request_id = nexus_ctx.nexus_context.request_id request_links = nexus_ctx._get_request_links() if _in_nexus_backing_start_context(): + req.request_id = nexus_ctx.nexus_context.request_id callbacks = nexus_ctx._get_callbacks( OperationToken( type=OperationTokenType.ACTIVITY, diff --git a/tests/nexus/test_link_propagation.py b/tests/nexus/test_link_propagation.py index 4f8c4323f..65273a172 100644 --- a/tests/nexus/test_link_propagation.py +++ b/tests/nexus/test_link_propagation.py @@ -497,7 +497,8 @@ async def test_activity_start_forwards_inbound_links() -> None: assert len(req.links) == 1 assert req.links[0] == _inbound_nexus_link() - assert req.request_id == "req-id" + assert req.request_id + assert req.request_id != "req-id" assert len(req.completion_callbacks) == 0 From dfd9a49131e96e047b3897eb76c8fda51f7ecc64 Mon Sep 17 00:00:00 2001 From: Alex Mazzeo Date: Thu, 20 Aug 2026 14:01:05 -0700 Subject: [PATCH 2/2] Test multiple activity starts from a Nexus handler --- tests/nexus/test_temporal_operation.py | 75 ++++++++++++++++++++++++++ 1 file changed, 75 insertions(+) diff --git a/tests/nexus/test_temporal_operation.py b/tests/nexus/test_temporal_operation.py index 85948deb3..952c3f217 100644 --- a/tests/nexus/test_temporal_operation.py +++ b/tests/nexus/test_temporal_operation.py @@ -1139,6 +1139,81 @@ async def test_temporal_operation_start_activity( assert result == "test" +async def test_temporal_operation_can_start_multiple_activities_with_same_id( + client: Client, env: WorkflowEnvironment +): + if env.supports_time_skipping: + pytest.skip( + "Standalone Nexus Operation tests don't work with time-skipping server" + ) + + task_queue = str(uuid.uuid4()) + endpoint_name = make_nexus_endpoint_name(task_queue) + activity_id = str(uuid.uuid4()) + await env.create_nexus_endpoint(endpoint_name, task_queue) + + @service_handler + class MultipleActivitiesHandler: + def __init__(self) -> None: + self.first_activity_run_id: str | None = None + self.start_completed = asyncio.Event() + + @nexus.temporal_operation + async def start_activities( + self, + _ctx: nexus.TemporalStartOperationContext, + client: nexus.TemporalNexusClient, + _input: None, + ) -> nexus.TemporalOperationResult[str]: + try: + first_handle = await nexus.client().start_activity( + echo_activity, + Input(value="first", task_queue=task_queue), + id=activity_id, + task_queue=task_queue, + schedule_to_close_timeout=timedelta(seconds=5), + ) + assert first_handle.run_id + self.first_activity_run_id = first_handle.run_id + await first_handle.result() + + return await client.start_activity( + echo_activity, + Input(value="second", task_queue=task_queue), + id=activity_id, + schedule_to_close_timeout=timedelta(seconds=5), + ) + finally: + self.start_completed.set() + + handler = MultipleActivitiesHandler() + async with Worker( + env.client, + task_queue=task_queue, + nexus_service_handlers=[handler], + activities=[echo_activity], + ): + nexus_client = client.create_nexus_client( + MultipleActivitiesHandler, endpoint_name + ) + operation_handle = await nexus_client.start_operation( + MultipleActivitiesHandler.start_activities, + None, + id=str(uuid.uuid4()), + ) + await asyncio.wait_for(handler.start_completed.wait(), timeout=15) + + assert handler.first_activity_run_id + first_result = await client.get_activity_handle( + activity_id, run_id=handler.first_activity_run_id, result_type=str + ).result() + second_result = await client.get_activity_handle( + activity_id, result_type=str + ).result() + assert [first_result, second_result] == ["first", "second"] + assert await operation_handle.result() == second_result + + async def test_temporal_operation_backing_activity_does_not_duplicate_links( client: Client, env: WorkflowEnvironment ):