From a7eebae8c0660bb5f15b6eb455b5c7025bb7efaf Mon Sep 17 00:00:00 2001 From: Diogo Martins Date: Sun, 4 Oct 2026 14:38:33 +0100 Subject: [PATCH 1/3] http2: a streamed response learns from FlushAsync that its peer reset the stream, on both stacks FlushAsync returns a FlushResult, IsCompleted once the peer has reset the stream or the connection is gone - what a TCP pipe writer already reports for a closed peer. ioxide.http2 kept the reset stream's window, so it went on sending DATA after RST_STREAM and a writer out of credit parked until the connection ended. RST_STREAM now drops the window and wakes the writer, and nothing more is sent on the stream. nghttp2 forgets a reset stream, so the next write failed - and failed the connection, with every other stream on it. That write now ends only its own stream. --- .../ioxide.http2/Http2Connection.Frames.cs | 7 + .../ioxide.http2/Http2Connection.Streamed.cs | 3 + .../ioxide.http2/Http2ResponseWriter.cs | 41 +++-- .../Connection/Nghttp2Connection.Callbacks.cs | 4 + .../Connection/Nghttp2Connection.Streamed.cs | 21 ++- .../Connection/Nghttp2Connection.cs | 1 + .../Http/Nghttp2ResponseWriter.cs | 46 +++++- tests/Ioxide.Tests.E2E/Program.cs | 1 + .../Protocols/StreamedPeerGoneTests.cs | 145 ++++++++++++++++++ 9 files changed, 248 insertions(+), 21 deletions(-) create mode 100644 tests/Ioxide.Tests.E2E/Protocols/StreamedPeerGoneTests.cs 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..6b4636cd 100644 --- a/src/protocols/ioxide.http2/Http2Connection.Streamed.cs +++ b/src/protocols/ioxide.http2/Http2Connection.Streamed.cs @@ -110,6 +110,9 @@ 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. diff --git a/src/protocols/ioxide.http2/Http2ResponseWriter.cs b/src/protocols/ioxide.http2/Http2ResponseWriter.cs index 30d1587e..fc29d524 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) { @@ -166,16 +180,17 @@ private async ValueTask FlushCore(bool endStream) int sent = 0; while (sent < _staged) { + if (!_connection.IsResponseLive(_streamId)) + { + _staged = 0; // reset by the peer, or the connection is gone: nobody to send it to + return PeerGone; + } + // 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 (_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 +222,11 @@ private async ValueTask FlushCore(bool endStream) _staged = 0; _sinceRealFlush += sent; + if (!_connection.IsResponseLive(_streamId)) + { + return PeerGone; // nothing was staged, or no stream is left to end + } + 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.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/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..43b3a8e1 --- /dev/null +++ b/tests/Ioxide.Tests.E2E/Protocols/StreamedPeerGoneTests.cs @@ -0,0 +1,145 @@ +using System.Buffers; +using System.IO.Pipelines; +using System.Net; +using ioxide; +using ioxide.http2; +using ioxide.nghttp2; +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. +/// +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 learned = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + int connections = 0; + int port = StartH2(ng, learned, () => Interlocked.Increment(ref connections)); + + 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(learned.Task.Wait(5_000), + "the handler never learned its stream was reset: FlushAsync went on accepting chunks, " + + "or stayed parked on credit the peer will never send"); + Assert.Equal("200: ok", Fetch(client, port, "/ok")); + Assert.True(Volatile.Read(ref connections) == 1, + $"{connections} connections served requests: the reset took the whole connection down " + + "with it, not just its own stream"); + }); + } + } + + // "/endless" writes 1 KiB every 5 ms until a flush reports the peer gone; "/ok" answers at once. + // Bounded, so a server that never learns cannot outlive the test. + private static int StartH2(bool native, TaskCompletionSource learned, Action onFirstRequest) + => TestServer.Start(async (reactor, conn) => + { + var timer = new RingTimer(reactor); + bool counted = false; + + async ValueTask Serve(string path, Func, ValueTask> writeAndFlush, Action headers) + { + if (!counted) + { + counted = true; + onFirstRequest(); + } + + 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); + } + } + + try + { + if (native) + { + await new Nghttp2Connection(conn).RunAsync((request, writer) => Serve( + System.Text.Encoding.ASCII.GetString(request.Path.Span), + body => + { + writer.Write(body.Span); + return writer.FlushAsync(); + }, + status => writer.WriteHeaders(new Nghttp2Response { Status = status }))); + } + else + { + await new Http2Connection(conn).RunAsync((request, writer) => Serve( + System.Text.Encoding.ASCII.GetString(request.Path.Span), + body => + { + writer.Write(body.Span); + return writer.FlushAsync(); + }, + status => writer.WriteHeaders(new Http2Response { Status = status }))); + } + } + finally + { + conn.DecRef(); + } + }); + + // 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}"; + } + } +} From 239063e3b3f144fd6df5394d630c9b7d6727ffba Mon Sep 17 00:00:00 2001 From: Diogo Martins Date: Sun, 4 Oct 2026 14:42:29 +0100 Subject: [PATCH 2/3] http3: a streamed response learns from FlushAsync that its peer stopped reading, on both stacks FlushAsync returns a FlushResult here too, IsCompleted once the peer has sent STOP_SENDING (or the stream closed under the handler) or the connection is gone. nghttp3 never pulls a stopped stream again, so a writer waiting for its chunk to be taken looped on PumpAsync - synchronously outside a read pass, which pinned the reactor thread at 100% and starved every connection on it. A stopped stream now ends that wait. The writer is also no longer pooled when its stream closes under a running handler: the next request on the connection could rent it while the old handler still wrote into it. It goes back when the second of its two owners lets go. ioxide.http3 kept sending into a stopped stream and never told the handler; a parked writer now wakes to find out. --- .../ioxide.http3/Http3Connection.Streamed.cs | 5 + src/protocols/ioxide.http3/Http3Connection.cs | 9 ++ .../ioxide.http3/Http3ResponseWriter.cs | 39 ++++-- .../Nghttp3Connection.PushToEngine.cs | 4 + .../Connection/Nghttp3Connection.Streamed.cs | 31 +++++ .../Http/Nghttp3ResponseWriter.cs | 39 +++++- .../Protocols/StreamedPeerGoneTests.cs | 123 ++++++++++++++---- tests/Ioxide.Tests.Harness/H3TestClient.cs | 52 ++++++-- 8 files changed, 254 insertions(+), 48 deletions(-) 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.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/Protocols/StreamedPeerGoneTests.cs b/tests/Ioxide.Tests.E2E/Protocols/StreamedPeerGoneTests.cs index 43b3a8e1..d373f503 100644 --- a/tests/Ioxide.Tests.E2E/Protocols/StreamedPeerGoneTests.cs +++ b/tests/Ioxide.Tests.E2E/Protocols/StreamedPeerGoneTests.cs @@ -3,7 +3,10 @@ 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; @@ -13,7 +16,7 @@ namespace Ioxide.Tests; /// 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. +/// same way, on every stack. /// internal static class StreamedPeerGoneTests { @@ -26,9 +29,8 @@ public static void Register(Runner runner) runner.Test($"h2 peer gone ({stack}): a reset stream is reported by FlushAsync, and the connection serves on", () => { - var learned = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - int connections = 0; - int port = StartH2(ng, learned, () => Interlocked.Increment(ref connections)); + var server = new EndlessServer(); + int port = StartH2(ng, server); using HttpClient client = H2Client(); HttpResponseMessage endless = client.SendAsync( @@ -44,31 +46,75 @@ public static void Register(Runner runner) Thread.Sleep(700); endless.Dispose(); - Assert.True(learned.Task.Wait(5_000), - "the handler never learned its stream was reset: FlushAsync went on accepting chunks, " - + "or stayed parked on credit the peer will never send"); + Assert.True(server.Learned.Task.Wait(5_000), NeverLearned); Assert.Equal("200: ok", Fetch(client, port, "/ok")); - Assert.True(Volatile.Read(ref connections) == 1, - $"{connections} connections served requests: the reset took the whole connection down " - + "with it, not just its own stream"); + 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)); }); } } - // "/endless" writes 1 KiB every 5 ms until a flush reports the peer gone; "/ok" answers at once. - // Bounded, so a server that never learns cannot outlive the test. - private static int StartH2(bool native, TaskCompletionSource learned, Action onFirstRequest) - => TestServer.Start(async (reactor, conn) => + 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; - async ValueTask Serve(string path, Func, ValueTask> writeAndFlush, Action headers) + return async (path, writeAndFlush, headers) => { if (!counted) { counted = true; - onFirstRequest(); + Interlocked.Increment(ref Connections); } headers(200); @@ -83,19 +129,24 @@ async ValueTask Serve(string path, Func, ValueTask TestServer.Start(async (reactor, conn) => + { + var serve = server.For(reactor); try { if (native) { - await new Nghttp2Connection(conn).RunAsync((request, writer) => Serve( - System.Text.Encoding.ASCII.GetString(request.Path.Span), + await new Nghttp2Connection(conn).RunAsync((request, writer) => serve(Path(request.Path), body => { writer.Write(body.Span); @@ -105,8 +156,7 @@ async ValueTask Serve(string path, Func, ValueTask Serve( - System.Text.Encoding.ASCII.GetString(request.Path.Span), + await new Http2Connection(conn).RunAsync((request, writer) => serve(Path(request.Path), body => { writer.Write(body.Span); @@ -121,6 +171,35 @@ async ValueTask Serve(string path, Func, ValueTask + { + 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()) { 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); From c869a9572944c85edef5e3e80b82e4888f2ff0c2 Mon Sep 17 00:00:00 2001 From: Diogo Martins Date: Sun, 4 Oct 2026 15:37:47 +0100 Subject: [PATCH 3/3] http2: read the stream's liveness from the credit lookup the flush already makes SendCredit's window lookup now doubles as the check: negative means the peer reset the stream or the connection is gone. The flush loop no longer pays a second dictionary lookup per DATA chunk for it - measured at -1.7% on a streamed response of 8 one-KiB chunks. --- .../ioxide.http2/Http2Connection.Streamed.cs | 20 +++++-------------- .../ioxide.http2/Http2ResponseWriter.cs | 16 +++++++-------- 2 files changed, 13 insertions(+), 23 deletions(-) diff --git a/src/protocols/ioxide.http2/Http2Connection.Streamed.cs b/src/protocols/ioxide.http2/Http2Connection.Streamed.cs index 6b4636cd..346942b6 100644 --- a/src/protocols/ioxide.http2/Http2Connection.Streamed.cs +++ b/src/protocols/ioxide.http2/Http2Connection.Streamed.cs @@ -115,27 +115,17 @@ internal void SendStreamedHeaders(int streamId, Http2Response response) /// /// 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 fc29d524..684248ba 100644 --- a/src/protocols/ioxide.http2/Http2ResponseWriter.cs +++ b/src/protocols/ioxide.http2/Http2ResponseWriter.cs @@ -180,16 +180,15 @@ private async ValueTask FlushCore(bool endStream) int sent = 0; while (sent < _staged) { - if (!_connection.IsResponseLive(_streamId)) + // 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) { _staged = 0; // reset by the peer, or the connection is gone: nobody to send it to return PeerGone; } - - // 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) { // 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 @@ -222,9 +221,10 @@ private async ValueTask FlushCore(bool endStream) _staged = 0; _sinceRealFlush += sent; - if (!_connection.IsResponseLive(_streamId)) + // Sending proved the stream live; with nothing sent it still has to be asked. + if (sent == 0 && !_connection.IsResponseLive(_streamId)) { - return PeerGone; // nothing was staged, or no stream is left to end + return PeerGone; } if (endStream && sent == 0)