From 01c3497c274ac620a4447150e95868eb619d3b17 Mon Sep 17 00:00:00 2001 From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com> Date: Wed, 30 Sep 2026 20:25:44 -0400 Subject: [PATCH 1/4] Count each lock_wait log line once and each lock wait once in the blocking fact --- .../Darling.Tests/PgTargetBlockingTests.cs | 10 +- .../PgTargetFactCollector.Blocking.cs | 115 ++++++++++++++---- .../PgTargetAdvice.Blocking.cs | 4 +- .../PgTargetScorer.Blocking.cs | 7 +- 4 files changed, 106 insertions(+), 30 deletions(-) diff --git a/Darling/Darling.Tests/PgTargetBlockingTests.cs b/Darling/Darling.Tests/PgTargetBlockingTests.cs index 8a60d6e63..a94ef022c 100644 --- a/Darling/Darling.Tests/PgTargetBlockingTests.cs +++ b/Darling/Darling.Tests/PgTargetBlockingTests.cs @@ -570,11 +570,17 @@ 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 < $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); diff --git a/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs b/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs index 60b251b86..ae6b8c0cb 100644 --- a/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs +++ b/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs @@ -102,11 +102,31 @@ AND c.role_name IS NULL AND c.name IN ('deadlock_timeout', 'log_lock_waits')"; /// - /// The window's lock_wait family of pg_log_events (#3601) — the engine's own record of every lock - /// wait that outlived deadlock_timeout, EVENT grain, four line shapes: still waiting for … after N ms - /// (one per wait, written when the wait crosses the timeout — the COUNT of waits), acquired … after N ms - /// (the same wait ending — the wait's true end-to-end length), avoided deadlock … after N ms and - /// detected deadlock … after N ms. + /// The window's lock_wait family of pg_log_events (#3601) — the engine's own record of lock waits + /// that outlived deadlock_timeout, EVENT grain, four line shapes: still waiting for … after N ms + /// (written when a wait crosses the timeout, and written AGAIN each time the backend's latch wakes while it is + /// still waiting — PostgreSQL 18's ProcSleep — so one wait can leave several lines, each with a larger + /// N), acquired … after N ms (the same wait ending — the wait's true end-to-end length), + /// avoided deadlock … after N ms and detected deadlock … after N ms. + /// + /// Each stored line counts once, in the window of its first sighting. The resumed log read + /// re-reads an overlap, so the same line (raw_line_hash) can be stored by two passes. in_window + /// keeps one row per distinct line collected in the window, and first_seen drops a line that an earlier + /// pass already stored: a line is never collected before it occurred, so an earlier sighting can only lie in + /// [occurred_at, $2), a range bounded by the line's own time (a NULL occurred_at searches all + /// history; it is rare and stays correct). This holds the same assumption DarlingPgLogEventReader makes: + /// the target's clock is not ahead of the store's. + /// + /// Each lock wait counts once, at the line that opens it. waits keeps a still waiting + /// line only when no EARLIER line of the same wait exists anywhere in the table: same backend (the pid the message + /// names), same lock text, a strictly smaller after N ms, and an occurred_at within + /// [L.occurred_at − (N + 1 s), L.occurred_at]. A backend waits on one lock at a time and every line of a + /// wait lies between its start and the line's own time, so the probe is bounded by the line's own N; the + /// one-second pad covers whole-second %t stamps. Two separate waits by one backend are counted twice, + /// because the second wait's lines start after the first ended. A line with no parsable lock text, duration or + /// time counts as its own wait. Residual, stated: a re-wait on the same lock text that starts within one + /// second of the previous wait's last logged line, when that wait ended without a lock_wait line (cancelled or + /// timed out), folds into it. /// /// The duration is read off the message, not the duration_ms column. Verified at source /// (PgLockWaitEventParser): the lock_wait parser stores the line with NO metrics lifted — duration_ms @@ -122,51 +142,100 @@ AND c.role_name IS NULL /// and written just before it is this pass's. Stated, bounded by one collector cadence, and the same rule every /// PostgreSQL family applies to its table. /// + /// The output. still_waiting is the count of WAITS; acquired, deadlocks, + /// lines, max_wait_ms, acquired_wait_ms and last_event_at are over the distinct lines + /// first sighted in the window; top_relation and top_fingerprint weigh each wait once, at its + /// opener. + /// /// The two denominators ride the row. log_captures — SUCCESS runs of the pg_log_events /// collector in the window, from collection_log — tells "no lines" apart from "nobody was reading the log" /// (the collector is optional and needs the RDS log API or file access); without it a silent collector would /// read as a quiet server. $1 server_id, $2/$3 window. /// public const string PgTargetLockWaitEventsSql = @" -WITH events AS ( +WITH in_window AS ( + SELECT DISTINCT ON (e.raw_line_hash) + e.raw_line_hash, + e.message, + e.occurred_at, + e.statement_fingerprint, + e.pid, + e.relation_name, + e.context, + e.duration_ms + FROM pg_log_events AS e + WHERE e.server_id = $1 + AND e.collection_time >= $2 + AND e.collection_time <= $3 + AND e.family = 'lock_wait' + AND e.raw_line_hash IS NOT NULL + ORDER BY e.raw_line_hash, e.collection_time +), +first_seen AS ( + SELECT w.* + FROM in_window AS w + WHERE w.occurred_at >= $2 + OR NOT EXISTS ( + SELECT 1 + FROM pg_log_events AS p + WHERE p.server_id = $1 + AND p.collection_time >= coalesce(w.occurred_at, '-infinity'::timestamp) + AND p.collection_time < $2 + AND p.raw_line_hash = w.raw_line_hash) +), +lines AS ( SELECT - message, - occurred_at, - statement_fingerprint, - coalesce(relation_name, substring(coalesce(context, '') from 'in relation ""([^""]+)""')) AS relation_name, - coalesce(duration_ms::DOUBLE PRECISION, - NULLIF(substring(message from 'after ([0-9]+(?:\.[0-9]+)?) ms'), '')::DOUBLE PRECISION) AS wait_ms - FROM pg_log_events - WHERE server_id = $1 - AND collection_time >= $2 - AND collection_time <= $3 - AND family = 'lock_wait' + f.raw_line_hash, + f.message, + f.occurred_at, + f.statement_fingerprint, + coalesce(f.relation_name, substring(coalesce(f.context, '') from 'in relation ""([^""]+)""')) AS relation_name, + coalesce(f.duration_ms::DOUBLE PRECISION, + NULLIF(substring(f.message from 'after ([0-9]+(?:\.[0-9]+)?) ms'), '')::DOUBLE PRECISION) AS wait_ms, + coalesce(substring(f.message from '^process ([0-9]+) ')::INT, f.pid) AS wait_pid, + substring(f.message from '^process [0-9]+ still waiting for (.+) after [0-9.]+ ms') AS lock_text + FROM first_seen AS f +), +waits AS ( + SELECT l.* + FROM lines AS l + WHERE l.message LIKE 'process % still waiting for %' + AND NOT EXISTS ( + SELECT 1 + FROM pg_log_events AS p + WHERE p.server_id = $1 + AND p.family = 'lock_wait' + AND p.collection_time >= l.occurred_at - (l.wait_ms + 1000) * INTERVAL '1 millisecond' + AND p.occurred_at >= l.occurred_at - (l.wait_ms + 1000) * INTERVAL '1 millisecond' + AND p.occurred_at <= l.occurred_at + AND p.raw_line_hash <> l.raw_line_hash + AND coalesce(substring(p.message from '^process ([0-9]+) ')::INT, p.pid) = l.wait_pid + AND substring(p.message from '^process [0-9]+ still waiting for (.+) after [0-9.]+ ms') = l.lock_text + AND NULLIF(substring(p.message from 'after ([0-9]+(?:\.[0-9]+)?) ms'), '')::DOUBLE PRECISION < l.wait_ms) ), shape AS ( SELECT - COUNT(*) FILTER (WHERE message LIKE 'process % still waiting for %') AS still_waiting, + (SELECT COUNT(*) FROM waits) AS still_waiting, COUNT(*) FILTER (WHERE message LIKE 'process % acquired %') AS acquired, COUNT(*) FILTER (WHERE message LIKE 'process % detected deadlock %') AS deadlocks, COUNT(*) AS lines, MAX(wait_ms) AS max_wait_ms, coalesce(SUM(wait_ms) FILTER (WHERE message LIKE 'process % acquired %'), 0) AS acquired_wait_ms, MAX(occurred_at) AS last_event_at - FROM events + FROM lines ), top_relation AS ( SELECT relation_name, COUNT(*) AS waits - FROM events + FROM waits WHERE relation_name IS NOT NULL - AND message LIKE 'process % still waiting for %' GROUP BY relation_name ORDER BY waits DESC, relation_name LIMIT 1 ), top_fingerprint AS ( SELECT statement_fingerprint, COUNT(*) AS waits - FROM events + FROM waits WHERE statement_fingerprint IS NOT NULL - AND message LIKE 'process % still waiting for %' GROUP BY statement_fingerprint ORDER BY waits DESC, statement_fingerprint LIMIT 1 diff --git a/PerformanceMonitor.Analysis/PgTargetAdvice.Blocking.cs b/PerformanceMonitor.Analysis/PgTargetAdvice.Blocking.cs index ff1e7dd1d..b5f157cda 100644 --- a/PerformanceMonitor.Analysis/PgTargetAdvice.Blocking.cs +++ b/PerformanceMonitor.Analysis/PgTargetAdvice.Blocking.cs @@ -59,8 +59,8 @@ public static partial class PgTargetAdvice private static readonly AdviceBlock s_lockWaitEventsStatic = new( Headline: "The engine logged lock waits past deadlock_timeout — the written record of contention between samples", Investigation: - "With log_lock_waits = on, PostgreSQL writes one 'process N still waiting for on after " + - "N ms' line for every lock wait that outlives deadlock_timeout (1 s by default), and a matching 'acquired " + + "With log_lock_waits = on, PostgreSQL writes a 'process N still waiting for on after " + + "N ms' line when a lock wait outlives deadlock_timeout (1 s by default), and may write it again while the wait continues, plus a matching 'acquired " + "… after N ms' line when the wait ends — the wait's true end-to-end length. pg_log_events stores those " + "lines as the lock_wait family, at EVENT grain: complete where the one-minute pg_blocking sample is not, " + "and one line deep where the sample has the blocker's statement. The fact counts the waits, rates them " + diff --git a/PerformanceMonitor.Analysis/PgTargetScorer.Blocking.cs b/PerformanceMonitor.Analysis/PgTargetScorer.Blocking.cs index a3ed3157c..96552405d 100644 --- a/PerformanceMonitor.Analysis/PgTargetScorer.Blocking.cs +++ b/PerformanceMonitor.Analysis/PgTargetScorer.Blocking.cs @@ -92,8 +92,9 @@ public static partial class PgTargetScorer /// engine default assumed. public const string BlockingDeadlockTimeoutFromSnapshotKey = "deadlock_timeout_from_snapshot"; - /// Metadata key: the count of still waiting lines — one per wait that outlived - /// deadlock_timeout — in the window's lock_wait events. + /// Metadata key: the count of lock WAITS that outlived deadlock_timeout in the window's + /// lock_wait events. A wait the engine re-logs while it continues (several still waiting lines) counts + /// once, and a line stored by two collection passes counts once, in the window of its first sighting. public const string LockWaitEventsCountKey = "wait_events"; /// Metadata key: the longest after N ms any lock_wait line in the window carried — the /// engine's own end-to-end length of a wait (the acquired … after N ms line), GRADED on the chain's bars. @@ -103,7 +104,7 @@ public static partial class PgTargetScorer public const string LockWaitEventsAcquiredMsKey = "acquired_wait_ms"; /// Metadata key: detected deadlock lines in the window's lock_wait family. public const string LockWaitEventsDeadlocksKey = "deadlocks_detected"; - /// Metadata key: still waiting lines per OBSERVED hour (context.ObservedDurationMs). + /// Metadata key: lock waits (see ) per OBSERVED hour (context.ObservedDurationMs). public const string LockWaitEventsPerHourKey = "wait_events_per_hour"; /// Metadata key: 1 when the family cannot know (see the two reason flags); the fact then makes no claim. public const string LockWaitUnavailableKey = "unavailable"; From 49b50c756f32ffc05eae72160072a84573a796f1 Mon Sep 17 00:00:00 2001 From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com> Date: Wed, 30 Sep 2026 20:29:51 -0400 Subject: [PATCH 2/4] Bound the lock-wait opener probe at the window end and pin the once-only counting live --- .../Darling.Tests/PgTargetBlockingTests.cs | 169 ++++++++++++++++++ .../PgTargetFactCollector.Blocking.cs | 6 +- 2 files changed, 174 insertions(+), 1 deletion(-) diff --git a/Darling/Darling.Tests/PgTargetBlockingTests.cs b/Darling/Darling.Tests/PgTargetBlockingTests.cs index a94ef022c..6f3b7b49e 100644 --- a/Darling/Darling.Tests/PgTargetBlockingTests.cs +++ b/Darling/Darling.Tests/PgTargetBlockingTests.cs @@ -578,6 +578,7 @@ public void TheCollectorsReads_CarryTheFloorAsAParameter_TheExclusions_TheMessag /* 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 <= $3", 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); @@ -1389,6 +1390,174 @@ 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 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(); + 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.000"); + var rows = new List<(DateTime, DateTime, string, int, string)> + { + (T(10, 40, 1), 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); + } + private static async Task DeleteRowsAsync(NpgsqlConnection connection, CancellationToken ct) { using var cleanup = new NpgsqlCommand( diff --git a/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs b/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs index ae6b8c0cb..0b88eaa7a 100644 --- a/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs +++ b/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs @@ -126,7 +126,10 @@ AND c.role_name IS NULL /// because the second wait's lines start after the first ended. A line with no parsable lock text, duration or /// time counts as its own wait. Residual, stated: a re-wait on the same lock text that starts within one /// second of the previous wait's last logged line, when that wait ended without a lock_wait line (cancelled or - /// timed out), folds into it. + /// timed out), folds into it. The probe reads what was collected by the window's end ($3), so the answer + /// for a window is stable. A wait whose earlier line was collected only after that end is counted here at its first + /// line collected in the window. That needs out-of-order collection, which a single log file's read doesn't + /// produce. /// /// The duration is read off the message, not the duration_ms column. Verified at source /// (PgLockWaitEventParser): the lock_wait parser stores the line with NO metrics lifted — duration_ms @@ -206,6 +209,7 @@ FROM pg_log_events AS p WHERE p.server_id = $1 AND p.family = 'lock_wait' AND p.collection_time >= l.occurred_at - (l.wait_ms + 1000) * INTERVAL '1 millisecond' + AND p.collection_time <= $3 AND p.occurred_at >= l.occurred_at - (l.wait_ms + 1000) * INTERVAL '1 millisecond' AND p.occurred_at <= l.occurred_at AND p.raw_line_hash <> l.raw_line_hash From b880c105ab00d5a48453bf5f5f0a2aec15ab8a0e Mon Sep 17 00:00:00 2001 From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com> Date: Wed, 30 Sep 2026 20:34:27 -0400 Subject: [PATCH 3/4] Bound the lock-wait opener probe at the line's own first sighting, not the window end --- Darling/Darling.Tests/PgTargetBlockingTests.cs | 3 ++- .../PgTargetFactCollector.Blocking.cs | 14 ++++++++------ 2 files changed, 10 insertions(+), 7 deletions(-) diff --git a/Darling/Darling.Tests/PgTargetBlockingTests.cs b/Darling/Darling.Tests/PgTargetBlockingTests.cs index 6f3b7b49e..cd46b8aa1 100644 --- a/Darling/Darling.Tests/PgTargetBlockingTests.cs +++ b/Darling/Darling.Tests/PgTargetBlockingTests.cs @@ -578,7 +578,8 @@ public void TheCollectorsReads_CarryTheFloorAsAParameter_TheExclusions_TheMessag /* 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 <= $3", 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); diff --git a/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs b/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs index 0b88eaa7a..e9bc39b3c 100644 --- a/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs +++ b/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs @@ -126,10 +126,10 @@ AND c.role_name IS NULL /// because the second wait's lines start after the first ended. A line with no parsable lock text, duration or /// time counts as its own wait. Residual, stated: a re-wait on the same lock text that starts within one /// second of the previous wait's last logged line, when that wait ended without a lock_wait line (cancelled or - /// timed out), folds into it. The probe reads what was collected by the window's end ($3), so the answer - /// for a window is stable. A wait whose earlier line was collected only after that end is counted here at its first - /// line collected in the window. That needs out-of-order collection, which a single log file's read doesn't - /// produce. + /// timed out), folds into it. The probe reads what was collected up to the opener candidate's own first + /// sighting (first_collected), because an earlier line of the same wait is read before it, in the same pass + /// or an earlier one. A wait whose earlier line was collected only after that is counted at its first line + /// collected in the window, which needs out-of-order collection; a single log file's read doesn't produce it. /// /// The duration is read off the message, not the duration_ms column. Verified at source /// (PgLockWaitEventParser): the lock_wait parser stores the line with NO metrics lifted — duration_ms @@ -165,7 +165,8 @@ SELECT DISTINCT ON (e.raw_line_hash) e.pid, e.relation_name, e.context, - e.duration_ms + e.duration_ms, + e.collection_time AS first_collected FROM pg_log_events AS e WHERE e.server_id = $1 AND e.collection_time >= $2 @@ -189,6 +190,7 @@ AND p.collection_time < $2 lines AS ( SELECT f.raw_line_hash, + f.first_collected, f.message, f.occurred_at, f.statement_fingerprint, @@ -209,7 +211,7 @@ FROM pg_log_events AS p WHERE p.server_id = $1 AND p.family = 'lock_wait' AND p.collection_time >= l.occurred_at - (l.wait_ms + 1000) * INTERVAL '1 millisecond' - AND p.collection_time <= $3 + AND p.collection_time <= l.first_collected AND p.occurred_at >= l.occurred_at - (l.wait_ms + 1000) * INTERVAL '1 millisecond' AND p.occurred_at <= l.occurred_at AND p.raw_line_hash <> l.raw_line_hash From 1f124e1a3908c94296f00cfaa465c546a2d9d2fa Mon Sep 17 00:00:00 2001 From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com> Date: Wed, 30 Sep 2026 21:06:15 -0400 Subject: [PATCH 4/4] State the lock-wait re-wait fold residual as the probe defines it, and make the two-waits pin depend on the time bound --- Darling/Darling.Tests/PgTargetBlockingTests.cs | 17 +++++++++++++++-- .../PgTargetFactCollector.Blocking.cs | 6 +++--- 2 files changed, 18 insertions(+), 5 deletions(-) diff --git a/Darling/Darling.Tests/PgTargetBlockingTests.cs b/Darling/Darling.Tests/PgTargetBlockingTests.cs index cd46b8aa1..32e122000 100644 --- a/Darling/Darling.Tests/PgTargetBlockingTests.cs +++ b/Darling/Darling.Tests/PgTargetBlockingTests.cs @@ -1542,10 +1542,10 @@ public async Task AWaitSplitAcrossTwoWindows_IsCountedInTheWindowOfItsOpeningLin public async Task TwoSeparateWaitsByOnePid_OnTheSameLock_CountTwice_EachLineStoredTwice(bool firstWaitAcquired) { var w1 = StillWaiting(4200, "1000.000"); - var w2 = StillWaiting(4200, "1000.000"); + var w2 = StillWaiting(4200, "1000.500"); var rows = new List<(DateTime, DateTime, string, int, string)> { - (T(10, 40, 1), T(10, 40, 0), "hw1", 4200, w1), (T(10, 50), T(10, 40, 0), "hw1", 4200, w1), + (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) @@ -1559,6 +1559,19 @@ public async Task TwoSeparateWaitsByOnePid_OnTheSameLock_CountTwice_EachLineStor 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( diff --git a/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs b/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs index e9bc39b3c..a5a62ecb6 100644 --- a/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs +++ b/Darling/PerformanceMonitor.Darling.Analysis/PgTargetFactCollector.Blocking.cs @@ -124,9 +124,9 @@ AND c.role_name IS NULL /// wait lies between its start and the line's own time, so the probe is bounded by the line's own N; the /// one-second pad covers whole-second %t stamps. Two separate waits by one backend are counted twice, /// because the second wait's lines start after the first ended. A line with no parsable lock text, duration or - /// time counts as its own wait. Residual, stated: a re-wait on the same lock text that starts within one - /// second of the previous wait's last logged line, when that wait ended without a lock_wait line (cancelled or - /// timed out), folds into it. The probe reads what was collected up to the opener candidate's own first + /// time counts as its own wait. Residual, stated: a re-wait by the same backend on the same lock text + /// that starts within one second after the previous wait's opening line folds into it, however that wait ended + /// (that previous wait then lasted under deadlock_timeout plus one second). The probe reads what was collected up to the opener candidate's own first /// sighting (first_collected), because an earlier line of the same wait is read before it, in the same pass /// or an earlier one. A wait whose earlier line was collected only after that is counted at its first line /// collected in the window, which needs out-of-order collection; a single log file's read doesn't produce it.