diff --git a/src/protocols/ioxide.http2/Http2Connection.Frames.cs b/src/protocols/ioxide.http2/Http2Connection.Frames.cs index 90264d30..a4fb06e8 100644 --- a/src/protocols/ioxide.http2/Http2Connection.Frames.cs +++ b/src/protocols/ioxide.http2/Http2Connection.Frames.cs @@ -504,6 +504,13 @@ private void HandleRstStream(in FrameHeader header) { pending.Dispose(); // the peer gave up; there is nobody to answer } + + // A response in flight may send nothing more on the stream (RFC 9113 5.1), and a writer parked + // on its credit has to wake to learn that - the WINDOW_UPDATE it waits for will never come. + if (_responseWindows.Remove(header.StreamId)) + { + ReleaseCreditWaiters(header.StreamId); + } } /// diff --git a/src/protocols/ioxide.http2/Http2Connection.Streamed.cs b/src/protocols/ioxide.http2/Http2Connection.Streamed.cs index c88f43d7..346942b6 100644 --- a/src/protocols/ioxide.http2/Http2Connection.Streamed.cs +++ b/src/protocols/ioxide.http2/Http2Connection.Streamed.cs @@ -110,29 +110,22 @@ private Http2ResponseWriter RentWriter(int streamId) internal void SendStreamedHeaders(int streamId, Http2Response response) => WriteHeaders(streamId, response, endStream: false); + /// False once the peer has reset this response's stream, or the connection is gone. + internal bool IsResponseLive(int streamId) => !IsBroken && _responseWindows.ContainsKey(streamId); + /// /// How many body bytes may be sent on this stream right now: the smaller of the connection - /// window, the stream window and the peer's maximum frame size. Zero means wait. + /// window, the stream window and the peer's maximum frame size. Zero means wait; negative means + /// never - the peer reset the stream, or the connection is gone. /// internal int SendCredit(int streamId) { - if (IsBroken) - { - return 0; - } - - int credit = Math.Min(_peerConnectionWindow, _peerMaxFrameSize); - - if (_responseWindows.TryGetValue(streamId, out int window)) - { - credit = Math.Min(credit, window); - } - else if (_streams.TryGetValue(streamId, out PendingRequest? pending)) + if (IsBroken || !_responseWindows.TryGetValue(streamId, out int window)) { - credit = Math.Min(credit, pending.SendWindow); + return -1; } - return Math.Max(credit, 0); + return Math.Max(Math.Min(Math.Min(_peerConnectionWindow, _peerMaxFrameSize), window), 0); } /// One DATA frame, already known to fit both windows. diff --git a/src/protocols/ioxide.http2/Http2ResponseWriter.cs b/src/protocols/ioxide.http2/Http2ResponseWriter.cs index 30d1587e..684248ba 100644 --- a/src/protocols/ioxide.http2/Http2ResponseWriter.cs +++ b/src/protocols/ioxide.http2/Http2ResponseWriter.cs @@ -1,4 +1,5 @@ using System.Buffers; +using System.IO.Pipelines; namespace ioxide.http2; @@ -72,7 +73,10 @@ public void WriteHeaders(Http2Response response) } _headersSent = true; - _connection.SendStreamedHeaders(_streamId, response); + if (_connection.IsResponseLive(_streamId)) + { + _connection.SendStreamedHeaders(_streamId, response); + } } /// @@ -105,7 +109,11 @@ public void Advance(int count) /// That wait is the backpressure: a peer that stops reading stops the producer rather than /// growing a queue behind it. /// - public ValueTask FlushAsync() => FlushCore(endStream: false); + /// + /// once the peer has reset the stream or the connection is + /// gone, as a TCP pipe writer reports a closed peer: nothing written from then on reaches anyone. + /// + public ValueTask FlushAsync() => FlushCore(endStream: false); /// /// Send what is left and mark END_STREAM. A handler that returns without calling this gets it @@ -153,10 +161,16 @@ internal async ValueTask FailAsync() _staging = []; } - await _connection.ResetStreamedAsync(_streamId); + // Not in answer to the peer's own reset (RFC 9113 5.4.2), and not on a connection that is gone. + if (_connection.IsResponseLive(_streamId)) + { + await _connection.ResetStreamedAsync(_streamId); + } } - private async ValueTask FlushCore(bool endStream) + private static readonly FlushResult PeerGone = new(isCanceled: false, isCompleted: true); + + private async ValueTask FlushCore(bool endStream) { if (!_headersSent) { @@ -169,13 +183,13 @@ private async ValueTask FlushCore(bool endStream) // Bounded by the frame size AND by both windows: exceeding either is a connection // error the peer would be right to hang up over. int credit = _connection.SendCredit(_streamId); - if (credit <= 0) + if (credit < 0) + { + _staged = 0; // reset by the peer, or the connection is gone: nobody to send it to + return PeerGone; + } + if (credit == 0) { - if (_connection.IsBroken) - { - return; - } - // Take the wait BEFORE flushing. The flush below is an await, and the WINDOW_UPDATE // it exists to provoke can arrive while this writer is still inside it - at which // point a release has no waiter to find, is dropped, and the writer then parks on a @@ -207,6 +221,12 @@ private async ValueTask FlushCore(bool endStream) _staged = 0; _sinceRealFlush += sent; + // Sending proved the stream live; with nothing sent it still has to be asked. + if (sent == 0 && !_connection.IsResponseLive(_streamId)) + { + return PeerGone; + } + if (endStream && sent == 0) { // Nothing left to send, but the stream still needs its end. @@ -218,11 +238,12 @@ private async ValueTask FlushCore(bool endStream) // limit, where it must write for real to yield. Outside a pass the flush is always real. if (_connection.InDispatchPass && _sinceRealFlush < Http2Connection.CoalesceLimit) { - return; + return default; } _sinceRealFlush = 0; await _connection.FlushOutboundAsync(); + return default; } private void EnsureStaging(int sizeHint) diff --git a/src/protocols/ioxide.http3/Http3Connection.Streamed.cs b/src/protocols/ioxide.http3/Http3Connection.Streamed.cs index ba4153be..c72e1680 100644 --- a/src/protocols/ioxide.http3/Http3Connection.Streamed.cs +++ b/src/protocols/ioxide.http3/Http3Connection.Streamed.cs @@ -9,6 +9,9 @@ namespace ioxide.http3; public sealed partial class Http3Connection { private readonly Stack _writerPool = new(); + + // Streamed responses in flight, so a STOP_SENDING from the peer can reach the one it ends. + private readonly Dictionary _writers = new(); private readonly List _capacityWaiters = []; /// True once the connection can no longer make progress; a parked writer gives up. @@ -61,6 +64,7 @@ private bool TryDispatchStreamedResponse(Http3Request request) } Http3ResponseWriter writer = RentWriter(request.StreamId); + _writers[request.StreamId] = writer; _ = ServeAsync(_streamedResponseHandler, request, writer); return true; } @@ -80,6 +84,7 @@ private async Task ServeAsync(Func } finally { + _writers.Remove(writer.StreamId); _writerPool.Push(writer); } } diff --git a/src/protocols/ioxide.http3/Http3Connection.cs b/src/protocols/ioxide.http3/Http3Connection.cs index 2fc5edcb..03a2fb99 100644 --- a/src/protocols/ioxide.http3/Http3Connection.cs +++ b/src/protocols/ioxide.http3/Http3Connection.cs @@ -231,6 +231,15 @@ private void Feed(in QuicRecvRing.Delivery item) ReleaseParseBuffers(dead); } _unis.Remove(item.StreamId); + + // A response in flight whose peer will read no more of it: its handler learns at the next + // flush, and one parked on send capacity wakes to find that out. + if (item.Kind is QuicStreamEvent.StopSending or QuicStreamEvent.Closed + && _writers.TryGetValue(item.StreamId, out Http3ResponseWriter? stopped)) + { + stopped.OnPeerGone(); + ReleaseCapacityWaiters(); + } return; } diff --git a/src/protocols/ioxide.http3/Http3ResponseWriter.cs b/src/protocols/ioxide.http3/Http3ResponseWriter.cs index fb66cf98..962dd551 100644 --- a/src/protocols/ioxide.http3/Http3ResponseWriter.cs +++ b/src/protocols/ioxide.http3/Http3ResponseWriter.cs @@ -1,4 +1,5 @@ using System.Buffers; +using System.IO.Pipelines; namespace ioxide.http3; @@ -35,6 +36,7 @@ public sealed class Http3ResponseWriter : IBufferWriter private bool _headersSent; private bool _completed; + private bool _gone; // the peer stopped reading, or the stream closed under the handler internal Http3ResponseWriter(Http3Connection connection, QuicConnection quic, long streamId) { @@ -72,7 +74,10 @@ public void WriteHeaders(Http3Response response) } _headersSent = true; - _connection.SendStreamedHeaders(_streamId, response); + if (!_gone) + { + _connection.SendStreamedHeaders(_streamId, response); + } } /// @@ -105,7 +110,17 @@ public void Advance(int count) /// send-retention high-water: that wait is the backpressure, and it is what keeps memory bound /// to one chunk rather than to the whole response. /// - public ValueTask FlushAsync() => FlushCore(fin: false); + /// + /// once the peer has stopped reading the stream or the + /// connection is gone, as a TCP pipe writer reports a closed peer: nothing written from then on + /// reaches anyone. + /// + public ValueTask FlushAsync() => FlushCore(fin: false); + + /// The peer will read no more of this stream: from here on nothing is sent on it. + internal void OnPeerGone() => _gone = true; + + private static readonly FlushResult PeerGone = new(isCanceled: false, isCompleted: true); /// /// Send what is left and close the stream. A handler that returns without calling this still @@ -146,7 +161,10 @@ internal ValueTask FailAsync() _completed = true; _staged = 0; - _quic.ResetStream(_streamId, H3InternalError); + if (!_gone) + { + _quic.ResetStream(_streamId, H3InternalError); // never on a stream the peer already stopped + } if (_staging.Length > 0) { @@ -156,7 +174,7 @@ internal ValueTask FailAsync() return ValueTask.CompletedTask; } - private async ValueTask FlushCore(bool fin) + private async ValueTask FlushCore(bool fin) { if (!_headersSent) { @@ -165,25 +183,26 @@ private async ValueTask FlushCore(bool fin) if (_staged == 0 && !fin) { - return; + return _gone || _connection.IsBroken ? PeerGone : default; } // The peer has stopped reading and the connection is holding all it is willing to. Wait // rather than queue: unbounded queueing here is exactly what streaming exists to avoid. - while (!_quic.CanQueueSend && !_connection.IsBroken) + while (!_quic.CanQueueSend && !_connection.IsBroken && !_gone) { await _connection.WaitForSendCapacityAsync(); } - if (_connection.IsBroken) + if (_connection.IsBroken || _gone) { - return; + _staged = 0; // nobody to send it to + return PeerGone; } if (_staged == 0) { _quic.SendStream(_streamId, ReadOnlySpan.Empty, fin: true); - return; + return default; } // [0x00][varint length][payload], sent as one call - the header is tiny and splitting it @@ -200,6 +219,7 @@ private async ValueTask FlushCore(bool fin) ArrayPool.Shared.Return(frame); _staged = 0; + return default; } private void EnsureStaging(int sizeHint) @@ -233,5 +253,6 @@ internal void Reset(long streamId) _staged = 0; _headersSent = false; _completed = false; + _gone = false; } } diff --git a/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.Callbacks.cs b/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.Callbacks.cs index 57f3d6ae..3f29c3d5 100644 --- a/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.Callbacks.cs +++ b/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.Callbacks.cs @@ -159,6 +159,10 @@ private static unsafe void CallbackStreamError(void* user, int streamId, uint er { pending.Dispose(); } + if (connection._writers.TryGetValue(streamId, out Nghttp2ResponseWriter? writer)) + { + writer.OnPeerReset(); // a response in flight: its handler learns at the next flush + } _ = errorCode; } catch (Exception e) diff --git a/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.Streamed.cs b/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.Streamed.cs index 186c23e4..d5e51ecc 100644 --- a/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.Streamed.cs +++ b/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.Streamed.cs @@ -14,6 +14,9 @@ namespace ioxide.nghttp2; public sealed partial class Nghttp2Connection { private readonly Stack _writerPool = new(); + + // Streamed responses in flight, so a RST_STREAM from the peer can reach the one it ends. + private readonly Dictionary _writers = new(); private Func? _streamedHandler; /// @@ -43,6 +46,7 @@ private bool TryDispatchStreamed(Nghttp2Request request, PendingRequest pending) } Nghttp2ResponseWriter writer = RentWriter(request.StreamId); + _writers[request.StreamId] = writer; _ = ServeStreamedAsync(_streamedHandler, request, writer, pending); return true; } @@ -74,6 +78,7 @@ private async Task ServeStreamedAsync(FuncHand a chunk to nghttp2 and wake the deferred stream. Copied natively. - internal unsafe void SendStreamedData(int streamId, ReadOnlySpan body) + /// False when nghttp2 no longer knows the stream: the peer reset it. + internal unsafe bool SendStreamedData(int streamId, ReadOnlySpan body) { if (_handle == 0 || body.IsEmpty) { - return; + return true; } fixed (byte* bodyBytes = body) { - if (Nghttp2.ih2_stream_write(_handle, streamId, bodyBytes, (nuint)body.Length) != 0) + int result = Nghttp2.ih2_stream_write(_handle, streamId, bodyBytes, (nuint)body.Length); + + // A stream closed by the peer is dropped natively, so the write finds nothing - one + // stream's end, not a fault in the connection every other stream is riding. + if (result == ErrInvalidArgument) + { + return false; + } + if (result != 0) { _failed = true; } } + return true; } /// A streamed response its handler could not finish: RST_STREAM INTERNAL_ERROR. diff --git a/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.cs b/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.cs index 1ef6a51d..c5ab04a3 100644 --- a/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.cs +++ b/src/protocols/ioxide.nghttp2/Connection/Nghttp2Connection.cs @@ -37,6 +37,7 @@ public sealed partial class Nghttp2Connection : IDisposable // RFC 9113 INTERNAL_ERROR: a stream this server could not finish. private const uint InternalError = 0x2; + private const int ErrInvalidArgument = -501; // NGHTTP2_ERR_INVALID_ARGUMENT private readonly IDuplexPipe _pipe; private readonly Nghttp2Options _options; diff --git a/src/protocols/ioxide.nghttp2/Http/Nghttp2ResponseWriter.cs b/src/protocols/ioxide.nghttp2/Http/Nghttp2ResponseWriter.cs index 90f73976..517f7e51 100644 --- a/src/protocols/ioxide.nghttp2/Http/Nghttp2ResponseWriter.cs +++ b/src/protocols/ioxide.nghttp2/Http/Nghttp2ResponseWriter.cs @@ -1,4 +1,5 @@ using System.Buffers; +using System.IO.Pipelines; namespace ioxide.nghttp2; @@ -26,6 +27,7 @@ public sealed class Nghttp2ResponseWriter : IBufferWriter private bool _headersSent; private bool _completed; + private bool _gone; // the peer reset the stream; nghttp2 has forgotten it internal Nghttp2ResponseWriter(Nghttp2Connection connection, int streamId) { @@ -47,7 +49,10 @@ public void WriteHeaders(Nghttp2Response response) } _headersSent = true; - _connection.SendStreamedHeaders(StreamId, response); + if (!_gone) + { + _connection.SendStreamedHeaders(StreamId, response); + } } // --- IBufferWriter --------------------------------------------------------------- @@ -73,22 +78,42 @@ public Span GetSpan(int sizeHint = 0) /// Headers are sent for you if the handler never called , because a /// body cannot precede them and a 200 is what the buffered path would have sent. /// - public async ValueTask FlushAsync() + /// + /// once the peer has reset the stream or the connection is + /// gone, as a TCP pipe writer reports a closed peer: nothing written from then on reaches anyone. + /// + public async ValueTask FlushAsync() { if (!_headersSent) { WriteHeaders(new Nghttp2Response { Status = 200 }); } - if (_staged > 0) + HandOverStaged(); + + if (_gone || _connection.IsBroken) { - _connection.SendStreamedData(StreamId, _staging.AsSpan(0, _staged)); - _staged = 0; + return PeerGone; } await _connection.FlushStreamedAsync(); + return default; } + private static readonly FlushResult PeerGone = new(isCanceled: false, isCompleted: true); + + private void HandOverStaged() + { + if (_staged > 0 && !_gone && !_connection.SendStreamedData(StreamId, _staging.AsSpan(0, _staged))) + { + _gone = true; + } + _staged = 0; + } + + /// The peer reset this stream: from here on nothing is sent on it. + internal void OnPeerReset() => _gone = true; + /// End the response: whatever is staged goes out, then END_STREAM. Idempotent. internal async ValueTask CompleteAsync() { @@ -105,10 +130,10 @@ internal async ValueTask CompleteAsync() WriteHeaders(new Nghttp2Response { Status = 200 }); } - if (_staged > 0) + HandOverStaged(); + if (_gone) { - _connection.SendStreamedData(StreamId, _staging.AsSpan(0, _staged)); - _staged = 0; + return; // reset by the peer: there is no end left to send } _connection.EndStreamedBody(StreamId); @@ -136,6 +161,10 @@ internal async ValueTask FailAsync() _completed = true; _staged = 0; + if (_gone) + { + return; // never in answer to the peer's own reset (RFC 9113 5.4.2) + } _connection.ResetStreamed(StreamId); await _connection.FlushStreamedAsync(); } @@ -164,6 +193,7 @@ internal void Reset(int streamId) _staged = 0; _headersSent = false; _completed = false; + _gone = false; } internal void Release() diff --git a/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.PushToEngine.cs b/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.PushToEngine.cs index bf497268..7f8f06e6 100644 --- a/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.PushToEngine.cs +++ b/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.PushToEngine.cs @@ -32,6 +32,10 @@ private unsafe void PushToEngine(in QuicRecvRing.Delivery item) // the pool. Without this the writer and its native staging block lived until the // CONNECTION closed - one per response, which is what made a streamed response cost // several times the memory of a buffered one. + if (item.Kind == QuicStreamEvent.StopSending && _writers.TryGetValue(item.StreamId, out Nghttp3ResponseWriter? stopped)) + { + stopped.OnPeerGone(); // never pulled again; a parked flush learns on this pass's drain + } if (item.Kind is not (QuicStreamEvent.Reset or QuicStreamEvent.StopSending)) { ReleaseWriter(item.StreamId); diff --git a/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.Streamed.cs b/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.Streamed.cs index 215fdf5c..075deab1 100644 --- a/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.Streamed.cs +++ b/src/protocols/ioxide.nghttp3/Connection/Nghttp3Connection.Streamed.cs @@ -217,6 +217,7 @@ private void DispatchStreamedResponse(Func private bool _headersSent; private bool _completed; // the handler has finished producing private bool _finReported; // nghttp3 has been told this is the end + private bool _gone; // the peer stopped reading, or the stream closed under the handler + + /// True from dispatch until the handler is done with this writer (see HandlerExited). + internal bool HandlerRunning; internal Nghttp3ResponseWriter(Nghttp3Connection connection, long streamId) { @@ -78,6 +83,8 @@ internal void Reset(long streamId) _headersSent = false; _completed = false; _finReported = false; + _gone = false; + HandlerRunning = false; // The scratch buffer survives - that is the point of pooling the writer - but a hand-out // from the previous stream must not be mistaken for one from this one. @@ -113,7 +120,10 @@ public void WriteHeaders(Nghttp3Response response) } _headersSent = true; - _connection.SubmitStreamedHeaders(_streamId, response, this); + if (!_gone) + { + _connection.SubmitStreamedHeaders(_streamId, response, this); + } } /// @@ -181,7 +191,12 @@ public unsafe void Advance(int count) /// taken, which is the backpressure: a peer that stops reading stops this returning, so a /// producer cannot outrun the connection. /// - public async ValueTask FlushAsync() + /// + /// once the peer has stopped reading the stream or the + /// connection is gone, as a TCP pipe writer reports a closed peer: nothing written from then on + /// reaches anyone. + /// + public async ValueTask FlushAsync() { if (!_headersSent) { @@ -189,27 +204,35 @@ public async ValueTask FlushAsync() } if (_staged == 0) { - return; + return _gone || _connection.IsFailed ? PeerGone : default; } // Wait for the previous chunk to be taken before replacing it - one is in flight at a time. - while (_inFlightLength > 0 && !_connection.IsFailed) + // A stream the peer stopped is never pulled again, so that ends the wait too. + while (_inFlightLength > 0 && !_gone && !_connection.IsFailed) { _connection.ResumeStreamedResponse(_streamId); await _connection.PumpAsync(); } - if (_connection.IsFailed) + if (_gone || _connection.IsFailed) { - return; + _staged = 0; // nobody to send it to + return PeerGone; } PromoteStagedChunk(); _connection.ResumeStreamedResponse(_streamId); _connection.PumpIfOutsidePass(); + return default; } + private static readonly FlushResult PeerGone = new(isCanceled: false, isCompleted: true); + + /// The peer will read no more of this stream: from here on nothing is sent on it. + internal void OnPeerGone() => _gone = true; + /// /// Flush what is left and mark the end of the body. /// @@ -238,6 +261,10 @@ public async ValueTask CompleteAsync() // pass pumps it out. The staged buffer stays alive until the stream closes, so there is // nothing to keep the handler around for. _completed = true; + if (_gone) + { + return; // the peer stopped the stream: there is no end left to send + } _connection.ResumeStreamedResponse(_streamId); _connection.PumpIfOutsidePass(); } diff --git a/tests/Ioxide.Tests.E2E/Program.cs b/tests/Ioxide.Tests.E2E/Program.cs index 7f2be229..28c1bfd8 100644 --- a/tests/Ioxide.Tests.E2E/Program.cs +++ b/tests/Ioxide.Tests.E2E/Program.cs @@ -46,6 +46,7 @@ private static int Main() H3AlpnTests.Register(runner); Http2BodyTests.Register(runner); Http2StreamedFaultTests.Register(runner); + StreamedPeerGoneTests.Register(runner); Http3BodyTests.Register(runner); return runner.Summary(); diff --git a/tests/Ioxide.Tests.E2E/Protocols/StreamedPeerGoneTests.cs b/tests/Ioxide.Tests.E2E/Protocols/StreamedPeerGoneTests.cs new file mode 100644 index 00000000..d373f503 --- /dev/null +++ b/tests/Ioxide.Tests.E2E/Protocols/StreamedPeerGoneTests.cs @@ -0,0 +1,224 @@ +using System.Buffers; +using System.IO.Pipelines; +using System.Net; +using ioxide; +using ioxide.http2; +using ioxide.http3; +using ioxide.nghttp2; +using ioxide.nghttp3; +using ioxide.ngtcp2; +using ioxide.timer; + +namespace Ioxide.Tests; + +/// +/// A streamed response whose peer has gone (#269). An event source writes until a write fails, so +/// unless the writer can say "nobody is reading" it produces for a closed tab forever - or, once +/// the stream window is spent, parks forever on credit that will never come. A TCP pipe writer +/// says it with ; the HTTP/2 and HTTP/3 writers now say it the +/// same way, on every stack. +/// +internal static class StreamedPeerGoneTests +{ + public static void Register(Runner runner) + { + foreach (bool native in new[] { false, true }) + { + bool ng = native; + string stack = ng ? "nghttp2" : "ioxide.http2"; + + runner.Test($"h2 peer gone ({stack}): a reset stream is reported by FlushAsync, and the connection serves on", () => + { + var server = new EndlessServer(); + int port = StartH2(ng, server); + + using HttpClient client = H2Client(); + HttpResponseMessage endless = client.SendAsync( + new HttpRequestMessage(HttpMethod.Get, $"http://127.0.0.1:{port}/endless") + { + Version = HttpVersion.Version20, // a hand-built request ignores the client's default + VersionPolicy = HttpVersionPolicy.RequestVersionExact, + }, + HttpCompletionOption.ResponseHeadersRead).GetAwaiter().GetResult(); + + // Unread for long enough to spend the 64 KiB stream window, so a writer that waits on + // credit is parked by now; then dropped, which is RST_STREAM CANCEL. + Thread.Sleep(700); + endless.Dispose(); + + Assert.True(server.Learned.Task.Wait(5_000), NeverLearned); + Assert.Equal("200: ok", Fetch(client, port, "/ok")); + Assert.True(Volatile.Read(ref server.Connections) == 1, WholeConnection(server)); + }); + } + + foreach (bool native in new[] { false, true }) + { + bool ng = native; + string stack = ng ? "nghttp3" : "ioxide.http3"; + + runner.Test($"h3 peer gone ({stack}): a cancelled stream is reported by FlushAsync, and the connection serves on", () => + { + (string certPath, string keyPath) = TestCert.Ensure(); + using var engine = new QuicEngine(certPath, keyPath, cidLength: 8, alpn: ["h3"]); + var server = new EndlessServer(); + int port = StartH3(ng, engine, server); + + using var client = new H3TestClient("127.0.0.1", port); + client.Connect(); + Assert.True(client.CompleteHandshake(timeoutMs: 5_000), "handshake did not complete"); + + int received = client.RequestThenCancel("/endless", readMs: 500); + Assert.True(received > 0, "none of the endless response arrived before the cancel, so no stream was in flight to stop"); + + // The server learns from the peer's frames, so the connection has to keep turning. + long until = Environment.TickCount64 + 5_000; + while (!server.Learned.Task.IsCompleted && Environment.TickCount64 < until) + { + client.Pump(50); + } + + Assert.True(server.Learned.Task.IsCompleted, NeverLearned); + (int status, string body) = client.Get("/ok", timeoutMs: 5_000); + Assert.True(status == 200 && body == "ok", $"the next request on the connection got {status} [{body}]"); + Assert.True(Volatile.Read(ref server.Connections) == 1, WholeConnection(server)); + }); + } + } + + private const string NeverLearned = + "the handler never learned its stream was abandoned: FlushAsync went on accepting chunks, or stayed " + + "parked waiting for a peer that will never read them"; + + private static string WholeConnection(EndlessServer server) + => $"{server.Connections} connections served requests: the abandoned stream took the whole connection " + + "down with it"; + + /// + /// "/endless" writes 1 KiB every 5 ms until a flush reports the peer gone; anything else answers + /// "ok". Bounded, so a server that never learns cannot outlive the test. + /// + private sealed class EndlessServer + { + public readonly TaskCompletionSource Learned = new(TaskCreationOptions.RunContinuationsAsynchronously); + public int Connections; + + /// One per connection: counts it once it serves a request. + public Func, ValueTask>, Action, ValueTask> For(Reactor reactor) + { + var timer = new RingTimer(reactor); + bool counted = false; + + return async (path, writeAndFlush, headers) => + { + if (!counted) + { + counted = true; + Interlocked.Increment(ref Connections); + } + + headers(200); + if (path != "/endless") + { + await writeAndFlush("ok"u8.ToArray()); + return; + } + + byte[] chunk = new byte[1024]; + for (int i = 0; i < 6_000; i++) + { + if ((await writeAndFlush(chunk)).IsCompleted) + { + Learned.TrySetResult(i); + return; + } + await timer.DelayAsync(5); + } + }; + } + } + + private static int StartH2(bool native, EndlessServer server) + => TestServer.Start(async (reactor, conn) => + { + var serve = server.For(reactor); + try + { + if (native) + { + await new Nghttp2Connection(conn).RunAsync((request, writer) => serve(Path(request.Path), + body => + { + writer.Write(body.Span); + return writer.FlushAsync(); + }, + status => writer.WriteHeaders(new Nghttp2Response { Status = status }))); + } + else + { + await new Http2Connection(conn).RunAsync((request, writer) => serve(Path(request.Path), + body => + { + writer.Write(body.Span); + return writer.FlushAsync(); + }, + status => writer.WriteHeaders(new Http2Response { Status = status }))); + } + } + finally + { + conn.DecRef(); + } + }); + + private static int StartH3(bool native, QuicEngine engine, EndlessServer server) + { + (_, int udpPort) = TestServer.StartDatagram( + onDatagram: null, + quicFactory: engine.CreateFactory(), + quicHandle: (reactor, conn) => + { + var serve = server.For(reactor); + return native + ? new Nghttp3Connection(conn).RunStreamedResponseAsync((request, writer) => serve(Path(request.Path), + body => + { + writer.Write(body.Span); + return writer.FlushAsync(); + }, + status => writer.WriteHeaders(new Nghttp3Response { Status = status }))) + : new Http3Connection(conn).RunStreamedResponseAsync((request, writer) => serve(Path(request.Path), + body => + { + writer.Write(body.Span); + return writer.FlushAsync(); + }, + status => writer.WriteHeaders(new Http3Response { Status = status }))); + }); + return udpPort; + } + + private static string Path(ReadOnlyMemory path) => System.Text.Encoding.ASCII.GetString(path.Span); + + // Cleartext HTTP/2 with prior knowledge, one connection for every request of a test. + private static HttpClient H2Client() => new(new SocketsHttpHandler()) + { + DefaultRequestVersion = HttpVersion.Version20, + DefaultVersionPolicy = HttpVersionPolicy.RequestVersionExact, + Timeout = TimeSpan.FromSeconds(10), + }; + + private static string Fetch(HttpClient client, int port, string path) + { + try + { + using HttpResponseMessage response = client.GetAsync($"http://127.0.0.1:{port}{path}").GetAwaiter().GetResult(); + string body = response.Content.ReadAsStringAsync().GetAwaiter().GetResult(); + return $"{(int)response.StatusCode}: {body}"; + } + catch (HttpRequestException e) + { + return $"reset: {e.GetBaseException().Message}"; + } + } +} diff --git a/tests/Ioxide.Tests.Harness/H3TestClient.cs b/tests/Ioxide.Tests.Harness/H3TestClient.cs index 2a9f62b2..4918f87f 100644 --- a/tests/Ioxide.Tests.Harness/H3TestClient.cs +++ b/tests/Ioxide.Tests.Harness/H3TestClient.cs @@ -154,6 +154,46 @@ public bool CompleteHandshake(int timeoutMs) => Request(method, path, body, extraHeaders: null, timeoutMs); public (int Status, string Body) Request(string method, string path, byte[]? body, (string Name, string Value)[]? extraHeaders, int timeoutMs) + { + Submit(method, path, body, extraHeaders); + Pump(timeoutMs); + + // Status 0 = never answered, which is what a refused connection looks like from here. + return (_status, Encoding.UTF8.GetString(_body.ToArray())); + } + + /// + /// Open a GET and read its response for , then abandon it the way a + /// closed browser tab does: STOP_SENDING and RESET_STREAM, H3_REQUEST_CANCELLED. Returns how + /// many body bytes had arrived by then. + /// + public int RequestThenCancel(string path, int readMs) + { + Submit("GET", path, null, null); + Pump(readMs); + + int received = _body.Count; + Assert.True(iq_conn_shutdown_stream(_conn, _requestSid, H3RequestCancelled) == 0, "shutdown_stream failed"); + _requestSid = -1; // whatever still arrives for it belongs to nobody + FlushOut(); + return received; + } + + /// Keep the connection turning over for , or until the request in hand ends. + public void Pump(int ms) + { + long deadline = Environment.TickCount64 + ms; + while (Environment.TickCount64 < deadline && !_done && !_peerClosed) + { + DrainH3Out(); + FlushOut(); + PumpIn(); + } + } + + private const ulong H3RequestCancelled = 0x010c; + + private void Submit(string method, string path, byte[]? body, (string Name, string Value)[]? extraHeaders) { EnsureH3Session(); @@ -193,17 +233,6 @@ public bool CompleteHandshake(int timeoutMs) pb, (nuint)(body?.Length ?? 0)) == 0, "submit_request failed"); } - - long deadline = Environment.TickCount64 + timeoutMs; - while (Environment.TickCount64 < deadline && !_done && !_peerClosed) - { - DrainH3Out(); - FlushOut(); - PumpIn(); - } - - // Status 0 = never answered, which is what a refused connection looks like from here. - return (_status, Encoding.UTF8.GetString(_body.ToArray())); } /// @@ -531,6 +560,7 @@ private struct Ih3Callbacks [MarshalAs(UnmanagedType.LPUTF8Str)] string? keyPath, IqCallbacks cbs); [DllImport(QuicLib)] private static extern long iq_client_open_bidi(nint conn); [DllImport(QuicLib)] private static extern long iq_conn_open_uni(nint conn); + [DllImport(QuicLib)] private static extern int iq_conn_shutdown_stream(nint conn, long streamId, ulong appErrorCode); [DllImport(QuicLib)] private static extern ulong iq_conn_expiry(nint conn); [DllImport(QuicLib)] private static extern int iq_conn_handle_expiry(nint conn, ulong ts); [DllImport(QuicLib)] private static extern nint iq_conn_write(nint conn, byte* dest, nuint destLen, long streamId, byte* data, nuint dataLen, int fin, long* pConsumed, ulong ts);