From 579b616e5d317f52dabb5a2a5acfab4d03b46477 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Wed, 2 Sep 2026 01:42:10 -0500 Subject: [PATCH 1/5] Propagate workflow stream subscription cancellation Signed-off-by: 1fanwang <1fannnw@gmail.com> --- CHANGELOG.md | 2 + .../contrib/workflow_streams/_client.py | 6 +- .../workflow_streams/test_workflow_streams.py | 80 +++++++++++++++++-- 3 files changed, 81 insertions(+), 7 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3a7ef493d..3d77dac9c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -150,6 +150,8 @@ This file contains assembled releases only. - Cancelling an activity from a signal while the workflow itself is cancelled no longer causes a nondeterminism error from duplicate activity-cancellation commands. +- `WorkflowStreamClient.subscribe` now propagates task cancellation instead of + ending the subscription normally. - `StrandsPlugin` now disables Botocore retries for its default Bedrock model so model request retries are handled exclusively by Temporal. - `temporalio.contrib.openai_agents` now honors the `retry-after-ms` and diff --git a/temporalio/contrib/workflow_streams/_client.py b/temporalio/contrib/workflow_streams/_client.py index 202558fd3..8ed6cb5b5 100644 --- a/temporalio/contrib/workflow_streams/_client.py +++ b/temporalio/contrib/workflow_streams/_client.py @@ -585,7 +585,7 @@ async def subscribe( self._polled_run_id = handle.workflow_run_id result: PollResult = await handle.result() except asyncio.CancelledError: - return + raise except WorkflowUpdateFailedError as e: cause_type = getattr(e.cause, "type", None) if cause_type == TRUNCATED_OFFSET_ERROR_TYPE: @@ -611,7 +611,9 @@ async def subscribe( continue return raise - except WorkflowUpdateRPCTimeoutOrCancelledError: + except WorkflowUpdateRPCTimeoutOrCancelledError as e: + if isinstance(e.__cause__, asyncio.CancelledError): + raise e.__cause__ if await self._follow_continue_as_new(): continue return diff --git a/tests/contrib/workflow_streams/test_workflow_streams.py b/tests/contrib/workflow_streams/test_workflow_streams.py index a98d467c4..d48890bb9 100644 --- a/tests/contrib/workflow_streams/test_workflow_streams.py +++ b/tests/contrib/workflow_streams/test_workflow_streams.py @@ -50,7 +50,7 @@ ) from temporalio.contrib.workflow_streams._types import _encode_payload from temporalio.converter import DataConverter, PayloadCodec -from temporalio.exceptions import ApplicationError +from temporalio.exceptions import ActivityError, ApplicationError from temporalio.nexus import WorkflowRunOperationContext, workflow_run_operation from temporalio.service import RPCError, RPCStatusCode from temporalio.testing import WorkflowEnvironment @@ -103,6 +103,37 @@ async def run(self) -> None: await workflow.wait_condition(lambda: self._closed) +@workflow.defn +class CancelSubscriptionWorkflow: + @workflow.init + def __init__(self) -> None: + self.stream = WorkflowStream() + self._cancel_requested = False + + @workflow.signal + def cancel_subscription(self) -> None: + self._cancel_requested = True + + @workflow.run + async def run(self) -> str: + self.stream.topic("events", type=bytes).publish(b"seed") + handle = workflow.start_activity( + "subscribe_until_cancelled", + start_to_close_timeout=timedelta(seconds=30), + heartbeat_timeout=timedelta(seconds=1), + cancellation_type=workflow.ActivityCancellationType.WAIT_CANCELLATION_COMPLETED, + ) + await workflow.wait_condition(lambda: self._cancel_requested) + handle.cancel() + try: + result = await handle + except ActivityError as err: + result = type(err.cause).__name__ + self.stream.detach_pollers() + await workflow.wait_condition(workflow.all_handlers_finished) + return result + + @workflow.defn class ActivityPublishWorkflow: @workflow.init @@ -375,6 +406,28 @@ async def publish_items(count: int) -> None: client.topic("events", type=bytes).publish(f"item-{i}".encode()) +class CancellableSubscriber: + def __init__(self) -> None: + self.started = asyncio.Event() + + @activity.defn(name="subscribe_until_cancelled") + async def subscribe(self) -> str: + async def heartbeat() -> None: + while True: + activity.heartbeat() + await asyncio.sleep(0.1) + + heartbeat_task = asyncio.create_task(heartbeat()) + try: + stream = WorkflowStreamClient.from_within_activity() + async for _ in stream.subscribe(result_type=bytes): + self.started.set() + return "subscription-ended" + finally: + heartbeat_task.cancel() + await asyncio.gather(heartbeat_task, return_exceptions=True) + + @activity.defn(name="publish_multi_topic") async def publish_multi_topic(count: int) -> None: topics = ["a", "b", "c"] @@ -1091,7 +1144,7 @@ async def test_priority_flush(client: Client) -> None: @pytest.mark.asyncio async def test_iterator_cancellation(client: Client) -> None: """Cancelling a subscription iterator after it has yielded an item - completes cleanly.""" + propagates cancellation.""" async with new_worker( client, BasicWorkflowStreamWorkflow, @@ -1127,10 +1180,8 @@ async def subscribe_and_collect() -> None: async with _async_timeout(5): await first_item.wait() task.cancel() - try: + with pytest.raises(asyncio.CancelledError): await task - except asyncio.CancelledError: - pass assert len(items) == 1 assert items[0].data == b"seed" @@ -1138,6 +1189,25 @@ async def subscribe_and_collect() -> None: await handle.signal(BasicWorkflowStreamWorkflow.close) +@pytest.mark.asyncio +async def test_activity_subscription_propagates_cancellation(client: Client) -> None: + subscriber = CancellableSubscriber() + async with new_worker( + client, + CancelSubscriptionWorkflow, + activities=[subscriber.subscribe], + ) as worker: + handle = await client.start_workflow( + CancelSubscriptionWorkflow.run, + id=f"workflow-stream-activity-cancel-{uuid.uuid4()}", + task_queue=worker.task_queue, + ) + async with _async_timeout(5): + await subscriber.started.wait() + await handle.signal(CancelSubscriptionWorkflow.cancel_subscription) + assert await handle.result() == "CancelledError" + + @pytest.mark.asyncio async def test_context_manager_flushes_on_exit(client: Client) -> None: """Context manager exit flushes all buffered items.""" From e94f9e0ae8c699beaeeffdffa83d9a0aa6f2b1f2 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Wed, 2 Sep 2026 02:32:58 -0500 Subject: [PATCH 2/5] Prove cancellation at the subscription boundary Signed-off-by: 1fanwang <1fannnw@gmail.com> --- tests/contrib/workflow_streams/test_workflow_streams.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/tests/contrib/workflow_streams/test_workflow_streams.py b/tests/contrib/workflow_streams/test_workflow_streams.py index d48890bb9..919bce6bd 100644 --- a/tests/contrib/workflow_streams/test_workflow_streams.py +++ b/tests/contrib/workflow_streams/test_workflow_streams.py @@ -420,8 +420,11 @@ async def heartbeat() -> None: heartbeat_task = asyncio.create_task(heartbeat()) try: stream = WorkflowStreamClient.from_within_activity() - async for _ in stream.subscribe(result_type=bytes): - self.started.set() + try: + async for _ in stream.subscribe(result_type=bytes): + self.started.set() + except asyncio.CancelledError: + return "subscription-cancelled" return "subscription-ended" finally: heartbeat_task.cancel() @@ -1205,7 +1208,7 @@ async def test_activity_subscription_propagates_cancellation(client: Client) -> async with _async_timeout(5): await subscriber.started.wait() await handle.signal(CancelSubscriptionWorkflow.cancel_subscription) - assert await handle.result() == "CancelledError" + assert await handle.result() == "subscription-cancelled" @pytest.mark.asyncio From dfd6b0a9523bf24a2d47702b42ead3866a879b44 Mon Sep 17 00:00:00 2001 From: Brian Strauch Date: Wed, 2 Sep 2026 13:10:23 -0700 Subject: [PATCH 3/5] Trigger CLA check From 3f7a5b47e5916d23f7d986ea6cd1e4127dae281e Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 1 Oct 2026 14:18:22 -0700 Subject: [PATCH 4/5] Move changelog entry under Unreleased Signed-off-by: 1fanwang <1fannnw@gmail.com> --- CHANGELOG.md | 40 ++++++++++++++++++++++++++++++++++++++-- 1 file changed, 38 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3d77dac9c..875bbef78 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,44 @@ This file contains assembled releases only. # Changelog +## [Unreleased] + +### Added + +### Changed + +- Payload converters exposed by data converters and workflow/activity accessors + retain transfer type conversion, so direct use behaves consistently with SDK + serialization. + +### Deprecated + +### :boom: Breaking Changes + +- The OpenAI Agents integration has moved to the independently versioned + [`temporalio-openai-agents`](https://pypi.org/project/temporalio-openai-agents/) + package. The existing `temporalio[openai-agents]` extra now installs that + package, and the old public `temporalio.contrib.openai_agents` imports + remain available at runtime and retain their static type information. + New code should depend on `temporalio-openai-agents` directly and import + `temporalio.openai_agents`. + +### Fixed + +- `WorkflowStreamClient.subscribe` now propagates task cancellation instead of + ending the subscription normally. +- Encoding a datetime search attribute without a timezone now raises + `ValueError("Timezone must be present on all search attribute dates")` on + the typed path, matching the deprecated untyped encoder, instead of sending + a naive ISO string that the server rejects with `BadSearchAttributes`. + +- `temporalio.contrib.opentelemetry`: `TracingInterceptor` and `OpenTelemetryInterceptor` no longer + log `Failed to detach context` when a context is torn down on a different thread while + OpenTelemetry's threading instrumentation (enabled by strands, among others) is active; a + context is now detached exactly when its token is still valid in the current + `contextvars.Context`, which it stays when a workflow resumes on another pool thread. +### Security + ## [1.34.0] - 2026-09-30 ### Added @@ -150,8 +188,6 @@ This file contains assembled releases only. - Cancelling an activity from a signal while the workflow itself is cancelled no longer causes a nondeterminism error from duplicate activity-cancellation commands. -- `WorkflowStreamClient.subscribe` now propagates task cancellation instead of - ending the subscription normally. - `StrandsPlugin` now disables Botocore retries for its default Bedrock model so model request retries are handled exclusively by Temporal. - `temporalio.contrib.openai_agents` now honors the `retry-after-ms` and From f5249d8eff32b283ae69795085076622b55fe0d7 Mon Sep 17 00:00:00 2001 From: Tim Conley Date: Mon, 5 Oct 2026 15:54:16 -0700 Subject: [PATCH 5/5] Move release note to a changelog fragment --- CHANGELOG.md | 38 -------------------------- changelog/fixed/somersaulting-sloth.md | 1 + 2 files changed, 1 insertion(+), 38 deletions(-) create mode 100644 changelog/fixed/somersaulting-sloth.md diff --git a/CHANGELOG.md b/CHANGELOG.md index 875bbef78..3a7ef493d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,44 +7,6 @@ This file contains assembled releases only. # Changelog -## [Unreleased] - -### Added - -### Changed - -- Payload converters exposed by data converters and workflow/activity accessors - retain transfer type conversion, so direct use behaves consistently with SDK - serialization. - -### Deprecated - -### :boom: Breaking Changes - -- The OpenAI Agents integration has moved to the independently versioned - [`temporalio-openai-agents`](https://pypi.org/project/temporalio-openai-agents/) - package. The existing `temporalio[openai-agents]` extra now installs that - package, and the old public `temporalio.contrib.openai_agents` imports - remain available at runtime and retain their static type information. - New code should depend on `temporalio-openai-agents` directly and import - `temporalio.openai_agents`. - -### Fixed - -- `WorkflowStreamClient.subscribe` now propagates task cancellation instead of - ending the subscription normally. -- Encoding a datetime search attribute without a timezone now raises - `ValueError("Timezone must be present on all search attribute dates")` on - the typed path, matching the deprecated untyped encoder, instead of sending - a naive ISO string that the server rejects with `BadSearchAttributes`. - -- `temporalio.contrib.opentelemetry`: `TracingInterceptor` and `OpenTelemetryInterceptor` no longer - log `Failed to detach context` when a context is torn down on a different thread while - OpenTelemetry's threading instrumentation (enabled by strands, among others) is active; a - context is now detached exactly when its token is still valid in the current - `contextvars.Context`, which it stays when a workflow resumes on another pool thread. -### Security - ## [1.34.0] - 2026-09-30 ### Added diff --git a/changelog/fixed/somersaulting-sloth.md b/changelog/fixed/somersaulting-sloth.md new file mode 100644 index 000000000..28d82fb6b --- /dev/null +++ b/changelog/fixed/somersaulting-sloth.md @@ -0,0 +1 @@ +`WorkflowStreamClient.subscribe` now propagates task cancellation instead of ending the subscription normally.