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); + } +}