Skip to content

Commit 5cdfeef

Browse files
andystaplesCopilot
andauthored
Fix orchestration discriminator in single-instance purge requests (#279)
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>
1 parent 866cc78 commit 5cdfeef

7 files changed

Lines changed: 202 additions & 6 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
99

1010
FIXED
1111

12+
- Single-instance `purge_orchestration()` requests now explicitly target
13+
orchestrations rather than entities in both synchronous and asynchronous clients.
1214
- Timer callbacks no longer schedule additional long-timer chunks, retry
1315
activities or sub-orchestrations, or resume orchestrator code after completion,
1416
failure, termination, or continue-as-new. Work scheduled before the terminal

‎azure-functions-durable/CHANGELOG.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
99

1010
FIXED
1111

12+
- With the corresponding core `durabletask` SDK fix, single-instance
13+
`purge_orchestration()` requests from `DurableFunctionsClient` and
14+
`SyncDurableFunctionsClient` now explicitly target orchestrations rather than
15+
entities. This also applies to the `purge_instance_history()` compatibility alias.
1216
- With the corresponding core `durabletask` SDK fix, timer callbacks no longer
1317
schedule additional long-timer chunks, retry activities or sub-orchestrations,
1418
or resume orchestrator code after completion, failure, termination, or

‎durabletask-azuremanaged/CHANGELOG.md‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,9 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
99

1010
FIXED
1111

12+
- With the corresponding core `durabletask` SDK fix, single-instance
13+
`purge_orchestration()` requests from synchronous and asynchronous Azure Managed
14+
clients now explicitly target orchestrations rather than entities.
1215
- With the corresponding core `durabletask` SDK fix, timer callbacks no longer
1316
retry activities or sub-orchestrations, or resume orchestrator code after
1417
completion, failure, termination, or continue-as-new. Azure Managed uses native

‎durabletask/client.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -857,7 +857,7 @@ def restart_orchestration(self, instance_id: str, *,
857857
return res.instanceId
858858

859859
def purge_orchestration(self, instance_id: str, recursive: bool = True) -> PurgeInstancesResult:
860-
req = pb.PurgeInstancesRequest(instanceId=instance_id, recursive=recursive)
860+
req = pb.PurgeInstancesRequest(instanceId=instance_id, recursive=recursive, isOrchestration=True)
861861
self._logger.info(f"Purging instance '{instance_id}'.")
862862
resp: pb.PurgeInstancesResponse = self._stub.PurgeInstances(req)
863863
return new_purge_instances_result(resp)
@@ -1402,7 +1402,7 @@ async def restart_orchestration(self, instance_id: str, *,
14021402
return res.instanceId
14031403

14041404
async def purge_orchestration(self, instance_id: str, recursive: bool = True) -> PurgeInstancesResult:
1405-
req = pb.PurgeInstancesRequest(instanceId=instance_id, recursive=recursive)
1405+
req = pb.PurgeInstancesRequest(instanceId=instance_id, recursive=recursive, isOrchestration=True)
14061406
self._logger.info(f"Purging instance '{instance_id}'.")
14071407
resp: pb.PurgeInstancesResponse = await self._get_stub().PurgeInstances(req)
14081408
return new_purge_instances_result(resp)

‎tests/azure-functions-durable/test_client_compat.py‎

Lines changed: 44 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@
2020
replace_url_origin,
2121
)
2222
from durabletask import history as dt_history, task as dt_task
23-
from durabletask.client import AsyncTaskHubGrpcClient, OrchestrationStatus
23+
from durabletask.client import AsyncTaskHubGrpcClient, OrchestrationStatus, PurgeInstancesResult
2424
from durabletask.entities import EntityInstanceId
2525
from durabletask.internal import orchestrator_service_pb2 as pb
2626
from durabletask.task import RetryPolicy
@@ -1288,14 +1288,54 @@ async def test_get_status_by_returns_wrapped_list():
12881288
# Return-type shims: PurgeHistoryResult
12891289
# ---------------------------------------------------------------------------
12901290

1291+
@pytest.mark.parametrize("recursive", [None, True, False])
1292+
def test_sync_purge_orchestration_request(recursive: bool | None) -> None:
1293+
client = df.SyncDurableFunctionsClient(_CLIENT_CONFIG)
1294+
stub = Mock()
1295+
stub.PurgeInstances.return_value = pb.PurgeInstancesResponse(deletedInstanceCount=3)
1296+
try:
1297+
with patch.object(client, "_stub", stub):
1298+
if recursive is None:
1299+
result = client.purge_orchestration("abc")
1300+
else:
1301+
result = client.purge_orchestration("abc", recursive=recursive)
1302+
stub.PurgeInstances.assert_called_once_with(pb.PurgeInstancesRequest(
1303+
instanceId="abc", recursive=True if recursive is None else recursive,
1304+
isOrchestration=True))
1305+
assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=None)
1306+
finally:
1307+
client.close()
1308+
1309+
1310+
@pytest.mark.parametrize("recursive", [None, True, False])
1311+
async def test_async_purge_orchestration_request(recursive: bool | None) -> None:
1312+
client = _make_client()
1313+
stub = Mock()
1314+
stub.PurgeInstances = AsyncMock(return_value=pb.PurgeInstancesResponse(deletedInstanceCount=3))
1315+
try:
1316+
with patch.object(client, "_get_stub", return_value=stub):
1317+
if recursive is None:
1318+
result = await client.purge_orchestration("abc")
1319+
else:
1320+
result = await client.purge_orchestration("abc", recursive=recursive)
1321+
stub.PurgeInstances.assert_awaited_once_with(pb.PurgeInstancesRequest(
1322+
instanceId="abc", recursive=True if recursive is None else recursive,
1323+
isOrchestration=True))
1324+
assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=None)
1325+
finally:
1326+
await client.close()
1327+
1328+
12911329
async def test_purge_instance_history_returns_purge_history_result():
12921330
client = _make_client()
1331+
stub = Mock()
1332+
stub.PurgeInstances = AsyncMock(return_value=pb.PurgeInstancesResponse(deletedInstanceCount=3))
12931333
try:
1294-
result = SimpleNamespace(deleted_instance_count=3, is_complete=True)
1295-
with patch.object(client, "purge_orchestration",
1296-
new=AsyncMock(return_value=result)):
1334+
with patch.object(client, "_get_stub", return_value=stub):
12971335
with pytest.warns(DeprecationWarning):
12981336
purge = await client.purge_instance_history("abc")
1337+
stub.PurgeInstances.assert_awaited_once_with(pb.PurgeInstancesRequest(
1338+
instanceId="abc", recursive=True, isOrchestration=True))
12991339
assert purge.instances_deleted == 3
13001340
finally:
13011341
await client.close()
Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
# Copyright (c) Microsoft Corporation.
2+
# Licensed under the MIT License.
3+
4+
from unittest.mock import AsyncMock, MagicMock, patch
5+
6+
import pytest
7+
8+
from durabletask.azuremanaged.client import (
9+
AsyncDurableTaskSchedulerClient,
10+
DurableTaskSchedulerClient,
11+
)
12+
from durabletask.client import PurgeInstancesResult
13+
from durabletask.internal import orchestrator_service_pb2 as pb
14+
15+
16+
@pytest.mark.parametrize("recursive", [None, True, False])
17+
def test_dts_purge_orchestration_request(recursive: bool | None) -> None:
18+
stub = MagicMock()
19+
stub.PurgeInstances.return_value = pb.PurgeInstancesResponse(deletedInstanceCount=3)
20+
21+
with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub):
22+
with DurableTaskSchedulerClient(
23+
host_address="localhost:4001", taskhub="hub",
24+
token_credential=None, channel=MagicMock()) as client:
25+
if recursive is None:
26+
result = client.purge_orchestration("instance")
27+
else:
28+
result = client.purge_orchestration("instance", recursive=recursive)
29+
30+
stub.PurgeInstances.assert_called_once_with(pb.PurgeInstancesRequest(
31+
instanceId="instance", recursive=True if recursive is None else recursive,
32+
isOrchestration=True))
33+
assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=None)
34+
35+
36+
@pytest.mark.asyncio
37+
@pytest.mark.parametrize("recursive", [None, True, False])
38+
async def test_async_dts_purge_orchestration_request(recursive: bool | None) -> None:
39+
stub = MagicMock()
40+
stub.PurgeInstances = AsyncMock(return_value=pb.PurgeInstancesResponse(deletedInstanceCount=3))
41+
42+
with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub):
43+
async with AsyncDurableTaskSchedulerClient(
44+
host_address="localhost:4001", taskhub="hub",
45+
token_credential=None, channel=MagicMock()) as client:
46+
if recursive is None:
47+
result = await client.purge_orchestration("instance")
48+
else:
49+
result = await client.purge_orchestration("instance", recursive=recursive)
50+
51+
stub.PurgeInstances.assert_awaited_once_with(pb.PurgeInstancesRequest(
52+
instanceId="instance", recursive=True if recursive is None else recursive,
53+
isOrchestration=True))
54+
assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=None)

