From 7ed9d29281978c3eaaf03ef4f0792fb4dcb0feef Mon Sep 17 00:00:00 2001 From: Diogo Martins Date: Sat, 3 Oct 2026 19:10:25 +0100 Subject: [PATCH] reactor: the wait's timeout replaces the 250 ms timer, and wakes for a due QUIC deadline The loop's wait is bounded by the next ticker run and the earliest engine deadline, so an idle reactor fires a retransmit, ACK or pacing step when it is due. The tickers run from the loop once their interval is up - checked every pass, since a busy reactor's wait never runs out - and the IORING_OP_TIMEOUT that used to drive them is gone. --- src/ioxide/Native/Native.IoUring.cs | 15 +- .../Loop/Reactor.Loop.DispatchCompletions.cs | 20 --- .../Reactor/Loop/Reactor.Loop.Incremental.cs | 12 +- .../Reactor/Loop/Reactor.Loop.SharedRing.cs | 12 +- src/ioxide/Reactor/Loop/Reactor.Timer.cs | 60 ++++--- src/ioxide/Reactor/Reactor.Runner.cs | 5 - src/ioxide/Reactor/Reactor.cs | 4 +- .../Reactor/Transport/Quic/Reactor.Quic.cs | 6 +- src/ioxide/io_uring/Ring.cs | 18 ++- .../QuicEngineConnection.KeepAlive.cs | 3 +- .../Connection/QuicEngineConnection.cs | 4 +- .../Ioxide.Tests.E2E/Core/TcpTimeoutTests.cs | 62 ++++++++ .../Protocols/QuicTimerTests.cs | 149 +++++++++++++++++- 13 files changed, 290 insertions(+), 80 deletions(-) diff --git a/src/ioxide/Native/Native.IoUring.cs b/src/ioxide/Native/Native.IoUring.cs index a2c01ed7..d35afc00 100644 --- a/src/ioxide/Native/Native.IoUring.cs +++ b/src/ioxide/Native/Native.IoUring.cs @@ -31,6 +31,7 @@ public static unsafe partial class Native { // CQE (IORING_CQE_F_NOTIF) once the buffer can be reused. public const byte IORING_OP_SEND_ZC = 47; public const uint IORING_ENTER_GETEVENTS = 1u << 0; + public const uint IORING_ENTER_EXT_ARG = 1u << 3; // the enter's arg is an io_uring_getevents_arg: a wait with a timeout public const long IORING_OFF_SQ_RING = 0; public const long IORING_OFF_SQES = 0x10000000; @@ -99,8 +100,11 @@ public static int io_uring_setup(uint entries, IoUringParams* p) } public static int io_uring_enter(int fd, uint toSubmit, uint minComplete, uint flags) + => io_uring_enter(fd, toSubmit, minComplete, flags, null, 0); + + public static int io_uring_enter(int fd, uint toSubmit, uint minComplete, uint flags, void* arg, nuint argSize) { - long rc = syscall6(SYS_IO_URING_ENTER, (uint)fd, toSubmit, minComplete, flags, null, 0); + long rc = syscall6(SYS_IO_URING_ENTER, (uint)fd, toSubmit, minComplete, flags, arg, argSize); return rc < 0 ? -Marshal.GetLastPInvokeError() : (int)rc; } @@ -178,4 +182,13 @@ public struct io_uring_buf_reg { // include/uapi/linux/time_types.h - the timespec IORING_OP_TIMEOUT reads. [StructLayout(LayoutKind.Sequential)] public struct __kernel_timespec { public long tv_sec; public long tv_nsec; } + + // io_uring_enter's argument under IORING_ENTER_EXT_ARG; ts points at a __kernel_timespec. + [StructLayout(LayoutKind.Sequential)] + public struct io_uring_getevents_arg { + public ulong sigmask; + public uint sigmask_sz; + public uint min_wait_usec; + public ulong ts; + } } diff --git a/src/ioxide/Reactor/Loop/Reactor.Loop.DispatchCompletions.cs b/src/ioxide/Reactor/Loop/Reactor.Loop.DispatchCompletions.cs index 73222fad..72adb52f 100644 --- a/src/ioxide/Reactor/Loop/Reactor.Loop.DispatchCompletions.cs +++ b/src/ioxide/Reactor/Loop/Reactor.Loop.DispatchCompletions.cs @@ -101,26 +101,6 @@ private void OnWakeCompletion(bool more) #endregion -#region Timer - - private void OnTimerTick() - { - for (int i = 0; i < _tickers.Count; i++) - { - try - { - _tickers[i](); - } - catch (Exception e) - { - Console.Error.WriteLine($"[r{_id}] ticker faulted: {e.Message}"); - } - } - ArmTimer(); // single-shot timer; re-arm for the next interval - } - -#endregion - #region Client private void OnClientCompletion(int slot, int result) diff --git a/src/ioxide/Reactor/Loop/Reactor.Loop.Incremental.cs b/src/ioxide/Reactor/Loop/Reactor.Loop.Incremental.cs index ca734f21..9ef288e8 100644 --- a/src/ioxide/Reactor/Loop/Reactor.Loop.Incremental.cs +++ b/src/ioxide/Reactor/Loop/Reactor.Loop.Incremental.cs @@ -185,8 +185,10 @@ private void LoopIncremental() // Anything else is a lifecycle or programming error. Throwing rather than breaking, // because a reactor that vanishes while the process reports healthy is the worst of // both; whether the process then dies is the host's call, via Reactor.OnFault. - int rc = _ring.SubmitAndWait(1); - if (rc < 0 && rc != -EINTR && rc != -EAGAIN && rc != -EBUSY) + // + // ETIME is not an error at all: the wait's own bound ran out (WaitForCompletions). + int rc = WaitForCompletions(); + if (rc < 0 && rc != -EINTR && rc != -EAGAIN && rc != -EBUSY && rc != -ETIME) { throw new InvalidOperationException( $"[r{_id}] io_uring_enter failed with errno {-rc}; this reactor cannot continue"); @@ -200,6 +202,8 @@ private void LoopIncremental() DispatchIncremental(in _ring.CqeAt(i)); } _ring.CqAdvance(ready); + + RunDueTickers(); } } @@ -240,10 +244,6 @@ private void DispatchIncremental(in IoUringCqe cqe) OnWakeCompletion(more); return; - case KindTimer: - OnTimerTick(); - return; - case KindCancel: return; } diff --git a/src/ioxide/Reactor/Loop/Reactor.Loop.SharedRing.cs b/src/ioxide/Reactor/Loop/Reactor.Loop.SharedRing.cs index 4fc12827..28c669b8 100644 --- a/src/ioxide/Reactor/Loop/Reactor.Loop.SharedRing.cs +++ b/src/ioxide/Reactor/Loop/Reactor.Loop.SharedRing.cs @@ -62,8 +62,10 @@ private void LoopSharedRing() // Anything else is a lifecycle or programming error. Throwing rather than breaking, // because a reactor that vanishes while the process reports healthy is the worst of // both; whether the process then dies is the host's call, via Reactor.OnFault. - int rc = _ring.SubmitAndWait(1); - if (rc < 0 && rc != -EINTR && rc != -EAGAIN && rc != -EBUSY) + // + // ETIME is not an error at all: the wait's own bound ran out (WaitForCompletions). + int rc = WaitForCompletions(); + if (rc < 0 && rc != -EINTR && rc != -EAGAIN && rc != -EBUSY && rc != -ETIME) { throw new InvalidOperationException( $"[r{_id}] io_uring_enter failed with errno {-rc}; this reactor cannot continue"); @@ -77,6 +79,8 @@ private void LoopSharedRing() DispatchSharedRing(in _ring.CqeAt(i)); } _ring.CqAdvance(ready); + + RunDueTickers(); } } @@ -120,10 +124,6 @@ private void DispatchSharedRing(in IoUringCqe cqe) OnWakeCompletion(more); return; - case KindTimer: - OnTimerTick(); - return; - case KindCancel: return; } diff --git a/src/ioxide/Reactor/Loop/Reactor.Timer.cs b/src/ioxide/Reactor/Loop/Reactor.Timer.cs index fa00a886..72217de7 100644 --- a/src/ioxide/Reactor/Loop/Reactor.Timer.cs +++ b/src/ioxide/Reactor/Loop/Reactor.Timer.cs @@ -1,20 +1,14 @@ -using System.Runtime.CompilerServices; -using System.Runtime.InteropServices; -using static ioxide.Native; - namespace ioxide; public sealed unsafe partial class Reactor { - // Periodic timer driving registered tickers (per-command timeout sweeps, pool replenishment). - // Single-shot, re-armed each fire. One timer in flight per reactor. - private __kernel_timespec* _timerTs; - private const long TimerIntervalNs = TickMs * 1_000_000L; + // Registered tickers (per-command timeout sweeps, pool replenishment), run by the loop every + // TickMs. No timer of their own: the loop's wait is bounded by the next run. + private long _nextTickMs; private readonly List _tickers = []; /// - /// The ticker's interval: the granularity of every sweep, and on an otherwise idle reactor of - /// QUIC engine timers too, which fire only when the loop wakes. + /// The ticker's interval: the granularity of every sweep. /// public const int TickMs = 250; @@ -24,23 +18,43 @@ public sealed unsafe partial class Reactor /// public void AddTicker(Action ticker) => _tickers.Add(ticker); - // Armed before the loop starts; the timespec is freed in Teardown, after the ring fd closes. + // The first run is one interval after the loop starts. private void StartTicker() { - _timerTs = (__kernel_timespec*)NativeMemory.Alloc((nuint)sizeof(__kernel_timespec)); - _timerTs->tv_sec = 0; - _timerTs->tv_nsec = TimerIntervalNs; - ArmTimer(); + NowMs = Environment.TickCount64; + _nextTickMs = NowMs + TickMs; + } + + // Checked on every pass, not left to the wait running out: a busy reactor's wait never does. + private void RunDueTickers() + { + if (NowMs < _nextTickMs) + { + return; + } + _nextTickMs = NowMs + TickMs; + + for (int i = 0; i < _tickers.Count; i++) + { + try + { + _tickers[i](); + } + catch (Exception e) + { + Console.Error.WriteLine($"[r{_id}] ticker faulted: {e.Message}"); + } + } } - private void ArmTimer() + // The loop's wait, bounded by the next ticker run and the earliest engine deadline. + private int WaitForCompletions() { - IoUringSqe* sqe = GetSqeOrFlush(); - Unsafe.InitBlockUnaligned(sqe, 0, 64); - sqe->opcode = IORING_OP_TIMEOUT; - sqe->addr = (ulong)_timerTs; - sqe->len = 1; - sqe->off = 0; // pure time-based (no completion-count trigger) - sqe->user_data = Tag(KindTimer, 0, 0); + // No QUIC connections left: the tracked deadline is the last one's leftover. + long quicDue = _quicConnSet.Count == 0 ? long.MaxValue : _quicNextTimeoutMs; + long wakeAt = Math.Min(_nextTickMs, quicDue); + + // +1: the ms clock can read just short of the deadline when the kernel's timer wakes us. + return _ring.SubmitAndWait(1, Math.Max(0, wakeAt - NowMs) + 1); } } diff --git a/src/ioxide/Reactor/Reactor.Runner.cs b/src/ioxide/Reactor/Reactor.Runner.cs index b4cbdcfd..44cee610 100644 --- a/src/ioxide/Reactor/Reactor.Runner.cs +++ b/src/ioxide/Reactor/Reactor.Runner.cs @@ -129,11 +129,6 @@ private void Teardown() } close(wakeFd); } - if (_timerTs != null) - { - NativeMemory.Free(_timerTs); - _timerTs = null; - } if (_opTimespecs != null) { NativeMemory.Free(_opTimespecs); diff --git a/src/ioxide/Reactor/Reactor.cs b/src/ioxide/Reactor/Reactor.cs index c6e39fd0..2bd31a75 100644 --- a/src/ioxide/Reactor/Reactor.cs +++ b/src/ioxide/Reactor/Reactor.cs @@ -44,7 +44,6 @@ public sealed unsafe partial class Reactor private const byte KindWake = 4; private const byte KindClient = 5; // low 32 bits = op slot (Reactor.RingHost.cs) private const byte KindCancel = 6; - private const byte KindTimer = 7; private const byte KindUdpRecv = 8; // low 32 bits = recv-slot index (Reactor.Udp.cs) private const byte KindUdpSend = 9; // low 32 bits = send-slot index (Reactor.Udp.cs) @@ -97,10 +96,11 @@ internal void ReturnWriteSlab(nint p) private readonly int _connBufRingEntries; private readonly uint _incRecvBufferSize; - // Transient io_uring_enter errnos. + // Transient io_uring_enter errnos, and ETIME: a bounded wait that ran out. private const int EINTR = 4; private const int EAGAIN = 11; private const int EBUSY = 16; + private const int ETIME = 62; public Reactor(int id, ServerConfig config) { diff --git a/src/ioxide/Reactor/Transport/Quic/Reactor.Quic.cs b/src/ioxide/Reactor/Transport/Quic/Reactor.Quic.cs index 00d96b65..dd89433e 100644 --- a/src/ioxide/Reactor/Transport/Quic/Reactor.Quic.cs +++ b/src/ioxide/Reactor/Transport/Quic/Reactor.Quic.cs @@ -270,7 +270,7 @@ public void QuicRemoveConnection(QuicConnection conn) } // Ticker callback (~250 ms): evict quiet connections. Engine deadlines are fired by - // QuicFireDueTimers at loop-pass granularity; this ticker's loop wake doubles as its floor. + // QuicFireDueTimers, and bound the loop's wait (WaitForCompletions). private void QuicSweep() { long now = Environment.TickCount64; @@ -294,8 +294,8 @@ private void QuicSweep() } // Earliest engine deadline across live conns; long.MaxValue = none. Checked at the top of every - // loop pass, so loss/PTO timers fire at completion-batch granularity (~RTT under load) instead - // of the 250 ms ticker - a retransmit that waits 250 ms per loss makes storms self-sustaining. + // loop pass and bounding the wait, so loss/PTO timers fire when due instead of on the 250 ms + // ticker - a retransmit that waits 250 ms per loss makes storms self-sustaining. private long _quicNextTimeoutMs = long.MaxValue; private void QuicFireDueTimers() diff --git a/src/ioxide/io_uring/Ring.cs b/src/ioxide/io_uring/Ring.cs index cc296b87..4ed9211e 100644 --- a/src/ioxide/io_uring/Ring.cs +++ b/src/ioxide/io_uring/Ring.cs @@ -164,7 +164,13 @@ public static Ring Create(uint entries) return &_sqes[slot]; } - public int SubmitAndWait(uint waitFor) + public int SubmitAndWait(uint waitFor) => SubmitAndWait(waitFor, -1); + + /// + /// As , with the wait bounded by + /// (negative: unbounded). A wait that runs out with nothing submitted returns -ETIME. + /// + public int SubmitAndWait(uint waitFor, long timeoutMs) { // liburing-style accounting: derive the submit count from the kernel-consumed head, so // SQEs published by an enter that consumed nothing (-EBUSY under CQ-overflow pressure) @@ -181,7 +187,15 @@ public int SubmitAndWait(uint waitFor) uint flags = waitFor > 0 ? IORING_ENTER_GETEVENTS : 0; - return io_uring_enter(_fd, toSubmit, waitFor, flags); + if (timeoutMs < 0) + { + return io_uring_enter(_fd, toSubmit, waitFor, flags); + } + + // The kernel copies both during the call, so the stack is fine. + var ts = new __kernel_timespec { tv_sec = timeoutMs / 1000, tv_nsec = timeoutMs % 1000 * 1_000_000 }; + var arg = new io_uring_getevents_arg { ts = (ulong)&ts }; + return io_uring_enter(_fd, toSubmit, waitFor, flags | IORING_ENTER_EXT_ARG, &arg, (nuint)sizeof(io_uring_getevents_arg)); } [MethodImpl(MethodImplOptions.AggressiveInlining)] diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.KeepAlive.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.KeepAlive.cs index 3efbc205..67990e59 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.KeepAlive.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.KeepAlive.cs @@ -37,8 +37,7 @@ private void ApplyKeepAlive() _keepAliveOn = want; _keepAliveStale = false; - // The sweep reaps on a tick, and an idle reactor sends the ping - and the ACK before it, - // which restarts ngtcp2's keep-alive clock - only on a tick. Keep two ticks clear. + // The sweep reaps on a tick; keep two ticks clear for the ping and the peer's ACK of it. int readMs = Math.Max(0, _reactor.QuicReadTimeoutMs); int boundMs = readMs - Math.Min(2 * Reactor.TickMs, readMs / 2); Ngtcp2.iq_conn_set_keep_alive(_conn, want ? 1 : 0, (ulong)boundMs * 1_000_000UL); diff --git a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.cs b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.cs index e4ce8b7f..42ae9dd2 100644 --- a/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.cs +++ b/src/protocols/ioxide.ngtcp2/Connection/QuicEngineConnection.cs @@ -7,8 +7,8 @@ namespace ioxide.ngtcp2; /// /// A live ngtcp2 server connection, bridging the reactor's QUIC transport to the native engine. /// Datagrams routed by CID arrive at and are fed to ngtcp2; the engine's -/// output is flushed back through the transport's Send; loss/idle deadlines ride the reactor -/// ticker via / . Everything runs on the owning +/// output is flushed back through the transport's Send; loss/idle deadlines reach the reactor +/// loop via / . Everything runs on the owning /// reactor thread, so the whole connection - transport half and engine half - is single-threaded. /// /// Application bytes flow through the read surface on : each decrypted diff --git a/tests/Ioxide.Tests.E2E/Core/TcpTimeoutTests.cs b/tests/Ioxide.Tests.E2E/Core/TcpTimeoutTests.cs index ba2dd8cb..92bc6a35 100644 --- a/tests/Ioxide.Tests.E2E/Core/TcpTimeoutTests.cs +++ b/tests/Ioxide.Tests.E2E/Core/TcpTimeoutTests.cs @@ -92,6 +92,68 @@ public static void Register(Runner runner) Assert.Equal(2, stream.Read(reply, 0, 2)); }); + runner.Test("tcp/read: a quiet connection is reaped while another keeps the reactor busy", () => + { + // A busy reactor's wait never runs out - a completion always ends it first - so a sweep + // that waited for an idle moment would leave every quiet connection open for as long as + // anyone else was talking. + int accepted = 0; + var quietClosed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + int port = StartWith(readMs: 500, sendMs: 0, async (reactor, conn) => + { + if (!await IsThisTestsConnection(conn)) + { + conn.DecRef(); + return; + } + + if (Interlocked.Increment(ref accepted) > 1) + { + conn.Write("ok"u8); + await conn.FlushAsync(); + conn.ResetRead(); + await EchoHandler(reactor, conn); // the busy one; releases its own ref + return; + } + + try + { + conn.ResetRead(); + RecvSnapshot snapshot = await conn.ReadAsync(); + quietClosed.TrySetResult(snapshot.IsClosed); + } + finally + { + conn.DecRef(); + } + }); + + using var quiet = new TcpClient(); + quiet.Connect("127.0.0.1", port); + quiet.GetStream().Write("hi"u8); + SpinWait.SpinUntil(() => Volatile.Read(ref accepted) == 1, 2_000); + + using var busy = new TcpClient(); + busy.Connect("127.0.0.1", port); + busy.ReceiveTimeout = 2_000; + NetworkStream stream = busy.GetStream(); + var reply = new byte[2]; + int exchanges = 0; + long until = Environment.TickCount64 + SweepGraceMs; + while (!quietClosed.Task.IsCompleted && Environment.TickCount64 < until) + { + stream.Write("ping"u8); + Assert.Equal(2, stream.Read(reply, 0, 2)); + exchanges++; + } + + Assert.True(quietClosed.Task.IsCompleted, + $"the quiet connection outlived its 500 ms read timeout by {SweepGraceMs} ms while another " + + $"made {exchanges} exchanges: the sweep never ran on the busy reactor"); + Assert.True(quietClosed.Task.Result, "the handler woke, but not with a closed snapshot"); + }); + runner.Test("tcp/read: 0 disables the clock", () => { int port = StartWith(readMs: 0, sendMs: 0, EchoHandler); diff --git a/tests/Ioxide.Tests.E2E/Protocols/QuicTimerTests.cs b/tests/Ioxide.Tests.E2E/Protocols/QuicTimerTests.cs index 5475683a..86edf14a 100644 --- a/tests/Ioxide.Tests.E2E/Protocols/QuicTimerTests.cs +++ b/tests/Ioxide.Tests.E2E/Protocols/QuicTimerTests.cs @@ -27,6 +27,7 @@ public static void Register(Runner runner) RegisterLossRecovery(runner); RegisterTimerFaults(runner); RegisterOutboundTimers(runner); + RegisterIdleReactor(runner); } // --- the expiry arithmetic itself ------------------------------------------------------- @@ -37,8 +38,8 @@ private static void RegisterExpiryArithmetic(Runner runner) { // GetNextTimeout converts ngtcp2's ns expiry into the sweep's TickCount64 ms frame by // subtracting the current ns clock. Both are UNSIGNED, and by the time the loop looks, - // the expiry has normally already passed - the loop parks in io_uring until a - // completion arrives, so it reads the clock milliseconds late, every time. Unguarded, + // the expiry has normally already passed - the loop wakes for it at millisecond + // granularity, or for some other completion, so it reads the clock after it. Unguarded, // that subtraction underflows and the connection's next deadline comes back roughly 584 // years out. It does not look like a crash: the connection works perfectly until the // first packet is lost, and then never recovers, for the rest of its life. @@ -72,16 +73,15 @@ private static void RegisterExpiryArithmetic(Runner runner) // Keep it turning over until the transport has actually been asked for a deadline that // had already passed - the state the guard exists for, and the only state in which this // test means anything. It happens within a second; the budget is generous only so a - // loaded machine cannot make it flaky. "Already passed" carries 2 ms of margin, so the - // answer cannot turn on the microseconds between the probe's sample and the engine's - // own re-read inside GetNextTimeout. + // loaded machine cannot make it flaky. Any amount past counts: the clock only moves + // forward, so an expiry the probe saw pass has passed for GetNextTimeout's re-read too. TimerProbeConnection.Observation[] alreadyDue; long deadline = Environment.TickCount64 + 20_000; do { client.Pump(100); alreadyDue = TimerProbeConnection.Snapshot() - .Where(o => o.ExpiryNs != ulong.MaxValue && o.ExpiryNs + 2_000_000 <= o.NowNs) + .Where(o => o.ExpiryNs != ulong.MaxValue && o.ExpiryNs < o.NowNs) .ToArray(); } while (alreadyDue.Length == 0 && Environment.TickCount64 < deadline); @@ -427,6 +427,99 @@ private static void RegisterOutboundTimers(Runner runner) }); } + // --- an idle reactor --------------------------------------------------------------------- + + private static void RegisterIdleReactor(Runner runner) + { + runner.Test("quic/timer: an idle reactor fires an engine deadline when it is due, not at the next tick", () => + { + // Nothing wakes an idle reactor but a completion, and the only one it can count on is the + // 250 ms ticker - so a retransmit due ~25 ms after a loss went out at the next tick. Here + // the answer is lost and the peer stays silent, so the retransmit timer fires again and + // again with backoff, each firing timed against the engine's own deadline. + (string certPath, string keyPath) = TestCert.Ensure(); + using var engine = new QuicEngine(certPath, keyPath, cidLength: 8); + var answered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + (_, int udpPort) = TestServer.StartDatagram( + onDatagram: null, + quicFactory: engine.CreateFactory(r => new TimerLatenessConnection(engine)), + quicHandle: RingDelayedEchoHandler(300, answered)); + + using var client = new LossyQuicClient(udpPort); + client.Connect(); + Assert.True(client.CompleteHandshake(5_000), "handshake did not complete"); + + client.Pump(500); + client.SendRequest(Encoding.ASCII.GetBytes("on-time")); + client.Pump(100); + + // From here the reactor has nothing to do but the answer, then its retransmits. + TimerLatenessConnection.Reset(); + client.Blackout(3_000); + + Assert.True(answered.Task.IsCompleted, "the handler had not answered by the end of the blackout"); + double[] lateMs = TimerLatenessConnection.Snapshot(); + Assert.True(lateMs.Length >= 3, + $"the engine timer fired {lateMs.Length} times in the blackout, too few to judge its timing"); + Assert.True(lateMs.Max() < 50, + $"engine deadlines fired up to {lateMs.Max():F0} ms late " + + $"([{string.Join(", ", lateMs.Select(l => l.ToString("F0")))}] ms): the reactor slept past " + + "them until the next tick"); + }); + + runner.Test("quic/timer: a reactor whose last connection is gone does not keep waking for it", () => + { + // The reactor tracks only the earliest engine deadline, and nothing resets it when the last + // connection leaves: the loop stops looking once there are none. A wait bounded by that + // leftover would wake an empty reactor every millisecond, for good. + (string certPath, string keyPath) = TestCert.Ensure(); + using var engine = new QuicEngine(certPath, keyPath, cidLength: 8); + + (_, int udpPort) = TestServer.StartDatagram( + onDatagram: null, + quicFactory: engine.CreateFactory(), + quicReadMs: 1_000, + quicHandle: EchoHandler); + + using (var client = new LossyQuicClient(udpPort)) + { + client.Connect(); + Assert.True(client.CompleteHandshake(5_000), "handshake did not complete"); + + // The answer is lost and the peer goes silent, so the connection is evicted for its read + // timeout with a retransmit still scheduled - a deadline that then passes. + client.SendRequest(Encoding.ASCII.GetBytes("leave-a-deadline")); + client.Blackout(3_000); + Assert.True(client.DroppedInbound > 0, "the server never answered, so it scheduled no retransmit"); + } + + string reactorThread = $"test-reactor-udp-{udpPort}"; + long before = VoluntarySwitches(reactorThread); + Thread.Sleep(1_000); + long wakes = VoluntarySwitches(reactorThread) - before; + + Assert.True(wakes < 50, + $"a reactor with no connections woke {wakes} times in a second; its ticker accounts for 4"); + }); + } + + // How often a thread has parked: each io_uring wait that sleeps is one voluntary switch. + private static long VoluntarySwitches(string threadName) + { + string comm = threadName.Length > 15 ? threadName[..15] : threadName; // the kernel keeps 15 bytes + foreach (string task in Directory.GetDirectories("/proc/self/task")) + { + if (File.ReadAllText(Path.Combine(task, "comm")).TrimEnd('\n') != comm) + { + continue; + } + string line = File.ReadLines(Path.Combine(task, "status")).First(l => l.StartsWith("voluntary_ctxt_switches:")); + return long.Parse(line["voluntary_ctxt_switches:".Length..].Trim()); + } + throw new InvalidOperationException($"no thread named {comm}"); + } + // --- handler ----------------------------------------------------------------------------- private static async Task EchoHandler(Reactor reactor, QuicConnection conn) @@ -519,7 +612,7 @@ private sealed class TimerProbeConnection(QuicEngine engine) : QuicEngineConnect { internal readonly record struct Observation(ulong ExpiryNs, ulong NowNs, long NowMs, long Deadline); - private static readonly System.Reflection.FieldInfo ConnHandle = + internal static readonly System.Reflection.FieldInfo ConnHandle = typeof(QuicEngineConnection).GetField("_conn", System.Reflection.BindingFlags.NonPublic | System.Reflection.BindingFlags.Instance) ?? throw new Exception("could not reflect QuicEngineConnection._conn"); @@ -545,7 +638,7 @@ public static Observation[] Snapshot() } // Same clock and the same call the engine makes; reactor thread, like everything else here. - private static ulong NowNs() => (ulong)(System.Diagnostics.Stopwatch.GetTimestamp() * + internal static ulong NowNs() => (ulong)(System.Diagnostics.Stopwatch.GetTimestamp() * (1_000_000_000.0 / System.Diagnostics.Stopwatch.Frequency)); public override long GetNextTimeout(long nowMs) @@ -570,6 +663,46 @@ public override void OnTimer(long nowMs) } } + /// + /// Records how late each firing was: the engine's own ns expiry against the ns clock, on entry to + /// OnTimer. A firing whose expiry is still ahead (the ms deadline rounds down) is not a late one. + /// + private sealed class TimerLatenessConnection(QuicEngine engine) : QuicEngineConnection(engine) + { + private static readonly List LateMs = []; + + public static void Reset() + { + lock (LateMs) + { + LateMs.Clear(); + } + } + + public static double[] Snapshot() + { + lock (LateMs) + { + return LateMs.ToArray(); + } + } + + public override void OnTimer(long nowMs) + { + nint conn = (nint)TimerProbeConnection.ConnHandle.GetValue(this)!; + ulong expiry = conn == 0 ? ulong.MaxValue : LossyQuicClient.Expiry(conn); + ulong now = TimerProbeConnection.NowNs(); + if (expiry <= now) + { + lock (LateMs) + { + LateMs.Add((now - expiry) / 1_000_000.0); + } + } + base.OnTimer(nowMs); + } + } + /// An engine connection whose deadlines are still armed but never acted on. private sealed class SuppressedTimerConnection(QuicEngine engine) : QuicEngineConnection(engine) {