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}"); + } }