diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.SendStream.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.SendStream.cs index 0cbab3b..d8b5640 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.SendStream.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.SendStream.cs @@ -257,8 +257,9 @@ private void PumpOut(long sid, OutStream os) os.FinSent = true; } } - else if (fin && len == 0 && n > 0) + else if (fin && len == 0 && consumed == 0) { + // Not n > 0: with the window full, a due ACK goes out alone and reports -1. os.FinSent = true; // bare-fin frame went out } if (n > 0) diff --git a/tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs b/tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs index 7d10495..0a95dda 100644 --- a/tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs +++ b/tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs @@ -158,6 +158,37 @@ public static void Register(Runner runner) Assert.Equal("answer-then-close", client.RequestEcho(sent, timeoutMs: 5000)); }); + runner.Test("quic: a FIN parked behind a full congestion window is still sent", () => + { + // With the window full and an ACK due, ngtcp2 sends the ACK alone and reports the stream + // frame unwritten (-1). A bare FIN was counted as sent on any packet, so that ACK + // swallowed it and nothing retried it: ~20 in 1.5M streamed h3 responses never ended. + (string certPath, string keyPath) = TestCert.Ensure(); + using var engine = new QuicEngine(certPath, keyPath, cidLength: 8); + + (_, int udpPort) = TestServer.StartDatagram( + onDatagram: null, + quicFactory: engine.CreateFactory(), + quicHandle: FinBehindFullWindowHandler); + + using var client = new QuicTestClient("127.0.0.1", udpPort); + client.Connect(); + Assert.True(client.CompleteHandshake(timeoutMs: 5000), "handshake did not complete"); + + long first = client.Send("a"u8.ToArray(), fin: true); + Assert.True(client.PumpUntil(() => client.BytesOn(first) == 1, timeoutMs: 5000), "the first stream was not answered"); + + // The server fills its window and parks the first stream's FIN behind it. A packet that + // arrives after a lost one makes its ACK due at once (RFC 9000 13.2.1), so the parked FIN + // is the next write - with the window still full, as nothing has been read and acked. + client.Send("fill"u8.ToArray(), fin: true); + client.Send("x"u8.ToArray(), fin: false, lost: true); + client.Send("y"u8.ToArray(), fin: false); + + Assert.True(client.PumpUntil(() => client.FinOn(first), timeoutMs: 5000), + "the first stream never ended: its FIN was counted as sent in an ACK-only packet"); + }); + runner.Test("quic: dual-pipe (PipeReader/PipeWriter) stream echo (loopback)", () => { (string certPath, string keyPath) = TestCert.Ensure(); @@ -287,6 +318,46 @@ private static Func SignalOnExit(TaskCompletionSo } }; + /// + /// Answers stream 0 without ending it. Stream 4 then fills the congestion window and is reset, + /// so that stream 0's bare FIN, sent right after, is the only thing waiting for the window. + /// + private static async Task FinBehindFullWindowHandler(Reactor reactor, QuicConnection conn) + { + try + { + while (true) + { + QuicRecvSnapshot snap = await conn.ReadAsync(); + + while (conn.TryGetDelivery(in snap, out QuicRecvRing.Delivery item)) + { + if (item.Kind == QuicStreamEvent.Data && item.Fin && item.StreamId == 0) + { + conn.SendStream(0, "A"u8, fin: false); + } + else if (item.Kind == QuicStreamEvent.Data && item.Fin && item.StreamId == 4) + { + conn.SendStream(4, new byte[1 << 20], fin: true); + conn.ResetStream(4, 0); + conn.SendStream(0, ReadOnlySpan.Empty, fin: true); + } + conn.ReturnBuffer(in item); + } + + if (snap.IsClosed) + { + break; + } + conn.ResetRead(); + } + } + finally + { + conn.DecRef(); + } + } + private static async Task EchoHandler(Reactor reactor, QuicConnection conn) { try @@ -328,6 +399,12 @@ internal sealed unsafe class QuicTestClient : IDisposable private long _streamId = -1; private readonly List _echo = []; private bool _echoFin; + private readonly Dictionary _bytesOn = []; + private readonly HashSet _finOn = []; + + public int BytesOn(long streamId) => _bytesOn.GetValueOrDefault(streamId); + + public bool FinOn(long streamId) => _finOn.Contains(streamId); private static ulong NowNs() => (ulong)(System.Diagnostics.Stopwatch.GetTimestamp() * (1_000_000_000.0 / System.Diagnostics.Stopwatch.Frequency)); @@ -444,6 +521,44 @@ public string RequestEcho(byte[] payload, int timeoutMs) private readonly byte[] _sendScratch = new byte[1452]; + /// Open a stream and send on it in its own packet, reading nothing. + /// A packet is built but never sent, so the next one arrives after a gap. + public long Send(byte[] data, bool fin, bool lost = false) + { + long streamId = iq_client_open_bidi(_conn); + Assert.True(streamId >= 0, "failed to open client stream"); + + long consumed; + fixed (byte* dest = _sendScratch) + fixed (byte* src = data) + { + nint n = iq_conn_write(_conn, dest, (nuint)_sendScratch.Length, streamId, + src, (nuint)data.Length, fin ? 1 : 0, &consumed, NowNs()); + Assert.True(n > 0 && consumed == data.Length, $"client write failed: {n}, consumed {consumed}"); + if (!lost) + { + _udp.Send(_sendScratch, (int)n); + } + } + return streamId; + } + + /// Read and ack until . + public bool PumpUntil(Func done, int timeoutMs) + { + long deadline = Environment.TickCount64 + timeoutMs; + while (!done()) + { + if (Environment.TickCount64 >= deadline) + { + return false; + } + FlushOut(); + PumpIn(); + } + return true; + } + // Drain the client engine's pending datagrams (ACKs, handshake) to the wire. private void FlushOut() { @@ -494,9 +609,11 @@ private static void OnClientStreamData(void* user, long streamId, byte* data, nu { var self = (QuicTestClient)GCHandle.FromIntPtr((nint)user).Target!; self._echo.AddRange(new ReadOnlySpan(data, (int)len).ToArray()); + self._bytesOn[streamId] = self.BytesOn(streamId) + (int)len; if (fin != 0) { self._echoFin = true; + self._finOn.Add(streamId); } }