diff --git a/CHANGELOG.md b/CHANGELOG.md
index 6900d0788..60f7789c0 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -2,6 +2,7 @@
## Unreleased
+- Add `Microsoft.DurableTask.LargePayloadPurge.Abstractions` with shared fetch/report and service-setting interfaces, and expose the existing blob purge tasks for host integration ([#805](https://github.com/microsoft/durabletask-dotnet/pull/805)).
- Preserve the original orchestration version when restarting through the orchestration-service client shim ([#463](https://github.com/microsoft/durabletask-dotnet/issues/463)).
- Support configurable scheduler token audiences and government defaults ([#806](https://github.com/microsoft/durabletask-dotnet/pull/806))
diff --git a/Microsoft.DurableTask.sln b/Microsoft.DurableTask.sln
index 380ee0a0e..40a4716b1 100644
--- a/Microsoft.DurableTask.sln
+++ b/Microsoft.DurableTask.sln
@@ -149,6 +149,10 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Extensions", "Extensions",
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "AzureBlobPayloads.Tests", "test\Extensions\AzureBlobPayloads.Tests\AzureBlobPayloads.Tests.csproj", "{3E509481-3CCC-4006-BCB2-9E8FA7C275F1}"
EndProject
+Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "LargePayloadPurge.Abstractions", "src\LargePayloadPurge.Abstractions\LargePayloadPurge.Abstractions.csproj", "{CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}"
+EndProject
+Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "LargePayloadPurge.Abstractions.Tests", "test\LargePayloadPurge.Abstractions.Tests\LargePayloadPurge.Abstractions.Tests.csproj", "{BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}"
+EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@@ -855,6 +859,30 @@ Global
{3E509481-3CCC-4006-BCB2-9E8FA7C275F1}.Release|x64.Build.0 = Release|Any CPU
{3E509481-3CCC-4006-BCB2-9E8FA7C275F1}.Release|x86.ActiveCfg = Release|Any CPU
{3E509481-3CCC-4006-BCB2-9E8FA7C275F1}.Release|x86.Build.0 = Release|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Debug|Any CPU.Build.0 = Debug|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Debug|x64.ActiveCfg = Debug|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Debug|x64.Build.0 = Debug|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Debug|x86.ActiveCfg = Debug|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Debug|x86.Build.0 = Debug|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Release|Any CPU.ActiveCfg = Release|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Release|Any CPU.Build.0 = Release|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Release|x64.ActiveCfg = Release|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Release|x64.Build.0 = Release|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Release|x86.ActiveCfg = Release|Any CPU
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28}.Release|x86.Build.0 = Release|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Debug|Any CPU.Build.0 = Debug|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Debug|x64.ActiveCfg = Debug|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Debug|x64.Build.0 = Debug|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Debug|x86.ActiveCfg = Debug|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Debug|x86.Build.0 = Debug|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Release|Any CPU.ActiveCfg = Release|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Release|Any CPU.Build.0 = Release|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Release|x64.ActiveCfg = Release|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Release|x64.Build.0 = Release|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Release|x86.ActiveCfg = Release|Any CPU
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE}.Release|x86.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@@ -929,6 +957,8 @@ Global
{3B8F957E-7773-4C0C-ACD7-91A1591D9312} = {5B448FF6-EC42-491D-A22E-1DC8B618E6D5}
{00205C88-F000-28F2-A910-C6FA00E065EE} = {E5637F81-2FB9-4CD7-900D-455363B142A7}
{3E509481-3CCC-4006-BCB2-9E8FA7C275F1} = {00205C88-F000-28F2-A910-C6FA00E065EE}
+ {CF0A6A55-1DEC-4027-8B8B-4B0DEE3C7A28} = {8AFC9781-F6F1-4696-BB4A-9ED7CA9D612B}
+ {BE1D641F-1F96-483E-A1A9-0F8A64E9BAFE} = {E5637F81-2FB9-4CD7-900D-455363B142A7}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {AB41CB55-35EA-4986-A522-387AB3402E71}
diff --git a/README.md b/README.md
index df785410f..75b22c694 100644
--- a/README.md
+++ b/README.md
@@ -200,6 +200,56 @@ For runnable DTS emulator examples that demonstrate versioning, see the [WorkerV
The [on-demand sandbox activities sample](samples/on-demand-sandbox/README.md) shows how to declare selected activities for Durable Task Scheduler (DTS)-managed on-demand sandbox execution and build the remote worker container image separately from the declarer app.
+### Blob auto-purge infrastructure integration
+
+The optional service contract is maintained in
+[`Microsoft.DurableTask.LargePayloadPurge.Abstractions`](src/LargePayloadPurge.Abstractions/README.md).
+It defines `DurableTask.LargePayloadPurge.IOrchestrationServiceLargePayloadPurgeClient` and
+`Microsoft.DurableTask.AzureBlobPayloads.ILargePayloadPurgeClient`, using the canonical SDK Client models
+without duplicating them. This package is versioned independently and is not BCL-only: its Client dependency
+transitively depends on SDK Abstractions and Durable Task Core. The Azure Blob implementation depends on these
+contracts; the contracts do not depend on Blob storage, gRPC or worker implementations. The service capability
+inherits the shared fetch/report interface and adds only the setting operation.
+
+`Microsoft.DurableTask.Extensions.AzureBlobPayloads` exposes reusable orchestration and activity
+implementations: `BlobPurgeJobOrchestrator`, `GetLargePayloadTombstonesActivity`, `DeleteExternalBlobActivity`,
+and `ReportLargePayloadPurgeResultsActivity`. Preserve their exact task names, empty version, input/output
+types, and retry/event/continue-as-new behavior. Keep the tasks registered even when auto-purge is disabled
+so existing work can finish. The existing client API manages the reserved per-task-hub orchestration instance
+and its configuration.
+
+The companion **.NET isolated Durable Functions** integration uses the optional
+`Microsoft.Azure.Functions.Worker.Extensions.DurableTask.AzureBlobPayloads` package. It supplies four ordinary
+`[Function]` methods that delegate to the shared tasks, plus worker-side payload-store configuration.
+The base Functions worker extension does not carry these function definitions. Referencing only this shared
+SDK package does not register Functions or enable auto-purge.
+
+The Functions Worker SDK discovers the compiled methods during the normal build and generates their metadata
+and invocation paths. The functions use ordinary trigger and `DurableClient` bindings in the isolated worker,
+passing the bound `TaskOrchestrationContext` and a `TaskActivityContext` with the invoking orchestration's
+instance ID to the shared tasks. Normal serialization and failure propagation preserve structured failure
+details and unprocessed events across continue-as-new.
+
+Construct the fetch/report activities with an `ILargePayloadPurgeClient` and their typed loggers, and the
+delete activity with the worker's configured `PayloadStore` and logger. Reuse `BlobPayloadStore` with
+`LargePayloadStorageOptions` for storage access rather than copying its ownership checks or deletion policy.
+Deletion runs in the language worker; purge RPCs do not carry storage credentials. The narrow purge client
+must honor UTC deadlines and preserve opaque tombstone tokens; the activities classify fetch/report gRPC failures.
+
+The Functions integration routes setting, fetch, and report operations through the bound client's local
+host endpoint to the provider's authenticated transport for the **same task hub**. The existing
+`LargePayloadPurge` gRPC service is separate from `TaskHubSidecarService`. The Functions client wrapper can
+forward the setting through the existing infrastructure `ILargePayloadAutoPurgeClient` interface so the
+original `client.SetLargePayloadAutoPurgeAsync(enabled, batchSize, cancellationToken)` extension retains
+ownership of bootstrap behavior. The integration owns client/store lifetimes, hub binding, and reconnection;
+this SDK surface alone does not supply Functions metadata or the local-host bridge.
+
+Enabling explicitly writes the setting, starts the reserved instance with live-status deduplication and an
+empty version, verifies its identity and Running status, then sends `SetBatchSize`. Disabling **only** writes
+the setting and ignores batch size. Calling neither leaves the setting untouched. These steps are not
+transactional; failures propagate without rollback. Existing standalone gRPC client and worker behavior
+is unchanged.
+
### Token audiences and Azure Government
`DurableTaskSchedulerClientOptions.ResourceId` and `DurableTaskSchedulerWorkerOptions.ResourceId`
diff --git a/doc/release_process.md b/doc/release_process.md
index 319a2cd49..116d71f9b 100644
--- a/doc/release_process.md
+++ b/doc/release_process.md
@@ -10,6 +10,23 @@ This repo publishes multiple NuGet packages. Most share a single version defined
We follow an approach of releasing everything together, even if a package has no changes — unless we intentionally hold a package back.
+`LargePayloadPurge.Abstractions` versions independently in its project file, starting at the planned
+`0.1.0` release. It uses source references to the SDK Client models. Publish the corresponding Client and
+Abstractions packages containing those models before publishing this package; released Client `1.26.0`
+does not contain them. The package uses the repository's MIT license and SDK strong-name key.
+Its package and assembly name is `Microsoft.DurableTask.LargePayloadPurge.Abstractions`, covered by the
+standard `Microsoft.DurableTask.*.dll` signing pattern. The existing source traversal, SBOM inclusion,
+NuGet signing and per-package approval-gated publication steps apply.
+
+Contract publication also waits for successful Client and Abstractions publication. If either prerequisite
+fails, including a duplicate-version upload failure, or is skipped or canceled, contract publication is
+skipped rather than treating that result as success. Other packages retain their independent publication jobs.
+
+Keep its `RELEASENOTES.md`: `eng/targets/Release.targets` reads it into NuGet package metadata, whereas the
+root `CHANGELOG.md` retains repository release history. The shared target also appends a link using the
+package's own version (`releases/tag/v0.1.0` for its initial release). Verify or create the corresponding
+release tag before publication; independent package versioning does not create that tag automatically.
+
### Versioning Scheme
We follow [semver](https://semver.org/) with optional pre-release tags:
diff --git a/eng/publish/publish.yml b/eng/publish/publish.yml
index a30094697..5f2c9a602 100644
--- a/eng/publish/publish.yml
+++ b/eng/publish/publish.yml
@@ -463,4 +463,30 @@ extends:
nuGetFeedType: external
publishFeedCredentials: 'DurableTask org NuGet API Key'
packagesToPush: '$(System.DefaultWorkingDirectory)/drop/Microsoft.DurableTask.Extensions.AzureBlobPayloads.*.nupkg;!$(System.DefaultWorkingDirectory)/**/*.symbols.nupkg' # Despite this being a custom command, we need to keep this for 1ES validation
- packageParentPath: $(System.DefaultWorkingDirectory) # This needs to be set to some prefix of the `packagesToPush` parameter. Apparently it helps with SDL tooling
\ No newline at end of file
+ packageParentPath: $(System.DefaultWorkingDirectory) # This needs to be set to some prefix of the `packagesToPush` parameter. Apparently it helps with SDL tooling
+
+ # NuGet release (Microsoft.DurableTask.LargePayloadPurge.Abstractions)
+ - job: nugetRelease_Microsoft_DurableTask_LargePayloadPurge_Abstractions
+ displayName: NuGet Release (Microsoft.DurableTask.LargePayloadPurge.Abstractions)
+ dependsOn:
+ - nugetApproval
+ - nugetRelease_Microsoft_DurableTask_Abstractions
+ - nugetRelease_Microsoft_DurableTask_Client
+ condition: succeeded()
+ templateContext:
+ type: releaseJob
+ isProduction: true
+ inputs:
+ - input: pipelineArtifact
+ pipeline: officialPipeline
+ artifactName: drop
+ targetPath: $(System.DefaultWorkingDirectory)/drop
+ steps:
+ - task: 1ES.PublishNuget@1
+ displayName: 'NuGet push (Microsoft.DurableTask.LargePayloadPurge.Abstractions)'
+ inputs:
+ command: push
+ nuGetFeedType: external
+ publishFeedCredentials: 'DurableTask org NuGet API Key'
+ packagesToPush: '$(System.DefaultWorkingDirectory)/drop/Microsoft.DurableTask.LargePayloadPurge.Abstractions.*.nupkg;!$(System.DefaultWorkingDirectory)/**/*.symbols.nupkg'
+ packageParentPath: $(System.DefaultWorkingDirectory)
\ No newline at end of file
diff --git a/list-nuget-packages-links.ps1 b/list-nuget-packages-links.ps1
index 2697869e3..bd6187bdf 100644
--- a/list-nuget-packages-links.ps1
+++ b/list-nuget-packages-links.ps1
@@ -72,6 +72,7 @@ $packages = @(
"Microsoft.DurableTask.Worker.Grpc",
"Microsoft.DurableTask.Client.OrchestrationServiceClientShim",
"Microsoft.DurableTask.Extensions.AzureBlobPayloads",
+ "Microsoft.DurableTask.LargePayloadPurge.Abstractions",
"Microsoft.DurableTask.Client.AzureManaged",
"Microsoft.DurableTask.Worker.AzureManaged",
"Microsoft.DurableTask.ScheduledTasks",
diff --git a/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/GetLargePayloadTombstonesActivity.cs b/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/GetLargePayloadTombstonesActivity.cs
index 2c3147215..90f3e74cc 100644
--- a/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/GetLargePayloadTombstonesActivity.cs
+++ b/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/GetLargePayloadTombstonesActivity.cs
@@ -4,8 +4,6 @@
using Grpc.Core;
using Microsoft.DurableTask.Client;
using Microsoft.Extensions.Logging;
-using static Microsoft.DurableTask.Protobuf.LargePayloads.LargePayloadPurge;
-using LP = Microsoft.DurableTask.Protobuf.LargePayloads;
namespace Microsoft.DurableTask.AzureBlobPayloads;
@@ -15,13 +13,17 @@ namespace Microsoft.DurableTask.AzureBlobPayloads;
///
/// The large-payload purge service client used to query the backend for tombstones.
/// The logger instance.
+///
+/// Infrastructure integration API for alternate .NET hosts. The supplied client must be bound to this
+/// worker's authenticated task hub. Its transport lifetime remains owned by the host.
+///
[DurableTask]
-internal sealed class GetLargePayloadTombstonesActivity(
- LargePayloadPurgeClient client,
+public sealed class GetLargePayloadTombstonesActivity(
+ ILargePayloadPurgeClient client,
ILogger logger)
: TaskActivity>
{
- readonly LargePayloadPurgeClient client = Check.NotNull(client);
+ readonly ILargePayloadPurgeClient client = Check.NotNull(client);
readonly ILogger logger = Check.NotNull(logger);
///
@@ -41,13 +43,10 @@ public override async Task> RunAsync(TaskActivityCon
nameof(input), input, $"Limit must be greater than 0 and less than or equal to {LargePayloadTombstone.MaxRequestLimit}.");
}
- LP.GetLargePayloadTombstonesResponse response;
+ List tombstones;
try
{
- using var call = this.client.GetLargePayloadTombstonesAsync(
- new LP.GetLargePayloadTombstonesRequest { Limit = input },
- deadline: DateTime.UtcNow.Add(this.RpcTimeout));
- response = await call;
+ tombstones = await this.client.GetLargePayloadTombstonesAsync(input, DateTime.UtcNow.Add(this.RpcTimeout));
}
catch (RpcException e) when (e.StatusCode == StatusCode.Cancelled)
{
@@ -76,12 +75,6 @@ public override async Task> RunAsync(TaskActivityCon
e);
}
- List tombstones = new(response.Tombstones.Count);
- foreach (LP.LargePayloadTombstone tombstone in response.Tombstones)
- {
- tombstones.Add(new LargePayloadTombstone(tombstone.TombstoneToken, tombstone.PayloadToken));
- }
-
this.logger.BlobPurgeFetchedTombstones(tombstones.Count);
return tombstones;
}
diff --git a/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/ReportLargePayloadPurgeResultsActivity.cs b/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/ReportLargePayloadPurgeResultsActivity.cs
index a2a7e05ab..1a86a03a5 100644
--- a/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/ReportLargePayloadPurgeResultsActivity.cs
+++ b/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/ReportLargePayloadPurgeResultsActivity.cs
@@ -4,8 +4,6 @@
using Grpc.Core;
using Microsoft.DurableTask.Client;
using Microsoft.Extensions.Logging;
-using static Microsoft.DurableTask.Protobuf.LargePayloads.LargePayloadPurge;
-using LP = Microsoft.DurableTask.Protobuf.LargePayloads;
namespace Microsoft.DurableTask.AzureBlobPayloads;
@@ -17,13 +15,17 @@ namespace Microsoft.DurableTask.AzureBlobPayloads;
///
/// The large-payload purge service client used to report purge results to the backend.
/// The logger instance.
+///
+/// Infrastructure integration API for alternate .NET hosts. The supplied client must be bound to this
+/// worker's authenticated task hub. Its transport lifetime remains owned by the host.
+///
[DurableTask]
-internal sealed class ReportLargePayloadPurgeResultsActivity(
- LargePayloadPurgeClient client,
+public sealed class ReportLargePayloadPurgeResultsActivity(
+ ILargePayloadPurgeClient client,
ILogger logger)
: TaskActivity, object?>
{
- readonly LargePayloadPurgeClient client = Check.NotNull(client);
+ readonly ILargePayloadPurgeClient client = Check.NotNull(client);
readonly ILogger logger = Check.NotNull(logger);
///
@@ -43,32 +45,9 @@ internal sealed class ReportLargePayloadPurgeResultsActivity(
return null;
}
- LP.ReportLargePayloadPurgeResultsRequest request = new();
- foreach (LargePayloadPurgeResult result in input)
- {
- request.Results.Add(new LP.LargePayloadPurgeResult
- {
- // Echoed back exactly as it was received. The SDK never parses or rebuilds this token, so a
- // change to what the backend puts in it needs no change here.
- TombstoneToken = result.TombstoneToken,
-
- // The managed disposition enum declares the same numeric values as its protobuf counterpart,
- // so it maps across by value. This is the only enum on the message and it only travels
- // outbound, so the SDK can never receive a value it does not know.
- Disposition = (LP.LargePayloadPurgeDisposition)result.Disposition,
- });
- }
-
- if (request.Results.Count == 0)
- {
- return null;
- }
-
try
{
- using var call = this.client.ReportLargePayloadPurgeResultsAsync(
- request, deadline: DateTime.UtcNow.Add(this.RpcTimeout));
- await call;
+ await this.client.ReportLargePayloadPurgeResultsAsync(input, DateTime.UtcNow.Add(this.RpcTimeout));
}
catch (RpcException e) when (e.StatusCode == StatusCode.Cancelled)
{
diff --git a/src/Extensions/AzureBlobPayloads/AutoPurge/Client/GrpcLargePayloadPurgeClient.cs b/src/Extensions/AzureBlobPayloads/AutoPurge/Client/GrpcLargePayloadPurgeClient.cs
new file mode 100644
index 000000000..3741f1360
--- /dev/null
+++ b/src/Extensions/AzureBlobPayloads/AutoPurge/Client/GrpcLargePayloadPurgeClient.cs
@@ -0,0 +1,51 @@
+// Copyright (c) Microsoft Corporation.
+// Licensed under the MIT License.
+
+using Microsoft.DurableTask.Client;
+using Proto = Microsoft.DurableTask.Protobuf.LargePayloads;
+
+namespace Microsoft.DurableTask.AzureBlobPayloads;
+
+///
+/// Adapts the worker's existing, rebindable purge transport without owning its lifetime.
+///
+sealed class GrpcLargePayloadPurgeClient(Proto.LargePayloadPurge.LargePayloadPurgeClient client) : ILargePayloadPurgeClient
+{
+ readonly Proto.LargePayloadPurge.LargePayloadPurgeClient client = Check.NotNull(client);
+
+ ///
+ public async Task> GetLargePayloadTombstonesAsync(
+ int limit, DateTime deadline, CancellationToken cancellationToken = default)
+ {
+ using var call = this.client.GetLargePayloadTombstonesAsync(
+ new Proto.GetLargePayloadTombstonesRequest { Limit = limit }, deadline: deadline, cancellationToken: cancellationToken);
+ Proto.GetLargePayloadTombstonesResponse response = await call;
+ List tombstones = new(response.Tombstones.Count);
+ foreach (Proto.LargePayloadTombstone tombstone in response.Tombstones)
+ {
+ tombstones.Add(new LargePayloadTombstone(tombstone.TombstoneToken, tombstone.PayloadToken));
+ }
+
+ return tombstones;
+ }
+
+ ///
+ public async Task ReportLargePayloadPurgeResultsAsync(
+ IReadOnlyList results, DateTime deadline, CancellationToken cancellationToken = default)
+ {
+ Proto.ReportLargePayloadPurgeResultsRequest request = new();
+ foreach (LargePayloadPurgeResult result in results)
+ {
+ request.Results.Add(new Proto.LargePayloadPurgeResult
+ {
+ // Echo the opaque correlation token unchanged. The managed and protobuf enums share values.
+ TombstoneToken = result.TombstoneToken,
+ Disposition = (Proto.LargePayloadPurgeDisposition)result.Disposition,
+ });
+ }
+
+ using var call = this.client.ReportLargePayloadPurgeResultsAsync(
+ request, deadline: deadline, cancellationToken: cancellationToken);
+ await call;
+ }
+}
diff --git a/src/Extensions/AzureBlobPayloads/AzureBlobPayloads.csproj b/src/Extensions/AzureBlobPayloads/AzureBlobPayloads.csproj
index 78d80fd7a..c793b77b8 100644
--- a/src/Extensions/AzureBlobPayloads/AzureBlobPayloads.csproj
+++ b/src/Extensions/AzureBlobPayloads/AzureBlobPayloads.csproj
@@ -18,6 +18,7 @@
+
diff --git a/src/Extensions/AzureBlobPayloads/DependencyInjection/DurableTaskWorkerBuilderExtensions.AzureBlobPayloads.cs b/src/Extensions/AzureBlobPayloads/DependencyInjection/DurableTaskWorkerBuilderExtensions.AzureBlobPayloads.cs
index dafe5784f..a55644317 100644
--- a/src/Extensions/AzureBlobPayloads/DependencyInjection/DurableTaskWorkerBuilderExtensions.AzureBlobPayloads.cs
+++ b/src/Extensions/AzureBlobPayloads/DependencyInjection/DurableTaskWorkerBuilderExtensions.AzureBlobPayloads.cs
@@ -8,6 +8,7 @@
using Microsoft.DurableTask.Worker.Grpc.Internal;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.DependencyInjection.Extensions;
+using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using static Microsoft.DurableTask.Protobuf.LargePayloads.LargePayloadPurge;
using P = Microsoft.DurableTask.Protobuf;
@@ -131,12 +132,14 @@ static IDurableTaskWorkerBuilder UseExternalizedPayloadsCore(IDurableTaskWorkerB
{
r.AddOrchestrator();
r.AddActivity(nameof(GetLargePayloadTombstonesActivity), sp =>
- ActivatorUtilities.CreateInstance(
- sp, sp.GetRequiredKeyedService(builder.Name)));
+ new GetLargePayloadTombstonesActivity(
+ new GrpcLargePayloadPurgeClient(sp.GetRequiredKeyedService(builder.Name)),
+ sp.GetRequiredService>()));
r.AddActivity();
r.AddActivity(nameof(ReportLargePayloadPurgeResultsActivity), sp =>
- ActivatorUtilities.CreateInstance(
- sp, sp.GetRequiredKeyedService(builder.Name)));
+ new ReportLargePayloadPurgeResultsActivity(
+ new GrpcLargePayloadPurgeClient(sp.GetRequiredKeyedService(builder.Name)),
+ sp.GetRequiredService>()));
});
return builder;
diff --git a/src/LargePayloadPurge.Abstractions/ILargePayloadPurgeClient.cs b/src/LargePayloadPurge.Abstractions/ILargePayloadPurgeClient.cs
new file mode 100644
index 000000000..54dab01b4
--- /dev/null
+++ b/src/LargePayloadPurge.Abstractions/ILargePayloadPurgeClient.cs
@@ -0,0 +1,40 @@
+// Copyright (c) Microsoft Corporation.
+// Licensed under the MIT License.
+
+using Microsoft.DurableTask.Client;
+
+namespace Microsoft.DurableTask.AzureBlobPayloads;
+
+///
+/// Provides transport operations bound to a task hub for integrating blob auto-purge with an alternate .NET host.
+///
+///
+/// This is an infrastructure integration API, not an application orchestration API. Implementations must
+/// use the same authenticated task hub as the associated orchestration client and preserve its authentication,
+/// metadata, reconnection and transport lifetime. The SDK does not own or dispose the supplied client.
+/// Fetch and report must propagate gRPC status exceptions unchanged: the activities own their cancellation,
+/// unsupported-backend and fetch-precondition handling. They must honor the supplied UTC deadline.
+/// Correlation tokens must be returned and reported exactly as received, without parsing or reconstruction.
+///
+public interface ILargePayloadPurgeClient
+{
+ ///
+ /// Fetches a bounded batch of due tombstones for this client's authenticated task hub.
+ ///
+ /// The requested maximum number of tombstones, from 1 through 1000.
+ /// The absolute UTC deadline for this backend attempt.
+ /// Cancels the fetch operation.
+ /// The fetched tombstones, with their opaque correlation and payload tokens unchanged.
+ Task> GetLargePayloadTombstonesAsync(
+ int limit, DateTime deadline, CancellationToken cancellationToken = default);
+
+ ///
+ /// Reports deletion outcomes for this client's authenticated task hub.
+ ///
+ /// The outcomes, each carrying the exact correlation token received during fetch.
+ /// The absolute UTC deadline for this backend attempt.
+ /// Cancels the report operation.
+ /// A task that completes when the backend acknowledges the results.
+ Task ReportLargePayloadPurgeResultsAsync(
+ IReadOnlyList results, DateTime deadline, CancellationToken cancellationToken = default);
+}
diff --git a/src/LargePayloadPurge.Abstractions/IOrchestrationServiceLargePayloadPurgeClient.cs b/src/LargePayloadPurge.Abstractions/IOrchestrationServiceLargePayloadPurgeClient.cs
new file mode 100644
index 000000000..c3c49ca6b
--- /dev/null
+++ b/src/LargePayloadPurge.Abstractions/IOrchestrationServiceLargePayloadPurgeClient.cs
@@ -0,0 +1,26 @@
+// Copyright (c) Microsoft Corporation.
+// Licensed under the MIT License.
+
+using Microsoft.DurableTask.AzureBlobPayloads;
+
+namespace DurableTask.LargePayloadPurge;
+
+///
+/// Optional orchestration service client capability for purging tombstoned large payloads.
+///
+///
+/// Extends the shared fetch and report transport contract with control of the task hub's auto-purge setting.
+///
+public interface IOrchestrationServiceLargePayloadPurgeClient : ILargePayloadPurgeClient
+{
+ ///
+ /// Records whether large payload auto-purge is enabled for the client's task hub.
+ ///
+ /// This operation does not start, stop, or wait for a purge runner.
+ /// Whether large payload auto-purge is enabled.
+ /// The caller's operation deadline in UTC, or
+ /// when the caller has not specified a deadline.
+ /// The token used to cancel the operation.
+ /// A task that represents the operation.
+ Task SetLargePayloadAutoPurgeAsync(bool enabled, DateTime deadlineUtc, CancellationToken cancellationToken);
+}
diff --git a/src/LargePayloadPurge.Abstractions/LargePayloadPurge.Abstractions.csproj b/src/LargePayloadPurge.Abstractions/LargePayloadPurge.Abstractions.csproj
new file mode 100644
index 000000000..4fac45e98
--- /dev/null
+++ b/src/LargePayloadPurge.Abstractions/LargePayloadPurge.Abstractions.csproj
@@ -0,0 +1,17 @@
+
+
+
+
+ netstandard2.0
+ DurableTask.LargePayloadPurge
+ Service and activity transport contracts for large payload purge using the Durable Task SDK client models.
+ 0.1.0
+
+ true
+
+
+
+
+
+
+
diff --git a/src/LargePayloadPurge.Abstractions/README.md b/src/LargePayloadPurge.Abstractions/README.md
new file mode 100644
index 000000000..5ac73b719
--- /dev/null
+++ b/src/LargePayloadPurge.Abstractions/README.md
@@ -0,0 +1,45 @@
+# Large payload purge contracts
+
+`Microsoft.DurableTask.LargePayloadPurge.Abstractions` provides two interfaces in the
+`Microsoft.DurableTask.LargePayloadPurge.Abstractions` assembly:
+
+- `Microsoft.DurableTask.AzureBlobPayloads.ILargePayloadPurgeClient` is the activity transport contract for
+ fetching tombstones with `GetLargePayloadTombstonesAsync` and reporting outcomes. It accepts a UTC deadline
+ and optional cancellation token.
+ Transport implementations bind it to the same authenticated task hub as the associated orchestration
+ client and preserve the documented gRPC status behavior without requiring this package to reference gRPC.
+- `DurableTask.LargePayloadPurge.IOrchestrationServiceLargePayloadPurgeClient` inherits that shared transport
+ contract and adds only `SetLargePayloadAutoPurgeAsync`. The setting operation requires a deadline and
+ cancellation token; `DateTime.MaxValue` represents an unspecified deadline. Setting the flag alone does
+ not start, stop or wait for a purge runner.
+
+The contract uses the canonical `LargePayloadTombstone`, `LargePayloadPurgeResult`, and
+`LargePayloadPurgeDisposition` types from `Microsoft.DurableTask.Client`. It does not copy, move, wrap or
+forward those types. Backend-issued tombstone tokens must be echoed unchanged.
+
+## Dependencies and release
+
+This project references the SDK Client project directly. Its packaged dependency graph is:
+
+```text
+Microsoft.DurableTask.LargePayloadPurge.Abstractions
+ -> Microsoft.DurableTask.Client
+ -> Microsoft.DurableTask.Abstractions
+ -> Microsoft.Azure.DurableTask.Core
+```
+
+It is not a BCL-only package. Core does not depend on this package or the SDK Client. The contract package
+contains no blob storage, gRPC transport, worker, or orchestration implementation.
+The Azure Blob implementation references this package, not the other way around.
+
+The initial version is planned as `0.1.0`, independently of the repository-wide SDK version. Release the
+matching SDK Client and Abstractions dependencies containing the purge models before this package.
+Published Client `1.26.0` predates those models and is not sufficient. Repository builds use source project
+references; no external Client-version bootstrap property is required.
+
+The assembly uses this repository's strong-name key. Consumers of the unreleased prototype package
+`Microsoft.Azure.DurableTask.LargePayloadPurge.Abstractions` must update the package reference and rebuild;
+the assembly name is now `Microsoft.DurableTask.LargePayloadPurge.Abstractions`. Source namespaces stay the same.
+Service implementations should rename `GetLargePayloadsToPurgeAsync` to the inherited
+`GetLargePayloadTombstonesAsync`, returning `Task>`; Report uses the existing
+shared transport signature. No type forwarder or duplicate DTO is provided.
diff --git a/src/LargePayloadPurge.Abstractions/RELEASENOTES.md b/src/LargePayloadPurge.Abstractions/RELEASENOTES.md
new file mode 100644
index 000000000..51224bf72
--- /dev/null
+++ b/src/LargePayloadPurge.Abstractions/RELEASENOTES.md
@@ -0,0 +1 @@
+- Initial service and activity transport contracts for large payload auto-purge, using the canonical Durable Task SDK Client models.
diff --git a/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/AlternateHostPurgeTaskTests.cs b/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/AlternateHostPurgeTaskTests.cs
new file mode 100644
index 000000000..46d270ea5
--- /dev/null
+++ b/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/AlternateHostPurgeTaskTests.cs
@@ -0,0 +1,264 @@
+// Copyright (c) Microsoft Corporation.
+// Licensed under the MIT License.
+
+using DurableTask.Core;
+using DurableTask.Core.Command;
+using DurableTask.Core.History;
+using Grpc.Core;
+using Microsoft.DurableTask.AzureBlobPayloads;
+using Microsoft.DurableTask.Client;
+using Microsoft.DurableTask.Converters;
+using Microsoft.DurableTask.Worker.Shims;
+using Microsoft.Extensions.Logging.Abstractions;
+
+namespace Microsoft.DurableTask.Extensions.AzureBlobPayloads.Tests.AutoPurge;
+
+public class AlternateHostPurgeTaskTests
+{
+ [Fact]
+ public async Task ActualTasks_ExecuteWithDtfXArguments_AndPreserveCorrelationAcrossReplayAsync()
+ {
+ // Arrange
+ const string FirstToken = "opaque:row/one+==";
+ const string SecondToken = "opaque:\"row two\"";
+ List tombstones =
+ [
+ new(FirstToken, "blob:v2:https://account.blob.core.windows.net/payloads/one"),
+ new(SecondToken, "blob:v2:https://account.blob.core.windows.net/payloads/two"),
+ ];
+ Mock purge = new(MockBehavior.Strict);
+ DateTime? fetchDeadline = null;
+ DateTime? reportDeadline = null;
+ purge.Setup(p => p.GetLargePayloadTombstonesAsync(37, It.IsAny(), default))
+ .Callback((_, deadline, _) => fetchDeadline = deadline)
+ .ReturnsAsync(tombstones);
+ List? reported = null;
+ purge.Setup(p => p.ReportLargePayloadPurgeResultsAsync(It.IsAny>(), It.IsAny(), default))
+ .Callback, DateTime, CancellationToken>((results, deadline, _) =>
+ {
+ reported = results.ToList();
+ reportDeadline = deadline;
+ }).Returns(Task.CompletedTask);
+ Mock store = new(MockBehavior.Strict);
+ store.Setup(s => s.DeleteAsync(tombstones[0].PayloadToken, It.IsAny())).ReturnsAsync(PayloadDeleteOutcome.Deleted);
+ store.Setup(s => s.DeleteAsync(tombstones[1].PayloadToken, It.IsAny())).ThrowsAsync(new TimeoutException());
+ Driver driver = new(new(37));
+ DateTime earliestDeadline = DateTime.UtcNow.AddSeconds(60);
+
+ // Act
+ ScheduleTaskOrchestratorAction fetch = driver.SingleActivity(nameof(GetLargePayloadTombstonesActivity));
+ Assert.Equal("[37]", fetch.Input);
+ await driver.RunActivityAsync(fetch, new GetLargePayloadTombstonesActivity(purge.Object, NullLogger.Instance));
+ ScheduleTaskOrchestratorAction delete = driver.SingleActivity(nameof(DeleteExternalBlobActivity));
+ Assert.StartsWith("[[", delete.Input);
+ await driver.RunActivityAsync(delete, new DeleteExternalBlobActivity(store.Object, NullLogger.Instance));
+ ScheduleTaskOrchestratorAction report = driver.SingleActivity(nameof(ReportLargePayloadPurgeResultsActivity));
+ Assert.StartsWith("[[", report.Input);
+ await driver.RunActivityAsync(report, new ReportLargePayloadPurgeResultsActivity(purge.Object, NullLogger.Instance));
+
+ // Assert
+ Assert.Equal(new[]
+ {
+ new LargePayloadPurgeResult(FirstToken, LargePayloadPurgeDisposition.Deleted),
+ new LargePayloadPurgeResult(SecondToken, LargePayloadPurgeDisposition.Retry),
+ }, reported);
+ Assert.Contains("\"PurgedCount\":1", driver.Result.CustomStatus);
+ Assert.Equal("[37]", driver.SingleActivity(nameof(GetLargePayloadTombstonesActivity)).Input);
+ Assert.All(new[] { fetchDeadline, reportDeadline }, deadline =>
+ {
+ DateTime actualDeadline = Assert.IsType(deadline);
+ Assert.Equal(DateTimeKind.Utc, actualDeadline.Kind);
+ Assert.InRange(actualDeadline, earliestDeadline, DateTime.UtcNow.AddSeconds(60));
+ });
+ Assert.Equal(driver.Snapshot(), driver.ReplaySnapshot());
+ store.VerifyAll();
+ purge.VerifyAll();
+ }
+
+ [Theory]
+ [InlineData(false)]
+ [InlineData(true)]
+ public async Task UnsupportedActivity_FailureDetailsReachOrchestrator_AndEventResumesWithoutRetryAsync(bool reportFailure)
+ {
+ // Arrange
+ Mock purge = new(MockBehavior.Strict);
+ RpcException unsupported = new(new Status(StatusCode.Unimplemented, "unsupported"));
+ purge.Setup(p => p.GetLargePayloadTombstonesAsync(It.IsAny(), It.IsAny(), default))
+ .ThrowsAsync(unsupported);
+ purge.Setup(p => p.ReportLargePayloadPurgeResultsAsync(It.IsAny>(), It.IsAny(), default))
+ .ThrowsAsync(unsupported);
+ Driver driver = new(new(20));
+ ScheduleTaskOrchestratorAction activity = driver.SingleActivity(nameof(GetLargePayloadTombstonesActivity));
+ ITaskActivity implementation = new GetLargePayloadTombstonesActivity(purge.Object, NullLogger.Instance);
+ if (reportFailure)
+ {
+ driver.Complete(activity, new[] { new LargePayloadTombstone("correlation", "blob:v2:payload") });
+ driver.Complete(driver.SingleActivity(nameof(DeleteExternalBlobActivity)), new[] { new BlobPurgeOutcome(LargePayloadPurgeDisposition.Deleted) });
+ activity = driver.SingleActivity(nameof(ReportLargePayloadPurgeResultsActivity));
+ implementation = new ReportLargePayloadPurgeResultsActivity(purge.Object, NullLogger.Instance);
+ }
+
+ // Act
+ Exception? error = await Record.ExceptionAsync(() => driver.InvokeActivityAsync(activity, implementation));
+ NotImplementedException failure = Assert.IsType(error);
+ driver.Fail(activity, failure);
+
+ // Assert
+ Assert.Same(unsupported, failure.InnerException);
+ Assert.Empty(driver.Result.Actions);
+ Assert.Contains("\"Status\":\"BackendUnsupported\"", driver.Result.CustomStatus);
+ Assert.Equal(driver.Snapshot(), driver.ReplaySnapshot());
+ driver.Turn(Driver.Configure(71));
+ Assert.Equal("[71]", driver.SingleActivity(nameof(GetLargePayloadTombstonesActivity)).Input);
+ }
+
+ [Fact]
+ public async Task TransientActivityFailure_KeepsDtfXRetryTimerAsync()
+ {
+ // Arrange
+ Mock purge = new(MockBehavior.Strict);
+ RpcException unavailable = new(new Status(StatusCode.Unavailable, "unavailable"));
+ purge.Setup(p => p.GetLargePayloadTombstonesAsync(It.IsAny(), It.IsAny(), default)).ThrowsAsync(unavailable);
+ Driver driver = new(new(20));
+ ScheduleTaskOrchestratorAction fetch = driver.SingleActivity(nameof(GetLargePayloadTombstonesActivity));
+ GetLargePayloadTombstonesActivity activity = new(purge.Object, NullLogger.Instance);
+
+ // Act
+ Exception? error = await Record.ExceptionAsync(() => driver.InvokeActivityAsync(fetch, activity));
+ driver.Fail(fetch, Assert.IsType(error));
+
+ // Assert
+ CreateTimerOrchestratorAction retry = Assert.IsType(Assert.Single(driver.Result.Actions));
+ Assert.Equal(driver.Now.AddSeconds(15), retry.FireAt);
+ Assert.Equal(driver.Snapshot(), driver.ReplaySnapshot());
+ }
+
+ [Fact]
+ public void ConfigurationAndContinueAsNew_PreserveBufferedEventsAndStateThroughDtfXReplay()
+ {
+ // Arrange
+ Driver driver = new(new(100, 9));
+ for (int i = 0; i < 5; i++)
+ {
+ driver.Complete(driver.SingleActivity(nameof(GetLargePayloadTombstonesActivity)), Array.Empty());
+ CreateTimerOrchestratorAction timer = Assert.IsType(Assert.Single(driver.Result.Actions));
+ if (i < 4)
+ {
+ driver.Turn(Driver.TimerFired(timer));
+ }
+ }
+
+ // Act
+ driver.Turn(Driver.Configure(700), Driver.Configure(800));
+ OrchestrationCompleteOrchestratorAction completed = Assert.IsType(Assert.Single(driver.Result.Actions));
+
+ // Assert
+ Assert.Equal(OrchestrationStatus.ContinuedAsNew, completed.OrchestrationStatus);
+ BlobPurgeJobRunRequest next = JsonDataConverter.Default.Deserialize(completed.Result)!;
+ Assert.Equal(new BlobPurgeJobRunRequest(700, 9), next);
+ EventRaisedEvent carried = Assert.IsType(Assert.Single(completed.CarryoverEvents));
+ Assert.Equal(BlobPurgeConstants.SetBatchSizeEvent, carried.Name);
+ Assert.Equal("800", carried.Input);
+ Assert.Equal(driver.Snapshot(), driver.ReplaySnapshot());
+ Driver nextDriver = new(next, carried);
+ nextDriver.Complete(nextDriver.SingleActivity(nameof(GetLargePayloadTombstonesActivity)), Array.Empty());
+ Assert.Equal("[800]", nextDriver.SingleActivity(nameof(GetLargePayloadTombstonesActivity)).Input);
+ }
+
+ sealed class Driver
+ {
+ readonly DurableTaskShimFactory factory = new();
+ readonly List history = [];
+ readonly OrchestrationInstance instance = new()
+ {
+ InstanceId = BlobPurgeConstants.OrchestratorInstanceId,
+ ExecutionId = "alternate-host",
+ };
+ List lastPast = [];
+ List lastNew = [];
+
+ public Driver(BlobPurgeJobRunRequest input, params HistoryEvent[] events)
+ {
+ this.Turn([new ExecutionStartedEvent(-1, JsonDataConverter.Default.Serialize(input))
+ {
+ Name = nameof(BlobPurgeJobOrchestrator),
+ Version = string.Empty,
+ OrchestrationInstance = this.instance,
+ }, .. events]);
+ }
+
+ public DateTime Now { get; private set; } = new(2026, 9, 1, 0, 0, 0, DateTimeKind.Utc);
+ public OrchestratorExecutionResult Result { get; private set; } = null!;
+
+ public static EventRaisedEvent Configure(int size) => new(-1, JsonDataConverter.Default.Serialize(size)) { Name = BlobPurgeConstants.SetBatchSizeEvent };
+ public static TimerFiredEvent TimerFired(CreateTimerOrchestratorAction timer) => new(-1, timer.FireAt) { TimerId = timer.Id };
+
+ public ScheduleTaskOrchestratorAction SingleActivity(string name)
+ {
+ ScheduleTaskOrchestratorAction task = Assert.IsType(Assert.Single(this.Result.Actions));
+ Assert.Equal(name, task.Name);
+ Assert.Equal(string.Empty, task.Version);
+ return task;
+ }
+
+ public Task InvokeActivityAsync(ScheduleTaskOrchestratorAction action, ITaskActivity implementation) =>
+ this.factory.CreateActivity(Assert.IsType(action.Name), implementation).RunAsync(
+ new TaskContext(this.instance, action.Name, action.Version, action.Id), action.Input);
+
+ public async Task RunActivityAsync(ScheduleTaskOrchestratorAction action, ITaskActivity implementation) =>
+ this.Turn(new TaskCompletedEvent(-1, action.Id, await this.InvokeActivityAsync(action, implementation)));
+
+ public void Complete(ScheduleTaskOrchestratorAction action, object result) =>
+ this.Turn(new TaskCompletedEvent(-1, action.Id, JsonDataConverter.Default.Serialize(result)));
+
+ public void Fail(ScheduleTaskOrchestratorAction action, Exception failure) =>
+ this.Turn(new TaskFailedEvent(-1, action.Id, failure.Message, null, new FailureDetails(failure)));
+
+ public string Snapshot() => Serialize(this.Result);
+ public string ReplaySnapshot() => Serialize(this.Replay());
+
+ public void Turn(params HistoryEvent[] events)
+ {
+ this.Now = events.OfType().Select(e => e.FireAt).Append(this.Now.AddSeconds(1)).Max();
+ this.lastPast = [.. this.history];
+ this.lastNew = [new OrchestratorStartedEvent(-1) { Timestamp = this.Now }];
+ foreach (HistoryEvent item in events)
+ {
+ item.Timestamp = this.Now;
+ this.lastNew.Add(item);
+ }
+
+ this.Result = this.Replay();
+ this.history.AddRange(this.lastNew);
+ foreach (OrchestratorAction action in this.Result.Actions)
+ {
+ if (action is ScheduleTaskOrchestratorAction task)
+ {
+ this.history.Add(new TaskScheduledEvent(task.Id, Assert.IsType(task.Name), task.Version, task.Input) { Timestamp = this.Now });
+ }
+ else if (action is CreateTimerOrchestratorAction timer)
+ {
+ this.history.Add(new TimerCreatedEvent(timer.Id, timer.FireAt) { Timestamp = this.Now });
+ }
+ }
+
+ this.history.Add(new OrchestratorCompletedEvent(-1) { Timestamp = this.Now });
+ }
+
+ static string Serialize(OrchestratorExecutionResult result) =>
+ Newtonsoft.Json.JsonConvert.SerializeObject(new { result.CustomStatus, Actions = result.Actions.ToArray() });
+
+ OrchestratorExecutionResult Replay()
+ {
+ OrchestrationRuntimeState state = new(this.lastPast);
+ foreach (HistoryEvent item in this.lastNew)
+ {
+ state.AddEvent(item);
+ }
+
+ TaskOrchestration task = this.factory.CreateOrchestration(nameof(BlobPurgeJobOrchestrator), new BlobPurgeJobOrchestrator());
+ TaskOrchestrationExecutor executor = new(state, task, BehaviorOnContinueAsNew.Carryover, ErrorPropagationMode.UseFailureDetails);
+ return executor.Execute();
+ }
+ }
+}
diff --git a/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/PurgeActivityBackendStatusTests.cs b/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/PurgeActivityBackendStatusTests.cs
index ddc3ea915..0a5ed6c8a 100644
--- a/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/PurgeActivityBackendStatusTests.cs
+++ b/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/PurgeActivityBackendStatusTests.cs
@@ -24,7 +24,7 @@ public async Task GetLargePayloadTombstones_WhenBackendUnimplemented_ThrowsNotIm
// Arrange - the backend rejects the fetch RPC because it does not implement it.
LargePayloadPurgeClient client = new(
new ThrowingCallInvoker(new RpcException(new Status(StatusCode.Unimplemented, "unknown method"))));
- GetLargePayloadTombstonesActivity activity = new(client, new TestLogger());
+ GetLargePayloadTombstonesActivity activity = new(new GrpcLargePayloadPurgeClient(client), new TestLogger());
// Act
Func act = () => activity.RunAsync(null!, 100);
@@ -42,7 +42,7 @@ public async Task ReportLargePayloadPurgeResults_WhenBackendUnimplemented_Throws
LargePayloadPurgeClient client = new(
new ThrowingCallInvoker(new RpcException(new Status(StatusCode.Unimplemented, "unknown method"))));
ReportLargePayloadPurgeResultsActivity activity =
- new(client, new TestLogger());
+ new(new GrpcLargePayloadPurgeClient(client), new TestLogger());
List results = new()
{
new LargePayloadPurgeResult("tombstone-token-1", LargePayloadPurgeDisposition.Deleted),
@@ -64,7 +64,7 @@ public async Task GetLargePayloadTombstones_WhenAutoPurgeDisabled_ReturnsEmptyAn
TestLogger logger = new();
LargePayloadPurgeClient client = new(
new ThrowingCallInvoker(new RpcException(new Status(StatusCode.FailedPrecondition, Detail))));
- GetLargePayloadTombstonesActivity activity = new(client, logger);
+ GetLargePayloadTombstonesActivity activity = new(new GrpcLargePayloadPurgeClient(client), logger);
// Act
List tombstones = await activity.RunAsync(null!, 100);
@@ -88,7 +88,7 @@ public async Task GetLargePayloadTombstones_WhenTaskHubBeingDeleted_IsNotMislabe
TestLogger logger = new();
LargePayloadPurgeClient client = new(
new ThrowingCallInvoker(new RpcException(new Status(StatusCode.FailedPrecondition, Detail))));
- GetLargePayloadTombstonesActivity activity = new(client, logger);
+ GetLargePayloadTombstonesActivity activity = new(new GrpcLargePayloadPurgeClient(client), logger);
// Act
List tombstones = await activity.RunAsync(null!, 100);
@@ -109,7 +109,7 @@ public async Task GetLargePayloadTombstones_WhenBackendCancels_ThrowsOperationCa
// Arrange - cancellation is unrelated to the new precondition path and must keep its own translation.
LargePayloadPurgeClient client = new(
new ThrowingCallInvoker(new RpcException(new Status(StatusCode.Cancelled, "canceled"))));
- GetLargePayloadTombstonesActivity activity = new(client, new TestLogger());
+ GetLargePayloadTombstonesActivity activity = new(new GrpcLargePayloadPurgeClient(client), new TestLogger());
// Act
Func act = () => activity.RunAsync(null!, 100);
@@ -126,7 +126,7 @@ public async Task GetLargePayloadTombstones_WhenBackendFailsOtherwise_Propagates
TestLogger logger = new();
LargePayloadPurgeClient client = new(
new ThrowingCallInvoker(new RpcException(new Status(StatusCode.Unavailable, "backend down"))));
- GetLargePayloadTombstonesActivity activity = new(client, logger);
+ GetLargePayloadTombstonesActivity activity = new(new GrpcLargePayloadPurgeClient(client), logger);
// Act
Func act = () => activity.RunAsync(null!, 100);
@@ -147,7 +147,7 @@ public async Task ReportLargePayloadPurgeResults_WhenBackendCancels_ThrowsOperat
LargePayloadPurgeClient client = new(
new ThrowingCallInvoker(new RpcException(new Status(StatusCode.Cancelled, "canceled"))));
ReportLargePayloadPurgeResultsActivity activity =
- new(client, new TestLogger());
+ new(new GrpcLargePayloadPurgeClient(client), new TestLogger());
List results = new()
{
new LargePayloadPurgeResult("tombstone-token-1", LargePayloadPurgeDisposition.Deleted),
@@ -172,7 +172,7 @@ public async Task ReportLargePayloadPurgeResults_WhenBackendFailsOtherwise_Propa
RpcException error = new(new Status(statusCode, "backend failure"));
LargePayloadPurgeClient client = new(new ThrowingCallInvoker(error));
TestLogger logger = new();
- ReportLargePayloadPurgeResultsActivity activity = new(client, logger);
+ ReportLargePayloadPurgeResultsActivity activity = new(new GrpcLargePayloadPurgeClient(client), logger);
List results = new()
{
new LargePayloadPurgeResult("tombstone-token-1", LargePayloadPurgeDisposition.Deleted),
@@ -196,7 +196,7 @@ public async Task GetLargePayloadTombstones_SetsDefaultUtcDeadlineAsync()
RpcException error = new(new Status(StatusCode.DeadlineExceeded, "deadline exceeded"));
ThrowingCallInvoker invoker = new(error);
LargePayloadPurgeClient client = new(invoker);
- GetLargePayloadTombstonesActivity activity = new(client, new TestLogger());
+ GetLargePayloadTombstonesActivity activity = new(new GrpcLargePayloadPurgeClient(client), new TestLogger());
DateTime earliestDeadline = DateTime.UtcNow.AddSeconds(60);
// Act
@@ -223,7 +223,7 @@ public async Task ReportLargePayloadPurgeResults_SetsDefaultUtcDeadlineAsync()
ThrowingCallInvoker invoker = new(error);
LargePayloadPurgeClient client = new(invoker);
ReportLargePayloadPurgeResultsActivity activity =
- new(client, new TestLogger());
+ new(new GrpcLargePayloadPurgeClient(client), new TestLogger());
List results = new()
{
new LargePayloadPurgeResult("tombstone-token-1", LargePayloadPurgeDisposition.Deleted),
diff --git a/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/PurgeActivityDeadlineTests.cs b/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/PurgeActivityDeadlineTests.cs
index d3bb89641..b48d1d64e 100644
--- a/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/PurgeActivityDeadlineTests.cs
+++ b/test/Extensions/AzureBlobPayloads.Tests/AutoPurge/PurgeActivityDeadlineTests.cs
@@ -24,7 +24,7 @@ public async Task GetLargePayloadTombstones_WhenDeadlineExpires_CancelsRequestAs
using BlockingHttpMessageHandler handler = new();
using GrpcChannel channel = GrpcChannel.ForAddress("http://localhost", new GrpcChannelOptions { HttpHandler = handler });
LargePayloadPurgeClient client = new(channel);
- GetLargePayloadTombstonesActivity activity = new(client, new TestLogger())
+ GetLargePayloadTombstonesActivity activity = new(new GrpcLargePayloadPurgeClient(client), new TestLogger())
{
RpcTimeout = TimeSpan.FromMilliseconds(200),
};
@@ -47,7 +47,7 @@ public async Task ReportLargePayloadPurgeResults_WhenDeadlineExpires_CancelsRequ
using BlockingHttpMessageHandler handler = new();
using GrpcChannel channel = GrpcChannel.ForAddress("http://localhost", new GrpcChannelOptions { HttpHandler = handler });
LargePayloadPurgeClient client = new(channel);
- ReportLargePayloadPurgeResultsActivity activity = new(client, new TestLogger())
+ ReportLargePayloadPurgeResultsActivity activity = new(new GrpcLargePayloadPurgeClient(client), new TestLogger())
{
RpcTimeout = TimeSpan.FromMilliseconds(200),
};
diff --git a/test/Extensions/AzureBlobPayloads.Tests/PayloadStore/BlobPayloadStoreDeleteResponseTests.cs b/test/Extensions/AzureBlobPayloads.Tests/PayloadStore/BlobPayloadStoreDeleteResponseTests.cs
index fa9be0eaf..a0f7dd786 100644
--- a/test/Extensions/AzureBlobPayloads.Tests/PayloadStore/BlobPayloadStoreDeleteResponseTests.cs
+++ b/test/Extensions/AzureBlobPayloads.Tests/PayloadStore/BlobPayloadStoreDeleteResponseTests.cs
@@ -51,7 +51,7 @@ public async Task DeleteResponse_ReportsOnlyConfirmedSuccessAsync(
TestLogger logger = new();
using GrpcChannel reportChannel = GrpcChannel.ForAddress("http://report.invalid", new() { HttpHandler = handler });
ReportLargePayloadPurgeResultsActivity report = new(
- new LP.LargePayloadPurge.LargePayloadPurgeClient(reportChannel),
+ new GrpcLargePayloadPurgeClient(new LP.LargePayloadPurge.LargePayloadPurgeClient(reportChannel)),
NullLogger.Instance);
// Act
diff --git a/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurge.Abstractions.Tests.csproj b/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurge.Abstractions.Tests.csproj
new file mode 100644
index 000000000..7414509cd
--- /dev/null
+++ b/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurge.Abstractions.Tests.csproj
@@ -0,0 +1,13 @@
+
+
+
+ net10.0
+ DurableTask.LargePayloadPurge.Tests
+
+
+
+
+
+
+
+
diff --git a/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurgeContractTests.cs b/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurgeContractTests.cs
new file mode 100644
index 000000000..f752e0842
--- /dev/null
+++ b/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurgeContractTests.cs
@@ -0,0 +1,160 @@
+// Copyright (c) Microsoft Corporation.
+// Licensed under the MIT License.
+
+using System.Reflection;
+using Microsoft.DurableTask.AzureBlobPayloads;
+using Microsoft.DurableTask.Client;
+using Xunit;
+
+namespace DurableTask.LargePayloadPurge.Tests;
+
+public class LargePayloadPurgeContractTests
+{
+ [Fact]
+ public void PackageOwnsBothInterfacesWithoutBlobDefinitionsOrForwarders()
+ {
+ // Arrange
+ Type contract = typeof(IOrchestrationServiceLargePayloadPurgeClient);
+ Type transport = typeof(ILargePayloadPurgeClient);
+ Assembly blob = typeof(GetLargePayloadTombstonesActivity).Assembly;
+
+ // Act
+ Type[] exported = contract.Assembly.GetExportedTypes();
+
+ // Assert
+ Assert.True(contract.IsInterface);
+ Assert.Equal("DurableTask.LargePayloadPurge", contract.Namespace);
+ Assert.Equal("Microsoft.DurableTask.LargePayloadPurge.Abstractions", contract.Assembly.GetName().Name);
+ Assert.Equal(new[] { contract, transport }.OrderBy(type => type.FullName), exported.OrderBy(type => type.FullName));
+ Assert.Same(contract.Assembly, transport.Assembly);
+ Assert.True(transport.IsInterface);
+ Assert.Equal("Microsoft.DurableTask.AzureBlobPayloads", transport.Namespace);
+ Assert.Equal([transport], contract.GetInterfaces());
+ Assert.Empty(transport.GetInterfaces());
+ Assert.Equal(nameof(IOrchestrationServiceLargePayloadPurgeClient.SetLargePayloadAutoPurgeAsync),
+ Assert.Single(contract.GetMethods()).Name);
+ Assert.Equal(2, transport.GetMethods().Length);
+ Assert.DoesNotContain(blob.GetTypes(), type => type.FullName == transport.FullName);
+ Assert.DoesNotContain(blob.GetForwardedTypes(), type => type.FullName == transport.FullName);
+ }
+
+ [Fact]
+ public void SetAcceptsExplicitChoiceAndCallerDeadlineAndCancellation()
+ {
+ // Arrange / Act / Assert
+ AssertSignature(
+ nameof(IOrchestrationServiceLargePayloadPurgeClient.SetLargePayloadAutoPurgeAsync),
+ typeof(Task),
+ [typeof(bool), typeof(DateTime), typeof(CancellationToken)],
+ ["enabled", "deadlineUtc", "cancellationToken"]);
+ }
+
+ [Fact]
+ public void GetReturnsCanonicalSdkTombstones()
+ {
+ // Arrange / Act / Assert
+ AssertTransportSignature(
+ nameof(ILargePayloadPurgeClient.GetLargePayloadTombstonesAsync),
+ typeof(Task>),
+ [typeof(int), typeof(DateTime), typeof(CancellationToken)],
+ ["limit", "deadline", "cancellationToken"]);
+ }
+
+ [Fact]
+ public void ReportAcceptsCanonicalSdkResults()
+ {
+ // Arrange / Act / Assert
+ AssertTransportSignature(
+ nameof(ILargePayloadPurgeClient.ReportLargePayloadPurgeResultsAsync),
+ typeof(Task),
+ [typeof(IReadOnlyList), typeof(DateTime), typeof(CancellationToken)],
+ ["results", "deadline", "cancellationToken"]);
+ }
+
+ [Fact]
+ public void ModelsComeFromSdkClientNotTheInterfacePackage()
+ {
+ // Arrange
+ Assembly client = typeof(DurableTaskClient).Assembly;
+
+ // Act
+ Assembly[] modelAssemblies =
+ [
+ typeof(LargePayloadTombstone).Assembly,
+ typeof(LargePayloadPurgeResult).Assembly,
+ typeof(LargePayloadPurgeDisposition).Assembly,
+ ];
+
+ // Assert
+ Assert.All(modelAssemblies, assembly => Assert.Same(client, assembly));
+ Assert.NotSame(client, typeof(IOrchestrationServiceLargePayloadPurgeClient).Assembly);
+ }
+
+ [Fact]
+ public void ContractUsesSdkSigningAndIndependentAssemblyVersion()
+ {
+ // Arrange
+ AssemblyName sdk = typeof(DurableTaskClient).Assembly.GetName();
+
+ // Act
+ AssemblyName contract = typeof(IOrchestrationServiceLargePayloadPurgeClient).Assembly.GetName();
+
+ // Assert
+ Assert.Equal(new Version(0, 1, 0, 0), contract.Version);
+ Assert.Equal("6A4C0315C2D1D937", Convert.ToHexString(contract.GetPublicKeyToken()!));
+ Assert.Equal(sdk.GetPublicKeyToken(), contract.GetPublicKeyToken());
+ }
+
+ [Fact]
+ public void ContractDependsOnSdkModelsWithoutAddingACoreReverseDependency()
+ {
+ // Arrange
+ Assembly contract = typeof(IOrchestrationServiceLargePayloadPurgeClient).Assembly;
+ Assembly core = typeof(Core.TaskHubClient).Assembly;
+ Assembly blob = typeof(GetLargePayloadTombstonesActivity).Assembly;
+
+ // Act
+ string?[] contractReferences = contract.GetReferencedAssemblies().Select(name => name.Name).ToArray();
+ string?[] coreReferences = core.GetReferencedAssemblies().Select(name => name.Name).ToArray();
+
+ // Assert
+ Assert.Contains("Microsoft.DurableTask.Client", contractReferences);
+ Assert.DoesNotContain("DurableTask.Core", contractReferences);
+ Assert.DoesNotContain("Microsoft.DurableTask.Worker", contractReferences);
+ Assert.DoesNotContain("Microsoft.DurableTask.Grpc", contractReferences);
+ Assert.DoesNotContain("Microsoft.DurableTask.Extensions.AzureBlobPayloads", contractReferences);
+ Assert.DoesNotContain(contractReferences, name => name!.StartsWith("Azure.", StringComparison.Ordinal));
+ Assert.DoesNotContain(contractReferences, name => name!.StartsWith("Grpc.", StringComparison.Ordinal));
+ Assert.DoesNotContain(contractReferences, name => name!.StartsWith("Microsoft.DurableTask.Worker", StringComparison.Ordinal));
+ Assert.Contains(blob.GetReferencedAssemblies(), name => name.Name == contract.GetName().Name);
+ Assert.DoesNotContain(contract.GetName().Name, coreReferences);
+ Assert.DoesNotContain("Microsoft.DurableTask.Client", coreReferences);
+ }
+
+ static void AssertSignature(string methodName, Type returnType, Type[] parameterTypes, string[] parameterNames)
+ {
+ MethodInfo method = typeof(IOrchestrationServiceLargePayloadPurgeClient).GetMethod(methodName)!;
+ Assert.NotNull(method);
+ Assert.Equal(returnType, method.ReturnType);
+ ParameterInfo[] parameters = method.GetParameters();
+ Assert.Equal(parameterTypes, parameters.Select(parameter => parameter.ParameterType));
+ Assert.Equal(parameterNames, parameters.Select(parameter => parameter.Name));
+ Assert.All(parameters, parameter => Assert.False(parameter.IsOptional));
+ }
+
+ static void AssertTransportSignature(string methodName, Type returnType, Type[] parameterTypes, string[] parameterNames)
+ {
+ MethodInfo method = typeof(ILargePayloadPurgeClient).GetMethod(methodName)!;
+ Assert.NotNull(method);
+ Type inheritedContract = Assert.Single(typeof(IOrchestrationServiceLargePayloadPurgeClient).GetInterfaces());
+ Assert.Equal(method, inheritedContract.GetMethod(methodName));
+ Assert.Equal(returnType, method.ReturnType);
+ ParameterInfo[] parameters = method.GetParameters();
+ Assert.Equal(parameterTypes, parameters.Select(parameter => parameter.ParameterType));
+ Assert.Equal(parameterNames, parameters.Select(parameter => parameter.Name));
+ Assert.All(parameters.Take(2), parameter => Assert.False(parameter.IsOptional));
+ Assert.True(parameters[2].IsOptional);
+ Assert.True(parameters[2].HasDefaultValue);
+ Assert.Null(parameters[2].DefaultValue);
+ }
+}