diff --git a/src/NetEvolve.Pulse/Interceptors/IdempotencyCommandInterceptor{TRequest,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/IdempotencyCommandInterceptor{TRequest,TResponse}.cs index cb25df4a..63e37b27 100644 --- a/src/NetEvolve.Pulse/Interceptors/IdempotencyCommandInterceptor{TRequest,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/IdempotencyCommandInterceptor{TRequest,TResponse}.cs @@ -27,9 +27,8 @@ /// Reservation before execution provides at-most-once semantics: a command whose handler fails /// keeps its key reserved, and retries with the same key are rejected with /// . Strict atomicity of the reservation itself depends on -/// the registered implementation of -/// ; the non-atomic default leaves a small window -/// between the existence check and the store operation. +/// the registered implementation of +/// , to which delegates. /// Registration: /// Use AddIdempotency() on the to register this interceptor. /// diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs index 9d9d94ec..7cc27375 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs @@ -27,6 +27,14 @@ IServiceInitializer databaseInitializer protected static DateTimeOffset TestDateTime { get; } = new DateTimeOffset(2025, 1, 1, 12, 0, 0, 0, TimeSpan.Zero); + private const int ConcurrentCalls = 20; + + /// + /// Gets a value indicating whether the provider reserves keys atomically, so that exactly one of several + /// concurrent reservations wins and a re-reserved expired key rejects later duplicates. + /// + protected virtual bool SupportsAtomicReservation => false; + protected async ValueTask RunAndVerify( Func testableCode, CancellationToken cancellationToken, @@ -444,6 +452,239 @@ await RunAndVerify( ) .ConfigureAwait(false); + [Test] + public async Task Should_Reserve_New_Key_Once(CancellationToken cancellationToken) => + await RunAndVerify( + async (services, token) => + { + var store = services.GetRequiredService(); + + var first = await store.TryReserveAsync("reserve-key", token).ConfigureAwait(false); + var second = await store.TryReserveAsync("reserve-key", token).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(first).IsTrue(); + _ = await Assert.That(second).IsFalse(); + } + }, + cancellationToken + ) + .ConfigureAwait(false); + + [Test] + public async Task Should_Reserve_Once_When_Reserving_Same_Key_Concurrently(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically."); + + await RunAndVerify( + async (services, token) => + { + // Open the provider connection up front, so the parallel calls race on the reservation only. + _ = await services + .GetRequiredService() + .ExistsAsync("warm-up", token) + .ConfigureAwait(false); + + var scopeFactory = services.GetRequiredService(); + + var results = await Task.WhenAll( + Enumerable + .Range(0, ConcurrentCalls) + .Select(_ => + Task.Run( + async () => + { + var scope = scopeFactory.CreateAsyncScope(); + await using (scope.ConfigureAwait(false)) + { + var store = + scope.ServiceProvider.GetRequiredService(); + return await store + .TryReserveAsync("concurrent-key", token) + .ConfigureAwait(false); + } + }, + token + ) + ) + ) + .ConfigureAwait(false); + + _ = await Assert.That(results.Count(reserved => reserved)).IsEqualTo(1); + }, + cancellationToken + ) + .ConfigureAwait(false); + } + + [Test] + public async Task Should_Run_Handler_Once_When_Sending_Same_Command_Concurrently( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically."); + + var counter = new InvocationCounter(); + + await RunAndVerify( + async (services, token) => + { + // Open the provider connection up front, so the parallel calls race on the reservation only. + _ = await services + .GetRequiredService() + .ExistsAsync("warm-up", token) + .ConfigureAwait(false); + + var scopeFactory = services.GetRequiredService(); + + var outcomes = await Task.WhenAll( + Enumerable + .Range(0, ConcurrentCalls) + .Select(_ => + Task.Run( + async () => + { + var scope = scopeFactory.CreateAsyncScope(); + await using (scope.ConfigureAwait(false)) + { + var mediator = scope.ServiceProvider.GetRequiredService(); + try + { + await mediator + .SendAsync(new CountingCommand("concurrent-command"), token) + .ConfigureAwait(false); + return true; + } + catch (IdempotencyConflictException) + { + return false; + } + } + }, + token + ) + ) + ) + .ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(counter.Count).IsEqualTo(1); + _ = await Assert.That(outcomes.Count(succeeded => succeeded)).IsEqualTo(1); + _ = await Assert.That(outcomes.Count(succeeded => !succeeded)).IsEqualTo(ConcurrentCalls - 1); + } + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(counter) + .AddSingleton, CountingCommandHandler>() + ) + .ConfigureAwait(false); + } + + [Test] + public async Task Should_Reserve_Expired_Key_Once_When_Reserving_Concurrently(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically."); + + var fakeTime = new FakeTimeProvider(); + fakeTime.AdjustTime(TestDateTime); + + await RunAndVerify( + async (services, token) => + { + var first = await services + .GetRequiredService() + .TryReserveAsync("concurrent-expired-key", token) + .ConfigureAwait(false); + + // Past the logical TTL, but within the physical expiry (TTL + 1h) of providers that have one. + fakeTime.Advance(TimeSpan.FromMinutes(90)); + + var scopeFactory = services.GetRequiredService(); + + var results = await Task.WhenAll( + Enumerable + .Range(0, ConcurrentCalls) + .Select(_ => + Task.Run( + async () => + { + var scope = scopeFactory.CreateAsyncScope(); + await using (scope.ConfigureAwait(false)) + { + var store = + scope.ServiceProvider.GetRequiredService(); + return await store + .TryReserveAsync("concurrent-expired-key", token) + .ConfigureAwait(false); + } + }, + token + ) + ) + ) + .ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(first).IsTrue(); + _ = await Assert.That(results.Count(reserved => reserved)).IsEqualTo(1); + } + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(fakeTime) + .Configure(o => o.TimeToLive = TimeSpan.FromHours(1)) + ) + .ConfigureAwait(false); + } + + [Test] + public async Task Should_Never_Reserve_Existing_Key_Again_When_TimeToLive_Is_Null( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var fakeTime = new FakeTimeProvider(); + fakeTime.AdjustTime(TestDateTime); + + await RunAndVerify( + async (services, token) => + { + var store = services.GetRequiredService(); + + var first = await store.TryReserveAsync("no-ttl-key", token).ConfigureAwait(false); + + fakeTime.Advance(TimeSpan.FromDays(365)); + + var second = await store.TryReserveAsync("no-ttl-key", token).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(first).IsTrue(); + _ = await Assert.That(second).IsFalse(); + } + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(fakeTime) + .Configure(o => o.TimeToLive = null) + ) + .ConfigureAwait(false); + } + private sealed record TestIdempotentVoidCommand(string IdempotencyKey) : IIdempotentCommand { public string? CausationId { get; set; } @@ -457,4 +698,33 @@ public Task HandleAsync( CancellationToken cancellationToken = default ) => Task.FromResult(Void.Completed); } + + private sealed class InvocationCounter + { + private int _count; + + public int Count => Volatile.Read(ref _count); + + public void Increment() => _ = Interlocked.Increment(ref _count); + } + + private sealed record CountingCommand(string IdempotencyKey) : IIdempotentCommand + { + public string? CausationId { get; set; } + public string? CorrelationId { get; set; } + } + + private sealed class CountingCommandHandler(InvocationCounter counter) : ICommandHandler + { + public async Task HandleAsync(CountingCommand command, CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + counter.Increment(); + + // Keep the handler busy so overlapping submissions hit the reservation while it runs. + await Task.Delay(50, cancellationToken).ConfigureAwait(false); + return Void.Completed; + } + } } diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs index 31aefb88..1b373f61 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs @@ -17,6 +17,8 @@ public class RedisIdempotencyTests(IServiceFixture databaseServiceFixture, IServiceInitializer databaseInitializer) : IdempotencyTestsBase(databaseServiceFixture, databaseInitializer) { + protected override bool SupportsAtomicReservation => true; + [Test] [Arguments("negative-offset-key", "2025-01-01T08:00:00.0000000-05:00")] [Arguments("positive-offset-key", "2025-01-01T14:30:00.0000000+02:00")]