From 1ed0f4c62a075090edc2fd4c1cec279ccd47a564 Mon Sep 17 00:00:00 2001 From: Andy Staples Date: Tue, 29 Sep 2026 11:54:24 -0600 Subject: [PATCH] Fix orchestration discriminator in single-instance purge requests Set isOrchestration on sync and async purge requests, cover inherited provider APIs and compatibility alias, and document the correction. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- CHANGELOG.md | 5 + azure-functions-durable/CHANGELOG.md | 7 ++ durabletask-azuremanaged/CHANGELOG.md | 6 ++ durabletask/client.py | 4 +- .../test_client_compat.py | 48 +++++++++- .../test_dts_purge.py | 54 +++++++++++ tests/durabletask/test_purge.py | 93 +++++++++++++++++++ 7 files changed, 211 insertions(+), 6 deletions(-) create mode 100644 tests/durabletask-azuremanaged/test_dts_purge.py create mode 100644 tests/durabletask/test_purge.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 9a320ea7..7d86ade0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ## Unreleased +FIXED + +- Single-instance `purge_orchestration()` requests now explicitly target +orchestrations rather than entities in both synchronous and asynchronous clients. + ## v1.11.0 FIXED diff --git a/azure-functions-durable/CHANGELOG.md b/azure-functions-durable/CHANGELOG.md index 5944f4bb..30c4543a 100644 --- a/azure-functions-durable/CHANGELOG.md +++ b/azure-functions-durable/CHANGELOG.md @@ -7,6 +7,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## Unreleased +FIXED + +- With the corresponding core `durabletask` SDK fix, single-instance +`purge_orchestration()` requests from `DurableFunctionsClient` and +`SyncDurableFunctionsClient` now explicitly target orchestrations rather than +entities. This also applies to the `purge_instance_history()` compatibility alias. + ## v2.0.0rc2 ADDED diff --git a/durabletask-azuremanaged/CHANGELOG.md b/durabletask-azuremanaged/CHANGELOG.md index d26c18ff..5ae7434d 100644 --- a/durabletask-azuremanaged/CHANGELOG.md +++ b/durabletask-azuremanaged/CHANGELOG.md @@ -7,6 +7,12 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ## Unreleased +FIXED + +- With the corresponding core `durabletask` SDK fix, single-instance +`purge_orchestration()` requests from synchronous and asynchronous Azure Managed +clients now explicitly target orchestrations rather than entities. + ## v1.11.0 ADDED diff --git a/durabletask/client.py b/durabletask/client.py index da45e9d6..4bec371a 100644 --- a/durabletask/client.py +++ b/durabletask/client.py @@ -857,7 +857,7 @@ def restart_orchestration(self, instance_id: str, *, return res.instanceId def purge_orchestration(self, instance_id: str, recursive: bool = True) -> PurgeInstancesResult: - req = pb.PurgeInstancesRequest(instanceId=instance_id, recursive=recursive) + req = pb.PurgeInstancesRequest(instanceId=instance_id, recursive=recursive, isOrchestration=True) self._logger.info(f"Purging instance '{instance_id}'.") resp: pb.PurgeInstancesResponse = self._stub.PurgeInstances(req) return new_purge_instances_result(resp) @@ -1402,7 +1402,7 @@ async def restart_orchestration(self, instance_id: str, *, return res.instanceId async def purge_orchestration(self, instance_id: str, recursive: bool = True) -> PurgeInstancesResult: - req = pb.PurgeInstancesRequest(instanceId=instance_id, recursive=recursive) + req = pb.PurgeInstancesRequest(instanceId=instance_id, recursive=recursive, isOrchestration=True) self._logger.info(f"Purging instance '{instance_id}'.") resp: pb.PurgeInstancesResponse = await self._get_stub().PurgeInstances(req) return new_purge_instances_result(resp) diff --git a/tests/azure-functions-durable/test_client_compat.py b/tests/azure-functions-durable/test_client_compat.py index 8ed5ca20..8cf96308 100644 --- a/tests/azure-functions-durable/test_client_compat.py +++ b/tests/azure-functions-durable/test_client_compat.py @@ -20,7 +20,7 @@ replace_url_origin, ) from durabletask import history as dt_history, task as dt_task -from durabletask.client import AsyncTaskHubGrpcClient, OrchestrationStatus +from durabletask.client import AsyncTaskHubGrpcClient, OrchestrationStatus, PurgeInstancesResult from durabletask.entities import EntityInstanceId from durabletask.internal import orchestrator_service_pb2 as pb from durabletask.task import RetryPolicy @@ -1288,14 +1288,54 @@ async def test_get_status_by_returns_wrapped_list(): # Return-type shims: PurgeHistoryResult # --------------------------------------------------------------------------- +@pytest.mark.parametrize("recursive", [None, True, False]) +def test_sync_purge_orchestration_request(recursive: bool | None) -> None: + client = df.SyncDurableFunctionsClient(_CLIENT_CONFIG) + stub = Mock() + stub.PurgeInstances.return_value = pb.PurgeInstancesResponse(deletedInstanceCount=3) + try: + with patch.object(client, "_stub", stub): + if recursive is None: + result = client.purge_orchestration("abc") + else: + result = client.purge_orchestration("abc", recursive=recursive) + stub.PurgeInstances.assert_called_once_with(pb.PurgeInstancesRequest( + instanceId="abc", recursive=True if recursive is None else recursive, + isOrchestration=True)) + assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=None) + finally: + client.close() + + +@pytest.mark.parametrize("recursive", [None, True, False]) +async def test_async_purge_orchestration_request(recursive: bool | None) -> None: + client = _make_client() + stub = Mock() + stub.PurgeInstances = AsyncMock(return_value=pb.PurgeInstancesResponse(deletedInstanceCount=3)) + try: + with patch.object(client, "_get_stub", return_value=stub): + if recursive is None: + result = await client.purge_orchestration("abc") + else: + result = await client.purge_orchestration("abc", recursive=recursive) + stub.PurgeInstances.assert_awaited_once_with(pb.PurgeInstancesRequest( + instanceId="abc", recursive=True if recursive is None else recursive, + isOrchestration=True)) + assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=None) + finally: + await client.close() + + async def test_purge_instance_history_returns_purge_history_result(): client = _make_client() + stub = Mock() + stub.PurgeInstances = AsyncMock(return_value=pb.PurgeInstancesResponse(deletedInstanceCount=3)) try: - result = SimpleNamespace(deleted_instance_count=3, is_complete=True) - with patch.object(client, "purge_orchestration", - new=AsyncMock(return_value=result)): + with patch.object(client, "_get_stub", return_value=stub): with pytest.warns(DeprecationWarning): purge = await client.purge_instance_history("abc") + stub.PurgeInstances.assert_awaited_once_with(pb.PurgeInstancesRequest( + instanceId="abc", recursive=True, isOrchestration=True)) assert purge.instances_deleted == 3 finally: await client.close() diff --git a/tests/durabletask-azuremanaged/test_dts_purge.py b/tests/durabletask-azuremanaged/test_dts_purge.py new file mode 100644 index 00000000..a5747e69 --- /dev/null +++ b/tests/durabletask-azuremanaged/test_dts_purge.py @@ -0,0 +1,54 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +from durabletask.azuremanaged.client import ( + AsyncDurableTaskSchedulerClient, + DurableTaskSchedulerClient, +) +from durabletask.client import PurgeInstancesResult +from durabletask.internal import orchestrator_service_pb2 as pb + + +@pytest.mark.parametrize("recursive", [None, True, False]) +def test_dts_purge_orchestration_request(recursive: bool | None) -> None: + stub = MagicMock() + stub.PurgeInstances.return_value = pb.PurgeInstancesResponse(deletedInstanceCount=3) + + with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub): + with DurableTaskSchedulerClient( + host_address="localhost:4001", taskhub="hub", + token_credential=None, channel=MagicMock()) as client: + if recursive is None: + result = client.purge_orchestration("instance") + else: + result = client.purge_orchestration("instance", recursive=recursive) + + stub.PurgeInstances.assert_called_once_with(pb.PurgeInstancesRequest( + instanceId="instance", recursive=True if recursive is None else recursive, + isOrchestration=True)) + assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=None) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("recursive", [None, True, False]) +async def test_async_dts_purge_orchestration_request(recursive: bool | None) -> None: + stub = MagicMock() + stub.PurgeInstances = AsyncMock(return_value=pb.PurgeInstancesResponse(deletedInstanceCount=3)) + + with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub): + async with AsyncDurableTaskSchedulerClient( + host_address="localhost:4001", taskhub="hub", + token_credential=None, channel=MagicMock()) as client: + if recursive is None: + result = await client.purge_orchestration("instance") + else: + result = await client.purge_orchestration("instance", recursive=recursive) + + stub.PurgeInstances.assert_awaited_once_with(pb.PurgeInstancesRequest( + instanceId="instance", recursive=True if recursive is None else recursive, + isOrchestration=True)) + assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=None) diff --git a/tests/durabletask/test_purge.py b/tests/durabletask/test_purge.py new file mode 100644 index 00000000..21ff5b73 --- /dev/null +++ b/tests/durabletask/test_purge.py @@ -0,0 +1,93 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +from unittest.mock import AsyncMock, MagicMock, patch + +import grpc +import pytest +from google.protobuf import wrappers_pb2 + +import durabletask.internal.orchestrator_service_pb2 as pb +from durabletask.client import AsyncTaskHubGrpcClient, PurgeInstancesResult, TaskHubGrpcClient + + +@pytest.mark.parametrize("recursive", [None, True, False]) +@pytest.mark.parametrize("is_complete", [None, True, False]) +def test_sync_purge_orchestration_request_and_result( + recursive: bool | None, is_complete: bool | None) -> None: + response = pb.PurgeInstancesResponse(deletedInstanceCount=3) + if is_complete is not None: + response.isComplete.CopyFrom(wrappers_pb2.BoolValue(value=is_complete)) + stub = MagicMock() + stub.PurgeInstances.return_value = response + + with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub): + with TaskHubGrpcClient(channel=MagicMock()) as client: + if recursive is None: + result = client.purge_orchestration("instance") + else: + result = client.purge_orchestration("instance", recursive=recursive) + + stub.PurgeInstances.assert_called_once() + request = stub.PurgeInstances.call_args.args[0] + assert request.instanceId == "instance" + assert request.recursive is (True if recursive is None else recursive) + assert request.isOrchestration is True + assert not request.HasField("purgeInstanceFilter") + assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=is_complete) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("recursive", [None, True, False]) +@pytest.mark.parametrize("is_complete", [None, True, False]) +async def test_async_purge_orchestration_request_and_result( + recursive: bool | None, is_complete: bool | None) -> None: + response = pb.PurgeInstancesResponse(deletedInstanceCount=3) + if is_complete is not None: + response.isComplete.CopyFrom(wrappers_pb2.BoolValue(value=is_complete)) + stub = MagicMock() + stub.PurgeInstances = AsyncMock(return_value=response) + + with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub): + async with AsyncTaskHubGrpcClient(channel=MagicMock()) as client: + if recursive is None: + result = await client.purge_orchestration("instance") + else: + result = await client.purge_orchestration("instance", recursive=recursive) + + stub.PurgeInstances.assert_awaited_once() + request = stub.PurgeInstances.call_args.args[0] + assert request.instanceId == "instance" + assert request.recursive is (True if recursive is None else recursive) + assert request.isOrchestration is True + assert not request.HasField("purgeInstanceFilter") + assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=is_complete) + + +def test_sync_purge_orchestration_propagates_rpc_error() -> None: + error = grpc.RpcError("purge failed") + stub = MagicMock() + stub.PurgeInstances.side_effect = error + + with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub): + with TaskHubGrpcClient(channel=MagicMock()) as client: + with pytest.raises(grpc.RpcError) as raised: + client.purge_orchestration("instance") + + assert raised.value is error + stub.PurgeInstances.assert_called_once() + + +@pytest.mark.asyncio +async def test_async_purge_orchestration_propagates_rpc_error() -> None: + error = grpc.RpcError("purge failed") + stub = MagicMock() + stub.PurgeInstances = AsyncMock(side_effect=error) + + with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub): + async with AsyncTaskHubGrpcClient(channel=MagicMock()) as client: + with pytest.raises(grpc.RpcError) as raised: + await client.purge_orchestration("instance") + + assert raised.value is error + stub.PurgeInstances.assert_awaited_once()