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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions Microsoft.DurableTask.sln
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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}
Expand Down
49 changes: 49 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
32 changes: 32 additions & 0 deletions azure-pipelines-release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
11 changes: 11 additions & 0 deletions doc/release_process.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<VersionPrefix>` and `<VersionSuffix>` 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:
Expand Down
25 changes: 24 additions & 1 deletion eng/publish/publish.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
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)
7 changes: 7 additions & 0 deletions eng/templates/build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -15,15 +14,30 @@ namespace Microsoft.DurableTask.AzureBlobPayloads;
/// </summary>
/// <param name="client">The large-payload purge service client used to query the backend for tombstones.</param>
/// <param name="logger">The logger instance.</param>
/// <remarks>
/// 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.
/// </remarks>
[DurableTask]
internal sealed class GetLargePayloadTombstonesActivity(
LargePayloadPurgeClient client,
public sealed class GetLargePayloadTombstonesActivity(
ILargePayloadPurgeClient client,
ILogger<GetLargePayloadTombstonesActivity> logger)
: TaskActivity<int, List<LargePayloadTombstone>>
{
readonly LargePayloadPurgeClient client = Check.NotNull(client);
readonly ILargePayloadPurgeClient client = Check.NotNull(client);
readonly ILogger<GetLargePayloadTombstonesActivity> logger = Check.NotNull(logger);

/// <summary>
/// Initializes a new instance of the <see cref="GetLargePayloadTombstonesActivity"/> class using the worker's transport.
/// </summary>
/// <param name="client">The worker's purge client.</param>
/// <param name="logger">The activity logger.</param>
internal GetLargePayloadTombstonesActivity(
LargePayloadPurgeClient client, ILogger<GetLargePayloadTombstonesActivity> logger)
: this(new GrpcLargePayloadPurgeClient(client), logger)
{
}

/// <summary>
/// Gets or sets the timeout for one backend RPC attempt.
/// </summary>
Expand All @@ -41,13 +55,10 @@ public override async Task<List<LargePayloadTombstone>> RunAsync(TaskActivityCon
nameof(input), input, $"Limit must be greater than 0 and less than or equal to {LargePayloadTombstone.MaxRequestLimit}.");
}

LP.GetLargePayloadTombstonesResponse response;
List<LargePayloadTombstone> 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)
{
Expand Down Expand Up @@ -76,12 +87,6 @@ public override async Task<List<LargePayloadTombstone>> RunAsync(TaskActivityCon
e);
}

List<LargePayloadTombstone> tombstones = new(response.Tombstones.Count);
Comment thread
YunchuWang marked this conversation as resolved.
foreach (LP.LargePayloadTombstone tombstone in response.Tombstones)
{
tombstones.Add(new LargePayloadTombstone(tombstone.TombstoneToken, tombstone.PayloadToken));
}

this.logger.BlobPurgeFetchedTombstones(tombstones.Count);
return tombstones;
}
Expand Down
Loading
Loading