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
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
117 changes: 117 additions & 0 deletions tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -287,6 +318,46 @@ private static Func<Reactor, QuicConnection, Task> SignalOnExit(TaskCompletionSo
}
};

/// <summary>
/// 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.
/// </summary>
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<byte>.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
Expand Down Expand Up @@ -328,6 +399,12 @@ internal sealed unsafe class QuicTestClient : IDisposable
private long _streamId = -1;
private readonly List<byte> _echo = [];
private bool _echoFin;
private readonly Dictionary<long, int> _bytesOn = [];
private readonly HashSet<long> _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));
Expand Down Expand Up @@ -444,6 +521,44 @@ public string RequestEcho(byte[] payload, int timeoutMs)

private readonly byte[] _sendScratch = new byte[1452];

/// <summary>Open a stream and send <paramref name="data"/> on it in its own packet, reading nothing.
/// A <paramref name="lost"/> packet is built but never sent, so the next one arrives after a gap.</summary>
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;
}

/// <summary>Read and ack until <paramref name="done"/>.</summary>
public bool PumpUntil(Func<bool> 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()
{
Expand Down Expand Up @@ -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<byte>(data, (int)len).ToArray());
self._bytesOn[streamId] = self.BytesOn(streamId) + (int)len;
if (fin != 0)
{
self._echoFin = true;
self._finOn.Add(streamId);
}
}

Expand Down
Loading