Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
306d19f
Add opt-in blob payload auto-purge job to AzureBlobPayloads
YunchuWang Jul 8, 2026
60f6637
Address PR #758 review feedback: naming, self-heal, poison-ack, start…
YunchuWang Jul 14, 2026
149c63a
Refine BlobPurgeJobStarter: pre-check bridge status before rescheduling
YunchuWang Jul 14, 2026
4d52005
Classify RequestFailedException 400 as permanent in DeleteExternalBlo…
YunchuWang Jul 14, 2026
780d743
Reuse shared PayloadStore and register purge starter conditionally on…
YunchuWang Jul 14, 2026
3a2215c
Register fallback PayloadStore in shared Core for both client and worker
YunchuWang Jul 14, 2026
3fbf061
Merge branch 'main' into yunchuwang-wangbill-blob-payload-autopurge-sdk
YunchuWang Jul 14, 2026
47651dc
Stop self-registering PayloadStore on the client; consume the shared …
YunchuWang Jul 14, 2026
7397fa6
Validate PayloadPurgeBatchSize once at specification (fail fast on ou…
YunchuWang Jul 15, 2026
4afeb8a
Translate gRPC Cancelled to OperationCanceledException in GetTombston…
YunchuWang Jul 15, 2026
e74f633
Raise auto-purge MaxBatchSize to 1000 (inclusive); relax gRPC GetTomb…
YunchuWang Jul 15, 2026
50ae944
Register PayloadStore in the client builder extension (symmetry with …
YunchuWang Jul 30, 2026
a680442
Merge branch 'main' into yunchuwang-wangbill-blob-payload-autopurge-sdk
YunchuWang Jul 31, 2026
a5ed298
Resolve v2 tokens in DeleteAsync and discard payloads in unreachable …
YunchuWang Jul 31, 2026
fff06b0
Align unreachable-account log, exception wording and v2 delete test n…
YunchuWang Jul 31, 2026
6e4f5d0
Gate blob auto-purge on v2 tokens and deleting stores
YunchuWang Jul 31, 2026
65e9cbb
Fix auto-purge starter registration and client resolution
YunchuWang Jul 31, 2026
8f436df
Back off on zero-ack purge cycles and document TokenPrefixV1
YunchuWang Jul 31, 2026
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
24 changes: 24 additions & 0 deletions src/Client/Core/DurableTaskClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -549,6 +549,30 @@ public virtual Task<Page<string>> ListInstanceIdsAsync(
$"{this.GetType()} does not support listing orchestration instance IDs filtered by completed time.");
}

/// <summary>
/// Gets a batch of tombstoned (soft-deleted) externalized payloads whose backing blobs should be deleted
/// by a credentialed caller before the backend hard-deletes the rows.
/// </summary>
/// <param name="limit">The maximum number of tombstoned payloads to request.</param>
/// <param name="cancellation">The cancellation token.</param>
/// <returns>The batch of tombstoned payloads whose blobs should be deleted.</returns>
/// <exception cref="NotSupportedException">Thrown if this implementation does not support the operation.</exception>
public virtual Task<List<TombstonedPayload>> GetTombstonedPayloadsAsync(
int limit, CancellationToken cancellation = default)
=> throw new NotSupportedException($"{this.GetType()} does not support retrieving tombstoned payloads.");

/// <summary>
/// Acknowledges tombstoned payloads whose backing blobs have been deleted so the backend can hard-delete
/// the corresponding rows.
/// </summary>
/// <param name="acks">The payloads whose blobs have been deleted.</param>
/// <param name="cancellation">The cancellation token.</param>
/// <returns>A task that completes when the acknowledgement has been sent.</returns>
/// <exception cref="NotSupportedException">Thrown if this implementation does not support the operation.</exception>
public virtual Task AckPurgedPayloadsAsync(
IEnumerable<PayloadPurgeAck> acks, CancellationToken cancellation = default)
=> throw new NotSupportedException($"{this.GetType()} does not support acknowledging purged payloads.");

// TODO: Create task hub

// TODO: Delete task hub
Expand Down
13 changes: 13 additions & 0 deletions src/Client/Core/PayloadPurgeAck.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
// Copyright (c) Microsoft Corporation.
// Licensed under the MIT License.

namespace Microsoft.DurableTask.Client;

/// <summary>
/// Serializable acknowledgement that the worker has deleted the blob for a tombstoned payload, so the
/// backend can hard-delete the soft-deleted row. Mirrors the <c>PayloadPurgeAck</c> protobuf message.
/// </summary>
/// <param name="PartitionId">The backend partition that owns the payload row.</param>
/// <param name="InstanceKey">The orchestration instance key the payload belongs to.</param>
/// <param name="PayloadId">The backend identifier of the soft-deleted payload row.</param>
public sealed record PayloadPurgeAck(int PartitionId, long InstanceKey, long PayloadId);
15 changes: 15 additions & 0 deletions src/Client/Core/TombstonedPayload.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
// Copyright (c) Microsoft Corporation.
// Licensed under the MIT License.

namespace Microsoft.DurableTask.Client;

/// <summary>
/// Serializable representation of a tombstoned payload the backend has soft-deleted and whose blob the
/// worker should delete. Mirrors the <c>TombstonedPayload</c> protobuf message but is safe to pass through
/// the orchestration/activity boundary.
/// </summary>
/// <param name="PartitionId">The backend partition that owns the payload row.</param>
/// <param name="InstanceKey">The orchestration instance key the payload belongs to.</param>
/// <param name="PayloadId">The backend identifier of the soft-deleted payload row.</param>
/// <param name="Token">The externalized payload token whose backing blob should be deleted.</param>
public sealed record TombstonedPayload(int PartitionId, long InstanceKey, long PayloadId, string Token);
66 changes: 66 additions & 0 deletions src/Client/Grpc/GrpcDurableTaskClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -624,6 +624,72 @@ public override async Task<IList<HistoryEvent>> GetOrchestrationHistoryAsync(
}
}

