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 @@ -12,6 +12,16 @@ public override void Close(ulong applicationErrorCode)
return;
}

// What this cycle queued is packetized first: after CONNECTION_CLOSE nothing more can be.
if (_inEngineCycle)
{
FlushEgress();
if (_closed)
{
return;
}
}

int farewell = 0;
if (_conn != 0)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ public void SendBatched(Action send)
/// </summary>
private void EndEngineCycle()
{
FlushEgress();
_inEngineCycle = false;
FlushGso();
ApplyKeepAlive();
Expand Down Expand Up @@ -114,16 +115,19 @@ private void FlushGso()

// --- engine egress pump -------------------------------------------------------------------

// Replay deferred stream bytes now that the window may have opened, then drain the engine's
// own frames. Runs after every inbound datagram (ACKs open the window) and every timer.
// Packetize every stream with bytes waiting - queued this cycle, or held back until the window
// opened - then drain the engine's own frames. Runs at the end of every engine cycle.
private void FlushEgress()
{
ReplayOut();
FlushConnection();
}

// Acks (processed just before this on the inbound path) freed retention: if a producer
// paused at the high-water, tell it to resume queueing. The read loop never sees these acks,
// so this is the resume trigger for a response larger than the retention window.
// Acks (processed just before this on the inbound path) freed retention: if a producer paused at
// the high-water, tell it to resume queueing. The read loop never sees these acks, so this is
// the resume trigger for a response larger than the retention window.
private void SignalSendCapacity()
{
if (_sendAtCapacity && !_closed && _outRetained < _maxSendRetention)
{
_sendAtCapacity = false;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ private void OnDatagramCore(ReadOnlySpan<byte> payload, byte tos, nint peerAddr,
HandshakeCompletedOnce();
}

FlushEgress();
SignalSendCapacity();
FireRecv();
FireHandshakeSignal();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ public override void OnTimer(long nowMs)
CloseFromEngine(rv);
return;
}
FlushEgress();
SignalSendCapacity();
FireRecv();
FireHandshakeSignal();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,8 @@ private void TakeSendRetention(Reactor reactor)
/// acks drain it (the egress pump re-runs on every inbound datagram).</summary>
public override bool CanQueueSend => !_closed && _outRetained < _maxSendRetention;

// Set when a producer filled retention to the high-water; FlushEgress fires the resume callback
// and clears it once acks bring retention back down.
// Set when a producer filled retention to the high-water; SignalSendCapacity fires the resume
// callback and clears it once acks bring retention back down.
private bool _sendAtCapacity;

/// <summary>Queue bytes on a stream and flush. streamId must come from a delivered item or
Expand Down Expand Up @@ -96,25 +96,41 @@ public override void SendStream(long streamId, ReadOnlySpan<byte> data, bool fin

if (_outRetained >= _maxSendRetention)
{
_sendAtCapacity = true; // arms the resume: FlushEgress fires it once acks drain below
_sendAtCapacity = true; // arms the resume, which SignalSendCapacity fires
}

if (_outRetained > _outRetainedCeiling)
{
Console.Error.WriteLine("[ioxide.ngtcp2] send retention backstop exceeded (producer ignored backpressure); closing connection.");

// What the cycle queued before the flood goes out ahead of the close, as it would have
// had it been packetized on the spot.
if (_inEngineCycle)
{
FlushEgress();
if (_closed)
{
return;
}
}

// Through Teardown, so the peer is told. The abort is ours, not the peer's, so it
// hears INTERNAL_ERROR rather than waiting out a timeout for silence.
Teardown(WriteTransportFarewell(Ngtcp2.NGTCP2_ERR_INTERNAL));
return;
}

PumpOut(streamId, os);
if (_closed)
// Inside an engine cycle nothing reaches the wire before the cycle ends, so the stream is
// packetized then, together with whatever else the cycle gives it.
if (!_inEngineCycle)
{
return;
PumpOut(streamId, os);
if (_closed)
{
return;
}
FlushConnection();
}
FlushConnection();
if (!OutDone(os) && !os.Pending)
{
os.Pending = true;
Expand Down Expand Up @@ -208,27 +224,36 @@ private void ReplayOut()
_outPending.RemoveRange(keep, _outPending.Count - keep);
}

// Chunks handed to one write. A packet carries one STREAM frame, and the frame can span them.
private const int MaxWriteChunks = 8;

// Feed one stream's unsent bytes (pointers into the retained chunks - stable until acked) into
// the engine, sending each produced datagram, until done or the engine can't take more.
private void PumpOut(long sid, OutStream os)
{
Ngtcp2.Vec* chunks = stackalloc Ngtcp2.Vec[MaxWriteChunks];
while (!_closed && !OutDone(os))
{
// Locate the first unsent byte inside the chunk chain.
byte* ptr = null;
int len = 0;
// The unsent bytes, from the first one on, as they lie in the chunk chain.
int count = 0;
long len = 0;
if (os.Sent < os.End)
{
long skip = os.Sent - os.Base;
foreach ((nint p, int l) in os.Chunks)
{
if (skip < l)
if (skip >= l)
{
skip -= l;
continue;
}
chunks[count++] = new Ngtcp2.Vec((byte*)p + skip, (nuint)(l - skip));
len += l - skip;
skip = 0;
if (count == MaxWriteChunks)
{
ptr = (byte*)p + skip;
len = (int)(l - skip);
break;
}
skip -= l;
}
}
bool fin = os.Fin && os.Sent + len == os.End;
Expand All @@ -237,8 +262,8 @@ private void PumpOut(long sid, OutStream os)
nint n;
fixed (byte* dest = _sendBuf)
{
n = Ngtcp2.iq_conn_write(_conn, dest, (nuint)_sendBuf.Length,
sid, ptr, (nuint)len, fin ? 1 : 0, &consumed, NowNs());
n = Ngtcp2.iq_conn_writev(_conn, dest, (nuint)_sendBuf.Length,
sid, count > 0 ? chunks : null, (nuint)count, fin ? 1 : 0, &consumed, NowNs());
}

int code = (int)n;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -291,14 +291,7 @@ public void Pump()
return;
}
_inEngineCycle = true;
try
{
FlushEgress();
}
finally
{
EndEngineCycle();
}
EndEngineCycle(); // where the cycle's egress is packetized and flushed
}

// Single handshake-done funnel (the engine callback and the OnDatagram poll can both detect
Expand Down
14 changes: 13 additions & 1 deletion src/protocols/ioxide.ngtcp2/Interop/Ngtcp2.cs
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,7 @@ internal struct Callbacks
[DllImport(Lib)] internal static extern uint iq_abi();

/// <summary>What this managed binding was written against. Bump both together.</summary>
internal const uint Abi = 5;
internal const uint Abi = 6;

/// <summary>
/// Refuse to run against a native library this binding was not built for.
Expand Down Expand Up @@ -146,6 +146,18 @@ internal static void RequireAbi()
long streamId, byte* data, nuint dataLen, int fin,
long* pConsumed, ulong ts);

/// <summary>ngtcp2_vec.</summary>
internal readonly struct Vec(byte* data, nuint length)
{
public readonly byte* Base = data;
public readonly nuint Len = length;
}

[DllImport(Lib)] internal static extern nint iq_conn_writev(
nint conn, byte* dest, nuint destLen,
long streamId, Vec* data, nuint dataCount, int fin,
long* pConsumed, ulong ts);

[DllImport(Lib)] internal static extern nint iq_conn_close(
nint conn, ulong appErrorCode, byte* dest, nuint destLen, ulong ts);

Expand Down
23 changes: 21 additions & 2 deletions src/protocols/ioxide.ngtcp2/native/ioxide_ngtcp2_shim.c
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
* conn = iq_accept(engine, addrs, first_pkt, ...) validates + creates the server conn
* iq_conn_read(...) feed one UDP datagram
* iq_conn_write(...) produce one UDP datagram (loop until 0)
* iq_conn_writev(...) the same, from several stream buffers
* iq_conn_open_uni / iq_conn_get_alpn H3 plumbing (uni streams, negotiated proto)
* iq_conn_expiry / iq_conn_handle_expiry ns-precision engine deadlines
* iq_conn_free / iq_engine_free
Expand Down Expand Up @@ -43,8 +44,9 @@
* 2 - iq_accept gained shard / shard_count for connection-id steering
* 3 - iq_conn_set_keep_alive added
* 4 - iq_conn_shutdown_stream added
* 5 - iq_client_engine_set_idle_timeout added */
#define IQ_ABI 5
* 5 - iq_client_engine_set_idle_timeout added
* 6 - iq_conn_writev added */
#define IQ_ABI 6

/* ---- callback table into C# ------------------------------------------------------------- */

Expand Down Expand Up @@ -1607,6 +1609,23 @@ ngtcp2_ssize iq_conn_write(iq_conn *c, uint8_t *dest, size_t destlen,
return n;
}

/* iq_conn_write with the stream bytes in several buffers. A packet carries one STREAM frame and the
* frame can span all of them, so a stream's queued writes - a response's headers, body and FIN -
* share a packet instead of taking one each. */
ngtcp2_ssize iq_conn_writev(iq_conn *c, uint8_t *dest, size_t destlen,
int64_t stream_id, const ngtcp2_vec *datav, size_t datavcnt, int fin,
int64_t *pconsumed, uint64_t ts)
{
ngtcp2_ssize consumed = -1;

ngtcp2_ssize n = ngtcp2_conn_writev_stream(
c->conn, &c->path, NULL, dest, destlen, &consumed,
fin ? NGTCP2_WRITE_STREAM_FLAG_FIN : 0, stream_id, datav, datavcnt, ts);

*pconsumed = consumed;
iq_sync_path(c); /* as in iq_conn_write: the destination ngtcp2 chose for this datagram */
return n;
}

/* App-initiated close: record an APPLICATION error (e.g. H3_NO_ERROR for graceful shutdown) and
* build the CONNECTION_CLOSE datagram carrying it. The caller sends it and tears the conn down. */
Expand Down
Binary file not shown.
61 changes: 61 additions & 0 deletions tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,29 @@ public static void Register(Runner runner)
"the first stream never ended: its FIN was counted as sent in an ACK-only packet");
});

runner.Test("quic: writes to a stream in one engine cycle arrive as one STREAM frame", () =>
{
// A streamed response's headers, body and end are three writes made in one pass. They
// were three packets - each write was packetized on the spot - where a buffered response
// is one.
(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: ThreeWritesHandler);

using var client = new QuicTestClient("127.0.0.1", udpPort);
client.Connect();
Assert.True(client.CompleteHandshake(timeoutMs: 5000), "handshake did not complete");

long stream = client.Send("go"u8.ToArray(), fin: true);
Assert.True(client.PumpUntil(() => client.FinOn(stream), timeoutMs: 5000), "the response never ended");
Assert.Equal(12, client.BytesOn(stream));
Assert.Equal(1, client.FramesOn(stream));
});

runner.Test("quic: dual-pipe (PipeReader/PipeWriter) stream echo (loopback)", () =>
{
(string certPath, string keyPath) = TestCert.Ensure();
Expand Down Expand Up @@ -318,6 +341,39 @@ private static Func<Reactor, QuicConnection, Task> SignalOnExit(TaskCompletionSo
}
};

/// <summary>Answers each request with three writes in the same cycle: head, body, then a bare FIN.</summary>
private static async Task ThreeWritesHandler(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)
{
conn.SendStream(item.StreamId, "head"u8, fin: false);
conn.SendStream(item.StreamId, "the body"u8, fin: false);
conn.SendStream(item.StreamId, ReadOnlySpan<byte>.Empty, fin: true);
}
conn.ReturnBuffer(in item);
}

if (snap.IsClosed)
{
break;
}
conn.ResetRead();
}
}
finally
{
conn.DecRef();
}
}

/// <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.
Expand Down Expand Up @@ -400,10 +456,14 @@ internal sealed unsafe class QuicTestClient : IDisposable
private readonly List<byte> _echo = [];
private bool _echoFin;
private readonly Dictionary<long, int> _bytesOn = [];
private readonly Dictionary<long, int> _framesOn = [];
private readonly HashSet<long> _finOn = [];

public int BytesOn(long streamId) => _bytesOn.GetValueOrDefault(streamId);

/// <summary>STREAM frames received on the stream: ngtcp2 delivers each in-order frame on its own.</summary>
public int FramesOn(long streamId) => _framesOn.GetValueOrDefault(streamId);

public bool FinOn(long streamId) => _finOn.Contains(streamId);

private static ulong NowNs() => (ulong)(System.Diagnostics.Stopwatch.GetTimestamp() *
Expand Down Expand Up @@ -610,6 +670,7 @@ 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;
self._framesOn[streamId] = self.FramesOn(streamId) + 1;
if (fin != 0)
{
self._echoFin = true;
Expand Down
Loading