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
Expand Up @@ -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;
Expand All @@ -21,8 +22,9 @@ namespace NetEvolve.Pulse.Idempotency;
/// 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>ISO-8601 Timestamps:</strong></para>
/// Stores <see cref="DateTimeOffset"/> values as ISO-8601 text strings, using SQLite's
/// native text affinity for reliable lexicographic ordering and TTL-based queries.
/// Stores <see cref="DateTimeOffset"/> 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 <c>julianday()</c> instead.
/// </remarks>
[SuppressMessage(
"Reliability",
Expand Down Expand Up @@ -71,6 +73,7 @@ public SQLiteIdempotencyKeyRepository(IOptions<IdempotencyKeyOptions> options)
_enableWalMode = opts.EnableWalMode;

var table = opts.FullTableName;
var createdAt = $"{table}.\"{IdempotencyKeySchema.Columns.CreatedAt}\"";

_existsSql = $"""
SELECT 1 FROM {table}
Expand All @@ -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;
""";

Expand All @@ -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, "<")};
""";
}

Expand All @@ -123,7 +126,7 @@ public async Task<bool> 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);
Expand All @@ -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);
}
Expand All @@ -176,17 +179,44 @@ public async Task<bool> 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;
}
}
}

/// <summary>
/// Formats <paramref name="value"/> as an ISO-8601 round-trip string in UTC (<c>+00:00</c> suffix),
/// so that stored values compare chronologically as text.
/// </summary>
/// <param name="value">The timestamp to format.</param>
/// <returns>The UTC round-trip representation of <paramref name="value"/>.</returns>
private static string ToUtcString(DateTimeOffset value) =>
value.ToUniversalTime().ToString("O", CultureInfo.InvariantCulture);

/// <summary>
/// Builds a SQL condition that compares the stored <paramref name="column"/> with <c>@validFrom</c> as points in time.
/// </summary>
/// <param name="column">The qualified <c>CreatedAt</c> column.</param>
/// <param name="comparison">The SQL comparison operator.</param>
/// <returns>The SQL condition.</returns>
/// <remarks>
/// UTC rows (<c>+00:00</c> suffix) are compared as text, which is exact to 100 ns.
/// Rows stored with another offset by earlier versions fall back to <c>julianday()</c>,
/// which is precise to the millisecond.
/// </remarks>
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)
""";

/// <summary>
/// Opens and returns a new <see cref="SqliteConnection"/> using the stored connection string.
/// Applies WAL mode once per repository instance when <see cref="IdempotencyKeyOptions.EnableWalMode"/> is
Expand Down
Original file line number Diff line number Diff line change
@@ -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<SQLiteDatabaseServiceFixture, SQLiteAdoNetIdempotencyInitializer>(
Shared = [SharedType.None, SharedType.None]
Expand All @@ -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<IIdempotencyKeyRepository>();

// 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<IIdempotencyKeyRepository>();

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<IIdempotencyKeyRepository>();

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<IIdempotencyKeyRepository>();

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<IIdempotencyKeyRepository>();

// 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<IOptions<IdempotencyKeyOptions>>().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);
}
}
}
}
Loading