@@ -570,11 +570,19 @@ public void TheCollectorsReads_CarryTheFloorAsAParameter_TheExclusions_TheMessag
570570 var events = PgTargetFactCollector . PgTargetLockWaitEventsSql ;
571571 Assert . Contains ( "family = 'lock_wait'" , events , StringComparison . Ordinal ) ;
572572 /* Verified at source: the lock_wait parser lifts no metrics, so the duration is the message's own figure. */
573- Assert . Contains ( "substring(message from 'after ([0-9]+(?:\\ .[0-9]+)?) ms')" , events , StringComparison . Ordinal ) ;
574- Assert . Contains ( "coalesce(duration_ms::DOUBLE PRECISION," , events , StringComparison . Ordinal ) ;
573+ Assert . Contains ( "substring(f. message from 'after ([0-9]+(?:\\ .[0-9]+)?) ms')" , events , StringComparison . Ordinal ) ;
574+ Assert . Contains ( "coalesce(f. duration_ms::DOUBLE PRECISION," , events , StringComparison . Ordinal ) ;
575575 Assert . Contains ( "message LIKE 'process % still waiting for %'" , events , StringComparison . Ordinal ) ;
576576 Assert . Contains ( "message LIKE 'process % acquired %'" , events , StringComparison . Ordinal ) ;
577577 Assert . Contains ( "message LIKE 'process % detected deadlock %'" , events , StringComparison . Ordinal ) ;
578+ /* Each line counts once, in the window of its first sighting; each wait counts once, at the line that opens it. */
579+ Assert . Contains ( "DISTINCT ON (e.raw_line_hash)" , events , StringComparison . Ordinal ) ;
580+ Assert . Contains ( "p.collection_time >= coalesce(w.occurred_at" , events , StringComparison . Ordinal ) ;
581+ Assert . Contains ( "p.collection_time <= l.first_collected" , events , StringComparison . Ordinal ) ;
582+ Assert . Contains ( "e.collection_time AS first_collected" , events , StringComparison . Ordinal ) ;
583+ Assert . Contains ( "p.collection_time < $2" , events , StringComparison . Ordinal ) ;
584+ Assert . Contains ( "* INTERVAL '1 millisecond'" , events , StringComparison . Ordinal ) ;
585+ Assert . Contains ( "(SELECT COUNT(*) FROM waits)" , events , StringComparison . Ordinal ) ;
578586 Assert . Contains ( "collector_name = 'pg_log_events'" , events , StringComparison . Ordinal ) ;
579587 Assert . Contains ( "status = 'SUCCESS'" , events , StringComparison . Ordinal ) ;
580588
@@ -1383,6 +1391,187 @@ INSERT INTO pg_log_events
13831391 await command . ExecuteNonQueryAsync ( ct ) ;
13841392 }
13851393
1394+ /* ───────────── gated: each lock_wait line counts once, each wait counts once ───────────── */
1395+
1396+ private static readonly DateTime LockWinStart = new ( 2026 , 6 , 1 , 10 , 0 , 0 , DateTimeKind . Unspecified ) ;
1397+ private static readonly DateTime LockWinEnd = new ( 2026 , 6 , 1 , 11 , 0 , 0 , DateTimeKind . Unspecified ) ;
1398+ private const string LockOrders = "while updating tuple (0,7) in relation \" orders\" " ;
1399+
1400+ private sealed record LockWaitShape ( long StillWaiting , long Acquired , long Deadlocks , long Lines , double ? MaxWaitMs , double AcquiredWaitMs , long TopRelationWaits ) ;
1401+
1402+ private static string StillWaiting ( int pid , string ms ) => $ "process { pid } still waiting for ShareLock on transaction 900 after { ms } ms";
1403+
1404+ private static async Task < LockWaitShape [ ] > RunLockWaitScenarioAsync (
1405+ ( DateTime Collected , DateTime Occurred , string Hash , int Pid , string Message ) [ ] lines ,
1406+ ( DateTime Start , DateTime End ) [ ] windows ,
1407+ CancellationToken ct )
1408+ {
1409+ var cs = Environment . GetEnvironmentVariable ( "DARLING_TEST_PG" ) ;
1410+ Assert . SkipWhen ( string . IsNullOrEmpty ( cs ) , "Set DARLING_TEST_PG to a Postgres connection string to run the lock-wait counting pins." ) ;
1411+
1412+ using var connection = new NpgsqlConnection ( cs ) ;
1413+ await connection . OpenAsync ( ct ) ;
1414+ await PgMigrations . MigrateAsync ( connection , ct ) ;
1415+ await DeleteRowsAsync ( connection , ct ) ;
1416+
1417+ var bodySucceeded = false ;
1418+ try
1419+ {
1420+ await PgTargetFactCollectorTests . RegisterServerAsync ( connection , ServerId , ServerName , "postgres" , 18 , ct ) ;
1421+ foreach ( var ( collected , occurred , hash , pid , message ) in lines )
1422+ {
1423+ using var insert = new NpgsqlCommand ( @"
1424+ INSERT INTO pg_log_events
1425+ (collection_id, collection_time, server_id, server_name, occurred_at, family, severity, sqlstate, database_name, user_name, application_name,
1426+ pid, message, detail, context, statement_fingerprint, raw_line_hash, relation_name, duration_ms)
1427+ VALUES ($1, $2, $3, $4, $5, 'lock_wait', 'LOG', '00000', 'appdb', 'app', 'web',
1428+ $6, $7, 'Process holding the lock: 9000.', $8, 'fp-orders-update', $9, NULL, NULL)" , connection ) ;
1429+ insert . Parameters . AddWithValue ( CollectionIdGenerator . Next ( ) ) ;
1430+ insert . Parameters . AddWithValue ( collected ) ;
1431+ insert . Parameters . AddWithValue ( ServerId ) ;
1432+ insert . Parameters . AddWithValue ( ServerName ) ;
1433+ insert . Parameters . AddWithValue ( occurred ) ;
1434+ insert . Parameters . AddWithValue ( pid ) ;
1435+ insert . Parameters . AddWithValue ( message ) ;
1436+ insert . Parameters . AddWithValue ( LockOrders ) ;
1437+ insert . Parameters . AddWithValue ( hash ) ;
1438+ await insert . ExecuteNonQueryAsync ( ct ) ;
1439+ }
1440+
1441+ var shapes = new List < LockWaitShape > ( ) ;
1442+ foreach ( var ( start , end ) in windows )
1443+ {
1444+ using var read = new NpgsqlCommand ( PgTargetFactCollector . PgTargetLockWaitEventsSql , connection ) ;
1445+ read . Parameters . AddWithValue ( ServerId ) ;
1446+ read . Parameters . AddWithValue ( start ) ;
1447+ read . Parameters . AddWithValue ( end ) ;
1448+ using var reader = await read . ExecuteReaderAsync ( ct ) ;
1449+ Assert . True ( await reader . ReadAsync ( ct ) ) ;
1450+ shapes . Add ( new LockWaitShape (
1451+ reader . GetInt64 ( 0 ) , reader . GetInt64 ( 1 ) , reader . GetInt64 ( 2 ) , reader . GetInt64 ( 3 ) ,
1452+ reader . IsDBNull ( 4 ) ? null : reader . GetDouble ( 4 ) ,
1453+ reader . IsDBNull ( 5 ) ? 0 : reader . GetDouble ( 5 ) ,
1454+ reader . IsDBNull ( 8 ) ? 0 : reader . GetInt64 ( 8 ) ) ) ;
1455+ }
1456+
1457+ bodySucceeded = true ;
1458+ return shapes . ToArray ( ) ;
1459+ }
1460+ finally
1461+ {
1462+ await LiveStoreCleanup . RunAsync ( cs ! , bodySucceeded , async ( cleanup , cleanupCt ) =>
1463+ await DeleteRowsAsync ( cleanup , cleanupCt ) ) ;
1464+ }
1465+ }
1466+
1467+ private static DateTime T ( int h , int m , int s = 0 ) => new ( 2026 , 6 , 1 , h , m , s , DateTimeKind . Unspecified ) ;
1468+
1469+ [ Fact ]
1470+ public async Task ARepeatedLineHash_InsideTheWindow_CountsOnce ( )
1471+ {
1472+ var m = StillWaiting ( 4200 , "1000.000" ) ;
1473+ var shape = ( await RunLockWaitScenarioAsync (
1474+ [ ( T ( 10 , 11 ) , T ( 10 , 10 ) , "h1" , 4200 , m ) , ( T ( 11 , 0 ) , T ( 10 , 10 ) , "h1" , 4200 , m ) ] ,
1475+ [ ( LockWinStart , LockWinEnd ) ] , TestContext . Current . CancellationToken ) ) [ 0 ] ;
1476+ Assert . Equal ( 1 , shape . StillWaiting ) ;
1477+ Assert . Equal ( 1 , shape . Lines ) ;
1478+ }
1479+
1480+ [ Fact ]
1481+ public async Task ALine_FirstStoredBeforeTheWindow_AndRestoredInside_IsNotCounted ( )
1482+ {
1483+ var m = StillWaiting ( 4200 , "1000.000" ) ;
1484+ var shape = ( await RunLockWaitScenarioAsync (
1485+ [ ( T ( 9 , 1 ) , T ( 9 , 0 ) , "h2" , 4200 , m ) , ( T ( 10 , 5 ) , T ( 9 , 0 ) , "h2" , 4200 , m ) ] ,
1486+ [ ( LockWinStart , LockWinEnd ) ] , TestContext . Current . CancellationToken ) ) [ 0 ] ;
1487+ Assert . Equal ( 0 , shape . StillWaiting ) ;
1488+ Assert . Equal ( 0 , shape . Lines ) ;
1489+ }
1490+
1491+ [ Fact ]
1492+ public async Task ALine_ReSightedThreeHoursLater_IsNotCounted_WhereAFixedTwoHourLookbackWould ( )
1493+ {
1494+ var m = StillWaiting ( 4200 , "1000.000" ) ;
1495+ var shape = ( await RunLockWaitScenarioAsync (
1496+ [ ( T ( 7 , 1 ) , T ( 7 , 0 ) , "h3" , 4200 , m ) , ( T ( 10 , 20 ) , T ( 7 , 0 ) , "h3" , 4200 , m ) ] ,
1497+ [ ( LockWinStart , LockWinEnd ) ] , TestContext . Current . CancellationToken ) ) [ 0 ] ;
1498+ Assert . Equal ( 0 , shape . StillWaiting ) ;
1499+ Assert . Equal ( 0 , shape . Lines ) ;
1500+ }
1501+
1502+ [ Fact ]
1503+ public async Task OutcomeLines_StoredTwice_CountOnce ( )
1504+ {
1505+ var acquired = "process 4200 acquired ShareLock on transaction 900 after 30000.000 ms" ;
1506+ var deadlock = "process 4300 detected deadlock while waiting for ShareLock on transaction 901 after 1000.000 ms" ;
1507+ var shape = ( await RunLockWaitScenarioAsync (
1508+ [ ( T ( 10 , 11 ) , T ( 10 , 10 ) , "h4a" , 4200 , acquired ) , ( T ( 10 , 40 ) , T ( 10 , 10 ) , "h4a" , 4200 , acquired ) ,
1509+ ( T ( 10 , 21 ) , T ( 10 , 20 ) , "h4d" , 4300 , deadlock ) , ( T ( 10 , 41 ) , T ( 10 , 20 ) , "h4d" , 4300 , deadlock ) ] ,
1510+ [ ( LockWinStart , LockWinEnd ) ] , TestContext . Current . CancellationToken ) ) [ 0 ] ;
1511+ Assert . Equal ( 1 , shape . Acquired ) ;
1512+ Assert . Equal ( 30000 , shape . AcquiredWaitMs ) ;
1513+ Assert . Equal ( 1 , shape . Deadlocks ) ;
1514+ }
1515+
1516+ [ Fact ]
1517+ public async Task TwoLinesOfOneWait_CountOneWait ( )
1518+ {
1519+ var shape = ( await RunLockWaitScenarioAsync (
1520+ [ ( T ( 10 , 30 , 1 ) , T ( 10 , 30 , 0 ) , "hb1" , 4200 , StillWaiting ( 4200 , "1000.000" ) ) ,
1521+ ( T ( 10 , 30 , 5 ) , T ( 10 , 30 , 4 ) , "hb2" , 4200 , StillWaiting ( 4200 , "5000.000" ) ) ] ,
1522+ [ ( LockWinStart , LockWinEnd ) ] , TestContext . Current . CancellationToken ) ) [ 0 ] ;
1523+ Assert . Equal ( 1 , shape . StillWaiting ) ;
1524+ Assert . Equal ( 1 , shape . TopRelationWaits ) ;
1525+ Assert . Equal ( 5000 , shape . MaxWaitMs ) ;
1526+ }
1527+
1528+ [ Fact ]
1529+ public async Task AWaitSplitAcrossTwoWindows_IsCountedInTheWindowOfItsOpeningLineOnly ( )
1530+ {
1531+ var shapes = await RunLockWaitScenarioAsync (
1532+ [ ( T ( 9 , 59 , 58 ) , T ( 9 , 59 , 57 ) , "hc1" , 4200 , StillWaiting ( 4200 , "1000.000" ) ) ,
1533+ ( T ( 10 , 0 , 2 ) , T ( 10 , 0 , 1 ) , "hc2" , 4200 , StillWaiting ( 4200 , "5000.000" ) ) ] ,
1534+ [ ( LockWinStart , LockWinEnd ) , ( T ( 9 , 0 ) , LockWinStart ) ] , TestContext . Current . CancellationToken ) ;
1535+ Assert . Equal ( 0 , shapes [ 0 ] . StillWaiting ) ;
1536+ Assert . Equal ( 1 , shapes [ 1 ] . StillWaiting ) ;
1537+ }
1538+
1539+ [ Theory ]
1540+ [ InlineData ( true ) ]
1541+ [ InlineData ( false ) ]
1542+ public async Task TwoSeparateWaitsByOnePid_OnTheSameLock_CountTwice_EachLineStoredTwice ( bool firstWaitAcquired )
1543+ {
1544+ var w1 = StillWaiting ( 4200 , "1000.000" ) ;
1545+ var w2 = StillWaiting ( 4200 , "1000.500" ) ;
1546+ var rows = new List < ( DateTime , DateTime , string , int , string ) >
1547+ {
1548+ ( T ( 10 , 41 , 11 ) , T ( 10 , 40 , 0 ) , "hw1" , 4200 , w1 ) , ( T ( 10 , 50 ) , T ( 10 , 40 , 0 ) , "hw1" , 4200 , w1 ) ,
1549+ ( T ( 10 , 41 , 11 ) , T ( 10 , 41 , 10 ) , "hw2" , 4200 , w2 ) , ( T ( 10 , 51 ) , T ( 10 , 41 , 10 ) , "hw2" , 4200 , w2 ) ,
1550+ } ;
1551+ if ( firstWaitAcquired )
1552+ {
1553+ var acq = "process 4200 acquired ShareLock on transaction 900 after 30000.000 ms" ;
1554+ rows . Add ( ( T ( 10 , 40 , 30 ) , T ( 10 , 40 , 29 ) , "hwa" , 4200 , acq ) ) ;
1555+ rows . Add ( ( T ( 10 , 52 ) , T ( 10 , 40 , 29 ) , "hwa" , 4200 , acq ) ) ;
1556+ }
1557+
1558+ var shape = ( await RunLockWaitScenarioAsync ( rows . ToArray ( ) , [ ( LockWinStart , LockWinEnd ) ] , TestContext . Current . CancellationToken ) ) [ 0 ] ;
1559+ Assert . Equal ( 2 , shape . StillWaiting ) ;
1560+ }
1561+
1562+ // 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.
1563+ [ Fact ]
1564+ public async Task AReWaitStartingWithinOneSecondOfThePreviousOpener_FoldsIntoIt_TheStatedResidual ( )
1565+ {
1566+ var w1 = StillWaiting ( 4200 , "1000.000" ) ;
1567+ var w2 = StillWaiting ( 4200 , "1000.500" ) ;
1568+ var shape = ( await RunLockWaitScenarioAsync (
1569+ [ ( T ( 10 , 40 , 1 ) , T ( 10 , 40 , 0 ) , "hf1" , 4200 , w1 ) ,
1570+ ( T ( 10 , 40 , 1 ) . AddMilliseconds ( 1500 ) , T ( 10 , 40 , 0 ) . AddMilliseconds ( 1500 ) , "hf2" , 4200 , w2 ) ] ,
1571+ [ ( LockWinStart , LockWinEnd ) ] , TestContext . Current . CancellationToken ) ) [ 0 ] ;
1572+ Assert . Equal ( 1 , shape . StillWaiting ) ;
1573+ }
1574+
13861575 private static async Task DeleteRowsAsync ( NpgsqlConnection connection , CancellationToken ct )
13871576 {
13881577 using var cleanup = new NpgsqlCommand (
0 commit comments