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
54 changes: 53 additions & 1 deletion src/protocols/ioxide.http3/Http3Connection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@ public Task RunAsync(Func<Http3Request, ValueTask<Http3Response>> handler)
private async Task RunCoreAsync(Func<Http3Request, Http3Response>? buffered, Func<Http3Request, ValueTask<Http3Response>>? streaming)
{
_streaming = streaming is not null;
_quicConnection.OnSendCapacityAvailable ??= SendPendingBodies; // a streamed connection has its own
try
{
while (true)
Expand Down Expand Up @@ -182,6 +183,7 @@ private async Task RunCoreAsync(Func<Http3Request, Http3Response>? 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
Expand Down Expand Up @@ -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<byte> rest)
{
public readonly long StreamId = streamId;
public ReadOnlyMemory<byte> Rest = rest;
}

private readonly Queue<PendingBody> _pendingBodies = new();

private void SendBody(long streamId, ReadOnlyMemory<byte> 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<byte> SendWhileCapacity(long streamId, ReadOnlyMemory<byte> 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 --------------------------------------------------------------------
Expand Down
1 change: 1 addition & 0 deletions tests/Ioxide.Tests.E2E/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ private static int Main()
Http2BodyTests.Register(runner);
Http2StreamedFaultTests.Register(runner);
Http3BodyTests.Register(runner);
H3LargeResponseTests.Register(runner);

return runner.Summary();
}
Expand Down
114 changes: 114 additions & 0 deletions tests/Ioxide.Tests.E2E/Protocols/H3LargeResponseTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
using ioxide;
using ioxide.http3;
using ioxide.nghttp3;
using ioxide.ngtcp2;

namespace Ioxide.Tests;

/// <summary>
/// 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.
/// </summary>
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<Reactor, QuicConnection, Task> 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;
}
}
61 changes: 61 additions & 0 deletions tests/Ioxide.Tests.Harness/H3TestClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<long, Response> _concurrent = [];

/// <summary>
/// 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.
/// </summary>
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();
}

/// <summary>
/// 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
Expand Down Expand Up @@ -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;
}
}
}

Expand All @@ -469,6 +522,10 @@ private static void OnH3Data(void* user, long streamId, byte* data, nuint len)
{
self._body.AddRange(new ReadOnlySpan<byte>(data, (int)len).ToArray());
}
else if (self._concurrent.TryGetValue(streamId, out Response? response))
{
response.Bytes += (long)len;
}
}

[UnmanagedCallersOnly]
Expand All @@ -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()
Expand Down
Loading