From b7b7349e90c872a4f395430aec6d6b896b3832d7 Mon Sep 17 00:00:00 2001 From: Paul DeVito Date: Sun, 27 Sep 2026 18:38:32 +0300 Subject: [PATCH] test(hosting): fix two races in the row-by-row job output test The host starts the pump on a thread-pool thread, so the pump can subscribe to Wake-Up Hints after StartAsync returns. FaultableStore.Wake then threw because no pump had subscribed yet. WakeAsync now waits for the subscription before it sends the hint. The pump also writes the fence-out log (1208) after the store write that makes both jobs terminal. The test read the logs as soon as both jobs were terminal, so it could miss that log. The test now waits for the log before it asserts on it. --- .../BackWave.Hosting.Tests/FaultableStore.cs | 11 ++++++++-- .../JobOutputTooLargeTests.cs | 20 ++++++++++++++++++- 2 files changed, 28 insertions(+), 3 deletions(-) diff --git a/tests/BackWave.Hosting.Tests/FaultableStore.cs b/tests/BackWave.Hosting.Tests/FaultableStore.cs index 03b3949..273a8ee 100644 --- a/tests/BackWave.Hosting.Tests/FaultableStore.cs +++ b/tests/BackWave.Hosting.Tests/FaultableStore.cs @@ -14,6 +14,7 @@ public sealed class TransientStoreException() : DbException("forced transient st /// Wraps the In-Memory Store to force invariant violations on demand. public sealed class FaultableStore(IJobStore inner) : IJobStore, IWakeUpHintSource { + private readonly TaskCompletionSource _subscribed = new(TaskCreationOptions.RunContinuationsAsynchronously); private Action? _onHint; private int _transientClaimFaults; private int _relinquishCalls; @@ -103,17 +104,23 @@ public int TransientClaimFaults /// /// Wakes the subscribed pump for this Queue, as an adapter's Wake-Up Hint would. A test that must keep /// every poll tick out of the window it exercises sets a long poll interval and claims through this. + /// + /// The host starts the pump on a thread-pool thread, so the pump can subscribe after the host's + /// StartAsync returns. This waits for that subscription before it sends the hint. + /// /// - public void Wake(string queue) + public async Task WakeAsync(string queue) { + await _subscribed.Task.WaitAsync(TimeSpan.FromSeconds(10)).ConfigureAwait(false); var onHint = Volatile.Read(ref _onHint) - ?? throw new InvalidOperationException("No pump has subscribed to this store's Wake-Up Hints."); + ?? throw new InvalidOperationException("The pump's Wake-Up Hint subscription is already disposed."); onHint(queue); } public Task SubscribeAsync(Action onHint, CancellationToken cancellationToken = default) { Volatile.Write(ref _onHint, onHint); + _subscribed.TrySetResult(); return Task.FromResult(new HintSubscription(this)); } diff --git a/tests/BackWave.Hosting.Tests/JobOutputTooLargeTests.cs b/tests/BackWave.Hosting.Tests/JobOutputTooLargeTests.cs index 112f257..579ba6f 100644 --- a/tests/BackWave.Hosting.Tests/JobOutputTooLargeTests.cs +++ b/tests/BackWave.Hosting.Tests/JobOutputTooLargeTests.cs @@ -188,7 +188,7 @@ public async Task OversizedOutput_OnARowByRowStore_ReAppliesTheBatchOnceAndCount var healthy = await client.EnqueueAsync(new OutputJob("small-1"), dueTime: DateTimeOffset.UtcNow); var oversized = await client.EnqueueAsync(new OutputJob("big-1"), dueTime: DateTimeOffset.UtcNow); // The one claim pass this test needs: both jobs are due, so they land in the same batch. - store.Wake("default"); + await store.WakeAsync("default"); await AwaitTerminalAsync(monitor, oversized, healthy); @@ -204,6 +204,10 @@ public async Task OversizedOutput_OnARowByRowStore_ReAppliesTheBatchOnceAndCount Assert.Contains("Job Output", dead.TerminalCause); // The settled row fences out on the second pass. Benign, so: Debug, no counter, group untouched. + // The pump logs only after the store write that made both jobs terminal, so wait for the log itself. + await WaitForAsync( + () => logs.Entries.Any(entry => entry.EventId == FencedOutEventId), + "the fenced-out outcome to be logged at Debug"); var debug = Assert.Single(logs.Entries, entry => entry.EventId == FencedOutEventId); Assert.Contains(healthy.ToString(), debug.Message); Assert.Empty(violations.Measurements); @@ -302,4 +306,18 @@ private static async Task AwaitTerminalAsync(BackWaveMonitor monitor, params Gui await Task.Delay(25); } } + + private static async Task WaitForAsync(Func condition, string description) + { + var deadline = DateTimeOffset.UtcNow + TestTimeout; + while (DateTimeOffset.UtcNow < deadline) + { + if (condition()) + { + return; + } + await Task.Delay(25); + } + Assert.Fail($"Timed out waiting for: {description}"); + } }