‎tests/durabletask/test_purge.py‎

Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
1+
# Copyright (c) Microsoft Corporation.
2+
# Licensed under the MIT License.
3+
4+
from unittest.mock import AsyncMock, MagicMock, patch
5+
6+
import grpc
7+
import pytest
8+
from google.protobuf import wrappers_pb2
9+
10+
import durabletask.internal.orchestrator_service_pb2 as pb
11+
from durabletask.client import AsyncTaskHubGrpcClient, PurgeInstancesResult, TaskHubGrpcClient
12+
13+
14+
@pytest.mark.parametrize("recursive", [None, True, False])
15+
@pytest.mark.parametrize("is_complete", [None, True, False])
16+
def test_sync_purge_orchestration_request_and_result(
17+
recursive: bool | None, is_complete: bool | None) -> None:
18+
response = pb.PurgeInstancesResponse(deletedInstanceCount=3)
19+
if is_complete is not None:
20+
response.isComplete.CopyFrom(wrappers_pb2.BoolValue(value=is_complete))
21+
stub = MagicMock()
22+
stub.PurgeInstances.return_value = response
23+
24+
with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub):
25+
with TaskHubGrpcClient(channel=MagicMock()) as client:
26+
if recursive is None:
27+
result = client.purge_orchestration("instance")
28+
else:
29+
result = client.purge_orchestration("instance", recursive=recursive)
30+
31+
stub.PurgeInstances.assert_called_once()
32+
request = stub.PurgeInstances.call_args.args[0]
33+
assert request.instanceId == "instance"
34+
assert request.recursive is (True if recursive is None else recursive)
35+
assert request.isOrchestration is True
36+
assert not request.HasField("purgeInstanceFilter")
37+
assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=is_complete)
38+
39+
40+
@pytest.mark.asyncio
41+
@pytest.mark.parametrize("recursive", [None, True, False])
42+
@pytest.mark.parametrize("is_complete", [None, True, False])
43+
async def test_async_purge_orchestration_request_and_result(
44+
recursive: bool | None, is_complete: bool | None) -> None:
45+
response = pb.PurgeInstancesResponse(deletedInstanceCount=3)
46+
if is_complete is not None:
47+
response.isComplete.CopyFrom(wrappers_pb2.BoolValue(value=is_complete))
48+
stub = MagicMock()
49+
stub.PurgeInstances = AsyncMock(return_value=response)
50+
51+
with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub):
52+
async with AsyncTaskHubGrpcClient(channel=MagicMock()) as client:
53+
if recursive is None:
54+
result = await client.purge_orchestration("instance")
55+
else:
56+
result = await client.purge_orchestration("instance", recursive=recursive)
57+
58+
stub.PurgeInstances.assert_awaited_once()
59+
request = stub.PurgeInstances.call_args.args[0]
60+
assert request.instanceId == "instance"
61+
assert request.recursive is (True if recursive is None else recursive)
62+
assert request.isOrchestration is True
63+
assert not request.HasField("purgeInstanceFilter")
64+
assert result == PurgeInstancesResult(deleted_instance_count=3, is_complete=is_complete)
65+
66+
67+
def test_sync_purge_orchestration_propagates_rpc_error() -> None:
68+
error = grpc.RpcError("purge failed")
69+
stub = MagicMock()
70+
stub.PurgeInstances.side_effect = error
71+
72+
with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub):
73+
with TaskHubGrpcClient(channel=MagicMock()) as client:
74+
with pytest.raises(grpc.RpcError) as raised:
75+
client.purge_orchestration("instance")
76+
77+
assert raised.value is error
78+
stub.PurgeInstances.assert_called_once()
79+
80+
81+
@pytest.mark.asyncio
82+
async def test_async_purge_orchestration_propagates_rpc_error() -> None:
83+
error = grpc.RpcError("purge failed")
84+
stub = MagicMock()
85+
stub.PurgeInstances = AsyncMock(side_effect=error)
86+
87+
with patch("durabletask.client.stubs.TaskHubSidecarServiceStub", return_value=stub):
88+
async with AsyncTaskHubGrpcClient(channel=MagicMock()) as client:
89+
with pytest.raises(grpc.RpcError) as raised:
90+
await client.purge_orchestration("instance")
91+
92+
assert raised.value is error
93+
stub.PurgeInstances.assert_awaited_once()

0 commit comments

Comments
 (0)