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
193 changes: 191 additions & 2 deletions Darling/Darling.Tests/PgTargetBlockingTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -570,11 +570,19 @@ public void TheCollectorsReads_CarryTheFloorAsAParameter_TheExclusions_TheMessag
var events = PgTargetFactCollector.PgTargetLockWaitEventsSql;
Assert.Contains("family = 'lock_wait'", events, StringComparison.Ordinal);
/* Verified at source: the lock_wait parser lifts no metrics, so the duration is the message's own figure. */
Assert.Contains("substring(message from 'after ([0-9]+(?:\\.[0-9]+)?) ms')", events, StringComparison.Ordinal);
Assert.Contains("coalesce(duration_ms::DOUBLE PRECISION,", events, StringComparison.Ordinal);
Assert.Contains("substring(f.message from 'after ([0-9]+(?:\\.[0-9]+)?) ms')", events, StringComparison.Ordinal);
Assert.Contains("coalesce(f.duration_ms::DOUBLE PRECISION,", events, StringComparison.Ordinal);
Assert.Contains("message LIKE 'process % still waiting for %'", events, StringComparison.Ordinal);
Assert.Contains("message LIKE 'process % acquired %'", events, StringComparison.Ordinal);
Assert.Contains("message LIKE 'process % detected deadlock %'", events, StringComparison.Ordinal);
/* Each line counts once, in the window of its first sighting; each wait counts once, at the line that opens it. */
Assert.Contains("DISTINCT ON (e.raw_line_hash)", events, StringComparison.Ordinal);
Assert.Contains("p.collection_time >= coalesce(w.occurred_at", events, StringComparison.Ordinal);
Assert.Contains("p.collection_time <= l.first_collected", events, StringComparison.Ordinal);
Assert.Contains("e.collection_time AS first_collected", events, StringComparison.Ordinal);
Assert.Contains("p.collection_time < $2", events, StringComparison.Ordinal);
Assert.Contains("* INTERVAL '1 millisecond'", events, StringComparison.Ordinal);
Assert.Contains("(SELECT COUNT(*) FROM waits)", events, StringComparison.Ordinal);
Assert.Contains("collector_name = 'pg_log_events'", events, StringComparison.Ordinal);
Assert.Contains("status = 'SUCCESS'", events, StringComparison.Ordinal);

Expand Down Expand Up @@ -1383,6 +1391,187 @@ INSERT INTO pg_log_events
await command.ExecuteNonQueryAsync(ct);
}

/* ───────────── gated: each lock_wait line counts once, each wait counts once ───────────── */

private static readonly DateTime LockWinStart = new(2026, 6, 1, 10, 0, 0, DateTimeKind.Unspecified);
private static readonly DateTime LockWinEnd = new(2026, 6, 1, 11, 0, 0, DateTimeKind.Unspecified);
private const string LockOrders = "while updating tuple (0,7) in relation \"orders\"";

private sealed record LockWaitShape(long StillWaiting, long Acquired, long Deadlocks, long Lines, double? MaxWaitMs, double AcquiredWaitMs, long TopRelationWaits);

private static string StillWaiting(int pid, string ms) => $"process {pid} still waiting for ShareLock on transaction 900 after {ms} ms";

