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")]