diff --git a/decisions/2026-09-28-idempotency-refresh-expired-keys-on-reserve.md b/decisions/2026-09-28-idempotency-refresh-expired-keys-on-reserve.md new file mode 100644 index 00000000..e42932fc --- /dev/null +++ b/decisions/2026-09-28-idempotency-refresh-expired-keys-on-reserve.md @@ -0,0 +1,60 @@ +--- +authors: + - Martin Stühmer + +applyTo: + - "src/NetEvolve.Pulse.Extensibility/Idempotency/*.cs" + - "src/NetEvolve.Pulse/Idempotency/*.cs" + - "src/**/Idempotency/*.cs" + - "src/**/Scripts/IdempotencyKey.sql" + +created: 2026-09-28 + +lastModified: 2026-09-28 + +state: proposed + +instructions: | + IdempotencyStore MUST reserve and store keys through IIdempotencyKeyRepository.TryReserveAsync(key, createdAt, validFrom), never through ExistsAsync followed by StoreAsync. + Every IIdempotencyKeyRepository implementation MUST make TryReserveAsync one atomic operation: insert an absent key, overwrite the timestamp of a key created before validFrom, leave a key that has not expired untouched, and return true for exactly one of several concurrent callers. + With validFrom = null (TimeToLive = null), an existing key MUST never be modified. +--- + +# Decision: Refresh Expired Idempotency Keys on Reserve + +Idempotency keys that have outlived `IdempotencyKeyOptions.TimeToLive` are refreshed, not left in place, when they are reserved or stored again. The check and the refresh happen in one atomic repository operation. + +## Context + +With a `TimeToLive`, `ExistsAsync` treats a key older than the cutoff as absent, and the library never deletes rows. Reservation used the default `IIdempotencyStore.TryReserveAsync` (`ExistsAsync` followed by `StoreAsync`), and every backend implemented `StoreAsync` as insert-if-absent. An expired key therefore kept its old timestamp forever. Every later duplicate of that key ran the handler again (#814). Refreshing the key needs the cutoff, but `IIdempotencyKeyRepository.StoreAsync` never receives it. A separate exists check followed by a write is also not race-safe. + +## Decision + +- `IIdempotencyKeyRepository` gets `Task TryReserveAsync(string idempotencyKey, DateTimeOffset createdAt, DateTimeOffset? validFrom, CancellationToken)`. It has no default implementation, in line with the pre-1.0 interface evolution decision. +- `IdempotencyStore` overrides `TryReserveAsync`, and routes `StoreAsync` through it, so `ExistsAsync` returns `true` after either call. +- Each provider implements the operation atomically with the tools of its backend: + - SQL Server: the `usp_ReserveIdempotencyKey` stored procedure with `MERGE ... WITH (HOLDLOCK)` and `WHEN MATCHED AND CreatedAt < @validFrom THEN UPDATE`. + - PostgreSQL: the `fn_reserve_idempotency_key` function with `ON CONFLICT ... DO UPDATE ... WHERE created_at < p_valid_from`. It is a new function because `CREATE OR REPLACE` cannot change the return type of `fn_insert_idempotency_key`. + - SQLite: `ON CONFLICT ... DO UPDATE ... WHERE`. + - MySQL: `INSERT IGNORE` followed by a conditional `UPDATE`, and one more `INSERT IGNORE` when the `UPDATE` matches no row, because a cleanup job can delete the expired row between the two statements. `ON DUPLICATE KEY UPDATE` is not used, because with the driver's default found-rows mode it reports the same row count for an insert and for an untouched duplicate. + - Entity Framework Core: `ExecuteDeleteAsync` of the expired row followed by the regular insert, whose primary key conflict returns `false`. `ExecuteUpdateAsync` is not used, because the Oracle MySQL provider cannot bind converted `DateTimeOffset` values in setters. InMemory uses change tracking, since it supports neither bulk operation. + - Redis: `SET NX` without a cutoff, and a Lua script with a cutoff that compares UTC round-trip timestamps and resets the expiry. Values with a non-UTC offset are treated as present, so a live key is never overwritten. + +## Consequences + +- A retried command is rejected again after its key has been re-reserved, in every provider. +- External implementers of `IIdempotencyKeyRepository` must implement the new member. +- Deployments of the SQL Server and PostgreSQL providers must re-run `IdempotencyKey.sql` together with the package upgrade. +- `IIdempotencyKeyRepository.StoreAsync` is no longer called by the library. It stays on the interface for direct callers. +- The Redis provider needs the scripting commands (`EVAL`, `EVALSHA`) once a `TimeToLive` is set. Redis values with a non-UTC offset, written by earlier versions through direct `StoreAsync` calls, stay unreservable until their physical expiry. + +## Alternatives Considered + +- **Refresh inside `StoreAsync`:** needs the cutoff as a new parameter, which is also a signature change, and keeps the non-atomic exists-then-store reservation. +- **Default interface implementation of `TryReserveAsync`:** any default would be exists-then-store and would silently keep the bug for external implementers. +- **Delete expired rows in the library:** needs a background job and still leaves a window between cleanup runs. + +## Related Decisions + +- [Extensibility Interface Evolution Before 1.0](./2026-09-24-extensibility-interface-evolution-pre-1-0.md) - governs how the new repository member is added and announced. +- [DateTimeOffset and TimeProvider Usage](./2026-01-21-datetimeoffset-and-timeprovider-usage.md) - the timestamps and the cutoff come from the injected `TimeProvider`. diff --git a/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs b/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs index 1a64d542..9cdda6c0 100644 --- a/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs +++ b/src/NetEvolve.Pulse.EntityFramework/Idempotency/EntityFrameworkIdempotencyKeyRepository{TContext}.cs @@ -67,11 +67,86 @@ public async Task StoreAsync( ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + _ = await TryInsertAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); + } + + /// + /// + /// An expired key is deleted and then inserted again instead of being updated in place, because + /// the Oracle MySQL provider cannot bind converted values in + /// ExecuteUpdateAsync setters. Both steps are safe under concurrency: only one caller deletes + /// the expired row, and the primary key lets only one caller insert the new one. + /// + public async Task TryReserveAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + + if (validFrom.HasValue) + { + await DeleteExpiredAsync(idempotencyKey, validFrom.Value, cancellationToken).ConfigureAwait(false); + } + + return await TryInsertAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false); + } + + /// + /// Deletes the stored key when it was created before . + /// + private async Task DeleteExpiredAsync( + string idempotencyKey, + DateTimeOffset validFrom, + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var tracked = _context.IdempotencyKeys.Local.FirstOrDefault(k => k.Key == idempotencyKey); + if (tracked is not null && tracked.CreatedAt < validFrom) + { + _context.Entry(tracked).State = EntityState.Detached; + } + + var expired = _context.IdempotencyKeys.Where(k => k.Key == idempotencyKey && k.CreatedAt < validFrom); + + if (_context.Database.IsRelational()) + { + _ = await expired.ExecuteDeleteAsync(cancellationToken).ConfigureAwait(false); + return; + } + + // Non-relational providers (e.g. InMemory) do not support ExecuteDeleteAsync. + var entry = await expired.FirstOrDefaultAsync(cancellationToken).ConfigureAwait(false); + if (entry is not null) + { + _ = _context.IdempotencyKeys.Remove(entry); + _ = await _context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); + } + } + + /// + /// Inserts the key unless it already exists. + /// + /// if the key was inserted; if it already existed. + private async Task TryInsertAsync( + string idempotencyKey, + DateTimeOffset createdAt, + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + // Check the local change tracker first to avoid a duplicate-tracking exception // from EF Core when the same key is stored twice within the same DbContext scope. if (_context.IdempotencyKeys.Local.Any(k => k.Key == idempotencyKey)) { - return; + return false; } var entry = new IdempotencyKey { Key = idempotencyKey, CreatedAt = createdAt }; @@ -81,12 +156,14 @@ public async Task StoreAsync( try { _ = await _context.SaveChangesAsync(cancellationToken).ConfigureAwait(false); + return true; } catch (Exception ex) when (IsDuplicateKeyException(ex) && IsIdempotencyKeyConflict(ex, entry)) { // A concurrent request already stored the same key — this is idempotent and safe to ignore. // Detach the conflicting entry so the context remains in a clean state. _context.Entry(entry).State = EntityState.Detached; + return false; } } diff --git a/src/NetEvolve.Pulse.Extensibility/Idempotency/IIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.Extensibility/Idempotency/IIdempotencyKeyRepository.cs index 608bbdfb..5ea551d2 100644 --- a/src/NetEvolve.Pulse.Extensibility/Idempotency/IIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.Extensibility/Idempotency/IIdempotencyKeyRepository.cs @@ -50,7 +50,37 @@ Task ExistsAsync( /// A task representing the asynchronous store operation. /// /// Implementations MUST handle duplicate-key exceptions gracefully and treat them as - /// a successful (idempotent) store operation. + /// a successful (idempotent) store operation. This member never refreshes an expired key. + /// The built-in IdempotencyStore reserves and stores keys through ; + /// this member remains for direct callers only. /// Task StoreAsync(string idempotencyKey, DateTimeOffset createdAt, CancellationToken cancellationToken = default); + + /// + /// Atomically reserves an idempotency key: inserts it when absent, or refreshes its creation + /// timestamp when the stored key has expired. + /// + /// The idempotency key to reserve. + /// The timestamp to associate with the reserved key. + /// + /// When set, a stored key created before this timestamp is expired and is overwritten with + /// . When , keys never expire and an existing + /// key is never modified. + /// + /// A token to monitor for cancellation requests. + /// + /// if the key was inserted or an expired key was refreshed; + /// if a key that has not expired already exists (duplicate submission). + /// + /// + /// Implementations MUST perform the check and the write as one atomic operation, so that exactly + /// one of several concurrent callers for the same key receives , and MUST NOT + /// modify a key that has not expired. + /// + Task TryReserveAsync( + 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..f063d1f4 100644 --- a/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.MySql/Idempotency/MySqlIdempotencyKeyRepository.cs @@ -51,6 +51,9 @@ internal sealed class MySqlIdempotencyKeyRepository : IIdempotencyKeyRepository /// Cached SQL statement for inserting an idempotency key. private readonly string _insertSql; + /// Cached SQL statement for refreshing the timestamp of an expired idempotency key. + private readonly string _refreshExpiredSql; + /// /// Initializes a new instance of the class. /// @@ -85,6 +88,16 @@ INSERT IGNORE INTO {table} (`{IdempotencyKeySchema.Columns.IdempotencyKey}`, `{IdempotencyKeySchema.Columns.CreatedAt}`) VALUES (@key, @createdAtTicks) """; + + // Only an expired row matches, so the affected-row count is unambiguous even with the + // driver's default found-rows semantics (unlike INSERT ... ON DUPLICATE KEY UPDATE, which + // reports 1 for both an insert and an unchanged duplicate). + _refreshExpiredSql = $""" + UPDATE {table} + SET `{IdempotencyKeySchema.Columns.CreatedAt}` = @createdAtTicks + WHERE `{IdempotencyKeySchema.Columns.IdempotencyKey}` = @key + AND `{IdempotencyKeySchema.Columns.CreatedAt}` < @validFromTicks + """; } /// @@ -144,6 +157,71 @@ public async Task StoreAsync( } } + /// + /// + /// Runs INSERT IGNORE and, when is set, a conditional + /// UPDATE of an expired row. InnoDB's row lock makes a concurrent refresh re-read the + /// already refreshed row, so only one caller for the same key receives . + /// When the refresh matches no row, the INSERT IGNORE runs once more, because the expired row + /// may have been deleted by a cleanup job between the two statements. + /// + public async Task TryReserveAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + + var connection = await CreateConnectionAsync(cancellationToken).ConfigureAwait(false); + await using (connection.ConfigureAwait(false)) + { + var insert = new MySqlCommand(_insertSql, connection); + await using (insert.ConfigureAwait(false)) + { + _ = insert.Parameters.AddWithValue("@key", idempotencyKey); + _ = insert.Parameters.AddWithValue("@createdAtTicks", createdAt.UtcTicks); + + if (await insert.ExecuteNonQueryAsync(cancellationToken).ConfigureAwait(false) > 0) + { + return true; + } + } + + if (!validFrom.HasValue) + { + return false; + } + + var refresh = new MySqlCommand(_refreshExpiredSql, connection); + await using (refresh.ConfigureAwait(false)) + { + _ = refresh.Parameters.AddWithValue("@key", idempotencyKey); + _ = refresh.Parameters.AddWithValue("@createdAtTicks", createdAt.UtcTicks); + _ = refresh.Parameters.AddWithValue("@validFromTicks", validFrom.Value.UtcTicks); + + if (await refresh.ExecuteNonQueryAsync(cancellationToken).ConfigureAwait(false) > 0) + { + return true; + } + } + + // A cleanup job can delete the expired row between the two statements; the key is then + // absent, so one more INSERT IGNORE reserves it instead of reporting a false duplicate. + var retry = new MySqlCommand(_insertSql, connection); + await using (retry.ConfigureAwait(false)) + { + _ = retry.Parameters.AddWithValue("@key", idempotencyKey); + _ = retry.Parameters.AddWithValue("@createdAtTicks", createdAt.UtcTicks); + + return await retry.ExecuteNonQueryAsync(cancellationToken).ConfigureAwait(false) > 0; + } + } + } + /// /// Opens and returns a new using the stored connection string. /// The caller is responsible for disposing the connection. diff --git a/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs index fe8d30ff..07dd0a09 100644 --- a/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.PostgreSql/Idempotency/PostgreSqlIdempotencyKeyRepository.cs @@ -20,6 +20,7 @@ namespace NetEvolve.Pulse.Idempotency; /// Duplicate Key Handling: /// Uses ON CONFLICT DO NOTHING to handle duplicate key inserts gracefully. /// Concurrent inserts of the same key are idempotent and will not throw exceptions. +/// Reservations use ON CONFLICT ... DO UPDATE ... WHERE, which also refreshes the timestamp of an expired key. /// Performance: /// Leverages stored functions for efficient operations and index utilization. /// @@ -44,6 +45,9 @@ internal sealed class PostgreSqlIdempotencyKeyRepository : IIdempotencyKeyReposi /// Cached SQL for inserting an idempotency key. private readonly string _insertSql; + /// Cached SQL for atomically reserving or refreshing an idempotency key. + private readonly string _reserveSql; + /// /// Initializes a new instance of the class. /// @@ -62,6 +66,7 @@ public PostgreSqlIdempotencyKeyRepository(IOptions option _existsSql = $"SELECT \"{schema}\".fn_exists_idempotency_key(@idempotency_key, @valid_from)"; _insertSql = $"SELECT \"{schema}\".fn_insert_idempotency_key(@idempotency_key, @created_at)"; + _reserveSql = $"SELECT \"{schema}\".fn_reserve_idempotency_key(@idempotency_key, @created_at, @valid_from)"; } /// @@ -126,6 +131,39 @@ public async Task StoreAsync( } } + /// + public async Task TryReserveAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + + var connection = await CreateConnectionAsync(cancellationToken).ConfigureAwait(false); + await using (connection.ConfigureAwait(false)) + { + var command = new NpgsqlCommand(_reserveSql, connection); + await using (command.ConfigureAwait(false)) + { + _ = command.Parameters.AddWithValue("idempotency_key", idempotencyKey); + _ = command.Parameters.AddWithValue("created_at", createdAt); + _ = command.Parameters.Add( + new NpgsqlParameter("valid_from", NpgsqlTypes.NpgsqlDbType.TimestampTz) + { + Value = validFrom.HasValue ? validFrom.Value : DBNull.Value, + } + ); + + var result = await command.ExecuteScalarAsync(cancellationToken).ConfigureAwait(false); + return result is true; + } + } + } + /// /// Opens and returns a new using the stored connection string. /// The caller is responsible for disposing the connection. diff --git a/src/NetEvolve.Pulse.PostgreSql/README.md b/src/NetEvolve.Pulse.PostgreSql/README.md index 80c4c19e..e959fc07 100644 --- a/src/NetEvolve.Pulse.PostgreSql/README.md +++ b/src/NetEvolve.Pulse.PostgreSql/README.md @@ -91,6 +91,8 @@ When upgrading from an earlier version, **re-run `OutboxMessage.sql`** against e `get_dead_letter_outbox_messages` now orders dead letters with equal `UpdatedAt` by `Id` descending, so paging returns every dead letter exactly once. Re-run `OutboxMessage.sql` to apply the updated function; its signature is unchanged. +The idempotency store reserves keys through `fn_reserve_idempotency_key` (`ON CONFLICT ... DO UPDATE ... WHERE`), which also refreshes the `created_at` of a key that has outlived `IdempotencyKeyOptions.TimeToLive`. **Re-run `IdempotencyKey.sql`** together with the package upgrade; it keeps the table and its data and creates the new function. + ## Quick Start ```csharp diff --git a/src/NetEvolve.Pulse.PostgreSql/Scripts/IdempotencyKey.sql b/src/NetEvolve.Pulse.PostgreSql/Scripts/IdempotencyKey.sql index 1fc93bfd..103752f0 100644 --- a/src/NetEvolve.Pulse.PostgreSql/Scripts/IdempotencyKey.sql +++ b/src/NetEvolve.Pulse.PostgreSql/Scripts/IdempotencyKey.sql @@ -79,6 +79,30 @@ BEGIN END; $$; +-- fn_reserve_idempotency_key: Atomically inserts an idempotency key or refreshes an expired one. +-- Returns TRUE when the key was inserted or refreshed, FALSE when a key that has not expired already exists. +CREATE OR REPLACE FUNCTION ":schema_name".fn_reserve_idempotency_key( + p_idempotency_key VARCHAR(500), + p_created_at TIMESTAMP WITH TIME ZONE, + p_valid_from TIMESTAMP WITH TIME ZONE DEFAULT NULL +) +RETURNS BOOLEAN +LANGUAGE plpgsql +AS $$ +DECLARE + affected_count INTEGER; +BEGIN + INSERT INTO ":schema_name".":table_name" AS t ("idempotency_key", "created_at") + VALUES (p_idempotency_key, p_created_at) + ON CONFLICT ("idempotency_key") DO UPDATE + SET "created_at" = EXCLUDED."created_at" + WHERE p_valid_from IS NOT NULL AND t."created_at" < p_valid_from; + + GET DIAGNOSTICS affected_count = ROW_COUNT; + RETURN affected_count > 0; +END; +$$; + -- fn_delete_expired_idempotency_keys: Removes expired idempotency keys (cleanup maintenance) CREATE OR REPLACE FUNCTION ":schema_name".fn_delete_expired_idempotency_keys( p_valid_from TIMESTAMP WITH TIME ZONE diff --git a/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs index e2b5c06f..a3720f67 100644 --- a/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.Redis/Idempotency/RedisIdempotencyKeyRepository.cs @@ -26,6 +26,27 @@ internal sealed class RedisIdempotencyKeyRepository : IIdempotencyKeyRepository { private const int DefaultDatabase = -1; + /// + /// Sets the key (ARGV[1] = UTC "O" timestamp, ARGV[3] = expiry in milliseconds or empty) unless it + /// holds a timestamp at or after the cutoff (ARGV[2], UTC "O"). Values are compared as text, so only + /// UTC values (ending in +00:00) can be refreshed. Values that are not timestamps, or carry a + /// non-UTC offset (written by earlier versions through direct calls), are + /// treated as present until their physical expiry, so a live key is never overwritten. + /// Returns 1 when the key was set, otherwise 0. + /// + private const string ReserveScript = """ + local current = redis.call('GET', KEYS[1]) + if current and (not string.match(current, '^%d%d%d%d%-.*%+00:00$') or current >= ARGV[2]) then + return 0 + end + if ARGV[3] == '' then + redis.call('SET', KEYS[1], ARGV[1]) + else + redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[3]) + end + return 1 + """; + private readonly IConnectionMultiplexer _multiplexer; private readonly IOptions _options; @@ -88,7 +109,7 @@ out var createdAt } /// - public async Task StoreAsync( + public Task StoreAsync( string idempotencyKey, DateTimeOffset createdAt, CancellationToken cancellationToken = default @@ -98,23 +119,59 @@ public async Task StoreAsync( cancellationToken.ThrowIfCancellationRequested(); + return TryReserveAsync(idempotencyKey, createdAt, null, cancellationToken); + } + + /// + /// + /// Without this is a plain SET NX. With it, a Lua script replaces an + /// expired value and resets its expiry atomically on the server. + /// + public async Task TryReserveAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + + cancellationToken.ThrowIfCancellationRequested(); + 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 timestamp = createdAt.ToString("O", CultureInfo.InvariantCulture); + var key = GetPrefixedKey(idempotencyKey); + var timestamp = FormatTimestamp(createdAt); - // Returns true when the key was set; false when the key already existed. - // Both outcomes are valid — no exception is thrown for duplicates. cancellationToken.ThrowIfCancellationRequested(); - _ = await database - .StringSetAsync(GetPrefixedKey(idempotencyKey), timestamp, physicalTtl, When.NotExists) + if (!validFrom.HasValue) + { + // Returns true when the key was set; false when the key already existed. + return await database.StringSetAsync(key, timestamp, physicalTtl, When.NotExists).ConfigureAwait(false); + } + + var expiry = physicalTtl.HasValue + ? ((long)physicalTtl.Value.TotalMilliseconds).ToString(CultureInfo.InvariantCulture) + : string.Empty; + + var result = await database + .ScriptEvaluateAsync(ReserveScript, [key], [timestamp, FormatTimestamp(validFrom.Value), expiry]) .ConfigureAwait(false); + + return (long)result == 1; } + /// + /// Formats a timestamp as UTC round-trip text, so that stored values compare lexicographically. + /// + private static string FormatTimestamp(DateTimeOffset value) => + value.ToUniversalTime().ToString("O", CultureInfo.InvariantCulture); + 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..4b19a913 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 single round-trip operations for high-throughput, distributed idempotency enforcement: `SET NX` when no `TimeToLive` is set, and a server-side Lua script that also refreshes expired keys when a `TimeToLive` is set. ## Features -- Atomic `SET key value NX` with expiry — single round-trip, no race conditions +- Without a `TimeToLive`: atomic `SET key value NX` in a single round-trip +- With a `TimeToLive`: one atomic Lua script (`EVAL`) that reserves an absent key, or refreshes an expired key and resets its `PX` expiry - 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()` @@ -74,4 +75,8 @@ No configuration section is bound automatically. To use `appsettings.json`, bind | `TimeToLive` | `null` | Logical expiry. When `null`, keys never expire logically. When set, it must be greater than zero and at most `TimeSpan.MaxValue` minus one hour (the physical expiry adds one hour). | The physical Redis expiry of each key is `TimeToLive` plus one hour. When `TimeToLive` is `null`, keys are stored without a Redis expiry. They are only removed if the server's `maxmemory-policy` evicts non-volatile keys (`allkeys-*`), and doing so breaks duplicate detection. Under `volatile-*` or `noeviction` policies the key space grows until the server runs out of memory. Set a `TimeToLive` if the key space must stay bounded. + +When `TimeToLive` is set, reservation runs a Lua script. The Redis user must be allowed to run scripting commands (`EVAL` and `EVALSHA`, for example through the ACL category `+@scripting`). Managed tiers or ACLs that block scripting make reservation fail once a `TimeToLive` is configured. + +The script compares stored timestamps as UTC round-trip text. The provider always writes UTC values. A value with a non-UTC offset, written by an earlier version through a direct `IIdempotencyKeyRepository.StoreAsync` call, is treated as present until its physical Redis expiry removes it. A value stored while `TimeToLive` was `null` has no physical expiry, so delete such keys manually if they must become reservable again. Invalid options cause an `OptionsValidationException` at startup or on first resolution of the options. diff --git a/src/NetEvolve.Pulse.Redis/RedisIdempotencyMediatorBuilderExtensions.cs b/src/NetEvolve.Pulse.Redis/RedisIdempotencyMediatorBuilderExtensions.cs index af3d5534..9db33b24 100644 --- a/src/NetEvolve.Pulse.Redis/RedisIdempotencyMediatorBuilderExtensions.cs +++ b/src/NetEvolve.Pulse.Redis/RedisIdempotencyMediatorBuilderExtensions.cs @@ -14,7 +14,8 @@ 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 using atomic SET NX (no TTL) or a + /// server-side Lua script that refreshes expired keys (TTL set). /// /// The mediator configurator. /// An optional action to configure . diff --git a/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs index 9cecb95e..fd0c729f 100644 --- a/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs @@ -19,6 +19,7 @@ namespace NetEvolve.Pulse.Idempotency; /// Duplicate Key Handling: /// Uses INSERT OR IGNORE to handle duplicate key inserts gracefully. /// Concurrent inserts of the same key are idempotent and will not throw exceptions. +/// Reservations use ON CONFLICT ... DO UPDATE ... WHERE, which also refreshes the timestamp of an expired key. /// ISO-8601 Timestamps: /// Stores values as ISO-8601 text strings, using SQLite's /// native text affinity for reliable lexicographic ordering and TTL-based queries. @@ -53,6 +54,9 @@ internal sealed class SQLiteIdempotencyKeyRepository : IIdempotencyKeyRepository /// Cached SQL statement for inserting an idempotency key. private readonly string _insertSql; + /// Cached SQL statement for atomically reserving or refreshing an idempotency key. + private readonly string _reserveSql; + /// /// Initializes a new instance of the class. /// @@ -86,6 +90,15 @@ INSERT OR IGNORE INTO {table} ("{IdempotencyKeySchema.Columns.IdempotencyKey}", "{IdempotencyKeySchema.Columns.CreatedAt}") VALUES (@key, @createdAt); """; + + _reserveSql = $""" + INSERT INTO {table} + ("{IdempotencyKeySchema.Columns.IdempotencyKey}", "{IdempotencyKeySchema.Columns.CreatedAt}") + VALUES (@key, @createdAt) + ON CONFLICT ("{IdempotencyKeySchema.Columns.IdempotencyKey}") DO UPDATE + SET "{IdempotencyKeySchema.Columns.CreatedAt}" = excluded."{IdempotencyKeySchema.Columns.CreatedAt}" + WHERE @validFrom IS NOT NULL AND {table}."{IdempotencyKeySchema.Columns.CreatedAt}" < @validFrom; + """; } /// @@ -144,6 +157,36 @@ public async Task StoreAsync( } } + /// + public async Task TryReserveAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + + var connection = await CreateConnectionAsync(cancellationToken).ConfigureAwait(false); + await using (connection.ConfigureAwait(false)) + { + var command = new SqliteCommand(_reserveSql, connection); + await using (command.ConfigureAwait(false)) + { + _ = command.Parameters.AddWithValue("@key", idempotencyKey); + _ = command.Parameters.AddWithValue("@createdAt", createdAt.ToString("O")); + _ = command.Parameters.AddWithValue( + "@validFrom", + validFrom.HasValue ? validFrom.Value.ToString("O") : DBNull.Value + ); + + return await command.ExecuteNonQueryAsync(cancellationToken).ConfigureAwait(false) > 0; + } + } + } + /// /// Opens and returns a new using the stored connection string. /// Applies WAL mode once per repository instance when is diff --git a/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs index 3756410d..7bbb7971 100644 --- a/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.SqlServer/Idempotency/SqlServerIdempotencyKeyRepository.cs @@ -21,6 +21,7 @@ namespace NetEvolve.Pulse.Idempotency; /// Duplicate Key Handling: /// Uses stored procedures with MERGE statement to handle duplicate key inserts gracefully. /// Concurrent inserts of the same key are idempotent and will not throw exceptions. +/// Reservations use MERGE ... WITH (HOLDLOCK), which also refreshes the timestamp of an expired key. /// Performance: /// Leverages stored procedures for efficient operations and index utilization. /// @@ -45,6 +46,9 @@ internal sealed class SqlServerIdempotencyKeyRepository : IIdempotencyKeyReposit /// Cached stored procedure name for inserting an idempotency key. private readonly string _insertSql; + /// Cached stored procedure name for atomically reserving or refreshing an idempotency key. + private readonly string _reserveSql; + /// /// Initializes a new instance of the class. /// @@ -63,6 +67,7 @@ public SqlServerIdempotencyKeyRepository(IOptions options _existsSql = $"[{schema}].[usp_ExistsIdempotencyKey]"; _insertSql = $"[{schema}].[usp_InsertIdempotencyKey]"; + _reserveSql = $"[{schema}].[usp_ReserveIdempotencyKey]"; } /// @@ -136,6 +141,51 @@ public async Task StoreAsync( } } + /// + public async Task TryReserveAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); + + var connection = await CreateConnectionAsync(cancellationToken).ConfigureAwait(false); + await using (connection.ConfigureAwait(false)) + { + var command = new SqlCommand(_reserveSql, connection) { CommandType = CommandType.StoredProcedure }; + await using (command.ConfigureAwait(false)) + { + _ = command.Parameters.Add( + new SqlParameter("@idempotencyKey", SqlDbType.NVarChar, 500) { Value = idempotencyKey } + ); + _ = command.Parameters.Add( + new SqlParameter("@createdAt", SqlDbType.DateTimeOffset) { Value = createdAt } + ); + _ = command.Parameters.Add( + new SqlParameter("@validFrom", SqlDbType.DateTimeOffset) + { + Value = validFrom.HasValue ? validFrom.Value : DBNull.Value, + } + ); + + try + { + var result = await command.ExecuteScalarAsync(cancellationToken).ConfigureAwait(false); + return result is true; + } + catch (SqlException ex) when (IsDuplicateKeyException(ex)) + { + // A concurrent request inserted the same key first — it owns the reservation. + return false; + } + } + } + } + /// /// Creates and opens a new SQL Server connection. /// diff --git a/src/NetEvolve.Pulse.SqlServer/README.md b/src/NetEvolve.Pulse.SqlServer/README.md index d7bda10e..4a57687f 100644 --- a/src/NetEvolve.Pulse.SqlServer/README.md +++ b/src/NetEvolve.Pulse.SqlServer/README.md @@ -152,7 +152,11 @@ sqlcmd -S your-server -d your-database -i IdempotencyKey.sql The script creates: - The `[IdempotencyKey]` table with `IdempotencyKey` (PK) and `CreatedAt` columns -- Stored procedures: `usp_ExistsIdempotencyKey`, `usp_InsertIdempotencyKey`, `usp_DeleteExpiredIdempotencyKeys` +- Stored procedures: `usp_ExistsIdempotencyKey`, `usp_InsertIdempotencyKey`, `usp_ReserveIdempotencyKey`, `usp_DeleteExpiredIdempotencyKeys` + +`usp_ReserveIdempotencyKey` reserves a key atomically (`MERGE ... WITH (HOLDLOCK)`). When `IdempotencyKeyOptions.TimeToLive` is set, it also refreshes the `CreatedAt` of an expired key, so duplicates are rejected again for the new window. + +**Upgrading:** re-run `IdempotencyKey.sql` together with the package upgrade. The script keeps the table and its data and recreates the stored procedures; the new package version calls `usp_ReserveIdempotencyKey`, which older scripts do not create. #### Using Idempotent Commands diff --git a/src/NetEvolve.Pulse.SqlServer/Scripts/IdempotencyKey.sql b/src/NetEvolve.Pulse.SqlServer/Scripts/IdempotencyKey.sql index e81e689b..6f5297a9 100644 --- a/src/NetEvolve.Pulse.SqlServer/Scripts/IdempotencyKey.sql +++ b/src/NetEvolve.Pulse.SqlServer/Scripts/IdempotencyKey.sql @@ -105,6 +105,57 @@ BEGIN END GO +-- usp_ReserveIdempotencyKey: Atomically inserts an idempotency key or refreshes an expired one. +-- Parameters: +-- @idempotencyKey Required. The idempotency key to reserve. +-- @createdAt Required. The creation timestamp stored for a new or refreshed key. +-- @validFrom Optional. Keys created before this cutoff count as expired and are refreshed; +-- NULL means keys never expire and an existing key is never modified. +-- Returns: one row with the BIT column [Reserved], 1 when the key was inserted or refreshed, +-- 0 when a key that has not expired already exists. +-- Errors: 50000 when @idempotencyKey or @createdAt is NULL; other errors are rethrown unchanged. +-- HOLDLOCK makes the MERGE serializable for the key range, so concurrent reservations cannot both win. +IF EXISTS (SELECT 1 FROM sys.objects WHERE [object_id] = OBJECT_ID(N'[$(SchemaName)].[usp_ReserveIdempotencyKey]') AND [type] = N'P') +BEGIN + DROP PROCEDURE [$(SchemaName)].[usp_ReserveIdempotencyKey]; +END +GO + +CREATE PROCEDURE [$(SchemaName)].[usp_ReserveIdempotencyKey] + @idempotencyKey NVARCHAR(500), + @createdAt DATETIMEOFFSET, + @validFrom DATETIMEOFFSET = NULL +AS +BEGIN + SET NOCOUNT ON; + + IF @idempotencyKey IS NULL OR @createdAt IS NULL + BEGIN + THROW 50000, N'usp_ReserveIdempotencyKey: @idempotencyKey and @createdAt must not be NULL.', 1; + END + + DECLARE @reserved BIT; + + BEGIN TRY + MERGE INTO [$(SchemaName)].[$(TableName)] WITH (HOLDLOCK) AS target + USING (SELECT @idempotencyKey AS [IdempotencyKey], @createdAt AS [CreatedAt]) AS source + ON (target.[IdempotencyKey] = source.[IdempotencyKey]) + WHEN MATCHED AND @validFrom IS NOT NULL AND target.[CreatedAt] < @validFrom THEN + UPDATE SET [CreatedAt] = source.[CreatedAt] + WHEN NOT MATCHED THEN + INSERT ([IdempotencyKey], [CreatedAt]) + VALUES (source.[IdempotencyKey], source.[CreatedAt]); + + SET @reserved = CASE WHEN @@ROWCOUNT > 0 THEN 1 ELSE 0 END; + END TRY + BEGIN CATCH + THROW; + END CATCH + + SELECT @reserved AS [Reserved]; +END +GO + -- usp_DeleteExpiredIdempotencyKeys: Removes expired idempotency keys (cleanup maintenance) IF EXISTS (SELECT 1 FROM sys.objects WHERE [object_id] = OBJECT_ID(N'[$(SchemaName)].[usp_DeleteExpiredIdempotencyKeys]') AND [type] = N'P') BEGIN diff --git a/src/NetEvolve.Pulse/Idempotency/IdempotencyKeyOptions.cs b/src/NetEvolve.Pulse/Idempotency/IdempotencyKeyOptions.cs index 77417952..65a59c20 100644 --- a/src/NetEvolve.Pulse/Idempotency/IdempotencyKeyOptions.cs +++ b/src/NetEvolve.Pulse/Idempotency/IdempotencyKeyOptions.cs @@ -40,7 +40,8 @@ public class IdempotencyKeyOptions /// /// /// TTL-based cleanup (physical row deletion) is out of scope and must be handled externally. - /// This option only controls whether expired keys are logically treated as absent. + /// When set, storing or reserving an expired key refreshes its creation timestamp, so duplicates are + /// rejected again for a new TTL window. When , existing keys are never modified. /// public TimeSpan? TimeToLive { get; set; } diff --git a/src/NetEvolve.Pulse/Idempotency/IdempotencyStore.cs b/src/NetEvolve.Pulse/Idempotency/IdempotencyStore.cs index 7da9cbb0..c398b607 100644 --- a/src/NetEvolve.Pulse/Idempotency/IdempotencyStore.cs +++ b/src/NetEvolve.Pulse/Idempotency/IdempotencyStore.cs @@ -14,7 +14,8 @@ namespace NetEvolve.Pulse.Idempotency; /// Time-to-Live: /// When is set, keys older than the TTL /// are treated as absent by . Physical deletion is not performed; -/// expired keys are logically ignored by passing a cutoff timestamp to the repository. +/// expired keys are logically ignored by passing a cutoff timestamp to the repository, and +/// and refresh the timestamp of an expired key. /// internal sealed class IdempotencyStore : IIdempotencyStore { @@ -50,20 +51,33 @@ public Task ExistsAsync(string idempotencyKey, CancellationToken cancellat ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); - DateTimeOffset? cutoff = _options.TimeToLive.HasValue - ? _timeProvider.GetUtcNow() - _options.TimeToLive.Value - : null; - - return _repository.ExistsAsync(idempotencyKey, cutoff, cancellationToken); + return _repository.ExistsAsync(idempotencyKey, GetCutoff(), cancellationToken); } /// - public Task StoreAsync(string idempotencyKey, CancellationToken cancellationToken = default) + /// + /// An expired key is refreshed, so that returns again afterwards. + /// + public Task StoreAsync(string idempotencyKey, CancellationToken cancellationToken = default) => + TryReserveAsync(idempotencyKey, cancellationToken); + + /// + /// + /// Delegates to the atomic , which also refreshes + /// the timestamp of a key that has outlived . + /// + public Task TryReserveAsync(string idempotencyKey, CancellationToken cancellationToken = default) { cancellationToken.ThrowIfCancellationRequested(); ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey); - return _repository.StoreAsync(idempotencyKey, _timeProvider.GetUtcNow(), cancellationToken); + var now = _timeProvider.GetUtcNow(); + return _repository.TryReserveAsync(idempotencyKey, now, GetCutoff(now), cancellationToken); } + + private DateTimeOffset? GetCutoff() => GetCutoff(_timeProvider.GetUtcNow()); + + private DateTimeOffset? GetCutoff(DateTimeOffset now) => + _options.TimeToLive.HasValue ? now - _options.TimeToLive.Value : null; } diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs index 37c5b807..9d9d94ec 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/IdempotencyTestsBase.cs @@ -262,6 +262,158 @@ await RunAndVerify( .ConfigureAwait(false); } + [Test] + public async Task Should_Reject_Duplicate_After_Expired_Reserve(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var fakeTime = new FakeTimeProvider(); + fakeTime.AdjustTime(TestDateTime); + + await RunAndVerify( + async (services, token) => + { + var store = services.GetRequiredService(); + + await store.StoreAsync("re-reserve-key", token).ConfigureAwait(false); + + fakeTime.Advance(TimeSpan.FromHours(2)); + + var reserved = await store.TryReserveAsync("re-reserve-key", token).ConfigureAwait(false); + var duplicate = await store.TryReserveAsync("re-reserve-key", token).ConfigureAwait(false); + var exists = await store.ExistsAsync("re-reserve-key", token).ConfigureAwait(false); + + _ = await Assert.That(reserved).IsTrue(); + _ = await Assert.That(duplicate).IsFalse(); + _ = await Assert.That(exists).IsTrue(); + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(fakeTime) + .Configure(o => o.TimeToLive = TimeSpan.FromHours(1)) + ) + .ConfigureAwait(false); + } + + [Test] + public async Task Should_Reject_Expired_Reserve_Dup_Across_Scopes(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var fakeTime = new FakeTimeProvider(); + fakeTime.AdjustTime(TestDateTime); + + await RunAndVerify( + async (services, token) => + { + var store = services.GetRequiredService(); + + await store.StoreAsync("re-reserve-scope-key", token).ConfigureAwait(false); + + fakeTime.Advance(TimeSpan.FromHours(2)); + + var scopeFactory = services.GetRequiredService(); + + var scope2 = scopeFactory.CreateAsyncScope(); + await using (scope2.ConfigureAwait(false)) + { + var store2 = scope2.ServiceProvider.GetRequiredService(); + var reserved = await store2 + .TryReserveAsync("re-reserve-scope-key", token) + .ConfigureAwait(false); + + _ = await Assert.That(reserved).IsTrue(); + } + + var scope3 = scopeFactory.CreateAsyncScope(); + await using (scope3.ConfigureAwait(false)) + { + var store3 = scope3.ServiceProvider.GetRequiredService(); + var duplicate = await store3 + .TryReserveAsync("re-reserve-scope-key", token) + .ConfigureAwait(false); + + _ = await Assert.That(duplicate).IsFalse(); + } + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(fakeTime) + .Configure(o => o.TimeToLive = TimeSpan.FromHours(1)) + ) + .ConfigureAwait(false); + } + + [Test] + public async Task Should_Report_Key_Present_After_Expired_Store(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var fakeTime = new FakeTimeProvider(); + fakeTime.AdjustTime(TestDateTime); + + await RunAndVerify( + async (services, token) => + { + var store = services.GetRequiredService(); + + await store.StoreAsync("re-store-key", token).ConfigureAwait(false); + + fakeTime.Advance(TimeSpan.FromHours(2)); + + await store.StoreAsync("re-store-key", token).ConfigureAwait(false); + + var result = await store.ExistsAsync("re-store-key", token).ConfigureAwait(false); + + _ = await Assert.That(result).IsTrue(); + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(fakeTime) + .Configure(o => o.TimeToLive = TimeSpan.FromHours(1)) + ) + .ConfigureAwait(false); + } + + [Test] + public async Task Should_Keep_CreatedAt_When_Reserving_Live_Key(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var fakeTime = new FakeTimeProvider(); + fakeTime.AdjustTime(TestDateTime); + + await RunAndVerify( + async (services, token) => + { + var store = services.GetRequiredService(); + + await store.StoreAsync("live-key", token).ConfigureAwait(false); + + fakeTime.Advance(TimeSpan.FromMinutes(30)); + + var duplicate = await store.TryReserveAsync("live-key", token).ConfigureAwait(false); + + // 75 minutes after the original store: only a refreshed CreatedAt would keep the key alive. + fakeTime.Advance(TimeSpan.FromMinutes(45)); + + var exists = await store.ExistsAsync("live-key", token).ConfigureAwait(false); + + _ = await Assert.That(duplicate).IsFalse(); + _ = await Assert.That(exists).IsFalse(); + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(fakeTime) + .Configure(o => o.TimeToLive = TimeSpan.FromHours(1)) + ) + .ConfigureAwait(false); + } + [Test] public async Task Should_Enforce_Idempotency_For_Void_Command_Through_Mediator( CancellationToken cancellationToken diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs index 84ecbd94..31aefb88 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/RedisIdempotencyTests.cs @@ -1,12 +1,57 @@ namespace NetEvolve.Pulse.Tests.Integration.Idempotency; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; +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 StackExchange.Redis; [ClassDataSource(Shared = [SharedType.None, SharedType.None])] [TestGroup("Redis")] [InheritsTests] public class RedisIdempotencyTests(IServiceFixture databaseServiceFixture, IServiceInitializer databaseInitializer) - : IdempotencyTestsBase(databaseServiceFixture, databaseInitializer); + : IdempotencyTestsBase(databaseServiceFixture, databaseInitializer) +{ + [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")] + public async Task Should_Keep_Live_Non_Utc_Value_On_Reserve( + string keyName, + string storedValue, + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var fakeTime = new FakeTimeProvider(); + fakeTime.AdjustTime(TestDateTime.AddHours(1)); + + await RunAndVerify( + async (services, token) => + { + var options = services.GetRequiredService>().Value; + var key = $"{options.Schema}:{options.TableName}:{keyName}"; + var database = services.GetRequiredService().GetDatabase(); + _ = await database.StringSetAsync(key, storedValue).ConfigureAwait(false); + + var store = services.GetRequiredService(); + var reserved = await store.TryReserveAsync(keyName, token).ConfigureAwait(false); + var value = await database.StringGetAsync(key).ConfigureAwait(false); + + _ = await Assert.That(reserved).IsFalse(); + _ = await Assert.That(value.ToString()).IsEqualTo(storedValue); + }, + cancellationToken, + configureServices: services => + services + .AddSingleton(fakeTime) + .Configure(o => o.TimeToLive = TimeSpan.FromHours(1)) + ) + .ConfigureAwait(false); + } +} diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Idempotency/IdempotencyStoreTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Idempotency/IdempotencyStoreTests.cs index b2db9163..50b97a0f 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Idempotency/IdempotencyStoreTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Idempotency/IdempotencyStoreTests.cs @@ -19,7 +19,12 @@ private static IdempotencyStore CreateStore( IIdempotencyKeyRepository repository, IdempotencyKeyOptions? options = null, TimeProvider? timeProvider = null - ) => new(repository, Options.Create(options ?? new IdempotencyKeyOptions()), timeProvider ?? TimeProvider.System); + ) => + new IdempotencyStore( + repository, + Options.Create(options ?? new IdempotencyKeyOptions()), + timeProvider ?? TimeProvider.System + ); [Test] public async Task Constructor_WithNullRepository_ThrowsArgumentNullException() => @@ -142,10 +147,84 @@ public async Task StoreAsync_PassesCurrentTimestampToRepository(CancellationToke _ = await Assert.That(repository.CapturedCreatedAt).IsEqualTo(expectedTimestamp); } + [Test] + public async Task TryReserveAsync_WithTtl_ReservesAtomicallyWithCutoff(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); + + var result = await store.TryReserveAsync("test-key", cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(result).IsTrue(); + _ = await Assert.That(repository.ReserveCount).IsEqualTo(1); + _ = await Assert.That(repository.CapturedCreatedAt).IsEqualTo(now); + _ = await Assert.That(repository.CapturedValidFrom).IsEqualTo(now - ttl); + } + } + + [Test] + public async Task TryReserveAsync_WithoutTtl_ReservesWithoutCutoff(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var repository = new TrackingIdempotencyKeyRepository(); + var store = CreateStore(repository, new IdempotencyKeyOptions { TimeToLive = null }); + + _ = await store.TryReserveAsync("test-key", cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(repository.ReserveCount).IsEqualTo(1); + _ = await Assert.That(repository.CapturedValidFrom).IsNull(); + } + } + + [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 StoreAsync_WithTtl_RefreshesExpiredKeyThroughReserve(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.StoreAsync("test-key", cancellationToken).ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(repository.ReserveCount).IsEqualTo(1); + _ = await Assert.That(repository.CapturedValidFrom).IsEqualTo(now - ttl); + } + } + private sealed class TrackingIdempotencyKeyRepository : IIdempotencyKeyRepository { public DateTimeOffset? CapturedValidFrom { get; private set; } = DateTimeOffset.MaxValue; public DateTimeOffset CapturedCreatedAt { get; private set; } + public int ReserveCount { get; private set; } public Task ExistsAsync( string idempotencyKey, @@ -170,5 +249,20 @@ public Task StoreAsync( CapturedCreatedAt = createdAt; return Task.CompletedTask; } + + public Task TryReserveAsync( + string idempotencyKey, + DateTimeOffset createdAt, + DateTimeOffset? validFrom = null, + CancellationToken cancellationToken = default + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + ReserveCount++; + CapturedCreatedAt = createdAt; + CapturedValidFrom = validFrom; + return Task.FromResult(true); + } } } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs index c31d216f..b0b7b85f 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Redis/RedisIdempotencyKeyRepositoryBehaviorTests.cs @@ -29,6 +29,7 @@ internal class FakeDatabase : DispatchProxy { public List StringSetCalls { get; } = new(); public List StringGetCalls { get; } = new(); + public List ScriptEvaluateCalls { get; } = new(); public Dictionary Storage { get; } = new(StringComparer.Ordinal); protected override object? Invoke(MethodInfo? targetMethod, object?[]? args) @@ -71,6 +72,15 @@ internal class FakeDatabase : DispatchProxy } throw new NotSupportedException($"Unexpected StringGetAsync overload: {targetMethod}"); } + case nameof(IDatabase.ScriptEvaluateAsync): + { + if (args is { Length: >= 3 } && args[0] is string && args[2] is RedisValue[] values) + { + ScriptEvaluateCalls.Add(values); + return Task.FromResult(RedisResult.Create((RedisValue)1)); + } + throw new NotSupportedException($"Unexpected ScriptEvaluateAsync overload: {targetMethod}"); + } default: throw new NotImplementedException($"FakeDatabase has no behavior for {targetMethod}"); } @@ -377,4 +387,60 @@ public async Task StoreAsync_With_cancelled_token_throws_without_calling_Redis() _ = await Assert.That(capture.StringSetCalls).IsEmpty(); } } + + // INVARIANT (#790): without a cutoff, TryReserveAsync is a plain SET NX and rejects an existing key. + [Test] + public async Task TryReserveAsync_Without_validFrom_uses_SET_NX(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + var repo = new RedisIdempotencyKeyRepository(mux, Options.Create(new IdempotencyKeyOptions())); + + var first = await repo.TryReserveAsync("k1", DateTimeOffset.UtcNow, null, cancellationToken) + .ConfigureAwait(false); + var second = await repo.TryReserveAsync("k1", DateTimeOffset.UtcNow, null, cancellationToken) + .ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(first).IsTrue(); + _ = await Assert.That(second).IsFalse(); + _ = await Assert.That(capture.StringSetCalls).HasCount(2); + _ = await Assert.That(capture.StringSetCalls[1].When).IsEqualTo(When.NotExists); + _ = await Assert.That(capture.ScriptEvaluateCalls).IsEmpty(); + } + } + + // INVARIANT (#814): with a cutoff, the refresh of an expired value runs as one server-side script + // that receives UTC timestamps and the physical expiry (TTL + 1h) in milliseconds. + [Test] + public async Task TryReserveAsync_With_validFrom_runs_reserve_script(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var (mux, capture) = BuildFakes(); + var options = Options.Create(new IdempotencyKeyOptions { TimeToLive = TimeSpan.FromHours(1) }); + var repo = new RedisIdempotencyKeyRepository(mux, options); + var createdAt = new DateTimeOffset(2025, 1, 1, 12, 0, 0, TimeSpan.FromHours(2)); + + var reserved = await repo.TryReserveAsync("k1", createdAt, createdAt.AddHours(-1), cancellationToken) + .ConfigureAwait(false); + + using (Assert.Multiple()) + { + _ = await Assert.That(reserved).IsTrue(); + _ = await Assert.That(capture.StringSetCalls).IsEmpty(); + _ = await Assert.That(capture.ScriptEvaluateCalls).HasCount(1); +#pragma warning disable S8969 // RedisValue's implicit string conversion is annotated nullable; the value is never null here + _ = await Assert + .That((string)capture.ScriptEvaluateCalls[0][0]!) + .IsEqualTo("2025-01-01T10:00:00.0000000+00:00"); + _ = await Assert + .That((string)capture.ScriptEvaluateCalls[0][1]!) + .IsEqualTo("2025-01-01T09:00:00.0000000+00:00"); + _ = await Assert.That((string)capture.ScriptEvaluateCalls[0][2]!).IsEqualTo("7200000"); +#pragma warning restore S8969 + } + } }