Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,21 @@ public async Task EnqueueAgainstStalledDrain_RaisesHighWaterPastSoftThreshold()
await using var stall = new StallStream();
await using var writer = new Http2FrameWriter(stall);

// First write blocks the drain on StallStream; subsequent frames pile up in the channel.
// Block the drain inside WriteAsync first. Enqueueing the backlog before that races the
// drain: it can coalesce a partial batch, stall, and leave HighWater under 64.
var primer = ArrayPool<byte>.Shared.Rent(32);
writer.EnqueueRented(primer, 32);
var writeStarted = stall.WaitForWriteStarted();
var started = await Task.WhenAny(writeStarted, Task.Delay(TimeSpan.FromSeconds(2)));
Assert.AreEqual(writeStarted, started, "drain did not reach WriteAsync");

for (var i = 0; i < 80; i++)
{
var rented = ArrayPool<byte>.Shared.Rent(32);
writer.EnqueueRented(rented, 32);
}

// Drain coalesces up to 64 frames into one WriteAsync that stalls; pending drops for those
// frames while HighWater still recorded the backlog peak.
Assert.IsTrue(SpinWait.SpinUntil(() => writer.HighWaterMark >= 64, TimeSpan.FromSeconds(2)),
Assert.IsTrue(writer.HighWaterMark >= 64,
$"expected high-water ≥ 64, got {writer.HighWaterMark}");
Assert.IsTrue(writer.HighWaterMark >= writer.PendingFrameCount);

Expand All @@ -38,12 +43,17 @@ private sealed class StallStream : Stream
{
private readonly TaskCompletionSource<bool> gate =
new(TaskCreationOptions.RunContinuationsAsynchronously);
private readonly TaskCompletionSource started =
new(TaskCreationOptions.RunContinuationsAsynchronously);

public void Release() => gate.TrySetResult(true);

public Task WaitForWriteStarted() => started.Task;

public override async ValueTask WriteAsync(ReadOnlyMemory<byte> buffer,
CancellationToken cancellationToken = default)
{
started.TrySetResult();
await gate.Task.WaitAsync(cancellationToken);
}

Expand Down
Loading