diff --git a/src/protocols/ioxide.http3/Http3Connection.cs b/src/protocols/ioxide.http3/Http3Connection.cs index 2fc5edc..d6712d3 100644 --- a/src/protocols/ioxide.http3/Http3Connection.cs +++ b/src/protocols/ioxide.http3/Http3Connection.cs @@ -112,6 +112,7 @@ public Task RunAsync(Func> handler) private async Task RunCoreAsync(Func? buffered, Func>? streaming) { _streaming = streaming is not null; + _quicConnection.OnSendCapacityAvailable ??= SendPendingBodies; // a streamed connection has its own try { while (true) @@ -182,6 +183,7 @@ private async Task RunCoreAsync(Func? buffered, Fun } FireBodyWakes(); _requests.Clear(); + _pendingBodies.Clear(); // Tell the peer why. A protocol error used to end the handler and leave the connection // registered and routable: no H3 error code ever reached the client, and the connection @@ -722,7 +724,57 @@ private void SendSubmitted(long streamId, Http3Response resp, byte[] head, int h } _quicConnection.SendStream(streamId, head.AsSpan(0, headLen), fin: false); - _quicConnection.SendStream(streamId, resp.Body.Span, fin: true); + SendBody(streamId, resp.Body); + } + + // A large body goes in pieces while the connection takes them, and the rest as acks free + // retention: handed over whole, 32 one-MiB responses on one connection crossed the backstop and + // lost the connection. A piece overshoots the high-water by at most its own size. + private const int BodyPiece = 64 * 1024; + + private sealed class PendingBody(long streamId, ReadOnlyMemory rest) + { + public readonly long StreamId = streamId; + public ReadOnlyMemory Rest = rest; + } + + private readonly Queue _pendingBodies = new(); + + private void SendBody(long streamId, ReadOnlyMemory body) + { + if (_pendingBodies.Count == 0) + { + body = SendWhileCapacity(streamId, body); + if (body.IsEmpty) + { + return; + } + } + _pendingBodies.Enqueue(new PendingBody(streamId, body)); // after the ones already waiting + } + + private ReadOnlyMemory SendWhileCapacity(long streamId, ReadOnlyMemory body) + { + while (!body.IsEmpty && _quicConnection.CanQueueSend) + { + int n = Math.Min(body.Length, BodyPiece); + _quicConnection.SendStream(streamId, body.Span[..n], fin: n == body.Length); + body = body[n..]; + } + return body; + } + + private void SendPendingBodies() + { + while (_pendingBodies.TryPeek(out PendingBody? next) && !IsBroken) + { + next.Rest = SendWhileCapacity(next.StreamId, next.Rest); + if (!next.Rest.IsEmpty) + { + return; // at the high-water again: the next capacity signal carries on + } + _pendingBodies.Dequeue(); + } } // --- streaming plumbing -------------------------------------------------------------------- diff --git a/tests/Ioxide.Tests.E2E/Program.cs b/tests/Ioxide.Tests.E2E/Program.cs index 7f2be22..b11adfc 100644 --- a/tests/Ioxide.Tests.E2E/Program.cs +++ b/tests/Ioxide.Tests.E2E/Program.cs @@ -47,6 +47,7 @@ private static int Main() Http2BodyTests.Register(runner); Http2StreamedFaultTests.Register(runner); Http3BodyTests.Register(runner); + H3LargeResponseTests.Register(runner); return runner.Summary(); } diff --git a/tests/Ioxide.Tests.E2E/Protocols/H3LargeResponseTests.cs b/tests/Ioxide.Tests.E2E/Protocols/H3LargeResponseTests.cs new file mode 100644 index 0000000..4f7f771 --- /dev/null +++ b/tests/Ioxide.Tests.E2E/Protocols/H3LargeResponseTests.cs @@ -0,0 +1,114 @@ +using ioxide; +using ioxide.http3; +using ioxide.nghttp3; +using ioxide.ngtcp2; + +namespace Ioxide.Tests; + +/// +/// HTTP/3 responses bigger than a connection's send retention: 16 MiB unacknowledged before a +/// producer has to wait, 32 MiB before the connection is closed as one that ignored the wait. A +/// buffered handler hands its body over whole, so the connection has to feed it out as the peer +/// acks it rather than take it all into retention. +/// +internal static class H3LargeResponseTests +{ + private static readonly byte[] Big = CreateBody(40 << 20); + private static readonly byte[] OneMiB = CreateBody(1 << 20); + + private static byte[] CreateBody(int length) + { + var body = new byte[length]; + body.AsSpan().Fill((byte)'x'); + return body; + } + + private static readonly string[] Stacks = ["ioxide.http3", "ioxide.http3 async", "nghttp3"]; + + private static Func ServeBuffered(string stack, byte[] body) => stack switch + { + "ioxide.http3" => (_, conn) => new Http3Connection(conn).RunAsync(_ => new Http3Response { Body = body }), + "ioxide.http3 async" => (_, conn) => new Http3Connection(conn).RunAsync(_ => ValueTask.FromResult(new Http3Response { Body = body })), + _ => (_, conn) => new Nghttp3Connection(conn).RunBufferedAsync(_ => new Nghttp3Response { Body = body }), + }; + + public static void Register(Runner runner) + { + foreach (string stack in Stacks) + { + string s = stack; + + runner.Test($"h3 buffered ({s}): a response past the send-retention ceiling is served whole", () => + { + using H3TestClient client = Connect(s, Big, out QuicEngine engine); + using QuicEngine _ = engine; + + (int status, string body) = client.Get("/big", timeoutMs: 60_000); + Assert.True(status == 200 && body.Length == Big.Length, + $"got status {status} and {body.Length} of {Big.Length} bytes: the connection was closed at the backstop"); + }); + + runner.Test($"h3 buffered ({s}): 40 one-MiB responses at once on one connection are all served", () => + { + // The shape that found it: h2load's 32 streams of 1 MiB each, handed over together. + using H3TestClient client = Connect(s, OneMiB, out QuicEngine engine); + using QuicEngine _ = engine; + + (int Status, long BodyLength)[] responses = client.GetConcurrent( + Enumerable.Range(0, 40).Select(i => $"/part/{i}").ToArray(), timeoutMs: 60_000); + + int whole = responses.Count(r => r.Status == 200 && r.BodyLength == OneMiB.Length); + Assert.True(whole == responses.Length, + $"{whole} of {responses.Length} responses arrived whole: the connection was closed at the backstop"); + }); + } + + runner.Test("h3 streamed (ioxide.http3): a response past the send-retention high-water completes", () => + { + // The buffered path's capacity callback must not displace the streamed one: a writer + // parked at the high-water is resumed by it. + (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 Http3Connection(conn).RunStreamedResponseAsync( + static async (_, writer) => + { + writer.WriteHeaders(new Http3Response { 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: 5_000), "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"); + }); + } + + private static H3TestClient Connect(string stack, byte[] body, out QuicEngine engine) + { + (string certPath, string keyPath) = TestCert.Ensure(); + engine = new QuicEngine(certPath, keyPath, cidLength: 8, alpn: ["h3"]); + + (_, int udpPort) = TestServer.StartDatagram( + onDatagram: null, + quicFactory: engine.CreateFactory(), + quicHandle: ServeBuffered(stack, body)); + + var client = new H3TestClient("127.0.0.1", udpPort); + client.Connect(); + Assert.True(client.CompleteHandshake(timeoutMs: 5_000), "handshake did not complete"); + return client; + } +} diff --git a/tests/Ioxide.Tests.Harness/H3TestClient.cs b/tests/Ioxide.Tests.Harness/H3TestClient.cs index 2a9f62b..8a16389 100644 --- a/tests/Ioxide.Tests.Harness/H3TestClient.cs +++ b/tests/Ioxide.Tests.Harness/H3TestClient.cs @@ -206,6 +206,55 @@ public bool CompleteHandshake(int timeoutMs) return (_status, Encoding.UTF8.GetString(_body.ToArray())); } + // GetConcurrent's responses, by stream. Bodies are counted, not kept. + private sealed class Response + { + public int Status = -1; + public long Bytes; + public bool Done; + } + + private readonly Dictionary _concurrent = []; + + /// + /// GET every path at once on this connection, each on its own stream, and wait for all of the + /// responses. Returns each one's status and body length, in path order. + /// + public (int Status, long BodyLength)[] GetConcurrent(string[] paths, int timeoutMs) + { + EnsureH3Session(); + for (int settle = 0; settle < 10; settle++) + { + DrainH3Out(); + FlushOut(); + PumpIn(); + } + + var streams = new long[paths.Length]; + for (int i = 0; i < paths.Length; i++) + { + streams[i] = iq_client_open_bidi(_conn); + Assert.True(streams[i] >= 0, "failed to open request stream"); + _concurrent[streams[i]] = new Response(); + + byte[] headers = PackHeaders([(":method", "GET"), (":scheme", "https"), (":authority", "localhost"), (":path", paths[i])]); + fixed (byte* p = headers) + { + Assert.True(ih3_submit_request(_h3, streams[i], p, (nuint)headers.Length, null, 0) == 0, "submit_request failed"); + } + } + + long deadline = Environment.TickCount64 + timeoutMs; + while (Environment.TickCount64 < deadline && !_peerClosed && streams.Any(sid => !_concurrent[sid].Done)) + { + DrainH3Out(); + FlushOut(); + PumpIn(); + } + + return streams.Select(sid => (_concurrent[sid].Status, _concurrent[sid].Bytes)).ToArray(); + } + /// /// The H3 session - client conn plus its control and QPACK streams - stood up once per /// CONNECTION, not per request. HTTP/3 allows one control stream per peer and RFC 9114 6.2.1 @@ -455,6 +504,10 @@ private static void OnH3Header(void* user, long streamId, byte* name, nuint name if (n == ":status") { self._status = int.Parse(Encoding.ASCII.GetString(value, (int)valueLen)); + if (self._concurrent.TryGetValue(streamId, out Response? response)) + { + response.Status = self._status; + } } } @@ -469,6 +522,10 @@ private static void OnH3Data(void* user, long streamId, byte* data, nuint len) { self._body.AddRange(new ReadOnlySpan(data, (int)len).ToArray()); } + else if (self._concurrent.TryGetValue(streamId, out Response? response)) + { + response.Bytes += (long)len; + } } [UnmanagedCallersOnly] @@ -479,6 +536,10 @@ private static void OnH3EndStream(void* user, long streamId) { self._done = true; } + else if (self._concurrent.TryGetValue(streamId, out Response? response)) + { + response.Done = true; + } } public void Dispose()