diff --git a/decisions/2026-09-24-outbox-unresolvable-event-types.md b/decisions/2026-09-24-outbox-unresolvable-event-types.md index d46fc9d2..acaf8042 100644 --- a/decisions/2026-09-24-outbox-unresolvable-event-types.md +++ b/decisions/2026-09-24-outbox-unresolvable-event-types.md @@ -12,7 +12,7 @@ applyTo: created: 2026-09-24 -lastModified: 2026-09-26 +lastModified: 2026-09-28 state: accepted @@ -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. diff --git a/decisions/2026-09-28-entityframework-and-cosmosdb-outbox-lease-reclaim.md b/decisions/2026-09-28-entityframework-and-cosmosdb-outbox-lease-reclaim.md new file mode 100644 index 00000000..0b678c4a --- /dev/null +++ b/decisions/2026-09-28-entityframework-and-cosmosdb-outbox-lease-reclaim.md @@ -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`. diff --git a/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptions.cs b/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptions.cs index 0a828be9..0899fb6b 100644 --- a/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptions.cs +++ b/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptions.cs @@ -81,6 +81,20 @@ public sealed class CosmosDbOutboxOptions /// public int TtlSeconds { get; set; } = DefaultTtlSeconds; + /// + /// Gets or sets the maximum duration a claimed message may remain in the + /// status + /// before it becomes eligible for reclaiming by a subsequent pending poll. + /// Default: 5 minutes. + /// + /// + /// When a worker crashes or is cancelled after claiming a message but before completing it, + /// the message stays in the Processing 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. + /// + public TimeSpan ProcessingLeaseTimeout { get; set; } = TimeSpan.FromMinutes(5); + /// /// Gets the validation error for , or when it is /// . @@ -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}'."; + /// + /// Gets the validation error for , or when it is + /// greater than . + /// + internal string? ProcessingLeaseTimeoutError => + ProcessingLeaseTimeout > TimeSpan.Zero + ? null + : $"{nameof(CosmosDbOutboxOptions)}.{nameof(ProcessingLeaseTimeout)} must be greater than zero, but was '{ProcessingLeaseTimeout}'."; + /// /// Throws when is not . /// diff --git a/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptionsValidator.cs b/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptionsValidator.cs index 6562c648..e968a11a 100644 --- a/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptionsValidator.cs +++ b/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxOptionsValidator.cs @@ -9,5 +9,7 @@ internal sealed class CosmosDbOutboxOptionsValidator : IValidateOptions 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; } diff --git a/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxRepository.cs b/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxRepository.cs index 3d3fcba0..26dc7488 100644 --- a/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxRepository.cs +++ b/src/NetEvolve.Pulse.CosmosDb/Outbox/CosmosDbOutboxRepository.cs @@ -16,6 +16,11 @@ /// Concurrency: /// Uses ETag-based conditional patch to atomically claim pending messages /// and prevent duplicate processing by concurrent workers. +/// Lease-based reclaim: +/// also reclaims documents that stayed in the Processing status +/// longer than (for example, after a worker +/// crash, cancellation or shutdown), so they are delivered at least once. The updatedAt field +/// serves as the lease timestamp, and the ETag precondition keeps concurrent reclaims exclusive. /// TTL Support: /// When is , /// the ttl field is set on completed and dead-letter documents so the Cosmos DB @@ -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; /// /// Initializes a new instance of the class. @@ -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; } /// @@ -93,9 +101,10 @@ public async Task> 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) @@ -401,6 +410,13 @@ CancellationToken cancellationToken /// as an IfMatchEtag precondition, avoiding an additional point read per candidate. /// Messages that have been modified by another worker (ETag mismatch) are silently skipped. /// + /// + /// When a later patch fails (for example, throttled with 429) 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 Processing by them are reclaimed + /// once expires. + /// private async Task> ClaimMessagesAsync( IReadOnlyList candidates, int targetStatus, @@ -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. diff --git a/src/NetEvolve.Pulse.CosmosDb/README.md b/src/NetEvolve.Pulse.CosmosDb/README.md index 78cab9f8..d083099e 100644 --- a/src/NetEvolve.Pulse.CosmosDb/README.md +++ b/src/NetEvolve.Pulse.CosmosDb/README.md @@ -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 diff --git a/src/NetEvolve.Pulse.EntityFramework/Outbox/EntityFrameworkOutboxRepository{TContext}.cs b/src/NetEvolve.Pulse.EntityFramework/Outbox/EntityFrameworkOutboxRepository{TContext}.cs index e1d82b1a..a005f6f9 100644 --- a/src/NetEvolve.Pulse.EntityFramework/Outbox/EntityFrameworkOutboxRepository{TContext}.cs +++ b/src/NetEvolve.Pulse.EntityFramework/Outbox/EntityFrameworkOutboxRepository{TContext}.cs @@ -2,6 +2,7 @@ using System.Linq.Expressions; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Options; using NetEvolve.Pulse.Extensibility.Outbox; /// @@ -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. +/// Lease-based reclaim: +/// also reclaims messages that stayed in the Processing status +/// longer than (for example, after a worker crash, +/// cancellation or shutdown), so they are delivered at least once. +/// serves as the lease timestamp. /// /// The DbContext type that implements . internal sealed class EntityFrameworkOutboxRepository : IOutboxRepository, IDisposable @@ -33,20 +39,29 @@ internal sealed class EntityFrameworkOutboxRepository : IOutboxReposit /// (change-tracking + SaveChangesAsync vs. bulk ExecuteUpdate / ExecuteDelete). /// private readonly IOutboxRepositoryExecutor _executor; + + /// + /// The duration after which a message in the Processing status is eligible for reclaiming. + /// + private readonly TimeSpan _processingLeaseTimeout; private bool _disposedValue; /// /// Initializes a new instance of the class. /// /// The DbContext for database operations. + /// The outbox configuration options. /// The time provider for timestamps. - public EntityFrameworkOutboxRepository(TContext context, TimeProvider timeProvider) + public EntityFrameworkOutboxRepository(TContext context, IOptions 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. @@ -81,9 +96,13 @@ public async Task> 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> 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); diff --git a/src/NetEvolve.Pulse.EntityFramework/README.md b/src/NetEvolve.Pulse.EntityFramework/README.md index 276d9789..2567cb4c 100644 --- a/src/NetEvolve.Pulse.EntityFramework/README.md +++ b/src/NetEvolve.Pulse.EntityFramework/README.md @@ -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(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 diff --git a/src/NetEvolve.Pulse.Extensibility/Outbox/OutboxEventTypeResolver.cs b/src/NetEvolve.Pulse.Extensibility/Outbox/OutboxEventTypeResolver.cs index 9e8092e8..b720768b 100644 --- a/src/NetEvolve.Pulse.Extensibility/Outbox/OutboxEventTypeResolver.cs +++ b/src/NetEvolve.Pulse.Extensibility/Outbox/OutboxEventTypeResolver.cs @@ -99,8 +99,8 @@ public static Type Resolve(string typeName) /// call per distinct stored name, with the error Cannot resolve event type '<stored name>'. /// Dead-lettering is best-effort: if it fails or is cancelled, the exception is swallowed so the resolvable /// messages of the claimed batch are still returned. The unresolvable messages are left in - /// and are dead-lettered again once a provider that reclaims - /// expired processing leases fetches them again. + /// and are dead-lettered again once their processing lease + /// expires and the next pending poll reclaims them. /// [SuppressMessage( "Usage", diff --git a/src/NetEvolve.Pulse/Outbox/OutboxOptions.cs b/src/NetEvolve.Pulse/Outbox/OutboxOptions.cs index e2f75b45..65d5fea2 100644 --- a/src/NetEvolve.Pulse/Outbox/OutboxOptions.cs +++ b/src/NetEvolve.Pulse/Outbox/OutboxOptions.cs @@ -51,8 +51,8 @@ public sealed class OutboxOptions /// the message stays in the Processing 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. - /// This setting is honored by providers that implement lease-based reclaim of stuck - /// Processing messages. + /// The SQL Server, PostgreSQL, MySQL, SQLite and Entity Framework Core providers honor this setting. + /// The MongoDB and Cosmos DB providers use the ProcessingLeaseTimeout of their own options. /// public TimeSpan ProcessingLeaseTimeout { get; set; } = TimeSpan.FromMinutes(5); } diff --git a/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs b/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs index 94b74bc4..78505688 100644 --- a/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs +++ b/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs @@ -377,7 +377,9 @@ await Parallel /// does not prevent processing of subsequent messages. /// Cancellation Behavior: /// Processing stops immediately when the cancellation token is triggered, leaving - /// remaining messages unprocessed. These messages will be re-polled and processed in subsequent cycles. + /// remaining messages in the Processing status. The repository reclaims them once + /// (or the provider's own lease option) expires, + /// and they are processed in a subsequent cycle. /// /// The repository resolved for this work item. /// The ordered array of outbox messages to process sequentially. diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/EntityFrameworkOutboxClaimRaceTestsBase.cs b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/EntityFrameworkOutboxClaimRaceTestsBase.cs index 13cff34f..83e06067 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/EntityFrameworkOutboxClaimRaceTestsBase.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/EntityFrameworkOutboxClaimRaceTestsBase.cs @@ -50,7 +50,36 @@ await RunAndVerify( async (services, token) => { await PrepareDatabaseAsync(token).ConfigureAwait(false); - await AddMessagesAsync(services, OutboxMessageStatus.Pending, token).ConfigureAwait(false); + await AddMessagesAsync(services, OutboxMessageStatus.Pending, TimeSpan.Zero, token) + .ConfigureAwait(false); + + var (first, second, overlapped) = await ClaimConcurrentlyAsync( + services, + outbox => outbox.GetPendingAsync(MessageCount * 2, token), + null, + token + ) + .ConfigureAwait(false); + + _ = await Assert.That(overlapped).IsTrue(); + _ = await Assert.That(first.Count).IsEqualTo(MessageCount); + _ = await Assert.That(second.Select(m => m.Id).Intersect(first.Select(m => m.Id))).IsEmpty(); + }, + cancellationToken, + configureServices: DisableProcessing + ) + .ConfigureAwait(false); + + [Test] + public async Task GetPendingAsync_WhenReclaimsOverlap_ReturnsDisjointBatches(CancellationToken cancellationToken) => + await RunAndVerify( + async (services, token) => + { + await PrepareDatabaseAsync(token).ConfigureAwait(false); + + // Processing rows whose default five-minute lease expired long ago, as left by a crashed worker. + await AddMessagesAsync(services, OutboxMessageStatus.Processing, TimeSpan.FromMinutes(10), token) + .ConfigureAwait(false); var (first, second, overlapped) = await ClaimConcurrentlyAsync( services, @@ -77,7 +106,8 @@ await RunAndVerify( async (services, token) => { await PrepareDatabaseAsync(token).ConfigureAwait(false); - await AddMessagesAsync(services, OutboxMessageStatus.Failed, token).ConfigureAwait(false); + await AddMessagesAsync(services, OutboxMessageStatus.Failed, TimeSpan.Zero, token) + .ConfigureAwait(false); var (first, second, overlapped) = await ClaimConcurrentlyAsync( services, @@ -102,7 +132,8 @@ await RunAndVerify( async (services, token) => { await PrepareDatabaseAsync(token).ConfigureAwait(false); - await AddMessagesAsync(services, OutboxMessageStatus.Failed, token).ConfigureAwait(false); + await AddMessagesAsync(services, OutboxMessageStatus.Failed, TimeSpan.Zero, token) + .ConfigureAwait(false); var nextRetryAt = services.GetRequiredService().GetUtcNow().AddHours(1); @@ -137,12 +168,13 @@ private static void DisableProcessing(IServiceCollection services) => private static async Task AddMessagesAsync( IServiceProvider services, OutboxMessageStatus status, + TimeSpan age, CancellationToken cancellationToken ) { cancellationToken.ThrowIfCancellationRequested(); - var now = services.GetRequiredService().GetUtcNow(); + var now = services.GetRequiredService().GetUtcNow() - age; var outbox = services.GetRequiredService(); for (var i = 0; i < MessageCount; i++) diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/OutboxTestsBase.cs b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/OutboxTestsBase.cs index 83c0afa8..91a7c361 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Outbox/OutboxTestsBase.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Outbox/OutboxTestsBase.cs @@ -490,6 +490,91 @@ await RunAndVerify( ) .ConfigureAwait(false); + [Test] + public async Task Should_GetPendingAsync_Reclaim_Expired_Processing(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var timeProvider = new FakeTimeProvider(); + timeProvider.AdjustTime(TestDateTime); + + await RunAndVerify( + async (services, token) => + { + var mediator = services.GetRequiredService(); + + await PublishEventsAsync(mediator, 3, x => new TestEvent { Id = $"Test{x:D3}" }, token) + .ConfigureAwait(false); + + var outbox = services.GetRequiredService(); + var claimed = await outbox.GetPendingAsync(50, token).ConfigureAwait(false); + + _ = await Assert.That(claimed.Count).IsEqualTo(3); + + // The worker completes one message and stops (shutdown or crash); the rest stay in Processing. + await outbox.MarkAsCompletedAsync(claimed[0].Id, token).ConfigureAwait(false); + + // Past the default processing lease of five minutes. + timeProvider.Advance(TimeSpan.FromMinutes(10)); + + var reclaimed = await outbox.GetPendingAsync(50, token).ConfigureAwait(false); + var reclaimedAgain = await outbox.GetPendingAsync(50, token).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert + .That(reclaimed.Select(m => m.Id)) + .IsEquivalentTo(claimed.Skip(1).Select(m => m.Id)); + _ = await Assert.That(reclaimed).All(m => m.Status == OutboxMessageStatus.Processing); + _ = await Assert.That(reclaimed).All(m => m.UpdatedAt == timeProvider.GetUtcNow()); + _ = await Assert.That(reclaimedAgain).IsEmpty(); + } + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(timeProvider) + .Configure(options => options.DisableProcessing = true) + ) + .ConfigureAwait(false); + } + + [Test] + public async Task Should_GetPendingAsync_Keep_Processing_Within_Lease(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var timeProvider = new FakeTimeProvider(); + timeProvider.AdjustTime(TestDateTime); + + await RunAndVerify( + async (services, token) => + { + var mediator = services.GetRequiredService(); + + await mediator.PublishAsync(new TestEvent { Id = "Test001" }, token).ConfigureAwait(false); + + var outbox = services.GetRequiredService(); + var claimed = await outbox.GetPendingAsync(50, token).ConfigureAwait(false); + + _ = await Assert.That(claimed.Count).IsEqualTo(1); + + // Well within the default processing lease of five minutes. + timeProvider.Advance(TimeSpan.FromMinutes(1)); + + var reclaimed = await outbox.GetPendingAsync(50, token).ConfigureAwait(false); + + _ = await Assert.That(reclaimed).IsEmpty(); + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(timeProvider) + .Configure(options => options.DisableProcessing = true) + ) + .ConfigureAwait(false); + } + [Test] public async Task Should_GetFailedForRetry_Returns_Empty_When_NoFailedMessages( CancellationToken cancellationToken diff --git a/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxOptionsValidatorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxOptionsValidatorTests.cs index eaa9c2a8..40672235 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxOptionsValidatorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxOptionsValidatorTests.cs @@ -35,4 +35,21 @@ public async Task Validate_When_PartitionKeyPath_is_unsupported_fails(string par _ = await Assert.That(result.FailureMessage).Contains("/id"); } } + + [Test] + [Arguments(0)] + [Arguments(-1)] + public async Task Validate_When_ProcessingLeaseTimeout_is_not_positive_fails(int leaseSeconds) + { + var result = _validator.Validate( + null, + new CosmosDbOutboxOptions { ProcessingLeaseTimeout = TimeSpan.FromSeconds(leaseSeconds) } + ); + + using (Assert.Multiple()) + { + _ = await Assert.That(result.Failed).IsTrue(); + _ = await Assert.That(result.FailureMessage).Contains(nameof(CosmosDbOutboxOptions.ProcessingLeaseTimeout)); + } + } } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxRepositoryClaimTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxRepositoryClaimTests.cs index 7d05e655..6458daf5 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxRepositoryClaimTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxRepositoryClaimTests.cs @@ -2,11 +2,13 @@ namespace NetEvolve.Pulse.Tests.Unit.CosmosDb; using System; using System.Collections.Generic; +using System.Linq; using System.Net; using System.Threading; using System.Threading.Tasks; using Microsoft.Azure.Cosmos; using Microsoft.Extensions.Options; +using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Outbox; using TUnit.Core; @@ -88,6 +90,43 @@ CancellationToken cancellationToken _ = await Assert.That(claimed.Count).IsEqualTo(0); } + [Test] + public async Task GetPendingAsync_WhenReclaimPatchThrowsPreconditionFailed_SkipsExpiredProcessingCandidate( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var document = CreateDocument(Guid.NewGuid(), status: 1); + document.UpdatedAt = DateTimeOffset.UtcNow.AddMinutes(-10); + document.ETag = "\"stale-etag\""; + var capturedOptions = new List(); + + var container = new FakeCosmosContainer + { + OnQueryIterator = (_, _, _) => + new FakeFeedIterator([ + [document], + ]), + OnPatchItem = (_, _, _, options) => + { + capturedOptions.Add(options); + throw new CosmosException("conflict", HttpStatusCode.PreconditionFailed, 0, "activity", 0); + }, + }; + + var repository = CreateRepository(container); + + var claimed = await repository.GetPendingAsync(10, cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(claimed.Count).IsEqualTo(0); + _ = await Assert.That(capturedOptions.Count).IsEqualTo(1); + _ = await Assert.That(capturedOptions[0]?.IfMatchEtag).IsEqualTo("\"stale-etag\""); + } + } + [Test] public async Task GetPendingAsync_WhenPatchThrowsNotFound_SkipsCandidate(CancellationToken cancellationToken) { @@ -112,6 +151,108 @@ public async Task GetPendingAsync_WhenPatchThrowsNotFound_SkipsCandidate(Cancell _ = await Assert.That(claimed.Count).IsEqualTo(0); } + [Test] + public async Task GetPendingAsync_WhenLaterPatchIsThrottled_ReturnsMessagesClaimedSoFar( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var claimedId = Guid.NewGuid(); + var throttledId = Guid.NewGuid(); + var first = CreateDocument(claimedId, status: 0); + var second = CreateDocument(throttledId, status: 0); + + var container = new FakeCosmosContainer + { + OnQueryIterator = (_, _, _) => + new FakeFeedIterator([ + [first, second], + ]), + OnPatchItem = (id, _, _, _) => + id == claimedId.ToString() + ? new FakeItemResponse(CreateDocument(claimedId, status: 1)) + : throw new CosmosException("throttled", HttpStatusCode.TooManyRequests, 0, "activity", 0), + }; + + var repository = CreateRepository(container); + + var claimed = await repository.GetPendingAsync(10, cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(claimed.Count).IsEqualTo(1); + _ = await Assert.That(claimed[0].Id).IsEqualTo(claimedId); + } + } + + [Test] + public async Task GetPendingAsync_QueriesExpiredProcessingLeases(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var timeProvider = new FakeTimeProvider(new DateTimeOffset(2025, 6, 1, 12, 0, 0, TimeSpan.Zero)); + QueryDefinition? capturedQuery = null; + + var container = new FakeCosmosContainer + { + OnQueryIterator = (_, query, _) => + { + capturedQuery = query; + return new FakeFeedIterator([ + [], + ]); + }, + }; + + using var client = new FakeCosmosClient(container); + var repository = new CosmosDbOutboxRepository( + client, + Options.Create( + new CosmosDbOutboxOptions { DatabaseName = "TestDb", ProcessingLeaseTimeout = TimeSpan.FromMinutes(3) } + ), + timeProvider + ); + + _ = await repository.GetPendingAsync(10, cancellationToken).ConfigureAwait(false); + + var leaseExpiredBefore = capturedQuery!.GetQueryParameters().Single(p => p.Name == "@leaseExpiredBefore").Value; + + using (Assert.Multiple()) + { + _ = await Assert + .That(capturedQuery.QueryText) + .Contains("c.status = 1 AND c.updatedAt <= @leaseExpiredBefore"); + _ = await Assert.That(leaseExpiredBefore).IsEqualTo(timeProvider.GetUtcNow().AddMinutes(-3)); + } + } + + [Test] + public async Task GetPendingAsync_WhenFirstPatchIsThrottled_ThrowsCosmosException( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var document = CreateDocument(Guid.NewGuid(), status: 0); + + var container = new FakeCosmosContainer + { + OnQueryIterator = (_, _, _) => + new FakeFeedIterator([ + [document], + ]), + OnPatchItem = (_, _, _, _) => + throw new CosmosException("throttled", HttpStatusCode.TooManyRequests, 0, "activity", 0), + }; + + var repository = CreateRepository(container); + + _ = await Assert + .That(async () => await repository.GetPendingAsync(10, cancellationToken).ConfigureAwait(false)) + .Throws(); + } + [Test] public async Task GetFailedForRetryAsync_WithCandidates_ClaimsAndReturnsMessages( CancellationToken cancellationToken diff --git a/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxRepositoryTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxRepositoryTests.cs index 70edb8ee..31d9d67a 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxRepositoryTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/CosmosDb/CosmosDbOutboxRepositoryTests.cs @@ -55,6 +55,32 @@ public async Task Constructor_WithNullTimeProvider_ThrowsArgumentNullException() .Throws(); } + [Test] + [Arguments(0)] + [Arguments(-1)] + public async Task Constructor_WithNonPositiveProcessingLeaseTimeout_ThrowsArgumentOutOfRangeException( + int leaseSeconds + ) + { + using var client = new CosmosClient(EmulatorConnectionString); + + _ = await Assert + .That(() => + new CosmosDbOutboxRepository( + client, + Options.Create( + new CosmosDbOutboxOptions + { + DatabaseName = "TestDb", + ProcessingLeaseTimeout = TimeSpan.FromSeconds(leaseSeconds), + } + ), + TimeProvider.System + ) + ) + .Throws(); + } + [Test] public async Task Constructor_WithEmptyDatabaseName_ThrowsArgumentException() { diff --git a/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryClaimStatementTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryClaimStatementTests.cs index ed2201ac..04eff9e6 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryClaimStatementTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryClaimStatementTests.cs @@ -10,6 +10,7 @@ namespace NetEvolve.Pulse.Tests.Unit.EntityFramework; using Microsoft.Data.Sqlite; using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.Diagnostics; +using Microsoft.Extensions.Options; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Extensibility.Outbox; using NetEvolve.Pulse.Outbox; @@ -131,7 +132,11 @@ CancellationToken cancellationToken context.ChangeTracker.Clear(); recorder.Commands.Clear(); - using var repository = new EntityFrameworkOutboxRepository(context, TimeProvider.System); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + TimeProvider.System + ); await claim(repository, cancellationToken).ConfigureAwait(false); return recorder.Commands.Single(c => c.StartsWith("UPDATE", StringComparison.OrdinalIgnoreCase)); diff --git a/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryInvariantTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryInvariantTests.cs index b5f946cc..66b46e8a 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryInvariantTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryInvariantTests.cs @@ -4,6 +4,7 @@ namespace NetEvolve.Pulse.Tests.Unit.EntityFramework; using System.Threading; using System.Threading.Tasks; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Options; using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Extensibility.Outbox; @@ -72,7 +73,11 @@ public async Task GetPendingAsync_Excludes_messages_with_future_NextRetryAt(Canc _ = await context.OutboxMessages.AddAsync(notDueMessage, cancellationToken).ConfigureAwait(false); _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); - using var repository = new EntityFrameworkOutboxRepository(context, fakeTime); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + fakeTime + ); var pending = await repository.GetPendingAsync(10, cancellationToken).ConfigureAwait(false); @@ -105,7 +110,11 @@ public async Task GetPendingAsync_Includes_messages_with_past_NextRetryAt(Cancel _ = await context.OutboxMessages.AddAsync(pastMessage, cancellationToken).ConfigureAwait(false); _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); - using var repository = new EntityFrameworkOutboxRepository(context, fakeTime); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + fakeTime + ); var pending = await repository.GetPendingAsync(10, cancellationToken).ConfigureAwait(false); @@ -132,7 +141,11 @@ public async Task GetPendingAsync_Transitions_claimed_messages_to_Processing(Can _ = await context.OutboxMessages.AddAsync(msg, cancellationToken).ConfigureAwait(false); _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); - using var repository = new EntityFrameworkOutboxRepository(context, fakeTime); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + fakeTime + ); var claimed = await repository.GetPendingAsync(10, cancellationToken).ConfigureAwait(false); @@ -165,7 +178,11 @@ public async Task MarkAsCompletedAsync_Does_not_change_DeadLetter_to_Completed(C _ = await context.OutboxMessages.AddAsync(dlq, cancellationToken).ConfigureAwait(false); _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); - using var repository = new EntityFrameworkOutboxRepository(context, fakeTime); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + fakeTime + ); await repository.MarkAsCompletedAsync(dlq.Id, cancellationToken).ConfigureAwait(false); @@ -192,7 +209,11 @@ public async Task MarkAsFailedAsync_Does_not_change_Completed_to_Failed(Cancella _ = await context.OutboxMessages.AddAsync(completed, cancellationToken).ConfigureAwait(false); _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); - using var repository = new EntityFrameworkOutboxRepository(context, fakeTime); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + fakeTime + ); await repository.MarkAsFailedAsync(completed.Id, "boom", cancellationToken).ConfigureAwait(false); @@ -228,7 +249,11 @@ CancellationToken cancellationToken _ = await context.OutboxMessages.AddAsync(processing, cancellationToken).ConfigureAwait(false); _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); - using var repository = new EntityFrameworkOutboxRepository(context, fakeTime); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + fakeTime + ); await repository.MarkAsCompletedAsync(processing.Id, cancellationToken).ConfigureAwait(false); @@ -278,7 +303,11 @@ await context .ConfigureAwait(false); _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); - using var repository = new EntityFrameworkOutboxRepository(context, fakeTime); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + fakeTime + ); // Cutoff = 1 day - only the 5-day-old Completed message qualifies. var deleted = await repository @@ -321,7 +350,11 @@ await context .ConfigureAwait(false); _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); - using var repository = new EntityFrameworkOutboxRepository(context, fakeTime); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + fakeTime + ); var retried = await repository .GetFailedForRetryAsync(maxRetryCount: 5, batchSize: 10, cancellationToken) @@ -344,7 +377,11 @@ public async Task AddAsync_Persists_message_and_AddAsync_with_null_throws(Cancel var context = CreateContext(nameof(AddAsync_Persists_message_and_AddAsync_with_null_throws)); await using (context.ConfigureAwait(false)) { - using var repository = new EntityFrameworkOutboxRepository(context, TimeProvider.System); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + TimeProvider.System + ); // null path _ = await Assert @@ -368,7 +405,11 @@ public async Task GetPendingAsync_Honors_cancellation_token() var context = CreateContext(nameof(GetPendingAsync_Honors_cancellation_token)); await using (context.ConfigureAwait(false)) { - using var repository = new EntityFrameworkOutboxRepository(context, TimeProvider.System); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + TimeProvider.System + ); using var cts = new CancellationTokenSource(); await cts.CancelAsync().ConfigureAwait(false); @@ -378,4 +419,48 @@ public async Task GetPendingAsync_Honors_cancellation_token() .Throws(); } } + + // INVARIANT (#813): GetPendingAsync reclaims Processing messages only after the configured + // OutboxOptions.ProcessingLeaseTimeout, not after a hardcoded default. + [Test] + [Arguments(30, 45, 15)] + [Arguments(3600, 5400, 1800)] + public async Task GetPendingAsync_Uses_configured_ProcessingLeaseTimeout( + int leaseSeconds, + int expiredAgeSeconds, + int activeAgeSeconds, + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var fakeTime = new FakeTimeProvider( + DateTimeOffset.Parse("2025-01-01T12:00:00Z", System.Globalization.CultureInfo.InvariantCulture) + ); + var now = fakeTime.GetUtcNow(); + var context = CreateContext($"{nameof(GetPendingAsync_Uses_configured_ProcessingLeaseTimeout)}_{leaseSeconds}"); + await using (context.ConfigureAwait(false)) + { + var expired = CreateMessage(OutboxMessageStatus.Processing, now.AddSeconds(-expiredAgeSeconds)); + var active = CreateMessage(OutboxMessageStatus.Processing, now.AddSeconds(-activeAgeSeconds)); + _ = await context.OutboxMessages.AddAsync(expired, cancellationToken).ConfigureAwait(false); + _ = await context.OutboxMessages.AddAsync(active, cancellationToken).ConfigureAwait(false); + _ = await context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); + + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions { ProcessingLeaseTimeout = TimeSpan.FromSeconds(leaseSeconds) }), + fakeTime + ); + + var claimed = await repository.GetPendingAsync(10, cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(claimed).HasCount(1); + _ = await Assert.That(claimed[0].Id).IsEqualTo(expired.Id); + _ = await Assert.That(claimed[0].UpdatedAt).IsEqualTo(now); + } + } + } } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryTests.cs index 6fa69e76..c7c90f27 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/EntityFramework/EntityFrameworkOutboxRepositoryTests.cs @@ -5,6 +5,7 @@ using System.Threading; using System.Threading.Tasks; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Options; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Outbox; using TUnit.Core; @@ -19,7 +20,11 @@ public async Task Constructor_WithRelationalProvider_CreatesExecutorWithSinglePe var context = new TestDbContext(options); await using (context.ConfigureAwait(false)) { - using var repository = new EntityFrameworkOutboxRepository(context, TimeProvider.System); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + TimeProvider.System + ); var executor = typeof(EntityFrameworkOutboxRepository) .GetField("_executor", BindingFlags.NonPublic | BindingFlags.Instance)! @@ -37,7 +42,13 @@ public async Task Constructor_WithRelationalProvider_CreatesExecutorWithSinglePe [Test] public async Task Constructor_WithNullContext_ThrowsArgumentNullException() => _ = await Assert - .That(() => new EntityFrameworkOutboxRepository(null!, TimeProvider.System)) + .That(() => + new EntityFrameworkOutboxRepository( + null!, + Options.Create(new OutboxOptions()), + TimeProvider.System + ) + ) .Throws(); [Test] @@ -50,11 +61,59 @@ public async Task Constructor_WithNullTimeProvider_ThrowsArgumentNullException() await using (context.ConfigureAwait(false)) { _ = await Assert - .That(() => new EntityFrameworkOutboxRepository(context, null!)) + .That(() => + new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + null! + ) + ) .Throws(); } } + [Test] + public async Task Constructor_WithNullOptions_ThrowsArgumentNullException() + { + var options = new DbContextOptionsBuilder() + .UseInMemoryDatabase(nameof(Constructor_WithNullOptions_ThrowsArgumentNullException)) + .Options; + var context = new TestDbContext(options); + await using (context.ConfigureAwait(false)) + { + _ = await Assert + .That(() => new EntityFrameworkOutboxRepository(context, null!, TimeProvider.System)) + .Throws(); + } + } + + [Test] + [Arguments(0)] + [Arguments(-1)] + public async Task Constructor_WithNonPositiveProcessingLeaseTimeout_ThrowsArgumentOutOfRangeException( + int leaseSeconds + ) + { + var options = new DbContextOptionsBuilder() + .UseInMemoryDatabase( + nameof(Constructor_WithNonPositiveProcessingLeaseTimeout_ThrowsArgumentOutOfRangeException) + ) + .Options; + var context = new TestDbContext(options); + await using (context.ConfigureAwait(false)) + { + var outboxOptions = Options.Create( + new OutboxOptions { ProcessingLeaseTimeout = TimeSpan.FromSeconds(leaseSeconds) } + ); + + _ = await Assert + .That(() => + new EntityFrameworkOutboxRepository(context, outboxOptions, TimeProvider.System) + ) + .Throws(); + } + } + [Test] public async Task Constructor_WithValidArguments_CreatesInstance() { @@ -64,7 +123,11 @@ public async Task Constructor_WithValidArguments_CreatesInstance() var context = new TestDbContext(options); await using (context.ConfigureAwait(false)) { - using var repository = new EntityFrameworkOutboxRepository(context, TimeProvider.System); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + TimeProvider.System + ); _ = await Assert.That(repository).IsNotNull(); } @@ -81,7 +144,11 @@ public async Task AddAsync_WithNullMessage_ThrowsArgumentNullException(Cancellat var context = new TestDbContext(options); await using (context.ConfigureAwait(false)) { - using var repository = new EntityFrameworkOutboxRepository(context, TimeProvider.System); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + TimeProvider.System + ); _ = await Assert .That(async () => await repository.AddAsync(null!, cancellationToken).ConfigureAwait(false)) @@ -100,7 +167,11 @@ public async Task IsHealthyAsync_WithInMemoryProvider_ReturnsTrue(CancellationTo var context = new TestDbContext(options); await using (context.ConfigureAwait(false)) { - using var repository = new EntityFrameworkOutboxRepository(context, TimeProvider.System); + using var repository = new EntityFrameworkOutboxRepository( + context, + Options.Create(new OutboxOptions()), + TimeProvider.System + ); var result = await repository.IsHealthyAsync(cancellationToken).ConfigureAwait(false);