Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions tests/BackWave.Hosting.Tests/FaultableStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ public sealed class TransientStoreException() : DbException("forced transient st
/// <summary>Wraps the In-Memory Store to force invariant violations on demand.</summary>
public sealed class FaultableStore(IJobStore inner) : IJobStore, IWakeUpHintSource
{
private readonly TaskCompletionSource _subscribed = new(TaskCreationOptions.RunContinuationsAsynchronously);
private Action<string>? _onHint;
private int _transientClaimFaults;
private int _relinquishCalls;
Expand Down Expand Up @@ -103,17 +104,23 @@ public int TransientClaimFaults
/// <summary>
/// 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.
/// <para>
/// 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.
/// </para>
/// </summary>
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<IAsyncDisposable> SubscribeAsync(Action<string> onHint, CancellationToken cancellationToken = default)
{
Volatile.Write(ref _onHint, onHint);
_subscribed.TrySetResult();
return Task.FromResult<IAsyncDisposable>(new HintSubscription(this));
}

Expand Down
20 changes: 19 additions & 1 deletion tests/BackWave.Hosting.Tests/JobOutputTooLargeTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand All @@ -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);
Expand Down Expand Up @@ -302,4 +306,18 @@ private static async Task AwaitTerminalAsync(BackWaveMonitor monitor, params Gui
await Task.Delay(25);
}
}

private static async Task WaitForAsync(Func<bool> 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}");
}
}
Loading