Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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<bool> 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`.
Original file line number Diff line number Diff line change
Expand Up @@ -67,11 +67,86 @@ public async Task StoreAsync(

ArgumentException.ThrowIfNullOrWhiteSpace(idempotencyKey);

_ = await TryInsertAsync(idempotencyKey, createdAt, cancellationToken).ConfigureAwait(false);
}

/// <inheritdoc />
/// <remarks>
/// An expired key is deleted and then inserted again instead of being updated in place, because
/// the Oracle MySQL provider cannot bind converted <see cref="DateTimeOffset"/> values in
/// <c>ExecuteUpdateAsync</c> 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.
/// </remarks>
public async Task<bool> 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);
}

/// <summary>
/// Deletes the stored key when it was created before <paramref name="validFrom"/>.
/// </summary>
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);
}
}

/// <summary>
/// Inserts the key unless it already exists.
/// </summary>
/// <returns><see langword="true"/> if the key was inserted; <see langword="false"/> if it already existed.</returns>
private async Task<bool> 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 };
Expand All @@ -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;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,37 @@ Task<bool> ExistsAsync(
/// <returns>A task representing the asynchronous store operation.</returns>
/// <remarks>
/// 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 <c>IdempotencyStore</c> reserves and stores keys through <see cref="TryReserveAsync"/>;
/// this member remains for direct callers only.
/// </remarks>
Task StoreAsync(string idempotencyKey, DateTimeOffset createdAt, CancellationToken cancellationToken = default);

/// <summary>
/// Atomically reserves an idempotency key: inserts it when absent, or refreshes its creation
/// timestamp when the stored key has expired.
/// </summary>
/// <param name="idempotencyKey">The idempotency key to reserve.</param>
/// <param name="createdAt">The timestamp to associate with the reserved key.</param>
/// <param name="validFrom">
/// When set, a stored key created before this timestamp is expired and is overwritten with
/// <paramref name="createdAt"/>. When <see langword="null"/>, keys never expire and an existing
/// key is never modified.
/// </param>
/// <param name="cancellationToken">A token to monitor for cancellation requests.</param>
/// <returns>
/// <see langword="true"/> if the key was inserted or an expired key was refreshed;
/// <see langword="false"/> if a key that has not expired already exists (duplicate submission).
/// </returns>
/// <remarks>
/// 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 <see langword="true"/>, and MUST NOT
/// modify a key that has not expired.
/// </remarks>
Task<bool> TryReserveAsync(
string idempotencyKey,
DateTimeOffset createdAt,
DateTimeOffset? validFrom = null,
CancellationToken cancellationToken = default
);
}
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,9 @@ internal sealed class MySqlIdempotencyKeyRepository : IIdempotencyKeyRepository
/// <summary>Cached SQL statement for inserting an idempotency key.</summary>
private readonly string _insertSql;

/// <summary>Cached SQL statement for refreshing the timestamp of an expired idempotency key.</summary>
private readonly string _refreshExpiredSql;

/// <summary>
/// Initializes a new instance of the <see cref="MySqlIdempotencyKeyRepository"/> class.
/// </summary>
Expand Down Expand Up @@ -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
""";
}

/// <inheritdoc />
Expand Down Expand Up @@ -144,6 +157,71 @@ public async Task StoreAsync(
}
}

/// <inheritdoc />
/// <remarks>
/// Runs <c>INSERT IGNORE</c> and, when <paramref name="validFrom"/> is set, a conditional
/// <c>UPDATE</c> 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 <see langword="true"/>.
/// When the refresh matches no row, the <c>INSERT IGNORE</c> runs once more, because the expired row
/// may have been deleted by a cleanup job between the two statements.
/// </remarks>
public async Task<bool> 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;
}
}
}

/// <summary>
/// Opens and returns a new <see cref="MySqlConnection"/> using the stored connection string.
/// The caller is responsible for disposing the connection.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ namespace NetEvolve.Pulse.Idempotency;
/// <para><strong>Duplicate Key Handling:</strong></para>
/// Uses <c>ON CONFLICT DO NOTHING</c> to handle duplicate key inserts gracefully.
/// Concurrent inserts of the same key are idempotent and will not throw exceptions.
/// Reservations use <c>ON CONFLICT ... DO UPDATE ... WHERE</c>, which also refreshes the timestamp of an expired key.
/// <para><strong>Performance:</strong></para>
/// Leverages stored functions for efficient operations and index utilization.
/// </remarks>
Expand All @@ -44,6 +45,9 @@ internal sealed class PostgreSqlIdempotencyKeyRepository : IIdempotencyKeyReposi
/// <summary>Cached SQL for inserting an idempotency key.</summary>
private readonly string _insertSql;

/// <summary>Cached SQL for atomically reserving or refreshing an idempotency key.</summary>
private readonly string _reserveSql;

/// <summary>
/// Initializes a new instance of the <see cref="PostgreSqlIdempotencyKeyRepository"/> class.
/// </summary>
Expand All @@ -62,6 +66,7 @@ public PostgreSqlIdempotencyKeyRepository(IOptions<IdempotencyKeyOptions> 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)";
}

/// <inheritdoc />
Expand Down Expand Up @@ -126,6 +131,39 @@ public async Task StoreAsync(
}
}

/// <inheritdoc />
public async Task<bool> 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;
}
}
}

/// <summary>
/// Opens and returns a new <see cref="NpgsqlConnection"/> using the stored connection string.
/// The caller is responsible for disposing the connection.
Expand Down
2 changes: 2 additions & 0 deletions src/NetEvolve.Pulse.PostgreSql/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading