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
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ ioxide hands you raw bytes and stays out of HTTP; when you want a framework on t
`ioxide.Kestrel` swaps the transport under an existing ASP.NET Core app with your endpoints
unchanged.

> Linux 6.1+ · .NET 10 / .NET 11 · `0.4.169` - experimental
> Linux 6.1+ · .NET 10 / .NET 11 · `0.14.236` - experimental

**[Documentation](https://mda2av.github.io/ioxide/)** - architecture, guides, and every example as
runnable code side by side.
Expand Down
2 changes: 1 addition & 1 deletion src/clients/ioxide.file/ioxide.file.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
<RootNamespace>ioxide.file</RootNamespace>

<PackageId>ioxide.file</PackageId>
<Version>0.13.233</Version>
<Version>0.14.236</Version>
<Authors>MDA2AV</Authors>
<Description>File serving for the ioxide io_uring runtime: immutable asset snapshots with baked responses, pooled positional ring reads, atomic reloads.</Description>
<PackageLicenseExpression>MIT</PackageLicenseExpression>
Expand Down
2 changes: 1 addition & 1 deletion src/clients/ioxide.httpclient/ioxide.httpclient.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
<RootNamespace>ioxide.httpclient</RootNamespace>

<PackageId>ioxide.httpclient</PackageId>
<Version>0.13.233</Version>
<Version>0.14.236</Version>
<Authors>MDA2AV</Authors>
<Description>The ring-native HTTP/1.1 client for the ioxide io_uring runtime - the upstream leg between a proxy and an origin. Connections are opened on the reactor thread that will use them, so a request never crosses a thread on its way out or back, and every response resumes the awaiting handler inline on its own reactor. Includes client-side TLS (SNI, ALPN, certificate verification and client certificates for mutual TLS) for https:// origins. Depends on ioxide core alone: no protocol package, no native asset.</Description>
<PackageLicenseExpression>MIT</PackageLicenseExpression>
Expand Down
2 changes: 1 addition & 1 deletion src/clients/ioxide.pg/ioxide.pg.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
<RootNamespace>ioxide.pg</RootNamespace>

<PackageId>ioxide.pg</PackageId>
<Version>0.13.233</Version>
<Version>0.14.236</Version>
<Authors>MDA2AV</Authors>
<Description>Postgres driver for the ioxide io_uring runtime: pooled ring-native connections per reactor, ring-native connect and handshake, inline completion resume.</Description>
<PackageLicenseExpression>MIT</PackageLicenseExpression>
Expand Down
2 changes: 1 addition & 1 deletion src/clients/ioxide.redis/ioxide.redis.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
<RootNamespace>ioxide.redis</RootNamespace>

<PackageId>ioxide.redis</PackageId>
<Version>0.13.233</Version>
<Version>0.14.236</Version>
<Authors>MDA2AV</Authors>
<Description>Redis client for the ioxide io_uring runtime: pooled ring-native connections per reactor, full RESP2 protocol, a generic command API plus typed helpers (strings, keys, hashes, lists, sets, sorted sets, pub/sub, transactions, scripting), and pipelining. Inline completion resume.</Description>
<PackageLicenseExpression>MIT</PackageLicenseExpression>
Expand Down
2 changes: 1 addition & 1 deletion src/clients/ioxide.timer/ioxide.timer.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
<RootNamespace>ioxide.timer</RootNamespace>

<PackageId>ioxide.timer</PackageId>
<Version>0.13.233</Version>
<Version>0.14.236</Version>
<Authors>MDA2AV</Authors>
<Description>Deadlines for the ioxide io_uring runtime: waits submitted to the reactor's own ring with IORING_OP_TIMEOUT, completing inline on the reactor that owns the caller, with no timer thread and nothing allocated per wait.</Description>
<PackageLicenseExpression>MIT</PackageLicenseExpression>
Expand Down
19 changes: 19 additions & 0 deletions src/ioxide/Connection/Tcp/TcpConnection.Write.Flush.cs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,24 @@ public sealed unsafe partial class TcpConnection : IValueTaskSource
private int _flushArmed;
private int _flushInProgress;

/// <summary>
/// The reactor's cached clock at the moment a flush was handed over, read by the sweep against
/// <see cref="TcpOptions.SendTimeoutMs"/>.
///
/// Only meaningful while <see cref="FlushOutstanding"/> - it is written on every arm and never
/// cleared, because clearing it would put a store on CompleteFlush, which is the hottest path
/// in the server, to maintain a value nothing reads in that state.
/// </summary>
internal long FlushArmedMs;

/// <summary>
/// Whether a flush is outstanding - which is what tells the reactor's sweep that this
/// connection is sending rather than idle, so the send clock governs it and the idle one does
/// not. Read rather than <see cref="FlushArmedMs"/> being non-zero, so the stamp never has to
/// double as a flag.
/// </summary>
internal bool FlushOutstanding => Volatile.Read(ref _flushInProgress) != 0;

public ValueTask FlushAsync()
{
if (Volatile.Read(ref _closed) == 1)
Expand Down Expand Up @@ -123,6 +141,7 @@ private ValueTask FlushCore()

_flushSignal.Reset();
WriteInFlight = target;
Volatile.Write(ref FlushArmedMs, _reactor.NowMs);

// A segmented response that spilled past the primary slab is gathered into one SENDMSG.
_flushVectored = _inOverflow;
Expand Down
21 changes: 21 additions & 0 deletions src/ioxide/Connection/Tcp/TcpConnection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,24 @@ public sealed unsafe partial class TcpConnection
/// <summary>The listener port this connection was accepted on; set per accept.</summary>
public ushort ListenerPort { get; internal set; }

/// <summary>
/// Environment.TickCount64 at the last completion this connection saw in either direction -
/// stamped at accept and on every recv and send completion, read by the reactor's sweep
/// (Reactor.Tcp.Sweep.cs) against <see cref="TcpOptions.IdleTimeoutMs"/>.
/// </summary>
/// <remarks>
/// A coarse tick rather than a precise clock on purpose: the sweep runs at ~250 ms and the
/// timeouts it serves are second-scale, so the cheapest read that cannot fall back is enough.
/// Reactor thread only, like the rest of the connection.
/// </remarks>
internal long LastActivityMs;

/// <summary>
/// Set once the sweep has shut this connection down, so the next tick skips it instead of
/// re-issuing shutdown() every 250 ms until the teardown completions land.
/// </summary>
internal bool SweepClosed;

/// <summary>
/// Whether this connection sends with SEND_ZC (zero-copy). Bound at accept from
/// <see cref="TcpOptions.ZeroCopySend"/>; kTLS forces it back to plain via the
Expand Down Expand Up @@ -124,6 +142,9 @@ internal void Clear()
Volatile.Write(ref _closed, 0);
Volatile.Write(ref _flushArmed, 0);
Volatile.Write(ref _flushInProgress, 0);
Volatile.Write(ref FlushArmedMs, 0);
LastActivityMs = _reactor.NowMs;
SweepClosed = false;

WriteHead = 0;
WriteTail = 0;
Expand Down
2 changes: 1 addition & 1 deletion src/ioxide/IoxideRuntime.cs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ namespace ioxide;
/// </summary>
public static class IoxideRuntime
{
public const string Version = "0.0.17";
public const string Version = "0.14.236";

// Wiring (a builder API will eventually wrap this):
// var reactor = new Reactor(id, config); // implements IRingHost
Expand Down
6 changes: 6 additions & 0 deletions src/ioxide/Native/Native.Socket.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,9 @@ public static unsafe partial class Native {
public const int SO_RCVBUF = 8;
public const int SO_REUSEPORT = 15;

/// <summary>shutdown(2) how: both directions.</summary>
public const int SHUT_RDWR = 2;

public const int AF_INET6 = 10;
public const int IPPROTO_IPV6 = 41;
public const int IPV6_V6ONLY = 26;
Expand All @@ -27,6 +30,9 @@ public static unsafe partial class Native {
/// bind to port 0 (QUIC client sockets take an ephemeral port).
[DllImport("libc")] public static extern int getsockname(int fd, void* addr, uint* len);
[DllImport("libc")] public static extern int listen(int fd, int backlog);
/// End a connection at the socket: the peer gets a FIN and any outstanding io_uring recv or
/// send completes, which is what releases a connection the reactor still holds a ref to.
[DllImport("libc")] public static extern int shutdown(int fd, int how);
[DllImport("libc")] public static extern int setsockopt(int fd, int level, int optname, void* optval, uint optlen);
[DllImport("libc")] public static extern int getsockopt(int fd, int level, int optname, void* optval, uint* optlen);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ private void OnSendCompletion(int fd, ushort gen, int res, uint cqeFlags)
return;
}
conn.WriteHead += res;
conn.LastActivityMs = NowMs;

// A zero-copy send posts its data CQE with F_MORE and a notif will follow; hold the slab until
// that notif arrives. Plain SEND never sets F_MORE, so this is a no-op for it.
Expand Down
2 changes: 2 additions & 0 deletions src/ioxide/Reactor/Loop/Reactor.Loop.Incremental.cs
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,8 @@ private void LoopIncremental()
break;
}

NowMs = Environment.TickCount64; // one read per batch; see Reactor.Tcp.Sweep.cs

uint ready = _ring.CqReady();
for (uint i = 0; i < ready; i++)
{
Expand Down
2 changes: 2 additions & 0 deletions src/ioxide/Reactor/Loop/Reactor.Loop.SharedRing.cs
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,8 @@ private void LoopSharedRing()
break;
}

NowMs = Environment.TickCount64; // one read per batch; see Reactor.Tcp.Sweep.cs

uint ready = _ring.CqReady();
for (uint i = 0; i < ready; i++)
{
Expand Down
8 changes: 8 additions & 0 deletions src/ioxide/Reactor/Reactor.Runner.cs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,14 @@ public void Run()
AnnounceListening();
ArmTcpAccepts();
ArmWakePoll();

// After OnStart, so a reactor that serves no TCP - or has both clocks off - registers
// nothing and the sweep costs it not even a table walk.
if (TcpSweepEnabled)
{
AddTicker(TcpSweep);
}

StartTicker();

if (_incremental) LoopIncremental();
Expand Down
2 changes: 2 additions & 0 deletions src/ioxide/Reactor/Reactor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,8 @@ public Reactor(int id, ServerConfig config)
_incRecvBufferSize = (uint)inc.RecvBufferSize;
_pool = new Stack<TcpConnection>(_tcp.PoolMax);
_zeroCopySend = _tcp.ZeroCopySend;
_idleTimeoutMs = _tcp.IdleTimeoutMs;
_sendTimeoutMs = _tcp.SendTimeoutMs;
}

[MethodImpl(MethodImplOptions.AggressiveInlining)]
Expand Down
102 changes: 102 additions & 0 deletions src/ioxide/Reactor/Transport/Tcp/Reactor.Tcp.Sweep.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
using static ioxide.Native;

namespace ioxide;

/// <summary>
/// The clocks on a TCP connection's lifecycle, which until now had none: an idle connection is
/// reaped, and so is one whose send the peer stopped draining.
/// </summary>
/// <remarks>
/// Rides the reactor's existing ~250 ms ticker rather than arming anything of its own, which is
/// what <see cref="ioxide.tls.TlsService"/> does for handshakes and the QUIC transport does for
/// idle connections. Second-scale timeouts do not need better granularity than that, and a
/// connection closes at the first tick past its deadline rather than exactly on it.
/// </remarks>
public sealed unsafe partial class Reactor
{
private readonly int _idleTimeoutMs;
private readonly int _sendTimeoutMs;

/// <summary>
/// Environment.TickCount64, refreshed once per loop pass rather than read per completion.
///
/// The stamps this feeds are read by a sweep that runs four times a second, so a clock good to
/// one batch of completions is far finer than anything that consumes it - while reading the
/// real one per CQE put a vDSO call on both the recv and the send hot path, three per request,
/// and measured as a 4-9% throughput cost on the small-response samples. Refreshing here is one
/// read per io_uring_enter, amortised over the whole batch it returned.
/// </summary>
internal long NowMs = Environment.TickCount64;

private bool TcpSweepEnabled => _tcpEnabled && (_idleTimeoutMs > 0 || _sendTimeoutMs > 0);

/// <summary>
/// One pass over the connection table. Runs on the reactor thread from the ticker, so it owns
/// the table outright and can touch a connection directly.
/// </summary>
private void TcpSweep()
{
long now = Environment.TickCount64;
TcpConnection?[] conns = _connections;

for (int fd = 0; fd < conns.Length; fd++)
{
TcpConnection? conn = conns[fd];
if (conn is null || conn.SweepClosed)
{
continue;
}

// A connection with a flush outstanding is not idle, it is sending - so the send clock
// governs it and the idle one does not apply. Without this split a large response to a
// slow peer would be reaped for making no INBOUND progress while it was working
// perfectly: under MSG_WAITALL the whole flush is a single completion, so nothing
// refreshes the activity stamp for as long as the send legitimately takes.
if (conn.FlushOutstanding)
{
if (_sendTimeoutMs > 0 && now - Volatile.Read(ref conn.FlushArmedMs) > _sendTimeoutMs)
{
TcpSweepClose(conn);
}
continue;
}

if (_idleTimeoutMs > 0 && now - conn.LastActivityMs > _idleTimeoutMs)
{
TcpSweepClose(conn);
}
}
}

/// <summary>
/// End one connection the sweep has condemned - and nothing more than that.
/// </summary>
/// <remarks>
/// Both halves are needed and neither is redundant, exactly as in TlsService.SweepHandshakes.
///
/// shutdown() is what the PEER sees, and it is also what releases the connection: a
/// TcpConnection is held by two refs, the handler's and the reactor's, and the reactor's is
/// given up only when its outstanding operation completes. For an idle connection that is a
/// multishot recv against a peer saying nothing, which otherwise never completes; for a stalled
/// one it is a SEND the peer's closed window is holding, which otherwise never completes
/// either. Shutting the socket down ends both.
///
/// MarkClosed is what wakes the handler NOW - parked on a read, or on the very flush being
/// timed out - with the closed state its loop already knows how to handle, rather than one
/// io_uring round trip later.
///
/// What this deliberately does NOT do is clear the table slot, cancel, or DecRef. The teardown
/// those completions already run (CloseFromRecv, and the send path's res &lt;= 0 branch) is the
/// one that gets the refcount right, and it only runs once the kernel has finished with the
/// connection's slab. Releasing the reactor's ref here instead would let the connection reach
/// zero - and be recycled, with its slab freed or resized - while a SEND the kernel has not
/// given back still points into it.
/// </remarks>
private void TcpSweepClose(TcpConnection conn)
{
conn.SweepClosed = true; // one shutdown per connection, not one per tick until it lands

shutdown(conn.ClientFd, SHUT_RDWR);
conn.MarkClosed();
}
}
7 changes: 7 additions & 0 deletions src/ioxide/Reactor/Transport/Tcp/Reactor.Tcp.cs
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@ private void OnTcpRecvCompletionShared(int fd, ushort gen, int res, uint flags)
// return to the group.
if (conn != null)
{
conn.LastActivityMs = NowMs;
_recvStarved.Add(((ulong)gen << 32) | (uint)fd);
}
return;
Expand Down Expand Up @@ -141,6 +142,8 @@ private void OnTcpRecvCompletionShared(int fd, ushort gen, int res, uint flags)
return;
}

conn.LastActivityMs = NowMs;

byte* ptr = hasBuf ? _bufSlab + (nuint)bid * (nuint)_recvBufferSize : null;
if (!conn.Complete(res, bid, hasBuf, ptr))
{
Expand Down Expand Up @@ -172,6 +175,7 @@ private void OnTcpRecvCompletionIncremental(int fd, ushort gen, int res, uint fl
// still holds buffers (#93). Park; the loop re-arms once a buffer recycles.
if (conn != null)
{
conn.LastActivityMs = NowMs;
_recvStarved.Add(((ulong)gen << 32) | (uint)fd);
}
return;
Expand All @@ -192,6 +196,8 @@ private void OnTcpRecvCompletionIncremental(int fd, ushort gen, int res, uint fl
return; // stale CQE; its ring is already gone
}

conn.LastActivityMs = NowMs;

// Data lands at the buffer's running offset; the kernel keeps appending
// to this bid until the buffer is full (F_BUF_MORE clear).
byte* ptr = conn.BufSlab + (nuint)bid * (nuint)_incRecvBufferSize + (nuint)conn.CumOffset![bid];
Expand Down Expand Up @@ -243,6 +249,7 @@ private void OnTcpAcceptCompletion(int listenFd, int res, bool more)
Track(clientFd, conn);
conn.InitRefs();
conn.ListenerPort = PortOf(listenFd);
conn.LastActivityMs = NowMs; // the idle clock starts at accept

if (_incremental)
{
Expand Down
46 changes: 46 additions & 0 deletions src/ioxide/Reactor/Transport/Tcp/TcpOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -45,4 +45,50 @@ public sealed record TcpOptions

// Per-connection SPSC recv queue depth (power of two); overflow closes the connection.
public int RecvQueueEntries { get; init; } = 64;

/// <summary>
/// Close a connection that has neither received nor sent anything for this long. 0 disables.
///
/// What it defends: a peer that connects and goes quiet holds an fd, a pooled
/// <see cref="TcpConnection"/> with its native write slab, and a recv queue - and in
/// incremental mode a registered buffer ring plus a gid, which is capped, so an idle
/// connection at the cap converts directly into shed accepts.
///
/// Enforced on the reactor's ticker, so the granularity is the tick (~250 ms) and a connection
/// closes at the first tick after its deadline rather than exactly on it.
/// </summary>
/// <remarks>
/// A connection with a flush in flight is NOT idle - it is sending, and
/// <see cref="SendTimeoutMs"/> governs it instead. Otherwise a large response to a slow peer
/// would be reaped for making no INBOUND progress while it was working perfectly.
///
/// The shape to check before deploying this: a protocol that legitimately goes quiet for
/// longer than the timeout in both directions - an idle websocket, a long-poll - is closed by
/// it. Raise it past the protocol's own keep-alive interval, or set 0 and bound those
/// connections some other way.
/// </remarks>
public int IdleTimeoutMs { get; init; } = 60_000;

/// <summary>
/// Close a connection whose flush has been in flight for this long. 0 disables.
///
/// What it defends: a peer that stops reading. Its window shuts, the socket send buffer fills,
/// and the SEND never completes - so <c>FlushAsync</c> parks forever, holding the connection,
/// its slab and the handler's state. TCP will not end it either: a zero window is legitimate
/// and a peer can hold one indefinitely. Nothing else in the stack bounds this.
/// </summary>
/// <remarks>
/// This is the deadline for the WHOLE flush, not for progress within it, because MSG_WAITALL
/// (the default - see <see cref="TcpConnection.SendOpFlags"/>) coalesces a flush into a single
/// completion: there is no per-chunk signal to measure progress against. So set it against the
/// slowest legitimate full response, not against a stall - a large body over a slow link is
/// the false positive to watch for.
///
/// It is deliberately not folded into <see cref="IdleTimeoutMs"/>: a peer that keeps SENDING
/// while it has stopped READING refreshes the idle stamp on every inbound completion, so an
/// idle sweep never fires while that connection's send is wedged. Duplex protocols - a
/// websocket written from a background task is the reported case - need this clock and are not
/// covered by the other one.
/// </remarks>
public int SendTimeoutMs { get; init; } = 60_000;
}
Loading
Loading