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 2fc375471..d382acebe 100644 --- a/README.md +++ b/README.md @@ -200,6 +200,55 @@ 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.Azure.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. + +`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. + ## Obtaining the Protobuf definitions This project utilizes protobuf definitions from [durabletask-protobuf](https://github.com/microsoft/durabletask-protobuf), which are copied (vendored) into this repository under the `src/Grpc` directory. See the corresponding [README.md](./src/Grpc/README.md) for more information about how to update the protobuf definitions. diff --git a/azure-pipelines-release.yml b/azure-pipelines-release.yml index 3bfb6e832..2b084345c 100644 --- a/azure-pipelines-release.yml +++ b/azure-pipelines-release.yml @@ -95,6 +95,38 @@ steps: } ] +# The service contract preserves its assembly name outside the Microsoft.DurableTask prefix. +- task: SFP.build-tasks.custom-build-task-1.EsrpCodeSigning@2 + displayName: 'ESRP CodeSigning: Large payload purge abstraction' + inputs: + ConnectedServiceName: 'ESRP Service' + FolderPath: $(bin_dir) + Pattern: 'DurableTask.LargePayloadPurge.Abstractions.dll' + signConfigType: inlineSignParams + inlineOperation: | + [ + { + "KeyCode": "CP-230012", + "OperationCode": "SigntoolSign", + "Parameters": { + "OpusName": "Microsoft", + "OpusInfo": "http://www.microsoft.com", + "FileDigest": "/fd \"SHA256\"", + "PageHash": "/NPH", + "TimeStamp": "/tr \"http://rfc3161.gtm.corp.microsoft.com/TSS/HttpTspServer\" /td sha256" + }, + "ToolName": "sign", + "ToolVersion": "1.0" + }, + { + "KeyCode": "CP-230012", + "OperationCode": "SigntoolVerify", + "Parameters": {}, + "ToolName": "sign", + "ToolVersion": "1.0" + } + ] + # SBOM generator task for additional supply chain protection - task: AzureArtifacts.manifest-generator-task.manifest-generator-task.ManifestGeneratorTask@0 displayName: 'SBOM Manifest Generator' diff --git a/doc/release_process.md b/doc/release_process.md index 319a2cd49..e420cad3d 100644 --- a/doc/release_process.md +++ b/doc/release_process.md @@ -5,11 +5,22 @@ | Package prefix | Registry | |---|---| | `Microsoft.DurableTask.*` | [NuGet](https://www.nuget.org/profiles/durabletask) | +| `Microsoft.Azure.DurableTask.LargePayloadPurge.Abstractions` | [NuGet](https://www.nuget.org/profiles/durabletask) | This repo publishes multiple NuGet packages. Most share a single version defined in `eng/targets/Release.props`. Individual packages can version independently by adding `` and `` properties directly in their `.csproj`. 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 retains the service interface's Apache-2.0 notice and the activity +transport interface's MIT notice, includes both license texts, and declares `Apache-2.0 AND MIT`. +It uses the SDK strong-name key. +Its assembly name is `DurableTask.LargePayloadPurge.Abstractions`, so both assembly-signing pipelines +explicitly include it in addition to the `Microsoft.DurableTask.*` assemblies. The existing source traversal, +SBOM inclusion and NuGet signing/packing steps apply unchanged. + ### 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..f8ac44f53 100644 --- a/eng/publish/publish.yml +++ b/eng/publish/publish.yml @@ -463,4 +463,27 @@ 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.Azure.DurableTask.LargePayloadPurge.Abstractions) + - job: nugetRelease_Microsoft_Azure_DurableTask_LargePayloadPurge_Abstractions + displayName: NuGet Release (Microsoft.Azure.DurableTask.LargePayloadPurge.Abstractions) + dependsOn: nugetApproval + condition: succeeded('nugetApproval') + 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.Azure.DurableTask.LargePayloadPurge.Abstractions)' + inputs: + command: push + nuGetFeedType: external + publishFeedCredentials: 'DurableTask org NuGet API Key' + packagesToPush: '$(System.DefaultWorkingDirectory)/drop/Microsoft.Azure.DurableTask.LargePayloadPurge.Abstractions.*.nupkg;!$(System.DefaultWorkingDirectory)/**/*.symbols.nupkg' + packageParentPath: $(System.DefaultWorkingDirectory) \ No newline at end of file diff --git a/eng/templates/build.yml b/eng/templates/build.yml index bcd404f81..420514ca7 100644 --- a/eng/templates/build.yml +++ b/eng/templates/build.yml @@ -76,6 +76,13 @@ jobs: pattern: Microsoft.DurableTask.*.dll signType: dll + - template: ci/sign-files.yml@eng + parameters: + displayName: Sign large payload purge abstraction + folderPath: $(bin_dir) + pattern: DurableTask.LargePayloadPurge.Abstractions.dll + signType: dll + # Packaging needs to be a separate step from build. # This will automatically pick up the signed DLLs. - task: DotNetCoreCLI@2 diff --git a/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/GetLargePayloadTombstonesActivity.cs b/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/GetLargePayloadTombstonesActivity.cs index 2c3147215..59767fa23 100644 --- a/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/GetLargePayloadTombstonesActivity.cs +++ b/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/GetLargePayloadTombstonesActivity.cs @@ -5,7 +5,6 @@ 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,15 +14,30 @@ 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); + /// + /// Initializes a new instance of the class using the worker's transport. + /// + /// The worker's purge client. + /// The activity logger. + internal GetLargePayloadTombstonesActivity( + LargePayloadPurgeClient client, ILogger logger) + : this(new GrpcLargePayloadPurgeClient(client), logger) + { + } + /// /// Gets or sets the timeout for one backend RPC attempt. /// @@ -41,13 +55,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 +87,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..cf3501b67 100644 --- a/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/ReportLargePayloadPurgeResultsActivity.cs +++ b/src/Extensions/AzureBlobPayloads/AutoPurge/Activities/ReportLargePayloadPurgeResultsActivity.cs @@ -5,7 +5,6 @@ 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,15 +16,30 @@ 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); + /// + /// Initializes a new instance of the class using the worker's transport. + /// + /// The worker's purge client. + /// The activity logger. + internal ReportLargePayloadPurgeResultsActivity( + LargePayloadPurgeClient client, ILogger logger) + : this(new GrpcLargePayloadPurgeClient(client), logger) + { + } + /// /// Gets or sets the timeout for one backend RPC attempt. /// @@ -43,32 +57,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..a4311e969 --- /dev/null +++ b/src/Extensions/AzureBlobPayloads/AutoPurge/Client/GrpcLargePayloadPurgeClient.cs @@ -0,0 +1,52 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +using Microsoft.DurableTask.Client; +using static Microsoft.DurableTask.Protobuf.LargePayloads.LargePayloadPurge; +using LP = Microsoft.DurableTask.Protobuf.LargePayloads; + +namespace Microsoft.DurableTask.AzureBlobPayloads; + +/// +/// Adapts the worker's existing, rebindable purge transport without owning its lifetime. +/// +sealed class GrpcLargePayloadPurgeClient(LargePayloadPurgeClient client) : ILargePayloadPurgeClient +{ + readonly LargePayloadPurgeClient client = Check.NotNull(client); + + /// + public async Task> GetLargePayloadTombstonesAsync( + int limit, DateTime deadline, CancellationToken cancellationToken = default) + { + using var call = this.client.GetLargePayloadTombstonesAsync( + new LP.GetLargePayloadTombstonesRequest { Limit = limit }, deadline: deadline, cancellationToken: cancellationToken); + LP.GetLargePayloadTombstonesResponse response = await call; + List tombstones = new(response.Tombstones.Count); + foreach (LP.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) + { + LP.ReportLargePayloadPurgeResultsRequest request = new(); + foreach (LargePayloadPurgeResult result in results) + { + request.Results.Add(new LP.LargePayloadPurgeResult + { + // Echo the opaque correlation token unchanged. The managed and protobuf enums share values. + TombstoneToken = result.TombstoneToken, + Disposition = (LP.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..604d44b33 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( + sp.GetRequiredKeyedService(builder.Name), + sp.GetRequiredService>())); r.AddActivity(); r.AddActivity(nameof(ReportLargePayloadPurgeResultsActivity), sp => - ActivatorUtilities.CreateInstance( - sp, sp.GetRequiredKeyedService(builder.Name))); + new ReportLargePayloadPurgeResultsActivity( + 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..a545401ce --- /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 task-hub-bound transport operations 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..e42ad5f12 --- /dev/null +++ b/src/LargePayloadPurge.Abstractions/IOrchestrationServiceLargePayloadPurgeClient.cs @@ -0,0 +1,56 @@ +// ---------------------------------------------------------------------------------- +// Copyright Microsoft Corporation +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// http://www.apache.org/licenses/LICENSE-2.0 +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// ---------------------------------------------------------------------------------- +// Adapted to the Durable Task .NET SDK repository's namespace and using conventions. + +using Microsoft.DurableTask.Client; + +namespace DurableTask.LargePayloadPurge; + +/// +/// Optional orchestration service client capability for purging tombstoned large payloads. +/// +public interface IOrchestrationServiceLargePayloadPurgeClient +{ + /// + /// 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); + + /// + /// Gets tombstoned large payloads that are ready to be purged. + /// + /// The maximum number of tombstones to return. + /// The caller's operation deadline in UTC, or + /// when the caller has not specified a deadline. + /// The token used to cancel the operation. + /// The tombstones to process. + Task> GetLargePayloadsToPurgeAsync( + int limit, DateTime deadlineUtc, CancellationToken cancellationToken); + + /// + /// Reports the outcomes of attempts to purge tombstoned large payloads. + /// + /// The purge outcomes, including the unchanged tombstone correlation tokens. + /// 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 ReportLargePayloadPurgeResultsAsync( + IReadOnlyList results, DateTime deadlineUtc, CancellationToken cancellationToken); +} diff --git a/src/LargePayloadPurge.Abstractions/LICENSE b/src/LargePayloadPurge.Abstractions/LICENSE new file mode 100644 index 000000000..dd5b3a58a --- /dev/null +++ b/src/LargePayloadPurge.Abstractions/LICENSE @@ -0,0 +1,174 @@ + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. diff --git a/src/LargePayloadPurge.Abstractions/LargePayloadPurge.Abstractions.csproj b/src/LargePayloadPurge.Abstractions/LargePayloadPurge.Abstractions.csproj new file mode 100644 index 000000000..48476c6e5 --- /dev/null +++ b/src/LargePayloadPurge.Abstractions/LargePayloadPurge.Abstractions.csproj @@ -0,0 +1,24 @@ + + + + + netstandard2.0 + DurableTask.LargePayloadPurge.Abstractions + DurableTask.LargePayloadPurge + Microsoft.Azure.DurableTask.LargePayloadPurge.Abstractions + Service and activity transport contracts for large payload purge using the Durable Task SDK client models. + Apache-2.0 AND MIT + 0.1.0 + + true + + $(NoWarn);SA1636 + + + + + + + + + diff --git a/src/LargePayloadPurge.Abstractions/README.md b/src/LargePayloadPurge.Abstractions/README.md new file mode 100644 index 000000000..1c9e99dee --- /dev/null +++ b/src/LargePayloadPurge.Abstractions/README.md @@ -0,0 +1,51 @@ +# Large payload purge contracts + +`Microsoft.Azure.DurableTask.LargePayloadPurge.Abstractions` provides two interfaces in the +`DurableTask.LargePayloadPurge.Abstractions` assembly: + +- `DurableTask.LargePayloadPurge.IOrchestrationServiceLargePayloadPurgeClient` is the optional service + capability for explicit auto-purge state, tombstone fetches, and outcome reports. Each operation requires + a UTC deadline and cancellation token. `DateTime.MaxValue` represents an unspecified deadline. + Setting the flag alone does not start or stop a purge runner. +- `Microsoft.DurableTask.AzureBlobPayloads.ILargePayloadPurgeClient` is the activity transport contract for + fetching tombstones 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. + +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.Azure.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 an earlier local prototype signed with +a different key must rebuild against the SDK-owned package; there is no compatibility promise for those +unreleased prototype binaries. The activity transport interface also moved here from the unreleased +Azure Blob implementation while retaining its full namespace, methods and optional-parameter defaults. +Rebuild consumers against this assembly; no type forwarder is provided. + +## License + +The service interface originated in [Azure/durabletask](https://github.com/Azure/durabletask) and retains +its Apache-2.0 notice; see [LICENSE](LICENSE). The activity transport interface retains its MIT notice; +see the [SDK MIT license](https://github.com/microsoft/durabletask-dotnet/blob/main/LICENSE). Both license texts +are included in the package as `LICENSE` and `licenses/MIT/LICENSE`; the package license expression +is `Apache-2.0 AND MIT`. The referenced SDK model assemblies retain their own licenses. 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/LargePayloadPurge.Abstractions.Tests/LargePayloadPurge.Abstractions.Tests.csproj b/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurge.Abstractions.Tests.csproj new file mode 100644 index 000000000..2cb181172 --- /dev/null +++ b/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurge.Abstractions.Tests.csproj @@ -0,0 +1,14 @@ + + + + net10.0 + DurableTask.LargePayloadPurge.Abstractions.Tests + 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..574831272 --- /dev/null +++ b/test/LargePayloadPurge.Abstractions.Tests/LargePayloadPurgeContractTests.cs @@ -0,0 +1,178 @@ +// ---------------------------------------------------------------------------------- +// Copyright Microsoft Corporation +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// http://www.apache.org/licenses/LICENSE-2.0 +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// ---------------------------------------------------------------------------------- +// Adapted from the original contract tests for SDK xUnit, signing and dependency checks. + +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("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.Empty(contract.GetInterfaces()); + Assert.Empty(transport.GetInterfaces()); + Assert.Equal(3, contract.GetMethods().Length); + 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 + AssertSignature( + nameof(IOrchestrationServiceLargePayloadPurgeClient.GetLargePayloadsToPurgeAsync), + typeof(Task>), + [typeof(int), typeof(DateTime), typeof(CancellationToken)], + ["limit", "deadlineUtc", "cancellationToken"]); + AssertTransportSignature( + nameof(ILargePayloadPurgeClient.GetLargePayloadTombstonesAsync), + typeof(Task>), + [typeof(int), typeof(DateTime), typeof(CancellationToken)], + ["limit", "deadline", "cancellationToken"]); + } + + [Fact] + public void ReportAcceptsCanonicalSdkResults() + { + // Arrange / Act / Assert + AssertSignature( + nameof(IOrchestrationServiceLargePayloadPurgeClient.ReportLargePayloadPurgeResultsAsync), + typeof(Task), + [typeof(IReadOnlyList), typeof(DateTime), typeof(CancellationToken)], + ["results", "deadlineUtc", "cancellationToken"]); + 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); + 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); + } +}