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.