Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions src/protocols/ioxide.http2/Http2Connection.Frames.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}

/// <summary>
Expand Down
23 changes: 8 additions & 15 deletions src/protocols/ioxide.http2/Http2Connection.Streamed.cs
Original file line number Diff line number Diff line change
Expand Up @@ -110,29 +110,22 @@ private Http2ResponseWriter RentWriter(int streamId)
internal void SendStreamedHeaders(int streamId, Http2Response response)
=> WriteHeaders(streamId, response, endStream: false);

/// <summary>False once the peer has reset this response's stream, or the connection is gone.</summary>
internal bool IsResponseLive(int streamId) => !IsBroken && _responseWindows.ContainsKey(streamId);

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

/// <summary>One DATA frame, already known to fit both windows.</summary>
Expand Down
43 changes: 32 additions & 11 deletions src/protocols/ioxide.http2/Http2ResponseWriter.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System.Buffers;
using System.IO.Pipelines;

namespace ioxide.http2;

Expand Down Expand Up @@ -72,7 +73,10 @@ public void WriteHeaders(Http2Response response)
}

_headersSent = true;
_connection.SendStreamedHeaders(_streamId, response);
if (_connection.IsResponseLive(_streamId))
{
_connection.SendStreamedHeaders(_streamId, response);
}
}

/// <inheritdoc />
Expand Down Expand Up @@ -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.
/// </summary>
public ValueTask FlushAsync() => FlushCore(endStream: false);
/// <returns>
/// <see cref="FlushResult.IsCompleted"/> 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.
/// </returns>
public ValueTask<FlushResult> FlushAsync() => FlushCore(endStream: false);

/// <summary>
/// Send what is left and mark END_STREAM. A handler that returns without calling this gets it
Expand Down Expand Up @@ -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<FlushResult> FlushCore(bool endStream)
{
if (!_headersSent)
{
Expand All @@ -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
Expand Down Expand Up @@ -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.
Expand All @@ -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)
Expand Down
5 changes: 5 additions & 0 deletions src/protocols/ioxide.http3/Http3Connection.Streamed.cs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,9 @@ namespace ioxide.http3;
public sealed partial class Http3Connection
{
private readonly Stack<Http3ResponseWriter> _writerPool = new();

// Streamed responses in flight, so a STOP_SENDING from the peer can reach the one it ends.
private readonly Dictionary<long, Http3ResponseWriter> _writers = new();
private readonly List<TaskCompletionSource> _capacityWaiters = [];

/// <summary>True once the connection can no longer make progress; a parked writer gives up.</summary>
Expand Down Expand 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;
}
Expand All @@ -80,6 +84,7 @@ private async Task ServeAsync(Func<Http3Request, Http3ResponseWriter, ValueTask>
}
finally
{
_writers.Remove(writer.StreamId);
_writerPool.Push(writer);
}
}
Expand Down
9 changes: 9 additions & 0 deletions src/protocols/ioxide.http3/Http3Connection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

Expand Down
39 changes: 30 additions & 9 deletions src/protocols/ioxide.http3/Http3ResponseWriter.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System.Buffers;
using System.IO.Pipelines;

namespace ioxide.http3;

Expand Down Expand Up @@ -35,6 +36,7 @@ public sealed class Http3ResponseWriter : IBufferWriter<byte>

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)
{
Expand Down Expand Up @@ -72,7 +74,10 @@ public void WriteHeaders(Http3Response response)
}

_headersSent = true;
_connection.SendStreamedHeaders(_streamId, response);
if (!_gone)
{
_connection.SendStreamedHeaders(_streamId, response);
}
}

/// <inheritdoc />
Expand Down Expand Up @@ -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.
/// </summary>
public ValueTask FlushAsync() => FlushCore(fin: false);
/// <returns>
/// <see cref="FlushResult.IsCompleted"/> 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.
/// </returns>
public ValueTask<FlushResult> FlushAsync() => FlushCore(fin: false);

/// <summary>The peer will read no more of this stream: from here on nothing is sent on it.</summary>
internal void OnPeerGone() => _gone = true;

private static readonly FlushResult PeerGone = new(isCanceled: false, isCompleted: true);

/// <summary>
/// Send what is left and close the stream. A handler that returns without calling this still
Expand Down Expand Up @@ -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)
{
Expand All @@ -156,7 +174,7 @@ internal ValueTask FailAsync()
return ValueTask.CompletedTask;
}

private async ValueTask FlushCore(bool fin)
private async ValueTask<FlushResult> FlushCore(bool fin)
{
if (!_headersSent)
{
Expand All @@ -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<byte>.Empty, fin: true);
return;
return default;
}

// [0x00][varint length][payload], sent as one call - the header is tiny and splitting it
Expand All @@ -200,6 +219,7 @@ private async ValueTask FlushCore(bool fin)
ArrayPool<byte>.Shared.Return(frame);

_staged = 0;
return default;
}

private void EnsureStaging(int sizeHint)
Expand Down Expand Up @@ -233,5 +253,6 @@ internal void Reset(long streamId)
_staged = 0;
_headersSent = false;
_completed = false;
_gone = false;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@ namespace ioxide.nghttp2;
public sealed partial class Nghttp2Connection
{
private readonly Stack<Nghttp2ResponseWriter> _writerPool = new();

// Streamed responses in flight, so a RST_STREAM from the peer can reach the one it ends.
private readonly Dictionary<int, Nghttp2ResponseWriter> _writers = new();
private Func<Nghttp2Request, Nghttp2ResponseWriter, ValueTask>? _streamedHandler;

/// <summary>
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -74,6 +78,7 @@ private async Task ServeStreamedAsync(Func<Nghttp2Request, Nghttp2ResponseWriter
// happen here - the arena backs the request's memories and the writer holds a pooled
// staging buffer.
pending.Dispose();
_writers.Remove(writer.StreamId);
writer.Release();
_writerPool.Push(writer);
}
Expand Down Expand Up @@ -118,20 +123,30 @@ internal unsafe void SendStreamedHeaders(int streamId, Nghttp2Response response)
}

/// <summary>Hand a chunk to nghttp2 and wake the deferred stream. Copied natively.</summary>
internal unsafe void SendStreamedData(int streamId, ReadOnlySpan<byte> body)
/// <returns>False when nghttp2 no longer knows the stream: the peer reset it.</returns>
internal unsafe bool SendStreamedData(int streamId, ReadOnlySpan<byte> 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;
}

/// <summary>A streamed response its handler could not finish: RST_STREAM INTERNAL_ERROR.</summary>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading