From a05d006fea287cbb2829d4a3429ba5504476e677 Mon Sep 17 00:00:00 2001 From: Diogo Martins Date: Sun, 4 Oct 2026 17:07:11 +0100 Subject: [PATCH] nghttp3: a streamed response at the send-retention high-water waits for acks instead of spinning At the high-water PumpEgress stops pulling, so a writer's chunk is not taken until acks drain retention - and a writer flushing from the dispatch pass looped on PumpAsync, which outside a drain pumped once and returned at once. The reactor never got back to reading the acks: 0 of 512 one-MiB responses completed under h2load, both reactors at 100% CPU for good. PumpAsync now parks it there, and the capacity signal runs the full drain, which releases parked writers where PumpEgress alone left them waiting. --- .../Connection/Nghttp3Connection.Streamed.cs | 12 ++++++- tests/Ioxide.Tests.E2E/Protocols/H3Tests.cs | 32 +++++++++++++++++++ 2 files changed, 43 insertions(+), 1 deletion(-) diff --git a/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.Streamed.cs b/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.Streamed.cs index 215fdf5c..d5c182cf 100644 --- a/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.Streamed.cs +++ b/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.Streamed.cs @@ -52,6 +52,10 @@ public async Task RunStreamedResponseAsync(Func 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();