private static async Task<LockWaitShape[]> RunLockWaitScenarioAsync(
(DateTime Collected, DateTime Occurred, string Hash, int Pid, string Message)[] lines,
(DateTime Start, DateTime End)[] windows,
CancellationToken ct)
{
var cs = Environment.GetEnvironmentVariable("DARLING_TEST_PG");
Assert.SkipWhen(string.IsNullOrEmpty(cs), "Set DARLING_TEST_PG to a Postgres connection string to run the lock-wait counting pins.");

using var connection = new NpgsqlConnection(cs);
await connection.OpenAsync(ct);
await PgMigrations.MigrateAsync(connection, ct);
await DeleteRowsAsync(connection, ct);

var bodySucceeded = false;
try
{
await PgTargetFactCollectorTests.RegisterServerAsync(connection, ServerId, ServerName, "postgres", 18, ct);
foreach (var (collected, occurred, hash, pid, message) in lines)
{
using var insert = new NpgsqlCommand(@"
INSERT INTO pg_log_events
(collection_id, collection_time, server_id, server_name, occurred_at, family, severity, sqlstate, database_name, user_name, application_name,
pid, message, detail, context, statement_fingerprint, raw_line_hash, relation_name, duration_ms)
VALUES ($1, $2, $3, $4, $5, 'lock_wait', 'LOG', '00000', 'appdb', 'app', 'web',
$6, $7, 'Process holding the lock: 9000.', $8, 'fp-orders-update', $9, NULL, NULL)", connection);
insert.Parameters.AddWithValue(CollectionIdGenerator.Next());
insert.Parameters.AddWithValue(collected);
insert.Parameters.AddWithValue(ServerId);
insert.Parameters.AddWithValue(ServerName);
insert.Parameters.AddWithValue(occurred);
insert.Parameters.AddWithValue(pid);
insert.Parameters.AddWithValue(message);
insert.Parameters.AddWithValue(LockOrders);
insert.Parameters.AddWithValue(hash);
await insert.ExecuteNonQueryAsync(ct);
}

var shapes = new List<LockWaitShape>();
foreach (var (start, end) in windows)
{
using var read = new NpgsqlCommand(PgTargetFactCollector.PgTargetLockWaitEventsSql, connection);
read.Parameters.AddWithValue(ServerId);
read.Parameters.AddWithValue(start);
read.Parameters.AddWithValue(end);
using var reader = await read.ExecuteReaderAsync(ct);
Assert.True(await reader.ReadAsync(ct));
shapes.Add(new LockWaitShape(
reader.GetInt64(0), reader.GetInt64(1), reader.GetInt64(2), reader.GetInt64(3),
reader.IsDBNull(4) ? null : reader.GetDouble(4),
reader.IsDBNull(5) ? 0 : reader.GetDouble(5),
reader.IsDBNull(8) ? 0 : reader.GetInt64(8)));
}

bodySucceeded = true;
return shapes.ToArray();
}
finally
{
await LiveStoreCleanup.RunAsync(cs!, bodySucceeded, async (cleanup, cleanupCt) =>
await DeleteRowsAsync(cleanup, cleanupCt));
}
}

private static DateTime T(int h, int m, int s = 0) => new(2026, 6, 1, h, m, s, DateTimeKind.Unspecified);

[Fact]
public async Task ARepeatedLineHash_InsideTheWindow_CountsOnce()
{
var m = StillWaiting(4200, "1000.000");
var shape = (await RunLockWaitScenarioAsync(
[(T(10, 11), T(10, 10), "h1", 4200, m), (T(11, 0), T(10, 10), "h1", 4200, m)],
[(LockWinStart, LockWinEnd)], TestContext.Current.CancellationToken))[0];
Assert.Equal(1, shape.StillWaiting);
Assert.Equal(1, shape.Lines);
}

[Fact]
public async Task ALine_FirstStoredBeforeTheWindow_AndRestoredInside_IsNotCounted()
{
var m = StillWaiting(4200, "1000.000");
var shape = (await RunLockWaitScenarioAsync(
[(T(9, 1), T(9, 0), "h2", 4200, m), (T(10, 5), T(9, 0), "h2", 4200, m)],
[(LockWinStart, LockWinEnd)], TestContext.Current.CancellationToken))[0];
Assert.Equal(0, shape.StillWaiting);
Assert.Equal(0, shape.Lines);
}

[Fact]
public async Task ALine_ReSightedThreeHoursLater_IsNotCounted_WhereAFixedTwoHourLookbackWould()
{
var m = StillWaiting(4200, "1000.000");
var shape = (await RunLockWaitScenarioAsync(
[(T(7, 1), T(7, 0), "h3", 4200, m), (T(10, 20), T(7, 0), "h3", 4200, m)],
[(LockWinStart, LockWinEnd)], TestContext.Current.CancellationToken))[0];
Assert.Equal(0, shape.StillWaiting);
Assert.Equal(0, shape.Lines);
}

[Fact]
public async Task OutcomeLines_StoredTwice_CountOnce()
{
var acquired = "process 4200 acquired ShareLock on transaction 900 after 30000.000 ms";
var deadlock = "process 4300 detected deadlock while waiting for ShareLock on transaction 901 after 1000.000 ms";
var shape = (await RunLockWaitScenarioAsync(
[(T(10, 11), T(10, 10), "h4a", 4200, acquired), (T(10, 40), T(10, 10), "h4a", 4200, acquired),
(T(10, 21), T(10, 20), "h4d", 4300, deadlock), (T(10, 41), T(10, 20), "h4d", 4300, deadlock)],
[(LockWinStart, LockWinEnd)], TestContext.Current.CancellationToken))[0];
Assert.Equal(1, shape.Acquired);
Assert.Equal(30000, shape.AcquiredWaitMs);
Assert.Equal(1, shape.Deadlocks);
}

[Fact]
public async Task TwoLinesOfOneWait_CountOneWait()
{
var shape = (await RunLockWaitScenarioAsync(
[(T(10, 30, 1), T(10, 30, 0), "hb1", 4200, StillWaiting(4200, "1000.000")),
(T(10, 30, 5), T(10, 30, 4), "hb2", 4200, StillWaiting(4200, "5000.000"))],
[(LockWinStart, LockWinEnd)], TestContext.Current.CancellationToken))[0];
Assert.Equal(1, shape.StillWaiting);
Assert.Equal(1, shape.TopRelationWaits);
Assert.Equal(5000, shape.MaxWaitMs);
}

[Fact]
public async Task AWaitSplitAcrossTwoWindows_IsCountedInTheWindowOfItsOpeningLineOnly()
{
var shapes = await RunLockWaitScenarioAsync(
[(T(9, 59, 58), T(9, 59, 57), "hc1", 4200, StillWaiting(4200, "1000.000")),
(T(10, 0, 2), T(10, 0, 1), "hc2", 4200, StillWaiting(4200, "5000.000"))],
[(LockWinStart, LockWinEnd), (T(9, 0), LockWinStart)], TestContext.Current.CancellationToken);
Assert.Equal(0, shapes[0].StillWaiting);
Assert.Equal(1, shapes[1].StillWaiting);
}

[Theory]
[InlineData(true)]
[InlineData(false)]
public async Task TwoSeparateWaitsByOnePid_OnTheSameLock_CountTwice_EachLineStoredTwice(bool firstWaitAcquired)
{
var w1 = StillWaiting(4200, "1000.000");
var w2 = StillWaiting(4200, "1000.500");
var rows = new List<(DateTime, DateTime, string, int, string)>
{
(T(10, 41, 11), T(10, 40, 0), "hw1", 4200, w1), (T(10, 50), T(10, 40, 0), "hw1", 4200, w1),
(T(10, 41, 11), T(10, 41, 10), "hw2", 4200, w2), (T(10, 51), T(10, 41, 10), "hw2", 4200, w2),
};
if (firstWaitAcquired)
{
var acq = "process 4200 acquired ShareLock on transaction 900 after 30000.000 ms";
rows.Add((T(10, 40, 30), T(10, 40, 29), "hwa", 4200, acq));
rows.Add((T(10, 52), T(10, 40, 29), "hwa", 4200, acq));
}

var shape = (await RunLockWaitScenarioAsync(rows.ToArray(), [(LockWinStart, LockWinEnd)], TestContext.Current.CancellationToken))[0];
Assert.Equal(2, shape.StillWaiting);
}

// Pins the stated limitation of the probe, not a goal: wait 2 starts (1.5 s - 1000.5 ms) within one second of wait 1's opener, so it folds.
[Fact]
public async Task AReWaitStartingWithinOneSecondOfThePreviousOpener_FoldsIntoIt_TheStatedResidual()
{
var w1 = StillWaiting(4200, "1000.000");
var w2 = StillWaiting(4200, "1000.500");
var shape = (await RunLockWaitScenarioAsync(
[(T(10, 40, 1), T(10, 40, 0), "hf1", 4200, w1),
(T(10, 40, 1).AddMilliseconds(1500), T(10, 40, 0).AddMilliseconds(1500), "hf2", 4200, w2)],
[(LockWinStart, LockWinEnd)], TestContext.Current.CancellationToken))[0];
Assert.Equal(1, shape.StillWaiting);
}

private static async Task DeleteRowsAsync(NpgsqlConnection connection, CancellationToken ct)
{
using var cleanup = new NpgsqlCommand(
Expand Down
Loading
Loading