/// <inheritdoc/>
public override async Task<List<TombstonedPayload>> GetTombstonedPayloadsAsync(
int limit, CancellationToken cancellation = default)
{
if (limit <= 0 || limit > 1000)
{
throw new ArgumentOutOfRangeException(
nameof(limit), limit, "Limit must be greater than 0 and less than or equal to 1000.");
}

P.GetTombstonedPayloadsResponse response;
try
{
response = await this.sidecarClient.GetTombstonedPayloadsAsync(
new P.GetTombstonedPayloadsRequest { Limit = limit },
cancellationToken: cancellation);
}
catch (RpcException e) when (e.StatusCode == StatusCode.Cancelled)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Medium / compatibility] During a mixed rollout, an older backend returns gRPC Unimplemented for this new RPC. Only cancellation is translated here, so the activity/orchestrator retries the unsupported operation indefinitely without a clear terminal diagnostic. Please add capability negotiation or map Unimplemented to an explicit unsupported-backend state that stops/disables the purge job.

{
throw new OperationCanceledException(
$"The {nameof(this.GetTombstonedPayloadsAsync)} operation was canceled.", e, cancellation);
}

List<TombstonedPayload> result = new(response.Payloads.Count);
foreach (P.TombstonedPayload payload in response.Payloads)
{
result.Add(new TombstonedPayload(
payload.PartitionId, payload.InstanceKey, payload.PayloadId, payload.Token));
}

return result;
}

/// <inheritdoc/>
public override async Task AckPurgedPayloadsAsync(
IEnumerable<PayloadPurgeAck> acks, CancellationToken cancellation = default)
{
Check.NotNull(acks);

P.AckPurgedPayloadsRequest request = new();
foreach (PayloadPurgeAck ack in acks)
{
request.Acks.Add(new P.PayloadPurgeAck
{
PartitionId = ack.PartitionId,
InstanceKey = ack.InstanceKey,
PayloadId = ack.PayloadId,
});
}

if (request.Acks.Count == 0)
{
return;
}

try
{
await this.sidecarClient.AckPurgedPayloadsAsync(request, cancellationToken: cancellation);
}
catch (RpcException e) when (e.StatusCode == StatusCode.Cancelled)
{
throw new OperationCanceledException(
$"The {nameof(this.AckPurgedPayloadsAsync)} operation was canceled.", e, cancellation);
}
}

