diff --git a/Darling/Darling.Tests/PgLogRotationEvidence.cs b/Darling/Darling.Tests/PgLogRotationEvidence.cs
new file mode 100644
index 000000000..ffd3cf9ca
--- /dev/null
+++ b/Darling/Darling.Tests/PgLogRotationEvidence.cs
@@ -0,0 +1,91 @@
+using System;
+using System.Collections.Generic;
+using System.Globalization;
+using System.Linq;
+using System.Text;
+using PerformanceMonitor.Collectors;
+
+namespace Darling.Tests;
+
+///
+/// Renders the evidence behind a failed log-rotation count check, so a failure names what was read instead of only
+/// "expected 1, actual 2": every row of the wait's backend (with its hash, time and text), which of them the test's
+/// identity matched, the resume state going in and coming out, and the target's log directory. The caller builds it
+/// only when a check fails.
+///
+internal static class PgLogRotationEvidence
+{
+ private const int TextLimit = 200;
+
+ internal static string Describe(
+ string what,
+ IEnumerable rows,
+ int waitPid,
+ DateTime waitFloorUtc,
+ Func isTheWait,
+ IReadOnlyDictionary? carriedState,
+ IReadOnlyDictionary? newState,
+ string logDirListing)
+ {
+ var all = rows.ToList();
+ var sb = new StringBuilder();
+ sb.Append("--- ").Append(what).AppendLine(" ---");
+ sb.Append("wait: pid=").Append(waitPid.ToString(CultureInfo.InvariantCulture))
+ .Append(" floor=").Append(waitFloorUtc.ToString("O", CultureInfo.InvariantCulture)).AppendLine();
+ sb.Append("rows read: ").Append(all.Count.ToString(CultureInfo.InvariantCulture))
+ .Append(", matching the wait: ").Append(all.Count(isTheWait).ToString(CultureInfo.InvariantCulture))
+ .Append(", distinct hashes among all rows: ").Append(all.Select(r => r.RawLineHash).Distinct().Count().ToString(CultureInfo.InvariantCulture)).AppendLine();
+
+ sb.AppendLine("rows matching the wait's identity, or carrying its pid (M = matched):");
+ var shown = 0;
+ foreach (var row in all.Where(r => isTheWait(r) || r.Pid == waitPid))
+ {
+ shown++;
+ sb.Append(isTheWait(row) ? " M " : " ")
+ .Append("hash=").Append(row.RawLineHash)
+ .Append(" at=").Append(row.OccurredAtUtc.ToString("O", CultureInfo.InvariantCulture))
+ .Append(" pid=").Append(row.Pid.ToString(CultureInfo.InvariantCulture))
+ .Append(" severity=").Append(row.Severity)
+ .Append(" family=").Append(row.Family)
+ .Append(" message=").Append(Clip(row.Message))
+ .Append(" detail=").Append(Clip(row.Detail))
+ .Append(" context=").Append(Clip(row.Context))
+ .AppendLine();
+ }
+
+ if (shown == 0)
+ {
+ sb.AppendLine(" (none)");
+ }
+
+ var repeated = all.GroupBy(r => r.RawLineHash).Where(g => g.Count() > 1).ToList();
+ sb.Append("hashes that occur more than once: ").Append(repeated.Count.ToString(CultureInfo.InvariantCulture)).AppendLine();
+ foreach (var group in repeated)
+ {
+ sb.Append(" ").Append(group.Key).Append(" x").Append(group.Count().ToString(CultureInfo.InvariantCulture))
+ .Append(" message=").Append(Clip(group.First().Message)).AppendLine();
+ }
+
+ sb.Append("carried state: ").AppendLine(State(carriedState));
+ sb.Append("resulting state: ").AppendLine(State(newState));
+ sb.AppendLine("log directory (name, size, modification):");
+ sb.AppendLine(string.IsNullOrEmpty(logDirListing) ? " (not read)" : logDirListing);
+ return sb.ToString();
+ }
+
+ private static string State(IReadOnlyDictionary? state) =>
+ state is null || state.Count == 0
+ ? "(none)"
+ : string.Join("; ", state.OrderBy(p => p.Key, StringComparer.Ordinal).Select(p => p.Key + "=" + p.Value));
+
+ private static string Clip(string? text)
+ {
+ if (text is null)
+ {
+ return "";
+ }
+
+ var flat = text.Replace("\r", "\\r", StringComparison.Ordinal).Replace("\n", "\\n", StringComparison.Ordinal);
+ return flat.Length <= TextLimit ? "\"" + flat + "\"" : "\"" + flat[..TextLimit] + "\"...(" + flat.Length.ToString(CultureInfo.InvariantCulture) + " chars)";
+ }
+}
diff --git a/Darling/Darling.Tests/PgLogRotationEvidenceTests.cs b/Darling/Darling.Tests/PgLogRotationEvidenceTests.cs
new file mode 100644
index 000000000..a8daf6fe9
--- /dev/null
+++ b/Darling/Darling.Tests/PgLogRotationEvidenceTests.cs
@@ -0,0 +1,64 @@
+using System;
+using System.Collections.Generic;
+using PerformanceMonitor.Collectors;
+using Xunit;
+
+namespace Darling.Tests;
+
+///
+/// The rotation live tests' failure description has to name each hash, pid, time and message it was given; a
+/// description that dropped them would turn a duplicate back into "expected 1, actual 2".
+///
+public sealed class PgLogRotationEvidenceTests
+{
+ private static PgLogEvent Row(string hash, int pid, DateTime at, string message) => new(
+ OccurredAtUtc: at, Family: PgLogFamilies.Error, Severity: "LOG", SqlState: null, DatabaseName: null, UserName: null,
+ ApplicationName: null, Pid: pid, Message: message, Detail: "Process holding the lock: 99.", Context: "while updating tuple",
+ StatementFingerprint: null, RawLineHash: hash, Metrics: default);
+
+ [Fact]
+ public void Describe_NamesEveryHashPidTimeAndMessageOfTheMatchedRows()
+ {
+ var floor = new DateTime(2026, 9, 30, 12, 0, 0, DateTimeKind.Utc);
+ var rows = new List
+ {
+ Row("hash-aaa", 4242, floor.AddSeconds(1), "process 4242 still waiting for ShareLock after 100.123 ms"),
+ Row("hash-bbb", 4242, floor.AddSeconds(2), "process 4242 still waiting for ShareLock after 900.456 ms"),
+ Row("hash-ccc", 7, floor.AddSeconds(3), "unrelated entry"),
+ Row("hash-aaa", 4242, floor.AddSeconds(1), "process 4242 still waiting for ShareLock after 100.123 ms"),
+ };
+
+ var text = PgLogRotationEvidence.Describe(
+ "after rotation", rows, 4242, floor, r => r.Pid == 4242 && r.OccurredAtUtc >= floor,
+ new Dictionary { ["log_tail_json"] = "100|postgresql-a.json" },
+ new Dictionary { ["log_tail_json"] = "200|postgresql-b.json" },
+ " postgresql-a.json 100 2026-09-30 12:00:00+00");
+
+ foreach (var expected in new[]
+ {
+ "pid=4242", floor.ToString("O", System.Globalization.CultureInfo.InvariantCulture), "hash=hash-aaa", "hash=hash-bbb",
+ "after 100.123 ms", "after 900.456 ms", "severity=LOG", "family=" + PgLogFamilies.Error, "detail=\"Process holding the lock: 99.\"",
+ "context=\"while updating tuple\"", "hash-aaa x2", "log_tail_json=100|postgresql-a.json", "log_tail_json=200|postgresql-b.json",
+ "postgresql-a.json 100",
+ })
+ {
+ Assert.Contains(expected, text, StringComparison.Ordinal);
+ }
+
+ Assert.DoesNotContain("hash=hash-ccc", text, StringComparison.Ordinal);
+ Assert.Contains("matching the wait: 3", text, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void Describe_ClipsLongTextAndToleratesMissingParts()
+ {
+ var floor = DateTime.UtcNow;
+ var text = PgLogRotationEvidence.Describe(
+ "x", [Row("h", 1, floor, new string('m', 500))], 1, floor, _ => true, null, null, string.Empty);
+
+ Assert.Contains("...(500 chars)", text, StringComparison.Ordinal);
+ Assert.DoesNotContain(new string('m', 250), text, StringComparison.Ordinal);
+ Assert.Contains("carried state: (none)", text, StringComparison.Ordinal);
+ Assert.Contains("(not read)", text, StringComparison.Ordinal);
+ }
+}
diff --git a/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs b/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs
index 46d29159d..2a57fbece 100644
--- a/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs
+++ b/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs
@@ -76,6 +76,35 @@ private static async Task RotateAsync(NpgsqlConnection connection, CancellationT
private static IReadOnlyDictionary Carry(CollectorContext context) =>
new Dictionary(context.PendingState);
+ ///
+ /// The evidence for a failed count check (), with the target's log directory read
+ /// now. Called only when a check has already failed, so a passing run never pays for it.
+ ///
+ private static async Task DescribeAsync(
+ NpgsqlConnection connection, string what, IEnumerable rows, LoggedWait wait,
+ IReadOnlyDictionary? carriedState, IReadOnlyDictionary? newState, CancellationToken ct)
+ {
+ string listing;
+ try
+ {
+ var lines = new List();
+ await using var command = new NpgsqlCommand("SELECT name, size, modification FROM pg_ls_logdir() ORDER BY modification, name", connection);
+ await using var reader = await command.ExecuteReaderAsync(ct);
+ while (await reader.ReadAsync(ct))
+ {
+ lines.Add(" " + reader.GetString(0) + " " + reader.GetInt64(1) + " " + reader.GetDateTime(2).ToString("O", System.Globalization.CultureInfo.InvariantCulture));
+ }
+
+ listing = string.Join("\n", lines);
+ }
+ catch (Exception ex) when (ex is NpgsqlException or InvalidOperationException)
+ {
+ listing = " (the log directory could not be read: " + ex.Message + ")";
+ }
+
+ return PgLogRotationEvidence.Describe(what, rows, wait.Pid, wait.FloorUtc, r => IsTheWait(r, wait), carriedState, newState, listing);
+ }
+
private static string Marker() => "pm4699t" + Guid.NewGuid().ToString("N")[..10];
///
@@ -89,10 +118,11 @@ private static IReadOnlyDictionary Carry(CollectorContext contex
/// Provokes one lock-wait log entry (log_lock_waits, a short deadlock_timeout) against a table named
/// ; the returned wait identifies the entry in the rows read back.
///
- private static async Task LogAsync(NpgsqlConnection connection, bool json, string marker, CancellationToken ct)
+ private static async Task LogAsync(NpgsqlConnection connection, bool json, string marker, CancellationToken ct, bool pokeTheWaiter = false)
{
await using var holder = await OpenAsync(json, ct);
await using var waiter = await OpenAsync(json, ct);
+ await using var poker = pokeTheWaiter ? await OpenAsync(json, ct) : null;
/* The floor is the target's own clock, read before any of this wait exists. Its entry is logged at least
deadlock_timeout (100 ms) after that, and every earlier wait had finished before it. */
await using var clockCommand = new NpgsqlCommand("SELECT clock_timestamp()", connection);
@@ -103,7 +133,10 @@ private static async Task LogAsync(NpgsqlConnection connection, bool
await ExecAsync(holder, "UPDATE " + marker + " SET id = 1", ct);
await using var pidCommand = new NpgsqlCommand("SELECT pg_backend_pid()", waiter);
var pid = (int)(await pidCommand.ExecuteScalarAsync(ct))!;
- var blocked = ExecAsync(waiter, "SET lock_timeout = '800ms'; UPDATE " + marker + " SET id = 1", ct);
+ var blocked = ExecAsync(waiter, "SET lock_timeout = '" + (pokeTheWaiter ? "3000ms" : "800ms") + "'; UPDATE " + marker + " SET id = 1", ct);
+ /* PostgreSQL logs "still waiting" again on every latch wakeup after the first deadlock check, so a wake-up
+ between deadlock_timeout (100 ms) and lock_timeout (800 ms; 3000 ms when poked, so the poke always lands inside the wait) gives this one wait a second line. */
+ var poke = pokeTheWaiter ? PokeAsync(poker!, pid, ct) : Task.CompletedTask;
try
{
await blocked;
@@ -113,12 +146,106 @@ private static async Task LogAsync(NpgsqlConnection connection, bool
/* lock_timeout after the wait was logged */
}
+ await poke;
await ExecAsync(holder, "ROLLBACK", ct);
await ExecAsync(holder, "DROP TABLE " + marker, ct);
await ExecAsync(connection, "SELECT pg_sleep(0.3);", ct);
return new LoggedWait(pid, floor);
}
+ /// Wakes the waiter's latch about 350 ms into its wait, from a third connection.
+ private static async Task PokeAsync(NpgsqlConnection third, int waiterPid, CancellationToken ct)
+ {
+ await Task.Delay(350, ct);
+ await using var command = new NpgsqlCommand("SELECT pg_log_backend_memory_contexts(" + waiterPid.ToString(System.Globalization.CultureInfo.InvariantCulture) + ")", third);
+ await command.ExecuteScalarAsync(ct);
+ }
+
+ ///
+ /// The messages of the wait's "still waiting" entries, from ONE read of each of the route's log files from byte 0,
+ /// parsed with the product's own parser. A wait can log more than one such line (each with its own "after X ms"),
+ /// so the set, not a count of one, is what the tail must return exactly.
+ ///
+ private static async Task> ExpectedWaitLinesAsync(NpgsqlConnection connection, bool json, LoggedWait wait, CancellationToken ct)
+ {
+ var names = new List();
+ await using (var list = new NpgsqlCommand("SELECT name FROM pg_ls_logdir() WHERE name LIKE @pattern ORDER BY name", connection))
+ {
+ list.Parameters.AddWithValue("pattern", json ? "%.json" : "%.csv");
+ await using var reader = await list.ExecuteReaderAsync(ct);
+ while (await reader.ReadAsync(ct))
+ {
+ names.Add(reader.GetString(0));
+ }
+ }
+
+ var prefix = $"process {wait.Pid} still waiting for ";
+ var messages = new List();
+ foreach (var name in names)
+ {
+ string body;
+ await using (var read = new NpgsqlCommand("SELECT pg_read_file(current_setting('log_directory') || '/' || @name)", connection))
+ {
+ read.Parameters.AddWithValue("name", name);
+ body = (string)(await read.ExecuteScalarAsync(ct))!;
+ }
+
+ var entries = json ? PgServerLogJsonParser.Parse(body, out _) : PgServerLogCsvParser.Parse(body, out _);
+ messages.AddRange(entries
+ .Where(e => e.Pid == wait.Pid && e.OccurredAtUtc >= wait.FloorUtc && e.Message.StartsWith(prefix, StringComparison.Ordinal))
+ .Select(e => e.Message));
+ }
+
+ return messages;
+ }
+
+ ///
+ /// Asserts the rotation outcome for one wait: the resumed read holds exactly the expected "still waiting" messages
+ /// (each once, each with its own hash), the read without state holds none, and across both cycles each message is one identity.
+ ///
+ private static async Task AssertExactWaitAsync(
+ NpgsqlConnection connection, bool json, LoggedWait wait, List expected,
+ (List Rows, CollectorContext Context) first, (List Rows, CollectorContext Context) withState,
+ IReadOnlyDictionary carried, bool binary, CancellationToken ct)
+ {
+ var expectedText = "expected from the files: [" + string.Join(" | ", expected) + "]\n";
+ Assert.True(expected.Count > 0, "the log files hold no entry for the wait\n" + expectedText);
+ Assert.True(expected.Distinct().Count() == expected.Count, "the log files repeat a wait message\n" + expectedText);
+
+ var resumed = withState.Rows.Where(r => IsTheWait(r, wait)).ToList();
+ Assert.True(
+ resumed.Select(r => r.Message).OrderBy(m => m, StringComparer.Ordinal).SequenceEqual(expected.OrderBy(m => m, StringComparer.Ordinal))
+ && resumed.Select(r => r.RawLineHash).Distinct().Count() == resumed.Count,
+ $"the resumed read should hold exactly the wait's {expected.Count} expected line(s), each once, got {resumed.Count}\n" + expectedText
+ + await DescribeAsync(connection, "resumed read after the rotation", withState.Rows, wait, carried, withState.Context.PendingState, ct));
+
+ /* No state is today's read: the newest file only, which does not hold the lines. */
+ var noState = await CycleAsync(connection, null, binary, json, ct);
+ var noStateCount = Count(noState.Rows, wait);
+ Assert.True(
+ noStateCount == 0,
+ $"the read without state should not hold the wait's entry, got {noStateCount}\n"
+ + await DescribeAsync(connection, "read without state", noState.Rows, wait, null, noState.Context.PendingState, ct));
+
+ /* Within the cycle no raw_line_hash repeats. */
+ var distinctHashes = withState.Rows.Select(r => r.RawLineHash).Distinct().Count();
+ Assert.True(
+ withState.Rows.Count == distinctHashes,
+ $"a raw_line_hash repeats within the resumed read: {withState.Rows.Count} rows, {distinctHashes} distinct hashes\n"
+ + await DescribeAsync(connection, "resumed read after the rotation", withState.Rows, wait, carried, withState.Context.PendingState, ct));
+
+ /* Across the two cycles each expected message is one identity. */
+ var both = first.Rows.Concat(withState.Rows).Where(r => IsTheWait(r, wait)).ToList();
+ foreach (var message in expected)
+ {
+ var hashes = both.Where(r => r.Message == message).Select(r => r.RawLineHash).Distinct().Count();
+ Assert.True(
+ hashes == 1,
+ $"the wait's entry \"{message}\" should have one identity across both reads, got {hashes}\n"
+ + await DescribeAsync(connection, "both reads together", both, wait, carried, withState.Context.PendingState, ct));
+ }
+ }
+
/* A pid alone does not identify a wait's entry. The holder and waiter connections come from Npgsql's pool, so a
later wait often runs on an earlier wait's backend (the same pid), and the resumed read re-reads the previous
read on purpose (a 1 MiB overlap), so the earlier wait's entry is in the same rows. The floor tells them apart. */
@@ -148,19 +275,34 @@ public async Task ARotationBetweenReads_DoesNotLoseTheLinesWrittenBeforeIt(bool
var wait = await LogAsync(connection, json, Marker(), ct);
await RotateAsync(connection, ct);
- var withState = await CycleAsync(connection, Carry(first.Context), binary, json, ct);
- Assert.Equal(1, Count(withState.Rows, wait));
+ var carried = Carry(first.Context);
+ var withState = await CycleAsync(connection, carried, binary, json, ct);
+ var expected = await ExpectedWaitLinesAsync(connection, json, wait, ct);
+ await AssertExactWaitAsync(connection, json, wait, expected, first, withState, carried, binary, ct);
+ }
- /* No state is today's read: the newest file only, which does not hold the line. */
- var noState = await CycleAsync(connection, null, binary, json, ct);
- Assert.Equal(0, Count(noState.Rows, wait));
+ [Theory]
+ [InlineData(false, true)]
+ [InlineData(true, true)]
+ [InlineData(false, false)]
+ [InlineData(true, false)]
+ public async Task ALockWaitThatLogsTwice_IsReadAsTwoLinesEachOnce(bool json, bool binary)
+ {
+ Assert.SkipWhen(string.IsNullOrEmpty(TargetFor(json)), SkipReason);
+ var ct = TestContext.Current.CancellationToken;
+ await using var connection = await OpenAsync(json, ct);
- /* Within the cycle no raw_line_hash repeats. */
- Assert.Equal(withState.Rows.Count, withState.Rows.Select(r => r.RawLineHash).Distinct().Count());
+ _ = await LogAsync(connection, json, Marker(), ct);
+ var first = await CycleAsync(connection, null, binary, json, ct);
+
+ var wait = await LogAsync(connection, json, Marker(), ct, pokeTheWaiter: true);
+ await RotateAsync(connection, ct);
- /* Across the two cycles the line is one identity. */
- var all = first.Rows.Concat(withState.Rows).Where(r => IsTheWait(r, wait)).Select(r => r.RawLineHash).Distinct();
- Assert.Single(all);
+ var carried = Carry(first.Context);
+ var withState = await CycleAsync(connection, carried, binary, json, ct);
+ var expected = await ExpectedWaitLinesAsync(connection, json, wait, ct);
+ Assert.True(expected.Count == 2, $"the poke should make one wait log two lines, the files hold {expected.Count}: [{string.Join(" | ", expected)}]");
+ await AssertExactWaitAsync(connection, json, wait, expected, first, withState, carried, binary, ct);
}
[Theory]