Skip to content
Merged
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
4 changes: 2 additions & 2 deletions decisions/2026-09-24-outbox-unresolvable-event-types.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ applyTo:

created: 2026-09-24

lastModified: 2026-09-26
lastModified: 2026-09-28

state: accepted

Expand Down Expand Up @@ -55,7 +55,7 @@ The Entity Framework Core provider materializes `OutboxMessage.EventType` throug
* An unresolvable event type no longer blocks the other messages of a batch in any provider.
* The behavior is consistent across providers and covered by shared integration tests in `OutboxTestsBase`.
* `GetPendingAsync` and `GetFailedForRetryAsync` may return fewer messages than claimed.
* If dead-lettering fails or is interrupted after the claim, only the unresolvable messages stay in `Processing`; they are reclaimed after the processing lease expires by the providers that reclaim expired leases. The Entity Framework Core and Cosmos DB fetch queries do not reclaim expired leases today, a pre-existing gap outside this decision.
* If dead-lettering fails or is interrupted after the claim, only the unresolvable messages stay in `Processing`; every provider reclaims them after the processing lease expires (see [Entity Framework Core and Cosmos DB Outbox Lease Reclaim](./2026-09-28-entityframework-and-cosmosdb-outbox-lease-reclaim.md)).
* A batch whose messages are all unresolvable returns an empty list, so the processor waits one polling interval before the next fetch. A large backlog of unresolvable messages is therefore drained at one batch per polling interval.
* The placeholder compares equal only to placeholders with the same stored name, not to `typeof(object)`.
* `NetEvolve.Pulse.Extensibility` gains a public static class with `Resolve` and `DeadLetterUnresolvableAsync`. External implementers of `IOutboxRepository` SHOULD use it in the same way. No interface changes.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
---
authors:
- Martin Stühmer

applyTo:
- "src/NetEvolve.Pulse.EntityFramework/Outbox/EntityFrameworkOutboxRepository*.cs"
- "src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxRepository.cs"
- "src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptions.cs"

created: 2026-09-28

lastModified: 2026-09-28

state: proposed

instructions: |
MUST reclaim Processing outbox messages whose UpdatedAt is at or before now minus ProcessingLeaseTimeout in the Entity Framework Core and Cosmos DB GetPendingAsync, as part of the same claim predicate as Pending messages; MUST NOT add a lease column.
MUST read the Entity Framework Core lease from OutboxOptions.ProcessingLeaseTimeout and the Cosmos DB lease from CosmosDbOutboxOptions.ProcessingLeaseTimeout, and MUST reject values that are not greater than TimeSpan.Zero.
MUST return the Cosmos DB documents already claimed when a later claim patch fails with a CosmosException, and MUST rethrow when nothing was claimed yet.
---

# Decision: Entity Framework Core and Cosmos DB Outbox Lease Reclaim

The Entity Framework Core and Cosmos DB outbox repositories reclaim messages that stay in `Processing` longer than the processing lease. `UpdatedAt` serves as the lease timestamp. A Cosmos DB claim that fails part-way returns the documents it has already claimed.

## Context

