From 8abfb0f41db175ebf896ffcb6949ea8a71142c50 Mon Sep 17 00:00:00 2001 From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com> Date: Wed, 30 Sep 2026 19:50:02 -0400 Subject: [PATCH 1/3] Make the log-rotation live test's count checks describe what they read When a count check in the csv/json rotation test fails, the message now lists the rows matching the wait's identity or carrying its pid (hash, time, pid, severity, family, message, detail, context), the resume state going in and coming out, and the target's log directory. The checks are exactly as strict as before. --- .../Darling.Tests/PgLogRotationEvidence.cs | 91 +++++++++++++++++++ .../PgLogRotationEvidenceTests.cs | 64 +++++++++++++ ...PgServerLogTailCsvJsonRotationLiveTests.cs | 57 ++++++++++-- 3 files changed, 206 insertions(+), 6 deletions(-) create mode 100644 Darling/Darling.Tests/PgLogRotationEvidence.cs create mode 100644 Darling/Darling.Tests/PgLogRotationEvidenceTests.cs 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..ac78f1c0e 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]; /// @@ -148,19 +177,35 @@ 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 withStateCount = Count(withState.Rows, wait); + Assert.True( + withStateCount == 1, + $"the resumed read should hold the wait's entry exactly once, got {withStateCount}\n" + + 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 line. */ var noState = await CycleAsync(connection, null, binary, json, ct); - Assert.Equal(0, Count(noState.Rows, wait)); + 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. */ - Assert.Equal(withState.Rows.Count, withState.Rows.Select(r => r.RawLineHash).Distinct().Count()); + 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 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 all = first.Rows.Concat(withState.Rows).Where(r => IsTheWait(r, wait)).Select(r => r.RawLineHash).Distinct().ToList(); + Assert.True( + all.Count == 1, + $"the wait's entry should have one identity across both reads, got {all.Count}\n" + + await DescribeAsync(connection, "both reads together", first.Rows.Concat(withState.Rows), wait, carried, withState.Context.PendingState, ct)); } [Theory] From d3fb2815f1350b33fd226b1ce980858f2702149d Mon Sep 17 00:00:00 2001 From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com> Date: Wed, 30 Sep 2026 20:22:37 -0400 Subject: [PATCH 2/3] Make the log-rotation live test's wait identity an exact set read from the log files, and pin a wait that logs twice --- ...PgServerLogTailCsvJsonRotationLiveTests.cs | 147 +++++++++++++++--- 1 file changed, 122 insertions(+), 25 deletions(-) diff --git a/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs b/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs index ac78f1c0e..7337c2010 100644 --- a/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs +++ b/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs @@ -118,7 +118,7 @@ private static async Task DescribeAsync( /// 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); @@ -133,6 +133,9 @@ private static async Task LogAsync(NpgsqlConnection connection, bool 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); + /* 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) gives this one wait a second line. */ + var poke = pokeTheWaiter ? PokeAsync(json, pid, ct) : Task.CompletedTask; try { await blocked; @@ -142,12 +145,107 @@ 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(bool json, int waiterPid, CancellationToken ct) + { + await Task.Delay(350, ct); + await using var third = await OpenAsync(json, 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. */ @@ -179,33 +277,32 @@ public async Task ARotationBetweenReads_DoesNotLoseTheLinesWrittenBeforeIt(bool var carried = Carry(first.Context); var withState = await CycleAsync(connection, carried, binary, json, ct); - var withStateCount = Count(withState.Rows, wait); - Assert.True( - withStateCount == 1, - $"the resumed read should hold the wait's entry exactly once, got {withStateCount}\n" - + await DescribeAsync(connection, "resumed read after the rotation", withState.Rows, wait, carried, withState.Context.PendingState, 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); - 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)); + [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. */ - 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)); + _ = await LogAsync(connection, json, Marker(), ct); + var first = await CycleAsync(connection, null, binary, json, 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().ToList(); - Assert.True( - all.Count == 1, - $"the wait's entry should have one identity across both reads, got {all.Count}\n" - + await DescribeAsync(connection, "both reads together", first.Rows.Concat(withState.Rows), wait, carried, withState.Context.PendingState, ct)); + var wait = await LogAsync(connection, json, Marker(), ct, pokeTheWaiter: true); + await RotateAsync(connection, ct); + + 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] From 8f6cb14ae40c8ccdf7bf1f602484edb9d7bb1e60 Mon Sep 17 00:00:00 2001 From: Erik Darling <2136037+erikdarlingdata@users.noreply.github.com> Date: Wed, 30 Sep 2026 20:28:16 -0400 Subject: [PATCH 3/3] Open the poking connection before the wait and give the forced wait a 3 s lock_timeout --- .../PgServerLogTailCsvJsonRotationLiveTests.cs | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs b/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs index 7337c2010..2a57fbece 100644 --- a/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs +++ b/Darling/Darling.Tests/PgServerLogTailCsvJsonRotationLiveTests.cs @@ -122,6 +122,7 @@ private static async Task LogAsync(NpgsqlConnection connection, bool { 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); @@ -132,10 +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) gives this one wait a second line. */ - var poke = pokeTheWaiter ? PokeAsync(json, pid, ct) : Task.CompletedTask; + 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; @@ -153,10 +154,9 @@ between deadlock_timeout (100 ms) and lock_timeout (800 ms) gives this one wait } /// Wakes the waiter's latch about 350 ms into its wait, from a third connection. - private static async Task PokeAsync(bool json, int waiterPid, CancellationToken ct) + private static async Task PokeAsync(NpgsqlConnection third, int waiterPid, CancellationToken ct) { await Task.Delay(350, ct); - await using var third = await OpenAsync(json, ct); await using var command = new NpgsqlCommand("SELECT pg_log_backend_memory_contexts(" + waiterPid.ToString(System.Globalization.CultureInfo.InvariantCulture) + ")", third); await command.ExecuteScalarAsync(ct); }