static AsyncDisposable GetCallInvoker(GrpcDurableTaskClientOptions options, ILogger logger, out CallInvoker callInvoker)
{
Func<GrpcChannel, CancellationToken, Task<GrpcChannel>>? recreator = options.Internal.ChannelRecreator;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
// Copyright (c) Microsoft Corporation.
// Licensed under the MIT License.

using Microsoft.DurableTask.Client;
using Microsoft.Extensions.Logging;

namespace Microsoft.DurableTask.AzureBlobPayloads;

/// <summary>
/// Activity that acknowledges to the backend the payloads whose blobs the worker has deleted, so the backend
/// can hard-delete the soft-deleted rows.
/// </summary>
/// <param name="client">The Durable Task client used to acknowledge purged payloads to the backend.</param>
/// <param name="logger">The logger instance.</param>
[DurableTask]
public class AckPurgedPayloadsActivity(
DurableTaskClient client,
ILogger<AckPurgedPayloadsActivity> logger)
: TaskActivity<List<PayloadPurgeAck>, object?>
{
readonly DurableTaskClient client = Check.NotNull(client);
readonly ILogger<AckPurgedPayloadsActivity> logger = Check.NotNull(logger);

/// <inheritdoc/>
public override async Task<object?> RunAsync(TaskActivityContext context, List<PayloadPurgeAck> input)
{
if (input is null || input.Count == 0)
{
return null;
}

await this.client.AckPurgedPayloadsAsync(input, CancellationToken.None);
this.logger.BlobPurgeAckedPayloads(input.Count);
return null;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
// Copyright (c) Microsoft Corporation.
// Licensed under the MIT License.

using Azure;
using Microsoft.Extensions.Logging;

namespace Microsoft.DurableTask.AzureBlobPayloads;

/// <summary>
/// Activity that deletes a single externalized payload blob given its token. Deletion is idempotent, so
/// re-delivered tokens and concurrent workers are safe.
/// </summary>
/// <remarks>
/// Outcome classification, verified against the Azure.Storage.Blobs / Azure.Core exception model (not
/// assumed):
/// <list type="bullet">
/// <item>
/// The Azure SDK already retries transient failures internally (connection errors plus HTTP
/// 408/429/500/502/503/504, with exponential backoff), so any exception that escapes
/// <see cref="PayloadStore.DeleteAsync"/> means those built-in retries were already exhausted.
/// </item>
/// <item>
/// Legacy <c>blob:v1:</c> tokens are discarded without a delete attempt: a v1 token carries only a container
/// name, not the storage account, so a delete against the currently-configured account cannot be verified - if
/// the store has been repointed the delete would report success while the real blob survives elsewhere. The
/// backend ack protocol has no "skip" status (an un-acked row is re-served every cycle), so the token is acked
/// to keep the pipeline moving and logged at error level as the operator's recovery pointer.
/// </item>
/// <item>
/// Permanent failures are discarded (acked so the backend clears the row) because retrying can never succeed:
/// an <see cref="ArgumentException"/> from the store's token decode - a genuinely malformed or unrecognized
/// token (v1 tokens are gated out above and never reach the store); a
/// <see cref="RequestFailedException"/> with <see cref="RequestFailedException.Status"/> 400 (for example
/// InvalidUri / InvalidResourceName when the decoded blob name violates Azure naming rules); and a
/// <see cref="PayloadStorageException"/> when a v2 token points at a storage account the configured credential
/// cannot reach (connection-string / account-key auth is account-specific). The backend batch is cursor-less,
/// so an undroppable row would otherwise re-stream every cycle and block the pipeline head-of-line; the
/// account-unreachable case is logged at error level so an operator can reconcile it.
/// </item>
/// <item>
/// Everything else is treated as transient and leaves the payload tombstoned to retry on a later cycle:
/// throttling / 5xx that outlived the SDK's retries, 403 authorization failures (which need an operator
/// credential fix rather than dropping data), timeouts / cancellation, and a <see cref="NotSupportedException"/>
/// from a misconfigured store that cannot delete at all (retried, not dropped, so the work is recoverable once
/// a deleting store is registered). A blob is never dropped on an uncertain error, and a single bad token never
/// fails the whole batch.
/// </item>
/// </list>
/// </remarks>
/// <param name="store">The payload store used to delete blobs.</param>
/// <param name="logger">The logger instance.</param>
[DurableTask]
public class DeleteExternalBlobActivity(
Comment thread
YunchuWang marked this conversation as resolved.
PayloadStore store,
ILogger<DeleteExternalBlobActivity> logger)
: TaskActivity<string, BlobDeleteResult>
{
readonly PayloadStore store = Check.NotNull(store);
readonly ILogger<DeleteExternalBlobActivity> logger = Check.NotNull(logger);

/// <inheritdoc/>
public override async Task<BlobDeleteResult> RunAsync(TaskActivityContext context, string input)
{
Check.NotNullOrEmpty(input, nameof(input));

if (input.StartsWith(BlobPayloadStore.TokenPrefixV1, StringComparison.Ordinal))
{
// Auto-purge deliberately does not act on legacy v1 tokens. A v1 token carries only a container
// *name* - not the storage account - so a delete against the currently-configured account cannot be
// verified: if the store has since been repointed, DeleteIfExistsAsync returns false and the purge
// would silently report success while the real blob survives in the old account. Rather than delete
// on an unverifiable pointer, the token is discarded. The backend ack protocol carries no "skip"
// status (PayloadPurgeAck is just partition/instance/payload id, and an un-acked row is re-served by
// an uncursored TOP(N) query every cycle), so declining without acking would permanently block the
// pipeline. The full token is logged at error level so it remains a recoverable pointer.
this.logger.BlobPurgeDeleteV1TokenUnsupported(input);
return BlobDeleteResult.Discarded;
}

try
{
await this.store.DeleteAsync(input, CancellationToken.None);
return BlobDeleteResult.Deleted;
}
catch (ArgumentException ex)
Comment thread
YunchuWang marked this conversation as resolved.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[High / data retention] This also catches the ArgumentException that BlobPayloadStore.DeleteAsync throws when the token's container differs from the configured container. Returning Discarded causes the orchestrator to acknowledge and permanently remove the backend row even though the blob remains. A container rename or configuration mismatch can therefore orphan data. Use a distinct malformed-token exception; treat container mismatch as retryable/configuration-fatal and do not acknowledge it.

{
// The token is malformed or points at a different container; it can never succeed. Discard it so
// the backend clears the row instead of re-streaming the same poison token every cycle.
this.logger.BlobPurgeDeleteDiscarded(ex, input);
return BlobDeleteResult.Discarded;
}
catch (RequestFailedException ex) when (ex.Status == 400)
{
// Service rejected the request as permanently invalid (e.g. InvalidUri / InvalidResourceName - the
// decoded blob name violates Azure naming rules). Retrying can never succeed, so discard it like a
// poison token: ack so the backend clears the row instead of re-streaming it forever.
this.logger.BlobPurgeDeleteDiscarded(ex, input);
return BlobDeleteResult.Discarded;
}
catch (NotSupportedException ex)
{
// The registered store does not implement deletion (PayloadStore.DeleteAsync is virtual and its base
// implementation throws). This is a misconfiguration, not poison data: every payload would fail the
// same way, so acking would hard-delete the backend's entire record of what still needs cleanup while
// every blob survives. Keep the payload tombstoned so the work is recoverable once an operator
// registers a store that can delete.
this.logger.BlobPurgeDeleteNotSupported(ex, input);
return BlobDeleteResult.Retry;
}
catch (PayloadStorageException ex)
{
// The token is well-formed but points at a storage account this worker's credential cannot reach
// (cross-account without AAD). Retrying can never succeed from this process, and because the backend
// streams tombstones with an uncursored TOP(N) query, leaving it un-acked would re-serve the same
// token every cycle and permanently block the purge pipeline. Discard it so the row is cleared, and
// log at Error so an operator can reclaim the orphaned blob out-of-band.
this.logger.BlobPurgeDeleteUnreachable(ex, input);
return BlobDeleteResult.Discarded;
}
catch (Exception ex) when (ex is not OutOfMemoryException and not StackOverflowException)
{
// Transient failure: leave the payload tombstoned so a later purge cycle can retry it. A single
// bad token must not fail the whole batch.
this.logger.BlobPurgeDeleteFailed(ex, input);
return BlobDeleteResult.Retry;
}
Comment on lines +110 to +126
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
// Copyright (c) Microsoft Corporation.
// Licensed under the MIT License.

using Microsoft.DurableTask.Client;
using Microsoft.Extensions.Logging;

namespace Microsoft.DurableTask.AzureBlobPayloads;

/// <summary>
/// Activity that fetches a batch of tombstoned payloads from the backend for the auto-purge job to delete.
/// </summary>
/// <param name="client">The Durable Task client used to query the backend for tombstoned payloads.</param>
/// <param name="logger">The logger instance.</param>
[DurableTask]
public class GetTombstonedPayloadsActivity(
DurableTaskClient client,
ILogger<GetTombstonedPayloadsActivity> logger)
: TaskActivity<int, List<TombstonedPayload>>
{
readonly DurableTaskClient client = Check.NotNull(client);
readonly ILogger<GetTombstonedPayloadsActivity> logger = Check.NotNull(logger);

/// <inheritdoc/>
public override async Task<List<TombstonedPayload>> RunAsync(TaskActivityContext context, int input)
{
int limit = input;
List<TombstonedPayload> payloads =
await this.client.GetTombstonedPayloadsAsync(limit, CancellationToken.None);
this.logger.BlobPurgeFetchedTombstones(payloads.Count);
return payloads;
}
}
Loading
Loading