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
Original file line number Diff line number Diff line change
Expand Up @@ -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
/// <see cref="IdempotencyConflictException"/>. Strict atomicity of the reservation itself depends on
/// the registered <see cref="IIdempotencyStore"/> implementation of
/// <see cref="IIdempotencyStore.TryReserveAsync"/>; the non-atomic default leaves a small window
/// between the existence check and the store operation.
/// the registered <see cref="IIdempotencyKeyRepository"/> implementation of
/// <see cref="IIdempotencyKeyRepository.TryReserveAsync"/>, to which <see cref="IdempotencyStore"/> delegates.
/// <para><strong>Registration:</strong></para>
/// Use <c>AddIdempotency()</c> on the <see cref="IMediatorBuilder"/> to register this interceptor.
/// </remarks>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

/// <summary>
/// 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.
/// </summary>
protected virtual bool SupportsAtomicReservation => false;

protected async ValueTask RunAndVerify(
Func<IServiceProvider, CancellationToken, Task> testableCode,
CancellationToken cancellationToken,
Expand Down Expand Up @@ -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<IIdempotencyStore>();

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<IIdempotencyStore>()
.ExistsAsync("warm-up", token)
.ConfigureAwait(false);

var scopeFactory = services.GetRequiredService<IServiceScopeFactory>();

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<IIdempotencyStore>();
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<IIdempotencyStore>()
.ExistsAsync("warm-up", token)
.ConfigureAwait(false);

var scopeFactory = services.GetRequiredService<IServiceScopeFactory>();

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<IMediator>();
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<ICommandHandler<CountingCommand, Void>, 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<IIdempotencyStore>()
.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<IServiceScopeFactory>();

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<IIdempotencyStore>();
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<TimeProvider>(fakeTime)
.Configure<IdempotencyKeyOptions>(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<IIdempotencyStore>();

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<TimeProvider>(fakeTime)
.Configure<IdempotencyKeyOptions>(o => o.TimeToLive = null)
)
.ConfigureAwait(false);
}

private sealed record TestIdempotentVoidCommand(string IdempotencyKey) : IIdempotentCommand
{
public string? CausationId { get; set; }
Expand All @@ -457,4 +698,33 @@ public Task<Void> 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<CountingCommand, Void>
{
public async Task<Void> 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;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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")]
Expand Down
Loading