Skip to content
Closed
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
15 changes: 14 additions & 1 deletion src/ioxide/Native/Native.IoUring.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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;
}
}
6 changes: 4 additions & 2 deletions src/ioxide/Reactor/Loop/Reactor.Loop.Incremental.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
6 changes: 4 additions & 2 deletions src/ioxide/Reactor/Loop/Reactor.Loop.SharedRing.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
3 changes: 1 addition & 2 deletions src/ioxide/Reactor/Loop/Reactor.Timer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,7 @@ public sealed unsafe partial class Reactor
private readonly List<Action> _tickers = [];

/// <summary>
/// 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.
/// </summary>
public const int TickMs = 250;

Expand Down
3 changes: 2 additions & 1 deletion src/ioxide/Reactor/Reactor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -97,10 +97,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)
{
Expand Down
21 changes: 18 additions & 3 deletions src/ioxide/Reactor/Transport/Quic/Reactor.Quic.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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()
Expand Down Expand Up @@ -362,6 +362,21 @@ internal void QuicArmTimer(QuicConnection conn)
}
}

// The loop's wait, bounded by the earliest engine deadline: otherwise only a completion wakes an
// idle reactor, and the one it can count on is the 250 ms ticker.
private int WaitForCompletions()
{
// No connections left: the tracked deadline is the last one's leftover.
long due = _quicConnSet.Count == 0 ? long.MaxValue : _quicNextTimeoutMs;
if (due == long.MaxValue)
{
return _ring.SubmitAndWait(1);
}

// +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, due - NowMs) + 1);
}

private void TeardownQuic()
{
if (_quicOptions == null && _quicConnSet.Count == 0)
Expand Down
18 changes: 16 additions & 2 deletions src/ioxide/io_uring/Ring.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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);

/// <summary>
/// As <see cref="SubmitAndWait(uint)"/>, with the wait bounded by <paramref name="timeoutMs"/>
/// (negative: unbounded). A wait that runs out with nothing submitted returns -ETIME.
/// </summary>
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)
Expand All @@ -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)]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,8 @@ namespace ioxide.ngtcp2;
/// <summary>
/// A live ngtcp2 server connection, bridging the reactor's QUIC transport to the native engine.
/// Datagrams routed by CID arrive at <see cref="OnDatagram(System.ReadOnlySpan{byte}, byte)"/> and are fed to ngtcp2; the engine's
/// output is flushed back through the transport's <c>Send</c>; loss/idle deadlines ride the reactor
/// ticker via <see cref="GetNextTimeout"/> / <see cref="OnTimer"/>. Everything runs on the owning
/// output is flushed back through the transport's <c>Send</c>; loss/idle deadlines reach the reactor
/// loop via <see cref="GetNextTimeout"/> / <see cref="OnTimer"/>. 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 <see cref="QuicConnection"/>: each decrypted
Expand Down
149 changes: 141 additions & 8 deletions tests/Ioxide.Tests.E2E/Protocols/QuicTimerTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ public static void Register(Runner runner)
RegisterLossRecovery(runner);
RegisterTimerFaults(runner);
RegisterOutboundTimers(runner);
RegisterIdleReactor(runner);
}

// --- the expiry arithmetic itself -------------------------------------------------------
Expand All @@ -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.
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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");
Expand All @@ -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)
Expand All @@ -570,6 +663,46 @@ public override void OnTimer(long nowMs)
}
}

/// <summary>
/// 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.
/// </summary>
private sealed class TimerLatenessConnection(QuicEngine engine) : QuicEngineConnection(engine)
{
private static readonly List<double> 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);
}
}

/// <summary>An engine connection whose deadlines are still armed but never acted on.</summary>
private sealed class SuppressedTimerConnection(QuicEngine engine) : QuicEngineConnection(engine)
{
Expand Down
Loading