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
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;
}
}
20 changes: 0 additions & 20 deletions src/ioxide/Reactor/Loop/Reactor.Loop.DispatchCompletions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
12 changes: 6 additions & 6 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 All @@ -200,6 +202,8 @@ private void LoopIncremental()
DispatchIncremental(in _ring.CqeAt(i));
}
_ring.CqAdvance(ready);

RunDueTickers();
}
}

Expand Down Expand Up @@ -240,10 +244,6 @@ private void DispatchIncremental(in IoUringCqe cqe)
OnWakeCompletion(more);
return;

case KindTimer:
OnTimerTick();
return;

case KindCancel:
return;
}
Expand Down
12 changes: 6 additions & 6 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 All @@ -77,6 +79,8 @@ private void LoopSharedRing()
DispatchSharedRing(in _ring.CqeAt(i));
}
_ring.CqAdvance(ready);

RunDueTickers();
}
}

Expand Down Expand Up @@ -120,10 +124,6 @@ private void DispatchSharedRing(in IoUringCqe cqe)
OnWakeCompletion(more);
return;

case KindTimer:
OnTimerTick();
return;

case KindCancel:
return;
}
Expand Down
60 changes: 37 additions & 23 deletions src/ioxide/Reactor/Loop/Reactor.Timer.cs
Original file line number Diff line number Diff line change
@@ -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<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 All @@ -24,23 +18,43 @@ public sealed unsafe partial class Reactor
/// </summary>
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);
}
}
5 changes: 0 additions & 5 deletions src/ioxide/Reactor/Reactor.Runner.cs
Original file line number Diff line number Diff line change
Expand Up @@ -129,11 +129,6 @@ private void Teardown()
}
close(wakeFd);
}
if (_timerTs != null)
{
NativeMemory.Free(_timerTs);
_timerTs = null;
}
if (_opTimespecs != null)
{
NativeMemory.Free(_opTimespecs);
Expand Down
4 changes: 2 additions & 2 deletions src/ioxide/Reactor/Reactor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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)
{
Expand Down
6 changes: 3 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
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
62 changes: 62 additions & 0 deletions tests/Ioxide.Tests.E2E/Core/TcpTimeoutTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<bool>(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);
Expand Down
Loading
Loading