diff --git a/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs b/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs index fd0c729f..09a01596 100644 --- a/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs +++ b/src/NetEvolve.Pulse.SQLite/Idempotency/SQLiteIdempotencyKeyRepository.cs @@ -2,6 +2,7 @@ namespace NetEvolve.Pulse.Idempotency; using System; using System.Diagnostics.CodeAnalysis; +using System.Globalization; using System.Threading; using System.Threading.Tasks; using Microsoft.Data.Sqlite; @@ -21,8 +22,9 @@ namespace NetEvolve.Pulse.Idempotency; /// 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. +/// Stores values as ISO-8601 text strings (round-trip format, normalized to UTC), +/// so that fixed-width lexicographic ordering matches chronological ordering for TTL-based queries. +/// Rows written with a non-UTC offset by earlier versions are compared via julianday() instead. /// [SuppressMessage( "Reliability", @@ -71,6 +73,7 @@ public SQLiteIdempotencyKeyRepository(IOptions options) _enableWalMode = opts.EnableWalMode; var table = opts.FullTableName; + var createdAt = $"{table}.\"{IdempotencyKeySchema.Columns.CreatedAt}\""; _existsSql = $""" SELECT 1 FROM {table} @@ -81,7 +84,7 @@ SELECT 1 FROM {table} _existsWithTtlSql = $""" SELECT 1 FROM {table} WHERE "{IdempotencyKeySchema.Columns.IdempotencyKey}" = @key - AND "{IdempotencyKeySchema.Columns.CreatedAt}" >= @validFrom + AND {UtcCondition(createdAt, ">=")} LIMIT 1; """; @@ -97,7 +100,7 @@ INSERT INTO {table} 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; + WHERE @validFrom IS NOT NULL AND {UtcCondition(createdAt, "<")}; """; } @@ -123,7 +126,7 @@ public async Task ExistsAsync( if (validFrom.HasValue) { - _ = command.Parameters.AddWithValue("@validFrom", validFrom.Value.ToString("O")); + _ = command.Parameters.AddWithValue("@validFrom", ToUtcString(validFrom.Value)); } var result = await command.ExecuteScalarAsync(cancellationToken).ConfigureAwait(false); @@ -150,7 +153,7 @@ public async Task StoreAsync( await using (command.ConfigureAwait(false)) { _ = command.Parameters.AddWithValue("@key", idempotencyKey); - _ = command.Parameters.AddWithValue("@createdAt", createdAt.ToString("O")); + _ = command.Parameters.AddWithValue("@createdAt", ToUtcString(createdAt)); _ = await command.ExecuteNonQueryAsync(cancellationToken).ConfigureAwait(false); } @@ -176,10 +179,10 @@ public async Task TryReserveAsync( await using (command.ConfigureAwait(false)) { _ = command.Parameters.AddWithValue("@key", idempotencyKey); - _ = command.Parameters.AddWithValue("@createdAt", createdAt.ToString("O")); + _ = command.Parameters.AddWithValue("@createdAt", ToUtcString(createdAt)); _ = command.Parameters.AddWithValue( "@validFrom", - validFrom.HasValue ? validFrom.Value.ToString("O") : DBNull.Value + validFrom.HasValue ? ToUtcString(validFrom.Value) : DBNull.Value ); return await command.ExecuteNonQueryAsync(cancellationToken).ConfigureAwait(false) > 0; @@ -187,6 +190,33 @@ public async Task TryReserveAsync( } } + /// + /// Formats as an ISO-8601 round-trip string in UTC (+00:00 suffix), + /// so that stored values compare chronologically as text. + /// + /// The timestamp to format. + /// The UTC round-trip representation of . + private static string ToUtcString(DateTimeOffset value) => + value.ToUniversalTime().ToString("O", CultureInfo.InvariantCulture); + + /// + /// Builds a SQL condition that compares the stored with @validFrom as points in time. + /// + /// The qualified CreatedAt column. + /// The SQL comparison operator. + /// The SQL condition. + /// + /// UTC rows (+00:00 suffix) are compared as text, which is exact to 100 ns. + /// Rows stored with another offset by earlier versions fall back to julianday(), + /// which is precise to the millisecond. + /// + private static string UtcCondition(string column, string comparison) => + $""" + (CASE WHEN substr({column}, -6) = '+00:00' + THEN {column} {comparison} @validFrom + ELSE julianday({column}) {comparison} julianday(@validFrom) END) + """; + /// /// Opens and returns a new using the stored connection string. /// Applies WAL mode once per repository instance when is diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/SQLiteAdoNetIdempotencyTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/SQLiteAdoNetIdempotencyTests.cs index 877f70a5..382c3f08 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/SQLiteAdoNetIdempotencyTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Idempotency/SQLiteAdoNetIdempotencyTests.cs @@ -1,9 +1,17 @@ namespace NetEvolve.Pulse.Tests.Integration.Idempotency; +using Microsoft.Data.Sqlite; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; 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 TUnit.Assertions; +using TUnit.Assertions.Extensions; +using TUnit.Core; [ClassDataSource( Shared = [SharedType.None, SharedType.None] @@ -14,4 +22,159 @@ public class SQLiteAdoNetIdempotencyTests( IServiceFixture databaseServiceFixture, IServiceInitializer databaseInitializer -) : IdempotencyTestsBase(databaseServiceFixture, databaseInitializer); +) : IdempotencyTestsBase(databaseServiceFixture, databaseInitializer) +{ + private static readonly TimeSpan Offset = TimeSpan.FromHours(2); + + [Test] + public async Task Should_Treat_Offset_Key_Before_Cutoff_As_Absent(CancellationToken cancellationToken) => + await RunAndVerify( + async (services, token) => + { + var repository = services.GetRequiredService(); + + // 13:00+02:00 is 11:00Z, one hour before the 12:00Z cutoff. + await repository + .StoreAsync("offset-key", TestDateTime.AddHours(-1).ToOffset(Offset), token) + .ConfigureAwait(false); + + var exists = await repository.ExistsAsync("offset-key", TestDateTime, token).ConfigureAwait(false); + + _ = await Assert.That(exists).IsFalse(); + }, + cancellationToken + ) + .ConfigureAwait(false); + + [Test] + public async Task Should_Treat_Key_After_Offset_Cutoff_As_Present(CancellationToken cancellationToken) => + await RunAndVerify( + async (services, token) => + { + var repository = services.GetRequiredService(); + + await repository.StoreAsync("utc-key", TestDateTime, token).ConfigureAwait(false); + + // 13:00+02:00 is 11:00Z, one hour before the key was created. + var exists = await repository + .ExistsAsync("utc-key", TestDateTime.AddHours(-1).ToOffset(Offset), token) + .ConfigureAwait(false); + + _ = await Assert.That(exists).IsTrue(); + }, + cancellationToken + ) + .ConfigureAwait(false); + + [Test] + public async Task Should_Reserve_Offset_Key_Before_Cutoff(CancellationToken cancellationToken) => + await RunAndVerify( + async (services, token) => + { + var repository = services.GetRequiredService(); + + await repository + .StoreAsync("offset-key", TestDateTime.AddHours(-1).ToOffset(Offset), token) + .ConfigureAwait(false); + + var reserved = await repository + .TryReserveAsync("offset-key", TestDateTime.AddMinutes(30), TestDateTime, token) + .ConfigureAwait(false); + + _ = await Assert.That(reserved).IsTrue(); + }, + cancellationToken + ) + .ConfigureAwait(false); + + [Test] + public async Task Should_Not_Reserve_Key_After_Offset_Cutoff(CancellationToken cancellationToken) => + await RunAndVerify( + async (services, token) => + { + var repository = services.GetRequiredService(); + + await repository.StoreAsync("utc-key", TestDateTime, token).ConfigureAwait(false); + + var reserved = await repository + .TryReserveAsync( + "utc-key", + TestDateTime.AddMinutes(30).ToOffset(Offset), + TestDateTime.AddHours(-1).ToOffset(Offset), + token + ) + .ConfigureAwait(false); + + _ = await Assert.That(reserved).IsFalse(); + }, + cancellationToken + ) + .ConfigureAwait(false); + + [Test] + public async Task Should_Compare_Legacy_Offset_Row_In_Utc(CancellationToken cancellationToken) => + await RunAndVerify( + async (services, token) => + { + var repository = services.GetRequiredService(); + + // Rows written before the UTC normalization keep the caller's offset: 13:00+02:00 is 11:00Z. + await InsertRawRowAsync(services, "legacy-old", "2025-01-01T13:00:00.0000000+02:00", token) + .ConfigureAwait(false); + // 15:00+02:00 is 13:00Z, one hour after the 12:00Z cutoff. + await InsertRawRowAsync(services, "legacy-new", "2025-01-01T15:00:00.0000000+02:00", token) + .ConfigureAwait(false); + + var oldExists = await repository + .ExistsAsync("legacy-old", TestDateTime, token) + .ConfigureAwait(false); + var newExists = await repository + .ExistsAsync("legacy-new", TestDateTime, token) + .ConfigureAwait(false); + var newReserved = await repository + .TryReserveAsync("legacy-new", TestDateTime.AddHours(2), TestDateTime, token) + .ConfigureAwait(false); + var oldReserved = await repository + .TryReserveAsync("legacy-old", TestDateTime.AddHours(2), TestDateTime, token) + .ConfigureAwait(false); + + _ = await Assert.That(oldExists).IsFalse(); + _ = await Assert.That(newExists).IsTrue(); + _ = await Assert.That(newReserved).IsFalse(); + _ = await Assert.That(oldReserved).IsTrue(); + }, + cancellationToken + ) + .ConfigureAwait(false); + + private static async Task InsertRawRowAsync( + IServiceProvider services, + string key, + string createdAt, + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + var options = services.GetRequiredService>().Value; + + var connection = new SqliteConnection(options.ConnectionString); + await using (connection.ConfigureAwait(false)) + { + await connection.OpenAsync(cancellationToken).ConfigureAwait(false); + +#pragma warning disable CA2100, S2077 // Table name comes from the test method name, not user input. + var command = new SqliteCommand( + $"""INSERT INTO "{options.TableName}" ("IdempotencyKey", "CreatedAt") VALUES (@key, @createdAt);""", + connection + ); +#pragma warning restore CA2100, S2077 + await using (command.ConfigureAwait(false)) + { + _ = command.Parameters.AddWithValue("@key", key); + _ = command.Parameters.AddWithValue("@createdAt", createdAt); + _ = await command.ExecuteNonQueryAsync(cancellationToken).ConfigureAwait(false); + } + } + } +}