The SQL Server, PostgreSQL, MySQL, SQLite and MongoDB repositories reclaim `Processing` messages whose lease has expired. The Entity Framework Core and Cosmos DB repositories selected only `Pending` messages in `GetPendingAsync`. A message stayed in `Processing` forever after a graceful shutdown mid-batch, a cancellation between send and completion, a failure after the claim committed, or a crashed host. `IOutboxManagement` replays and dismisses only dead-letter messages, so recovery needed a manual database edit (issue #813).

The Cosmos DB claim patches the candidates one at a time. A `429 Too Many Requests` or `503 Service Unavailable` on a later candidate threw away the documents that the earlier patches had already moved to `Processing`.

## Decision

- `GetPendingAsync` selects `(Status = Pending AND (NextRetryAt IS NULL OR NextRetryAt <= now)) OR (Status = Processing AND UpdatedAt <= now - ProcessingLeaseTimeout)` in both providers.
- Entity Framework Core passes this predicate as the claim filter. The bulk executor re-checks it in the claiming `UPDATE`, as required by [Entity Framework Outbox Claim Concurrency](./2026-09-28-entityframework-outbox-claim-concurrency.md). The reclaim writes a fresh `UpdatedAt`, so a competing poller no longer matches the row. The change-tracking executors reject the same race through the `Status` and `UpdatedAt` concurrency tokens.
- Cosmos DB keeps the `IfMatchEtag` claim, which rejects a second reclaim of the same document.
- Entity Framework Core reads `OutboxOptions.ProcessingLeaseTimeout`, like the ADO.NET providers. Cosmos DB reads the new `CosmosDbOutboxOptions.ProcessingLeaseTimeout`, like MongoDB reads `MongoDbOutboxOptions.ProcessingLeaseTimeout`, because Cosmos DB applications configure only `CosmosDbOutboxOptions`. Both default to 5 minutes, and both repositories reject values that are not greater than `TimeSpan.Zero`.
- When a Cosmos DB claim patch fails with a `CosmosException` other than `412 Precondition Failed` or `404 Not Found` after at least one document was claimed, the claim stops and returns the documents claimed so far. When nothing was claimed yet, the exception propagates. Cancellation always propagates.

## Consequences

- The Entity Framework Core and Cosmos DB outboxes keep at-least-once delivery after shutdowns, cancellations and crashes. The shared tests in `OutboxTestsBase` cover the reclaim and the active lease for every provider, and `EntityFrameworkOutboxClaimRaceTestsBase` covers overlapping reclaims.
- The schema does not change. Documents and rows that were stranded in `Processing` before the upgrade are reclaimed on the first poll after the upgrade.
- A message whose dispatch takes longer than the lease can be dispatched twice. This matches the other providers. Choose the lease comfortably larger than the longest dispatch.
- `UpdatedAt` also changes on every other status transition, so the lease starts at the claim. A claim that is not renewed during a long dispatch expires after the lease timeout.
- `GetPendingCountAsync` still counts only `Pending` messages. Expired `Processing` messages do not appear in the pending count until they are reclaimed.
- The Cosmos DB pending query orders by `_ts` and filters on `status`, `nextRetryAt` and `updatedAt`. It needs no composite index under the default indexing policy.
- `CosmosDbOutboxOptions` gains a public property. External implementers of `IOutboxRepository` are not affected.

## Alternatives Considered

- **A dedicated lease column (for example `ProcessingStartedAt`, as in MongoDB).** It separates the lease from other updates, but requires a schema change for every existing installation. `UpdatedAt` is written on every claim and already serves as the claim token.
- **Cosmos DB reading `OutboxOptions.ProcessingLeaseTimeout`.** This adds no public API, but `AddCosmosDbOutbox` users never configure `OutboxOptions`, so the setting would sit apart from the rest of the provider configuration.
- **Reclaiming in `GetFailedForRetryAsync`.** The other providers reclaim in the pending poll. Doing the same keeps the processor behavior identical across providers.
- **Retrying throttled Cosmos DB patches in the claim.** The SDK already retries throttled requests. Returning the claimed documents keeps the claim simple, and the next poll picks up the remaining candidates.

## Related Decisions (Optional)

- [Entity Framework Outbox Claim Concurrency](./2026-09-28-entityframework-outbox-claim-concurrency.md) - The widened claim predicate reaches the claiming `UPDATE` through the claim filter defined there.
- [Outbox Unresolvable Event Types](./2026-09-24-outbox-unresolvable-event-types.md) - Unresolvable messages left in `Processing` are now reclaimed by every provider.
- [DateTimeOffset and TimeProvider Usage](./2026-01-21-datetimeoffset-and-timeprovider-usage.md) - The lease cutoff comes from `TimeProvider`.
23 changes: 23 additions & 0 deletions src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,20 @@ public sealed class CosmosDbOutboxOptions
/// </remarks>
public int TtlSeconds { get; set; } = DefaultTtlSeconds;

/// <summary>
/// Gets or sets the maximum duration a claimed message may remain in the
/// <see cref="Extensibility.Outbox.OutboxMessageStatus.Processing"/> status
/// before it becomes eligible for reclaiming by a subsequent pending poll.
/// Default: 5 minutes.
/// </summary>
/// <remarks>
/// When a worker crashes or is cancelled after claiming a message but before completing it,
/// the message stays in the <c>Processing</c> status. Once this lease expires, the next pending
/// poll claims the message again, preserving at-least-once delivery. Choose a value comfortably
/// larger than the longest expected message dispatch duration to avoid duplicate publishing.
/// </remarks>
public TimeSpan ProcessingLeaseTimeout { get; set; } = TimeSpan.FromMinutes(5);

/// <summary>
/// Gets the validation error for <see cref="PartitionKeyPath"/>, or <see langword="null"/> when it is
/// <see cref="DefaultPartitionKeyPath"/>.
Expand All @@ -95,6 +109,15 @@ public sealed class CosmosDbOutboxOptions
? null
: $"{nameof(CosmosDbOutboxOptions)}.{nameof(PartitionKeyPath)} '{PartitionKeyPath}' is not supported. The only supported value is '{DefaultPartitionKeyPath}', and the container must be partitioned on '{DefaultPartitionKeyPath}'.";

/// <summary>
/// Gets the validation error for <see cref="ProcessingLeaseTimeout"/>, or <see langword="null"/> when it is
/// greater than <see cref="TimeSpan.Zero"/>.
/// </summary>
internal string? ProcessingLeaseTimeoutError =>
ProcessingLeaseTimeout > TimeSpan.Zero
? null
: $"{nameof(CosmosDbOutboxOptions)}.{nameof(ProcessingLeaseTimeout)} must be greater than zero, but was '{ProcessingLeaseTimeout}'.";

/// <summary>
/// Throws when <see cref="PartitionKeyPath"/> is not <see cref="DefaultPartitionKeyPath"/>.
/// </summary>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,5 +9,7 @@ internal sealed class CosmosDbOutboxOptionsValidator : IValidateOptions<CosmosDb
{
/// <inheritdoc />
public ValidateOptionsResult Validate(string? name, CosmosDbOutboxOptions options) =>
options.PartitionKeyPathError is { } error ? ValidateOptionsResult.Fail(error) : ValidateOptionsResult.Success;
(options.PartitionKeyPathError ?? options.ProcessingLeaseTimeoutError) is { } error
? ValidateOptionsResult.Fail(error)
: ValidateOptionsResult.Success;
}
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,11 @@
/// <para><strong>Concurrency:</strong></para>
/// Uses ETag-based conditional patch to atomically claim pending messages
/// and prevent duplicate processing by concurrent workers.
/// <para><strong>Lease-based reclaim:</strong></para>
/// <see cref="GetPendingAsync"/> also reclaims documents that stayed in the <c>Processing</c> status
/// longer than <see cref="CosmosDbOutboxOptions.ProcessingLeaseTimeout"/> (for example, after a worker
/// crash, cancellation or shutdown), so they are delivered at least once. The <c>updatedAt</c> field
/// serves as the lease timestamp, and the ETag precondition keeps concurrent reclaims exclusive.
/// <para><strong>TTL Support:</strong></para>
/// When <see cref="CosmosDbOutboxOptions.EnableTimeToLive"/> is <see langword="true"/>,
/// the <c>ttl</c> field is set on completed and dead-letter documents so the Cosmos DB
Expand All @@ -42,6 +47,7 @@ internal sealed class CosmosDbOutboxRepository : IOutboxRepository
private readonly TimeProvider _timeProvider;
private readonly bool _enableTtl;
private readonly int _ttlSeconds;
private readonly TimeSpan _processingLeaseTimeout;

/// <summary>
/// Initializes a new instance of the <see cref="CosmosDbOutboxRepository"/> class.
Expand All @@ -63,11 +69,13 @@ TimeProvider timeProvider
ArgumentException.ThrowIfNullOrWhiteSpace(opts.DatabaseName);
ArgumentException.ThrowIfNullOrWhiteSpace(opts.ContainerName);
opts.ThrowIfPartitionKeyPathIsNotSupported();
ArgumentOutOfRangeException.ThrowIfLessThanOrEqual(opts.ProcessingLeaseTimeout, TimeSpan.Zero);

_container = cosmosClient.GetContainer(opts.DatabaseName, opts.ContainerName);
_timeProvider = timeProvider;
_enableTtl = opts.EnableTimeToLive;
_ttlSeconds = opts.TtlSeconds;
_processingLeaseTimeout = opts.ProcessingLeaseTimeout;
}

/// <inheritdoc />
Expand All @@ -93,9 +101,10 @@ public async Task<IReadOnlyList<OutboxMessage>> GetPendingAsync(
var now = _timeProvider.GetUtcNow();

var query = new QueryDefinition(
"SELECT * FROM c WHERE c.status = 0 AND (IS_NULL(c.nextRetryAt) OR c.nextRetryAt <= @now) ORDER BY c._ts ASC OFFSET 0 LIMIT @batchSize"
"SELECT * FROM c WHERE (c.status = 0 AND (IS_NULL(c.nextRetryAt) OR c.nextRetryAt <= @now)) OR (c.status = 1 AND c.updatedAt <= @leaseExpiredBefore) ORDER BY c._ts ASC OFFSET 0 LIMIT @batchSize"
)
.WithParameter("@now", now)
.WithParameter("@leaseExpiredBefore", now - _processingLeaseTimeout)
.WithParameter("@batchSize", batchSize);

var candidates = await ExecuteQueryAsync(query, CreateBatchQueryOptions(batchSize), cancellationToken)
Expand Down Expand Up @@ -401,6 +410,13 @@ CancellationToken cancellationToken
/// as an <c>IfMatchEtag</c> precondition, avoiding an additional point read per candidate.
/// Messages that have been modified by another worker (ETag mismatch) are silently skipped.
/// </summary>
/// <remarks>
/// When a later patch fails (for example, throttled with <c>429</c>) after earlier candidates were
/// claimed, the documents claimed so far are returned instead of the exception, so they are dispatched
/// right away rather than waiting for their processing lease to expire. A failure before the first
/// claim, and cancellation, still propagate; documents left in <c>Processing</c> by them are reclaimed
/// once <see cref="CosmosDbOutboxOptions.ProcessingLeaseTimeout"/> expires.
/// </remarks>
private async Task<IReadOnlyList<OutboxMessage>> ClaimMessagesAsync(
IReadOnlyList<CosmosDbOutboxDocument> candidates,
int targetStatus,
Expand Down Expand Up @@ -447,6 +463,11 @@ CancellationToken cancellationToken
{
// Document deleted between read and patch — skip.
}
catch (CosmosException) when (claimed.Count > 0)
{
// Hand out what is already claimed; the remaining candidates stay eligible for the next poll.
break;
}
}

// Candidates arrive in _ts order; the outbox contract requires CreatedAt order.
Expand Down
1 change: 1 addition & 0 deletions src/NetEvolve.Pulse.CosmosDb/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,7 @@ services.AddPulse(config => config
| `PartitionKeyPath` | `string` | `/id` | Only `/id` is supported; other values fail validation at startup. The container must use `/id` (see [Container Setup](#container-setup)). |
| `EnableTimeToLive` | `bool` | `false` | Sets the `ttl` property on documents that become `Completed` or `DeadLetter`, so the Cosmos DB TTL engine deletes them. Replaying a dead-letter message sets its `ttl` to `-1`, so it does not expire while pending. Requires `DefaultTimeToLive` on the container. |
| `TtlSeconds` | `int` | `86400` (24 hours) | TTL in seconds for completed and dead-letter documents. Only applies when `EnableTimeToLive` is `true`. |
| `ProcessingLeaseTimeout` | `TimeSpan` | 5 minutes | How long a claimed message may stay in `Processing` before the next pending poll reclaims it, for example after a crash or shutdown. Must be greater than zero; other values fail validation at startup. |

## Concurrency

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

using System.Linq.Expressions;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Options;
using NetEvolve.Pulse.Extensibility.Outbox;

/// <summary>
Expand All @@ -17,6 +18,11 @@
/// Bulk executors re-check the full eligibility predicate on every claimed row; change-tracking
/// executors rely on optimistic concurrency tokens. For high-throughput scenarios,
/// consider using the SQL Server ADO.NET provider with explicit locking.
/// <para><strong>Lease-based reclaim:</strong></para>
/// <see cref="GetPendingAsync"/> also reclaims messages that stayed in the <c>Processing</c> status
/// longer than <see cref="OutboxOptions.ProcessingLeaseTimeout"/> (for example, after a worker crash,
/// cancellation or shutdown), so they are delivered at least once. <see cref="OutboxMessage.UpdatedAt"/>
/// serves as the lease timestamp.
/// </remarks>
/// <typeparam name="TContext">The DbContext type that implements <see cref="IOutboxDbContext"/>.</typeparam>
internal sealed class EntityFrameworkOutboxRepository<TContext> : IOutboxRepository, IDisposable
Expand All @@ -33,20 +39,29 @@ internal sealed class EntityFrameworkOutboxRepository<TContext> : IOutboxReposit
/// (change-tracking + <c>SaveChangesAsync</c> vs. bulk <c>ExecuteUpdate</c> / <c>ExecuteDelete</c>).
/// </summary>
private readonly IOutboxRepositoryExecutor _executor;

/// <summary>
/// The duration after which a message in the <c>Processing</c> status is eligible for reclaiming.
/// </summary>
private readonly TimeSpan _processingLeaseTimeout;
private bool _disposedValue;

/// <summary>
/// Initializes a new instance of the <see cref="EntityFrameworkOutboxRepository{TContext}"/> class.
/// </summary>
/// <param name="context">The DbContext for database operations.</param>
/// <param name="options">The outbox configuration options.</param>
/// <param name="timeProvider">The time provider for timestamps.</param>
public EntityFrameworkOutboxRepository(TContext context, TimeProvider timeProvider)
public EntityFrameworkOutboxRepository(TContext context, IOptions<OutboxOptions> options, TimeProvider timeProvider)
{
ArgumentNullException.ThrowIfNull(context);
ArgumentNullException.ThrowIfNull(options);
ArgumentNullException.ThrowIfNull(timeProvider);
ArgumentOutOfRangeException.ThrowIfLessThanOrEqual(options.Value.ProcessingLeaseTimeout, TimeSpan.Zero);

_context = context;
_timeProvider = timeProvider;
_processingLeaseTimeout = options.Value.ProcessingLeaseTimeout;
_executor = context.Database.ProviderName switch
{
// InMemory does not support ExecuteUpdate/ExecuteDelete at all.
Expand Down Expand Up @@ -81,9 +96,13 @@ public async Task<IReadOnlyList<OutboxMessage>> GetPendingAsync(
cancellationToken.ThrowIfCancellationRequested();

var now = _timeProvider.GetUtcNow();
var leaseExpiredBefore = now - _processingLeaseTimeout;

// Bulk executors re-check this predicate in the claiming UPDATE, and a reclaim writes a fresh
// UpdatedAt, so a competing poller cannot reclaim the same expired message twice.
Expression<Func<OutboxMessage, bool>> claimFilter = m =>
m.Status == OutboxMessageStatus.Pending && (m.NextRetryAt == null || m.NextRetryAt <= now);
(m.Status == OutboxMessageStatus.Pending && (m.NextRetryAt == null || m.NextRetryAt <= now))
|| (m.Status == OutboxMessageStatus.Processing && m.UpdatedAt <= leaseExpiredBefore);

var baseQuery = _context.OutboxMessages.Where(claimFilter).OrderBy(m => m.CreatedAt).Take(batchSize);

Expand Down
13 changes: 13 additions & 0 deletions src/NetEvolve.Pulse.EntityFramework/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -306,6 +306,19 @@ services.AddPulse(config => config
);
```

## Processing Lease Reclaim

A message claimed by `GetPendingAsync` stays in `Processing` until it is completed or failed. If a worker crashes or shuts down in between, the next pending poll reclaims the message once its `UpdatedAt` is older than `OutboxOptions.ProcessingLeaseTimeout` (default: 5 minutes, must be greater than zero).

```csharp
services.AddPulse(config => config
.AddEntityFrameworkOutbox<ApplicationDbContext>(options =>
options.ProcessingLeaseTimeout = TimeSpan.FromMinutes(10))
);
```

Choose a value well above the longest expected dispatch. A dispatch that runs longer than the lease can be reclaimed by another poller and delivered twice.

## Requirements

- .NET 8.0, .NET 9.0, or .NET 10.0
Expand Down
Loading
Loading