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