Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

FIXED

- Single-instance `purge_orchestration()` requests now explicitly target
orchestrations rather than entities in both synchronous and asynchronous clients.
- `continue_as_new(..., save_events=True)` now preserves the global arrival
order of unconsumed buffered external events across different event names,
instead of grouping carryover events by name.
Expand Down
4 changes: 4 additions & 0 deletions azure-functions-durable/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

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.
- With the corresponding core `durabletask` SDK fix,
`continue_as_new(..., save_events=True)` preserves the global arrival order
of unconsumed buffered external events across different event names.
Expand Down
3 changes: 3 additions & 0 deletions durabletask-azuremanaged/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,9 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

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.
- With the corresponding core `durabletask` SDK fix,
`continue_as_new(..., save_events=True)` preserves the global arrival order
of unconsumed buffered external events across different event names.
Expand Down
4 changes: 2 additions & 2 deletions durabletask/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
48 changes: 44 additions & 4 deletions tests/azure-functions-durable/test_client_compat.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down
54 changes: 54 additions & 0 deletions tests/durabletask-azuremanaged/test_dts_purge.py
Original file line number Diff line number Diff line change
@@ -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)
93 changes: 93 additions & 0 deletions tests/durabletask/test_purge.py
Original file line number Diff line number Diff line change
@@ -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()
Loading