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);