From 21f003b638da3f9e663800564e76e8c2cacadaa1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Mon, 28 Sep 2026 22:58:39 +0200 Subject: [PATCH 1/7] test(redis): cover concurrent and expired idempotency key reservations --- .../Idempotency/RedisIdempotencyTests.cs | 229 +++++++++++++++++- 1 file changed, 227 insertions(+), 2 deletions(-) diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs index 84ecbd94..10f6ce40 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs @@ -1,12 +1,237 @@ -namespace NetEvolve.Pulse.Tests.Integration.Idempotency; +namespace NetEvolve.Pulse.Tests.Integration.Idempotency; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; +using NetEvolve.Pulse.Extensibility; +using NetEvolve.Pulse.Extensibility.Idempotency; +using NetEvolve.Pulse.Idempotency; using NetEvolve.Pulse.Tests.Integration.Internals; using NetEvolve.Pulse.Tests.Integration.Internals.Idempotency; using NetEvolve.Pulse.Tests.Integration.Internals.Services; +using TUnit.Assertions; +using TUnit.Assertions.Extensions; +using TUnit.Core; [ClassDataSource(Shared = [SharedType.None, SharedType.None])] [TestGroup("Redis")] [InheritsTests] public class RedisIdempotencyTests(IServiceFixture databaseServiceFixture, IServiceInitializer databaseInitializer) - : IdempotencyTestsBase(databaseServiceFixture, databaseInitializer); + : IdempotencyTestsBase(databaseServiceFixture, databaseInitializer) +{ + private const int ConcurrentCalls = 20; + + [Test] + public async Task Should_Reserve_Exactly_Once_When_Reserving_Same_Key_Concurrently( + CancellationToken cancellationToken + ) => + await RunAndVerify( + async (services, token) => + { + // Connect the shared multiplexer 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(); + + var counter = new InvocationCounter(); + + await RunAndVerify( + async (services, token) => + { + // Connect the shared multiplexer 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_Logically_Expired_Key_That_Is_Still_Physically_Present( + 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("expired-reserve-key", token).ConfigureAwait(false); + + // The physical Redis expiry is TTL + 1h, so the key is still present after 90 minutes. + fakeTime.Advance(TimeSpan.FromMinutes(90)); + + var second = await store.TryReserveAsync("expired-reserve-key", token).ConfigureAwait(false); + var third = await store.TryReserveAsync("expired-reserve-key", token).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(first).IsTrue(); + _ = await Assert.That(second).IsTrue(); + _ = await Assert.That(third).IsFalse(); + } + }, + 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 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; + } + } +} From 795be11ab0fe56305b16b80b32bd3644ba1f29b3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 00:26:36 +0200 Subject: [PATCH 2/7] feat(idempotency): add TryStoreAsync to IIdempotencyKeyRepository and reserve keys through it IdempotencyStore.TryReserveAsync now delegates to the repository, passing the current timestamp and the logical-expiry cutoff. All first-party providers implement the new member by composing ExistsAsync and StoreAsync, which keeps their current non-atomic behaviour. Refs #816 --- ...eworkIdempotencyKeyRepository{TContext}.cs | 22 ++++++ .../Idempotency/IIdempotencyKeyRepository.cs | 26 +++++++ .../MySqlIdempotencyKeyRepository.cs | 22 ++++++ .../PostgreSqlIdempotencyKeyRepository.cs | 22 ++++++ .../RedisIdempotencyKeyRepository.cs | 22 ++++++ .../SQLiteIdempotencyKeyRepository.cs | 22 ++++++ .../SqlServerIdempotencyKeyRepository.cs | 22 ++++++ .../Idempotency/IdempotencyStore.cs | 17 +++++ .../Idempotency/IdempotencyStoreTests.cs | 76 +++++++++++++++++++ 9 files changed, 251 insertions(+) diff --git a/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs b/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs index 1a64d542..8db29a3d 100644 --- a/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs +++ b/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs @@ -56,6 +56,28 @@ public Task ExistsAsync( return _context.IdempotencyKeys.AnyAsync(k => k.Key == idempotencyKey, cancellationToken); } + /// + /// + /// Composes and and is therefore not atomic. + /// + public async Task TryStoreAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + if (await ExistsAsync(idempotencyKey, validFrom, cancellationToken).ConfigureAwait(false)) + { + return false; + } + + await StoreAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); + return true; + } + /// public async Task StoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.Extensibility/Idempotency/IIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.Extensibility/Idempotency/IIdempotencyKeyRepository.cs index 608bbdfb..d98aabbd 100644 --- a/src/NetEvolve.Pulse.Extensibility/Idempotency/IIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.Extensibility/Idempotency/IIdempotencyKeyRepository.cs @@ -53,4 +53,30 @@ Task ExistsAsync( /// a successful (idempotent) store operation. /// Task StoreAsync(string idempotencyKey, DateTimeOffset createdAt, CancellationToken cancellationToken = default); + + /// + /// Attempts to store an idempotency key only if no valid key with the same value exists yet. + /// + /// The idempotency key to reserve. + /// The timestamp to associate with the stored key. + /// + /// When set, an existing key created before this timestamp is treated as absent and may be + /// replaced by this reservation. When , an existing key is never replaced. + /// + /// A token to monitor for cancellation requests. + /// + /// if this call stored the key; if a valid key already exists. + /// + /// + /// Implementations backed by storage that supports an atomic check-and-set (for example Redis SET NX) + /// SHOULD perform the check and the store as one atomic operation, so that of several concurrent calls + /// for the same key at most one returns . Implementations that compose + /// and are not atomic. + /// + Task TryStoreAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ); } diff --git a/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs index defd8d38..bb3644ad 100644 --- a/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs @@ -119,6 +119,28 @@ public async Task ExistsAsync( } } + /// + /// + /// Composes and and is therefore not atomic. + /// + public async Task TryStoreAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + if (await ExistsAsync(idempotencyKey, validFrom, cancellationToken).ConfigureAwait(false)) + { + return false; + } + + await StoreAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); + return true; + } + /// public async Task StoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs index fe8d30ff..675f9955 100644 --- a/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs @@ -93,6 +93,28 @@ public async Task ExistsAsync( } } + /// + /// + /// Composes and and is therefore not atomic. + /// + public async Task TryStoreAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + if (await ExistsAsync(idempotencyKey, validFrom, cancellationToken).ConfigureAwait(false)) + { + return false; + } + + await StoreAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); + return true; + } + /// public async Task StoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs index e2b5c06f..1e1edeec 100644 --- a/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs @@ -87,6 +87,28 @@ out var createdAt return true; } + /// + /// + /// Composes and and is therefore not atomic. + /// + public async Task TryStoreAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + if (await ExistsAsync(idempotencyKey, validFrom, cancellationToken).ConfigureAwait(false)) + { + return false; + } + + await StoreAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); + return true; + } + /// public async Task StoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs index 9cecb95e..ad898f8e 100644 --- a/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs @@ -119,6 +119,28 @@ public async Task ExistsAsync( } } + /// + /// + /// Composes and and is therefore not atomic. + /// + public async Task TryStoreAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + if (await ExistsAsync(idempotencyKey, validFrom, cancellationToken).ConfigureAwait(false)) + { + return false; + } + + await StoreAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); + return true; + } + /// public async Task StoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs index 3756410d..99e10221 100644 --- a/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs @@ -99,6 +99,28 @@ public async Task ExistsAsync( } } + /// + /// + /// Composes and and is therefore not atomic. + /// + public async Task TryStoreAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + if (await ExistsAsync(idempotencyKey, validFrom, cancellationToken).ConfigureAwait(false)) + { + return false; + } + + await StoreAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); + return true; + } + /// public async Task StoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse/Idempotency/IdempotencyStore.cs b/src/NetEvolve.Pulse/Idempotency/IdempotencyStore.cs index 7da9cbb0..194d8fd1 100644 --- a/src/NetEvolve.Pulse/Idempotency/IdempotencyStore.cs +++ b/src/NetEvolve.Pulse/Idempotency/IdempotencyStore.cs @@ -66,4 +66,21 @@ public Task StoreAsync(string idempotencyKey, CancellationToken cancellationToke return _repository.StoreAsync(idempotencyKey, _timeProvider.GetUtcNow(), cancellationToken); } + + /// + /// + /// Delegates to , so the reservation is atomic + /// whenever the registered repository implements it atomically. + /// + public Task TryReserveAsync(string idempotencyKey, CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + + var now = _timeProvider.GetUtcNow(); + DateTimeOffset? cutoff = _options.TimeToLive.HasValue ? now - _options.TimeToLive.Value : null; + + return _repository.TryStoreAsync(idempotencyKey, now, cutoff, cancellationToken); + } } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Idempotency/IdempotencyStoreTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Idempotency/IdempotencyStoreTests.cs index b2db9163..fd62c039 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Idempotency/IdempotencyStoreTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Idempotency/IdempotencyStoreTests.cs @@ -142,10 +142,72 @@ public async Task StoreAsync_PassesCurrentTimestampToRepository(CancellationToke _ = await Assert.That(repository.CapturedCreatedAt).IsEqualTo(expectedTimestamp); } + [Test] + public async Task TryReserveAsync_WithEmptyKey_ThrowsArgumentException(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var store = CreateStore(new TrackingIdempotencyKeyRepository()); + + _ = await Assert + .That(async () => await store.TryReserveAsync(string.Empty, cancellationToken).ConfigureAwait(false)) + .Throws(); + } + + [Test] + public async Task TryReserveAsync_WithTtl_PassesTimestampAndCutoffToRepository(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var fakeTime = new FakeTimeProvider(); + var now = fakeTime.GetUtcNow(); + var ttl = TimeSpan.FromMinutes(10); + + var repository = new TrackingIdempotencyKeyRepository(); + var store = CreateStore(repository, new IdempotencyKeyOptions { TimeToLive = ttl }, fakeTime); + + _ = await store.TryReserveAsync("test-key", cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(repository.CapturedCreatedAt).IsEqualTo(now); + _ = await Assert.That(repository.CapturedValidFrom).IsEqualTo(now - ttl); + } + } + + [Test] + public async Task TryReserveAsync_WithoutTtl_PassesNullCutoffToRepository(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var repository = new TrackingIdempotencyKeyRepository(); + var store = CreateStore(repository, new IdempotencyKeyOptions { TimeToLive = null }); + + _ = await store.TryReserveAsync("test-key", cancellationToken).ConfigureAwait(false); + + _ = await Assert.That(repository.CapturedValidFrom).IsNull(); + } + + [Test] + [Arguments(true)] + [Arguments(false)] + public async Task TryReserveAsync_ReturnsRepositoryResult(bool stored, CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var repository = new TrackingIdempotencyKeyRepository { TryStoreResult = stored }; + var store = CreateStore(repository); + + var result = await store.TryReserveAsync("test-key", cancellationToken).ConfigureAwait(false); + + _ = await Assert.That(result).IsEqualTo(stored); + } + private sealed class TrackingIdempotencyKeyRepository : IIdempotencyKeyRepository { public DateTimeOffset? CapturedValidFrom { get; private set; } = DateTimeOffset.MaxValue; public DateTimeOffset CapturedCreatedAt { get; private set; } + public bool TryStoreResult { get; init; } = true; public Task ExistsAsync( string idempotencyKey, @@ -170,5 +232,19 @@ public Task StoreAsync( CapturedCreatedAt = createdAt; return Task.CompletedTask; } + + public Task TryStoreAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + CapturedCreatedAt = createdAt; + CapturedValidFrom = validFrom; + return Task.FromResult(TryStoreResult); + } } } From f78f40233e10e98287ebae05eb6c0739fa0fc567 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 00:41:26 +0200 Subject: [PATCH 3/7] fix(redis): reserve idempotency keys atomically with SET NX and compare-and-set RedisIdempotencyKeyRepository.TryStoreAsync now returns the SET NX result instead of composing ExistsAsync and StoreAsync. A logically expired key that is still physically present is replaced through a WATCH/MULTI transaction conditioned on the observed value, so exactly one concurrent reservation wins and the refreshed timestamp rejects later duplicates. Without a TimeToLive an existing key is never reserved again. Closes #816 --- .../RedisIdempotencyKeyRepository.cs | 83 +++++++++++----- src/NetEvolve.Pulse.Redis/README.md | 5 +- ...disIdempotencyMediatorBuilderExtensions.cs | 6 +- ...isIdempotencyKeyRepositoryBehaviorTests.cs | 98 +++++++++++++++++++ 4 files changed, 167 insertions(+), 25 deletions(-) diff --git a/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs index 1e1edeec..0c0061ad 100644 --- a/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs @@ -18,6 +18,10 @@ namespace NetEvolve.Pulse.Idempotency; /// without a Redis expiry and are only removed if the server's maxmemory-policy evicts non-volatile /// keys (allkeys-*), which breaks duplicate detection. TTL-based logical expiry is handled by the /// wrapper using the injected , which makes it testable with fake clocks. +/// Reservation: +/// reserves a key atomically with SET NX. A key that is logically expired but +/// still physically present is replaced through a compare-and-set transaction, so of several concurrent +/// reservations for the same key exactly one succeeds. /// Prerequisites: /// must be registered in the DI container by the caller /// before using this provider. @@ -70,26 +74,16 @@ public async Task ExistsAsync( return true; } - // Parse the stored creation timestamp and check it is within the TTL window. - if ( - DateTimeOffset.TryParse( - value.ToString(), - CultureInfo.InvariantCulture, - DateTimeStyles.RoundtripKind, - out var createdAt - ) - ) - { - return createdAt >= validFrom.Value; - } - - // If the value cannot be parsed (e.g. legacy entry), treat it as present. - return true; + // An unparseable value (e.g. legacy entry) is not expired and therefore treated as present. + return !IsExpired(value, validFrom.Value); } /// /// - /// Composes and and is therefore not atomic. + /// Reserves the key with an atomic SET NX. When the key already exists and + /// is set, a stored timestamp older than is replaced through a transaction that + /// only commits while the key still holds the observed value, so of several concurrent calls at most one + /// returns . /// public async Task TryStoreAsync( string idempotencyKey, @@ -100,13 +94,45 @@ public async Task TryStoreAsync( { cancellationToken.ThrowIfCancellationRequested(); - if (await ExistsAsync(idempotencyKey, validFrom, cancellationToken).ConfigureAwait(false)) + ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + + var database = _multiplexer.GetDatabase(DefaultDatabase); + var key = GetPrefixedKey(idempotencyKey); + var timestamp = createdAt.ToString("O", CultureInfo.InvariantCulture); + var physicalTtl = GetPhysicalTimeToLive(); + + if (await database.StringSetAsync(key, timestamp, physicalTtl, When.NotExists).ConfigureAwait(false)) + { + return true; + } + + // Without a TTL an existing key never expires logically, so it is never reservable again. + if (!validFrom.HasValue) { return false; } - await StoreAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); - return true; + cancellationToken.ThrowIfCancellationRequested(); + + var existing = await database.StringGetAsync(key).ConfigureAwait(false); + + if (!existing.HasValue) + { + // The key expired physically between SET NX and GET; try the plain reservation once more. + return await database.StringSetAsync(key, timestamp, physicalTtl, When.NotExists).ConfigureAwait(false); + } + + if (!IsExpired(existing, validFrom.Value)) + { + return false; + } + + // Compare-and-set: replace the expired value only if no concurrent call changed it meanwhile. + var transaction = database.CreateTransaction(); + _ = transaction.AddCondition(Condition.StringEqual(key, existing)); + _ = transaction.StringSetAsync(key, timestamp, physicalTtl, When.Always); + + return await transaction.ExecuteAsync().ConfigureAwait(false); } /// @@ -122,9 +148,7 @@ public async Task StoreAsync( var database = _multiplexer.GetDatabase(DefaultDatabase); - // Physical expiry is TTL + 1h headroom; a null TTL means "never expire", so no expiry is set. - // Logical expiry is handled by the IdempotencyStore wrapper via TimeProvider. - var physicalTtl = _options.Value.TimeToLive + TimeSpan.FromHours(1); + var physicalTtl = GetPhysicalTimeToLive(); var timestamp = createdAt.ToString("O", CultureInfo.InvariantCulture); @@ -137,6 +161,21 @@ public async Task StoreAsync( .ConfigureAwait(false); } + /// + /// Returns the physical Redis expiry: TTL plus one hour headroom, or (no expiry) + /// when no TTL is configured. Logical expiry is handled by the IdempotencyStore wrapper via TimeProvider. + /// + private TimeSpan? GetPhysicalTimeToLive() => _options.Value.TimeToLive + TimeSpan.FromHours(1); + + private static bool IsExpired(RedisValue value, DateTimeOffset validFrom) => + DateTimeOffset.TryParse( + value.ToString(), + CultureInfo.InvariantCulture, + DateTimeStyles.RoundtripKind, + out var createdAt + ) + && createdAt < validFrom; + private string GetPrefixedKey(string idempotencyKey) => $"{_options.Value.Schema}:{_options.Value.TableName}:{idempotencyKey}"; } diff --git a/src/NetEvolve.Pulse.Redis/README.md b/src/NetEvolve.Pulse.Redis/README.md index 3ba77a79..8f0e2035 100644 --- a/src/NetEvolve.Pulse.Redis/README.md +++ b/src/NetEvolve.Pulse.Redis/README.md @@ -4,11 +4,12 @@ [![NuGet Downloads](https://img.shields.io/nuget/dt/NetEvolve.Pulse.Redis.svg)](https://www.nuget.org/packages/NetEvolve.Pulse.Redis/) [![License](https://img.shields.io/github/license/dailydevops/pulse.svg)](https://github.com/dailydevops/pulse/blob/main/LICENSE) -Redis idempotency provider for Pulse using `StackExchange.Redis`. Implements `IIdempotencyKeyRepository` with atomic `SET NX` operations (with expiry) for high-throughput, distributed idempotency enforcement without read-before-write round-trips. +Redis idempotency provider for Pulse using `StackExchange.Redis`. Implements `IIdempotencyKeyRepository` with atomic `SET NX` reservations (with expiry) for high-throughput, distributed idempotency enforcement. ## Features -- Atomic `SET key value NX` with expiry — single round-trip, no race conditions +- Atomic reservation with `SET key value NX` and expiry: of several concurrent submissions of the same key exactly one wins +- A logically expired key that is still physically present is replaced with a compare-and-set transaction (`WATCH`/`MULTI`); without a `TimeToLive` an existing key is never reserved again - Keys namespaced as `{Schema}:{TableName}:{idempotencyKey}` (default `pulse:IdempotencyKey:{idempotencyKey}`) - Logical TTL evaluated through `TimeProvider`, plus a physical Redis expiry for automatic cleanup - Startup validation via `ValidateOnStart()` diff --git a/src/NetEvolve.Pulse.Redis/RedisIdempotencyMediatorBuilderExtensions.cs b/src/NetEvolve.Pulse.Redis/RedisIdempotencyMediatorBuilderExtensions.cs index af3d5534..4605a304 100644 --- a/src/NetEvolve.Pulse.Redis/RedisIdempotencyMediatorBuilderExtensions.cs +++ b/src/NetEvolve.Pulse.Redis/RedisIdempotencyMediatorBuilderExtensions.cs @@ -14,7 +14,7 @@ namespace NetEvolve.Pulse; public static class RedisIdempotencyMediatorBuilderExtensions { /// - /// Adds a Redis-backed idempotency store using atomic SET NX EX operations. + /// Adds a Redis-backed idempotency store that reserves keys atomically with SET NX and a physical expiry. /// /// The mediator configurator. /// An optional action to configure . @@ -34,6 +34,10 @@ public static class RedisIdempotencyMediatorBuilderExtensions /// Key layout: /// Keys are stored as {Schema}:{TableName}:{idempotencyKey} with a physical Redis expiry of /// plus one hour, or without any expiry when no TTL is configured. + /// Reservation: + /// Of several concurrent reservations for the same key exactly one succeeds. A key that is logically expired + /// (older than ) but still physically present is replaced with a + /// compare-and-set transaction. Without a TTL an existing key is never reserved again. /// Note: /// Core idempotency services are registered automatically; calling /// before this method is optional but harmless. diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs index c31d216f..91e36f6e 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs @@ -377,4 +377,102 @@ public async Task StoreAsync_With_cancelled_token_throws_without_calling_Redis() _ = await Assert.That(capture.StringSetCalls).IsEmpty(); } } + + // INVARIANT (#816): TryStoreAsync reserves an absent key with a single atomic SET NX including the physical TTL. + [Test] + public async Task TryStoreAsync_Absent_key_is_reserved_with_SET_NX(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + var options = Options.Create(new IdempotencyKeyOptions { TimeToLive = TimeSpan.FromHours(6) }); + var repo = new RedisIdempotencyKeyRepository(mux, options); + var now = new DateTimeOffset(2025, 1, 1, 10, 0, 0, TimeSpan.Zero); + + var result = await repo.TryStoreAsync("k1", now, now.AddHours(-6), cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(result).IsTrue(); + _ = await Assert.That(capture.StringSetCalls).HasCount(1); + _ = await Assert.That(capture.StringSetCalls[0].When).IsEqualTo(When.NotExists); + _ = await Assert.That(capture.StringSetCalls[0].Expiry).IsEqualTo(TimeSpan.FromHours(7)); + _ = await Assert.That(capture.StringGetCalls).IsEmpty(); + } + } + + // INVARIANT (#790, #816): without a TTL an existing key is never reserved again. + [Test] + public async Task TryStoreAsync_Existing_key_without_validFrom_returns_false(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + capture.Storage["pulse:IdempotencyKey:k1"] = "2000-01-01T00:00:00.0000000+00:00"; + var repo = new RedisIdempotencyKeyRepository(mux, Options.Create(new IdempotencyKeyOptions())); + + var result = await repo.TryStoreAsync("k1", DateTimeOffset.UtcNow, null, cancellationToken) + .ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(result).IsFalse(); + _ = await Assert.That(capture.StringGetCalls).IsEmpty(); +#pragma warning disable S8969 // RedisValue's implicit string conversion is annotated nullable; the value is never null here + _ = await Assert + .That((string)capture.Storage["pulse:IdempotencyKey:k1"]!) + .IsEqualTo("2000-01-01T00:00:00.0000000+00:00"); +#pragma warning restore S8969 + } + } + + // INVARIANT (#816): a key created at or after validFrom is still valid and is not replaced. + [Test] + public async Task TryStoreAsync_Valid_existing_key_returns_false(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + capture.Storage["pulse:IdempotencyKey:k1"] = "2025-01-01T10:00:00.0000000+00:00"; + var repo = new RedisIdempotencyKeyRepository(mux, Options.Create(new IdempotencyKeyOptions())); + var validFrom = new DateTimeOffset(2025, 1, 1, 10, 0, 0, TimeSpan.Zero); + + var result = await repo.TryStoreAsync("k1", validFrom.AddMinutes(30), validFrom, cancellationToken) + .ConfigureAwait(false); + + _ = await Assert.That(result).IsFalse(); + } + + // INVARIANT (#816): an unparseable value fails closed, like ExistsAsync, and is not replaced. + [Test] + public async Task TryStoreAsync_Unparsable_existing_value_returns_false(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + capture.Storage["pulse:IdempotencyKey:k1"] = "not-a-timestamp"; + var repo = new RedisIdempotencyKeyRepository(mux, Options.Create(new IdempotencyKeyOptions())); + var now = new DateTimeOffset(2025, 1, 1, 10, 0, 0, TimeSpan.Zero); + + var result = await repo.TryStoreAsync("k1", now, now.AddHours(-1), cancellationToken).ConfigureAwait(false); + + _ = await Assert.That(result).IsFalse(); + } + + [Test] + public async Task TryStoreAsync_With_cancelled_token_throws_without_calling_Redis() + { + var (mux, capture) = BuildFakes(); + var repo = new RedisIdempotencyKeyRepository(mux, Options.Create(new IdempotencyKeyOptions())); + using var cts = new CancellationTokenSource(); + await cts.CancelAsync().ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert + .That(() => repo.TryStoreAsync("k1", DateTimeOffset.UtcNow, null, cts.Token)) + .Throws(); + _ = await Assert.That(capture.StringSetCalls).IsEmpty(); + } + } } From c0da7b45c76ac7e99fb1bb9331a716e4eeb0c7ed Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 01:22:25 +0200 Subject: [PATCH 4/7] test(idempotency): run reservation tests against all providers from the shared base Reservation and TTL tests now live in IdempotencyTestsBase. The concurrency tests and the refresh assertion only run for providers that reserve keys atomically (currently Redis); other providers keep the non-atomic behaviour tracked in #814. --- .../Idempotency/IdempotencyTestsBase.cs | 255 ++++++++++++++++++ .../Idempotency/RedisIdempotencyTests.cs | 226 +--------------- 2 files changed, 257 insertions(+), 224 deletions(-) diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs index 37c5b807..5d32ebc4 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, @@ -292,6 +300,224 @@ 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_Exactly_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_Logically_Expired_Key_That_Is_Still_Physically_Present( + 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("expired-reserve-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 second = await store.TryReserveAsync("expired-reserve-key", token).ConfigureAwait(false); + var third = await store.TryReserveAsync("expired-reserve-key", token).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(first).IsTrue(); + _ = await Assert.That(second).IsTrue(); + + // Non-atomic providers keep the old timestamp on re-reservation (tracked in #814). + if (SupportsAtomicReservation) + { + _ = await Assert.That(third).IsFalse(); + } + } + }, + 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; } @@ -305,4 +531,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 10f6ce40..937d2816 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs @@ -1,17 +1,9 @@ -namespace NetEvolve.Pulse.Tests.Integration.Idempotency; +namespace NetEvolve.Pulse.Tests.Integration.Idempotency; -using Microsoft.Extensions.DependencyInjection; -using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; -using NetEvolve.Pulse.Extensibility; -using NetEvolve.Pulse.Extensibility.Idempotency; -using NetEvolve.Pulse.Idempotency; using NetEvolve.Pulse.Tests.Integration.Internals; using NetEvolve.Pulse.Tests.Integration.Internals.Idempotency; using NetEvolve.Pulse.Tests.Integration.Internals.Services; -using TUnit.Assertions; -using TUnit.Assertions.Extensions; -using TUnit.Core; [ClassDataSource(Shared = [SharedType.None, SharedType.None])] [TestGroup("Redis")] @@ -19,219 +11,5 @@ namespace NetEvolve.Pulse.Tests.Integration.Idempotency; public class RedisIdempotencyTests(IServiceFixture databaseServiceFixture, IServiceInitializer databaseInitializer) : IdempotencyTestsBase(databaseServiceFixture, databaseInitializer) { - private const int ConcurrentCalls = 20; - - [Test] - public async Task Should_Reserve_Exactly_Once_When_Reserving_Same_Key_Concurrently( - CancellationToken cancellationToken - ) => - await RunAndVerify( - async (services, token) => - { - // Connect the shared multiplexer 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(); - - var counter = new InvocationCounter(); - - await RunAndVerify( - async (services, token) => - { - // Connect the shared multiplexer 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_Logically_Expired_Key_That_Is_Still_Physically_Present( - 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("expired-reserve-key", token).ConfigureAwait(false); - - // The physical Redis expiry is TTL + 1h, so the key is still present after 90 minutes. - fakeTime.Advance(TimeSpan.FromMinutes(90)); - - var second = await store.TryReserveAsync("expired-reserve-key", token).ConfigureAwait(false); - var third = await store.TryReserveAsync("expired-reserve-key", token).ConfigureAwait(false); - - using (Assert.Multiple()) - { - _ = await Assert.That(first).IsTrue(); - _ = await Assert.That(second).IsTrue(); - _ = await Assert.That(third).IsFalse(); - } - }, - 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 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; - } - } + protected override bool SupportsAtomicReservation => true; } From 92e1ce4f6f8bff009a168d9d57a940d0e18cd8fd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 01:29:14 +0200 Subject: [PATCH 5/7] docs(idempotency): point reservation atomicity at the key repository and track SQL/EF follow-up in #907 --- .../EntityFrameworkIdempotencyKeyRepository{TContext}.cs | 2 +- .../Idempotency/MySqlIdempotencyKeyRepository.cs | 2 +- .../Idempotency/PostgreSqlIdempotencyKeyRepository.cs | 2 +- .../Idempotency/SQLiteIdempotencyKeyRepository.cs | 2 +- .../Idempotency/SqlServerIdempotencyKeyRepository.cs | 2 +- .../IdempotencyCommandInterceptor{TRequest,TResponse}.cs | 7 ++++--- .../Idempotency/IdempotencyTestsBase.cs | 4 ++-- 7 files changed, 11 insertions(+), 10 deletions(-) diff --git a/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs b/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs index 8db29a3d..29c663c8 100644 --- a/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs +++ b/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs @@ -58,7 +58,7 @@ public Task ExistsAsync( /// /// - /// Composes and and is therefore not atomic. + /// Composes and and is therefore not atomic (tracked in #907). /// public async Task TryStoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs index bb3644ad..70285c81 100644 --- a/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs @@ -121,7 +121,7 @@ public async Task ExistsAsync( /// /// - /// Composes and and is therefore not atomic. + /// Composes and and is therefore not atomic (tracked in #907). /// public async Task TryStoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs index 675f9955..c1b6e9f6 100644 --- a/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs @@ -95,7 +95,7 @@ public async Task ExistsAsync( /// /// - /// Composes and and is therefore not atomic. + /// Composes and and is therefore not atomic (tracked in #907). /// public async Task TryStoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs index ad898f8e..19ddc60f 100644 --- a/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs @@ -121,7 +121,7 @@ public async Task ExistsAsync( /// /// - /// Composes and and is therefore not atomic. + /// Composes and and is therefore not atomic (tracked in #907). /// public async Task TryStoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs index 99e10221..b2842af3 100644 --- a/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs @@ -101,7 +101,7 @@ public async Task ExistsAsync( /// /// - /// Composes and and is therefore not atomic. + /// Composes and and is therefore not atomic (tracked in #907). /// public async Task TryStoreAsync( string idempotencyKey, diff --git a/src/NetEvolve.Pulse/Interceptors/IdempotencyCommandInterceptor{TRequest,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/IdempotencyCommandInterceptor{TRequest,TResponse}.cs index cb25df4a..7558a58c 100644 --- a/src/NetEvolve.Pulse/Interceptors/IdempotencyCommandInterceptor{TRequest,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/IdempotencyCommandInterceptor{TRequest,TResponse}.cs @@ -27,9 +27,10 @@ /// 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. +/// The Redis provider reserves atomically; the SQL and Entity Framework providers compose the existence +/// check and the store operation and leave a small window between them. /// 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 5d32ebc4..3ecb3382 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs @@ -327,7 +327,7 @@ CancellationToken cancellationToken { cancellationToken.ThrowIfCancellationRequested(); - Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically."); + Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically (tracked in #907)."); await RunAndVerify( async (services, token) => @@ -377,7 +377,7 @@ CancellationToken cancellationToken { cancellationToken.ThrowIfCancellationRequested(); - Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically."); + Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically (tracked in #907)."); var counter = new InvocationCounter(); From 690ba40990731f68cc44997699742c4907534fef Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 01:50:25 +0200 Subject: [PATCH 6/7] test(redis): cover the compare-and-set path for concurrent re-reservation of expired keys --- .../Idempotency/IdempotencyTestsBase.cs | 76 ++++++++- ...isIdempotencyKeyRepositoryBehaviorTests.cs | 158 +++++++++++++++++- 2 files changed, 231 insertions(+), 3 deletions(-) diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs index 3ecb3382..c96e7c11 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs @@ -327,7 +327,10 @@ CancellationToken cancellationToken { cancellationToken.ThrowIfCancellationRequested(); - Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically (tracked in #907)."); + Skip.When( + !SupportsAtomicReservation, + "The provider does not reserve idempotency keys atomically (tracked in #907)." + ); await RunAndVerify( async (services, token) => @@ -377,7 +380,10 @@ CancellationToken cancellationToken { cancellationToken.ThrowIfCancellationRequested(); - Skip.When(!SupportsAtomicReservation, "The provider does not reserve idempotency keys atomically (tracked in #907)."); + Skip.When( + !SupportsAtomicReservation, + "The provider does not reserve idempotency keys atomically (tracked in #907)." + ); var counter = new InvocationCounter(); @@ -482,6 +488,72 @@ await RunAndVerify( .ConfigureAwait(false); } + [Test] + public async Task Should_Reserve_Logically_Expired_Key_Exactly_Once_When_Reserving_Concurrently( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + Skip.When( + !SupportsAtomicReservation, + "The provider does not reserve idempotency keys atomically (tracked in #907)." + ); + + 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 diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs index 91e36f6e..13c1e674 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs @@ -13,7 +13,7 @@ namespace NetEvolve.Pulse.Tests.Unit.Redis; /// Behavioral invariants for . /// IConnectionMultiplexer/IDatabase are very large interfaces; this file uses /// to intercept the few methods the repository actually calls -/// (GetDatabase, StringSetAsync, StringGetAsync) and capture their +/// (GetDatabase, StringSetAsync, StringGetAsync, CreateTransaction) and capture their /// arguments. Every other call routes to so accidental /// new dependencies on Redis methods surface immediately. /// @@ -30,6 +30,13 @@ internal class FakeDatabase : DispatchProxy public List StringSetCalls { get; } = new(); public List StringGetCalls { get; } = new(); public Dictionary Storage { get; } = new(StringComparer.Ordinal); + public List Transactions { get; } = new(); + + // Simulates a concurrent writer that changes the key before EXEC, so the compare-and-set fails. + public bool TransactionConditionFails { get; set; } + + // Simulates the key expiring physically between SET NX and GET. + public bool RemoveKeyOnGet { get; set; } protected override object? Invoke(MethodInfo? targetMethod, object?[]? args) { @@ -66,17 +73,74 @@ internal class FakeDatabase : DispatchProxy { StringGetCalls.Add(key); #pragma warning disable S8969 // RedisKey's implicit string conversion is annotated nullable; the value is never null here + if (RemoveKeyOnGet) + { + _ = Storage.Remove((string)key!); + } + return Task.FromResult(Storage.TryGetValue((string)key!, out var v) ? v : RedisValue.Null); #pragma warning restore S8969 } throw new NotSupportedException($"Unexpected StringGetAsync overload: {targetMethod}"); } + case nameof(IDatabase.CreateTransaction): + { + var transactionProxy = DispatchProxy.Create(); + var transaction = (FakeTransaction)(object)transactionProxy; + transaction.Database = this; + Transactions.Add(transaction); + return transactionProxy; + } default: throw new NotImplementedException($"FakeDatabase has no behavior for {targetMethod}"); } } } + internal class FakeTransaction : DispatchProxy + { + public FakeDatabase? Database { get; set; } + public List Conditions { get; } = new(); + public List StringSetCalls { get; } = new(); + public bool? Executed { get; private set; } + + protected override object? Invoke(MethodInfo? targetMethod, object?[]? args) + { + if (targetMethod is null) + { + throw new InvalidOperationException("Null target method"); + } + + switch (targetMethod.Name) + { + case nameof(ITransaction.AddCondition) when args is [Condition condition]: + Conditions.Add(condition); + return null; + case nameof(ITransaction.StringSetAsync) + when args is { Length: >= 4 } && args[0] is RedisKey key && args[1] is RedisValue value: + StringSetCalls.Add( + new StringSetCall(key, value, (TimeSpan?)args[2], (When)(args[3] ?? When.Always)) + ); + return Task.FromResult(true); + case nameof(ITransaction.ExecuteAsync): + Executed = !Database!.TransactionConditionFails; + if (Executed.Value) + { + foreach (var call in StringSetCalls) + { +#pragma warning disable S8969 // RedisKey's implicit string conversion is annotated nullable; the value is never null here + Database.Storage[(string)call.Key!] = call.Value; +#pragma warning restore S8969 + } + } + + return Task.FromResult(Executed.Value); + default: + throw new NotImplementedException($"FakeTransaction has no behavior for {targetMethod}"); + } + } + } + internal class FakeMultiplexer : DispatchProxy { public IDatabase? Database { get; set; } @@ -459,6 +523,98 @@ public async Task TryStoreAsync_Unparsable_existing_value_returns_false(Cancella _ = await Assert.That(result).IsFalse(); } + // INVARIANT (#816): a logically expired key is replaced only through a compare-and-set on the observed value. + [Test] + public async Task TryStoreAsync_Expired_existing_key_is_replaced_by_compare_and_set( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + capture.Storage["pulse:IdempotencyKey:k1"] = "2025-01-01T08:00:00.0000000+00:00"; + var options = Options.Create(new IdempotencyKeyOptions { TimeToLive = TimeSpan.FromHours(1) }); + var repo = new RedisIdempotencyKeyRepository(mux, options); + var now = new DateTimeOffset(2025, 1, 1, 10, 0, 0, TimeSpan.Zero); + var expectedCondition = Condition.StringEqual("pulse:IdempotencyKey:k1", "2025-01-01T08:00:00.0000000+00:00"); + + var result = await repo.TryStoreAsync("k1", now, now.AddHours(-1), cancellationToken).ConfigureAwait(false); + + _ = await Assert.That(capture.Transactions).HasCount(1); + var transaction = capture.Transactions[0]; + + using (Assert.Multiple()) + { + _ = await Assert.That(result).IsTrue(); + _ = await Assert.That(transaction.Conditions).HasCount(1); + _ = await Assert.That(transaction.Conditions[0].ToString()).IsEqualTo(expectedCondition.ToString()); + _ = await Assert.That(transaction.StringSetCalls).HasCount(1); + _ = await Assert.That(transaction.StringSetCalls[0].When).IsEqualTo(When.Always); + _ = await Assert.That(transaction.StringSetCalls[0].Expiry).IsEqualTo(TimeSpan.FromHours(2)); + _ = await Assert.That(transaction.Executed).IsTrue(); +#pragma warning disable S8969 // RedisValue's implicit string conversion is annotated nullable; the value is never null here + _ = await Assert + .That((string)capture.Storage["pulse:IdempotencyKey:k1"]!) + .IsEqualTo("2025-01-01T10:00:00.0000000+00:00"); +#pragma warning restore S8969 + } + } + + // INVARIANT (#816): when a concurrent call changed the expired key first, the compare-and-set fails and + // this call does not reserve the key. + [Test] + public async Task TryStoreAsync_Expired_key_changed_concurrently_returns_false(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + capture.Storage["pulse:IdempotencyKey:k1"] = "2025-01-01T08:00:00.0000000+00:00"; + capture.TransactionConditionFails = true; + var repo = new RedisIdempotencyKeyRepository(mux, Options.Create(new IdempotencyKeyOptions())); + var now = new DateTimeOffset(2025, 1, 1, 10, 0, 0, TimeSpan.Zero); + + var result = await repo.TryStoreAsync("k1", now, now.AddHours(-1), cancellationToken).ConfigureAwait(false); + + _ = await Assert.That(capture.Transactions).HasCount(1); + + using (Assert.Multiple()) + { + _ = await Assert.That(result).IsFalse(); + _ = await Assert.That(capture.Transactions[0].Executed).IsFalse(); +#pragma warning disable S8969 // RedisValue's implicit string conversion is annotated nullable; the value is never null here + _ = await Assert + .That((string)capture.Storage["pulse:IdempotencyKey:k1"]!) + .IsEqualTo("2025-01-01T08:00:00.0000000+00:00"); +#pragma warning restore S8969 + } + } + + // INVARIANT (#816): a key that expires physically between SET NX and GET is reserved with a second SET NX. + [Test] + public async Task TryStoreAsync_Key_expiring_between_SET_NX_and_GET_is_reserved_with_SET_NX( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + capture.Storage["pulse:IdempotencyKey:k1"] = "2025-01-01T08:00:00.0000000+00:00"; + capture.RemoveKeyOnGet = true; + var repo = new RedisIdempotencyKeyRepository(mux, Options.Create(new IdempotencyKeyOptions())); + var now = new DateTimeOffset(2025, 1, 1, 10, 0, 0, TimeSpan.Zero); + + var result = await repo.TryStoreAsync("k1", now, now.AddHours(-1), cancellationToken).ConfigureAwait(false); + + _ = await Assert.That(capture.StringSetCalls).HasCount(2); + + using (Assert.Multiple()) + { + _ = await Assert.That(result).IsTrue(); + _ = await Assert.That(capture.StringSetCalls[1].When).IsEqualTo(When.NotExists); + _ = await Assert.That(capture.Transactions).IsEmpty(); + } + } + [Test] public async Task TryStoreAsync_With_cancelled_token_throws_without_calling_Redis() { From 9455ef34a99acab03e55d2a70d7e1721072b99f5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 02:17:01 +0200 Subject: [PATCH 7/7] test(idempotency): keep reservation test names within the MySQL 64-character table name limit --- .../Idempotency/IdempotencyTestsBase.cs | 12 +++--------- 1 file changed, 3 insertions(+), 9 deletions(-) diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs index c96e7c11..d7008563 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs @@ -321,9 +321,7 @@ await RunAndVerify( .ConfigureAwait(false); [Test] - public async Task Should_Reserve_Exactly_Once_When_Reserving_Same_Key_Concurrently( - CancellationToken cancellationToken - ) + public async Task Should_Reserve_Once_When_Reserving_Same_Key_Concurrently(CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); @@ -445,9 +443,7 @@ await mediator } [Test] - public async Task Should_Reserve_Logically_Expired_Key_That_Is_Still_Physically_Present( - CancellationToken cancellationToken - ) + public async Task Should_Reserve_Expired_Key_That_Is_Still_Physically_Present(CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); @@ -489,9 +485,7 @@ await RunAndVerify( } [Test] - public async Task Should_Reserve_Logically_Expired_Key_Exactly_Once_When_Reserving_Concurrently( - CancellationToken cancellationToken - ) + public async Task Should_Reserve_Expired_Key_Once_When_Reserving_Concurrently(CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested();