diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.Close.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.Close.cs index 038c366..92f6b1f 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.Close.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.Close.cs @@ -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) { diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.Egress.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.Egress.cs index 43cd7d7..1c4cda4 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.Egress.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.Egress.cs @@ -54,6 +54,7 @@ public void SendBatched(Action send) /// private void EndEngineCycle() { + FlushEgress(); _inEngineCycle = false; FlushGso(); ApplyKeepAlive(); @@ -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; diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.OnDatagram.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.OnDatagram.cs index c3e49cb..83505cd 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.OnDatagram.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.OnDatagram.cs @@ -50,7 +50,7 @@ private void OnDatagramCore(ReadOnlySpan payload, byte tos, nint peerAddr, HandshakeCompletedOnce(); } - FlushEgress(); + SignalSendCapacity(); FireRecv(); FireHandshakeSignal(); } diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.OnTimer.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.OnTimer.cs index 06a8473..964dd09 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.OnTimer.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.OnTimer.cs @@ -43,7 +43,7 @@ public override void OnTimer(long nowMs) CloseFromEngine(rv); return; } - FlushEgress(); + SignalSendCapacity(); FireRecv(); FireHandshakeSignal(); } diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.SendStream.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.SendStream.cs index 4db7f60..c78ad59 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.SendStream.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.SendStream.cs @@ -50,8 +50,8 @@ private void TakeSendRetention(Reactor reactor) /// acks drain it (the egress pump re-runs on every inbound datagram). 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; /// Queue bytes on a stream and flush. streamId must come from a delivered item or @@ -96,25 +96,41 @@ public override void SendStream(long streamId, ReadOnlySpan 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; @@ -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; @@ -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; diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.cs index 193ae42..644607a 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.cs @@ -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 diff --git a/src/protocols/ioxide.ngtcp2/Interop/Ngtcp2.cs b/src/protocols/ioxide.ngtcp2/Interop/Ngtcp2.cs index 979b889..30b72ed 100644 --- a/src/protocols/ioxide.ngtcp2/Interop/Ngtcp2.cs +++ b/src/protocols/ioxide.ngtcp2/Interop/Ngtcp2.cs @@ -96,7 +96,7 @@ internal struct Callbacks [DllImport(Lib)] internal static extern uint iq_abi(); /// What this managed binding was written against. Bump both together. - internal const uint Abi = 5; + internal const uint Abi = 6; /// /// Refuse to run against a native library this binding was not built for. @@ -146,6 +146,18 @@ internal static void RequireAbi() long streamId, byte* data, nuint dataLen, int fin, long* pConsumed, ulong ts); + /// ngtcp2_vec. + 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); diff --git a/src/protocols/ioxide.ngtcp2/native/ioxide_ngtcp2_shim.c b/src/protocols/ioxide.ngtcp2/native/ioxide_ngtcp2_shim.c index 1d46f87..93fb221 100644 --- a/src/protocols/ioxide.ngtcp2/native/ioxide_ngtcp2_shim.c +++ b/src/protocols/ioxide.ngtcp2/native/ioxide_ngtcp2_shim.c @@ -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 @@ -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# ------------------------------------------------------------- */ @@ -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. */ diff --git a/src/protocols/ioxide.ngtcp2/runtimes/linux-x64/native/libioxide_ngtcp2.so b/src/protocols/ioxide.ngtcp2/runtimes/linux-x64/native/libioxide_ngtcp2.so index bd3258e..1ad1e8a 100755 Binary files a/src/protocols/ioxide.ngtcp2/runtimes/linux-x64/native/libioxide_ngtcp2.so and b/src/protocols/ioxide.ngtcp2/runtimes/linux-x64/native/libioxide_ngtcp2.so differ diff --git a/tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs b/tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs index 0a95dda..0dbd0ec 100644 --- a/tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs +++ b/tests/Ioxide.Tests.E2E/Protocols/QuicEngineTests.cs @@ -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(); @@ -318,6 +341,39 @@ private static Func SignalOnExit(TaskCompletionSo } }; + /// Answers each request with three writes in the same cycle: head, body, then a bare FIN. + 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.Empty, fin: true); + } + conn.ReturnBuffer(in item); + } + + if (snap.IsClosed) + { + break; + } + conn.ResetRead(); + } + } + finally + { + conn.DecRef(); + } + } + /// /// 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. @@ -400,10 +456,14 @@ internal sealed unsafe class QuicTestClient : IDisposable private readonly List _echo = []; private bool _echoFin; private readonly Dictionary _bytesOn = []; + private readonly Dictionary _framesOn = []; private readonly HashSet _finOn = []; public int BytesOn(long streamId) => _bytesOn.GetValueOrDefault(streamId); + /// STREAM frames received on the stream: ngtcp2 delivers each in-order frame on its own. + 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() * @@ -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(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;