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);
+ }
+ }
+ }
+}