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
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,10 @@ public async Task RunStreamedResponseAsync(Func<Nghttp3Request, Nghttp3ResponseW
{
_streaming = true;

// Writers parked at the retention high-water resume when acks drain it: the drain pumps and
// releases them, where a bare PumpEgress would leave them parked.
_quicConnection.OnSendCapacityAvailable = DrainStreamed;

try
{
while (true)
Expand Down Expand Up @@ -344,7 +348,13 @@ internal Task PumpAsync()
if (!_inStreamedDrain)
{
DrainStreamed();
return Task.CompletedTask;

// At the high-water nothing more is taken until acks drain it, and acks are only read once
// the reactor is back in its loop: spinning here never let it get there.
if (_quicConnection.CanQueueSend)
{
return Task.CompletedTask;
}
}

// NOT RunContinuationsAsynchronously: everything on a reactor resumes inline on the
Expand Down
32 changes: 32 additions & 0 deletions tests/Ioxide.Tests.E2E/Protocols/H3Tests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,38 @@ static async (_, writer) =>
Assert.Equal("chunkchunkchunkchunk", body);
});

runner.Test("h3: a streamed response past the send-retention high-water completes", () =>
{
// 20 MiB in 64 KiB flushes, from the dispatch pass: past the 16 MiB high-water the writer
// has to wait for acks, and its flush spun on the reactor thread that reads them instead.
(string certPath, string keyPath) = TestCert.Ensure();
using var engine = new QuicEngine(certPath, keyPath, cidLength: 8, alpn: ["h3"]);

const int Chunks = 320;
(_, int udpPort) = TestServer.StartDatagram(
onDatagram: null,
quicFactory: engine.CreateFactory(),
quicHandle: static (_, conn) => new Nghttp3Connection(conn).RunStreamedResponseAsync(
static async (_, writer) =>
{
writer.WriteHeaders(new Nghttp3Response { Status = 200 });
for (int i = 0; i < Chunks; i++)
{
writer.GetSpan(64 * 1024)[..(64 * 1024)].Fill((byte)'x');
writer.Advance(64 * 1024);
await writer.FlushAsync();
}
}));

using var client = new H3TestClient("127.0.0.1", udpPort);
client.Connect();
Assert.True(client.CompleteHandshake(timeoutMs: 5000), "handshake did not complete");

(int status, string body) = client.Get("/big", timeoutMs: 30_000);
Assert.True(status == 200 && body.Length == Chunks * 64 * 1024,
$"got status {status} and {body.Length} of {Chunks * 64 * 1024} bytes: the response stalled at the high-water");
});

runner.Test("h3: buffered-async handler (whole body in req.Body, handler may await)", () =>
{
(string certPath, string keyPath) = TestCert.Ensure();
Expand Down
Loading