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]