## Problem
`RunEventLoop` span through its body regardless of whether there was anything to do — ~10% of a desktop core for an empty shard, and ~70% of a core on a 3 vCPU VPS. A process that never idles is exactly what burstable vCPU plans throttle, which is how this surfaced: lag spikes that went away when the operator bought more cores. The spin also denied the GC its natural pause points, so memory climbed until a world save forced a collection — alarming in task manager, harmless in practice, and a recurring source of "is my server leaking?" reports.
## Result
Windows desktop, real world of **190,728 items / 33,158 mobiles**, no players, saves and prebake off, three consecutive runs:
| | Legacy spin | Idle sleeping |
|---|---|---|
| **CPU** | 10.42 – 10.50% of one core | **0.78 – 1.00%** |
| **Tick lag** (peak/15s) | 4–10 ms | 5–11 ms |
**~10× less CPU with tick lag unchanged** — the CPU came free rather than being traded for latency. Slower hosts gain proportionally more. Spin mode (`server.eventLoopIdleWaitMs=0`) independently gained **7× the iterations per core** (1.19M → 8.3M cycles/sec) from the ring's AcceptEx rework.
## How
The loop blocks in `NetState.WaitForCompletion` whenever every queue it drains is empty (all the drains are bounded, so leftovers keep it awake). Receive completions, new connections, and cross-thread `LoopContext.Post` (via the ring's sticky `Wake()`) are all in the wait set, so sleeping adds no latency to any of them. Only timer-driven logic sees wheel lag, bounded by the idle wait.
**Health is measured at the only place sleeping can cause harm.** A sleep is bounded by the time to the next wheel turn, so a correctly honoured sleep can never miss a deadline — the only failure mode is the host returning the wait late. That overshoot is measured on every sleep (one extra timestamp read; production's entire accounting cost), and an escalating backoff suspends sleeping when it persists. By construction, server work — saves, heavy staff commands, deep timer callbacks — cannot trip it, so the warning means exactly one thing: *the host is not scheduling the process promptly*, with two known remedies (dedicated CPU, or `=0`). Hosts with no high-resolution wait mechanism at all are detected once at startup and spin instead.
**CPS is removed.** `Core.CyclesPerSecond`/`AverageCPS` measured nothing actionable before and became actively misleading once the loop sleeps (the rate is set by the sleep, not by shard health). The admin gump's Performance page now shows the verdict instead: `Healthy` / `Sleep suspended (host)` / `Spinning (configured)`.
## Configuration
| Setting | Default | Meaning |
|---|---|---|
| `server.eventLoopIdleWaitMs` | `2` | Longest idle block. Measured across 1/2/4/8 ms, 2 is where the trade stops being free. `0` = never sleep: ~98% of a core, zero scheduling overhead — for large shards on dedicated CPU. |
| `server.lateWakeThreshold` | `1` | Idle waits the host may return a full tick late, per second, before sleeping backs off. Raise for jittery hosts; very high disables the backoff. |
## Diagnostics (compiled out by default)
`dotnet build -p:EventLoopProfiling=true` compiles in `EventLoopProfiler` — every hook is `[Conditional("EVENT_LOOP_PROFILING")]`, so normal builds contain zero profiling IL. The profiling build decomposes each second of wall time into **work (per loop phase) / sleep / GC pause / stolen residual**, keeps ~15 minutes of history in a ring buffer, and the `[LoopStats` command prints the last minute and dumps the full history to CSV. `dev-docs/debugging-event-loop.md` is the diagnosis guide (for humans and AI): what production already tells you, when to flip the profiling build, the signature table for host-steal vs deep-processing vs GC vs wake bugs, why dotnet-trace comes last, and the GC/RAM "leak" misconception.
## Verification
- 815 Server.Tests green; both build configurations compile.
- Docker echo harness green on epoll and io_uring (ping-pong mode); kqueue verified manually on an M1 Max.
- A/B measurements and per-change numbers: `measure/event-loop` branch.
## Notes
The full measurement harness and vendored ring sources used to develop this live on the [`measure/event-loop`](https://github.com/modernuo/ModernUO/tree/measure/event-loop) branch, kept for future loop work.
624 lines
22 KiB
C#
624 lines
22 KiB
C#
/*************************************************************************
|
|
* ModernUO *
|
|
* Copyright 2019-2026 - ModernUO Development Team *
|
|
* Email: hi@modernuo.com *
|
|
* File: NetState.Network.cs *
|
|
* *
|
|
* This program is free software: you can redistribute it and/or modify *
|
|
* it under the terms of the GNU General Public License as published by *
|
|
* the Free Software Foundation, either version 3 of the License, or *
|
|
* (at your option) any later version. *
|
|
* *
|
|
* You should have received a copy of the GNU General Public License *
|
|
* along with this program. If not, see <http://www.gnu.org/licenses/>. *
|
|
*************************************************************************/
|
|
|
|
using System;
|
|
using System.Collections.Generic;
|
|
using System.Linq;
|
|
using System.Net;
|
|
using System.Net.NetworkInformation;
|
|
using System.Network;
|
|
using System.Numerics;
|
|
|
|
namespace Server.Network;
|
|
|
|
/// <summary>
|
|
/// Network infrastructure for IORingGroup-based socket I/O.
|
|
/// </summary>
|
|
public partial class NetState
|
|
{
|
|
// Buffer sizes
|
|
private const int RecvBufferSize = 1024 * 64; // 64KB recv buffers
|
|
private const int DefaultSendBufferSize = 1024 * 256; // 256KB send buffers
|
|
private const int MinSendBufferSize = 1024 * 64; // Platform allocation granularity
|
|
private const int MaxConnections = 4096; // Max concurrent connections
|
|
|
|
private static readonly Queue<NetState> _disposed = [];
|
|
private static readonly TimeSpan ConnectingSocketIdleLimit = TimeSpan.FromMilliseconds(5000); // 5 seconds
|
|
|
|
// Socket manager handles buffer pools, socket lifecycle, and I/O operations
|
|
private static RingSocketManager _socketManager;
|
|
|
|
// NetState storage indexed by RingSocket.Id
|
|
private static readonly NetState[] _netStates = new NetState[MaxConnections];
|
|
|
|
// Events buffer for ProcessCompletions. Bounded by one event per peeked completion
|
|
// (maxSockets), doubled for headroom. Undersizing drops DataReceived events whose bytes were
|
|
// already committed, leaving them unparsed until the next recv completes.
|
|
private static readonly RingSocketEvent[] _events = new RingSocketEvent[MaxConnections * 2];
|
|
|
|
// Listener management
|
|
private static nint[] _listeners = Array.Empty<nint>();
|
|
private static int _pendingAcceptCount;
|
|
private const int PendingAcceptsPerListener = 32;
|
|
|
|
private const long AliveCheckIntervalMs = 5000;
|
|
private static long _nextAliveCheck;
|
|
|
|
/// <summary>
|
|
/// Gets the IORingGroup instance for socket operations.
|
|
/// </summary>
|
|
public static IIORingGroup Ring => _socketManager?.Ring;
|
|
|
|
/// <summary>
|
|
/// Waits for network I/O completions or until the specified timeout expires.
|
|
/// Used by the game loop to sleep efficiently while remaining responsive to network events.
|
|
/// </summary>
|
|
/// <param name="timeoutMs">Maximum time to wait in milliseconds.</param>
|
|
public static void WaitForCompletion(int timeoutMs)
|
|
{
|
|
_socketManager?.WaitForCompletion(timeoutMs);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Wakes the game loop if it is blocked in <see cref="WaitForCompletion"/>. Safe from any
|
|
/// thread; a no-op before networking is configured or after teardown. The signal is sticky,
|
|
/// so a wake racing the loop's decision to sleep is not lost.
|
|
/// </summary>
|
|
public static void Wake()
|
|
{
|
|
_socketManager?.Ring?.Wake();
|
|
}
|
|
|
|
/// <summary>
|
|
/// True when no queued network work remains for the loop to drain. <see cref="Slice"/> defers
|
|
/// work in several places, so an empty completion queue alone is not enough.
|
|
/// </summary>
|
|
internal static bool IsIdle =>
|
|
_throttled.Count == 0 && _throttledPending.Count == 0 &&
|
|
_flushPending.Count == 0 && _disposed.Count == 0;
|
|
|
|
/// <summary>
|
|
/// Gets the listening addresses that the server is bound to.
|
|
/// </summary>
|
|
public static IPEndPoint[] ListeningAddresses { get; private set; }
|
|
|
|
private static IPRateLimiter _ipRateLimiter;
|
|
|
|
/// <summary>
|
|
/// Configures the IORingGroup and socket manager.
|
|
/// </summary>
|
|
private static void ConfigureNetwork()
|
|
{
|
|
// Skip if already configured
|
|
if (_socketManager != null)
|
|
{
|
|
return;
|
|
}
|
|
|
|
// Initialize IP rate limiter
|
|
_ipRateLimiter = new IPRateLimiter(10, 10000, 1000, 2.0, 3_600_000, Core.ClosingTokenSource.Token);
|
|
|
|
// Sends in flight per connection; honoured by RIO only (see IIORingGroup). Costs a
|
|
// request-queue and completion-queue slot per send, not another buffer. Worst-case added
|
|
// latency is roughly completion RTT / this value.
|
|
var maxOutstandingSends = ServerConfiguration.GetOrUpdateSetting("network.maxOutstandingSends", 32);
|
|
|
|
// Initialize IORingGroup
|
|
var ring = IORingGroup.Create(
|
|
queueSize: MaxConnections * 2,
|
|
maxConnections: MaxConnections,
|
|
maxOutstandingSends: maxOutstandingSends
|
|
);
|
|
|
|
// Per-connection send buffer: the lever for "send buffer exhausted" disconnects, and the
|
|
// per-connection memory ceiling.
|
|
var sendBufferSize = GetSendBufferSize();
|
|
|
|
// Create socket manager which handles buffer pools and socket lifecycle
|
|
_socketManager = new RingSocketManager(
|
|
ring,
|
|
maxSockets: MaxConnections,
|
|
recvBufferSize: RecvBufferSize,
|
|
sendBufferSize: sendBufferSize,
|
|
initialBufferSlabs: 8,
|
|
maxBufferSlabs: 32
|
|
);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Reads the configured send buffer size, coerced to a power of two of at least the platform
|
|
/// allocation granularity. IORingBuffer requires this and would otherwise throw at socket
|
|
/// creation rather than at startup.
|
|
/// </summary>
|
|
private static int GetSendBufferSize()
|
|
{
|
|
var configured = ServerConfiguration.GetOrUpdateSetting("network.sendBufferSize", DefaultSendBufferSize);
|
|
var size = Math.Max(MinSendBufferSize, configured);
|
|
|
|
if (!BitOperations.IsPow2(size))
|
|
{
|
|
size = (int)BitOperations.RoundUpToPowerOf2((uint)size);
|
|
}
|
|
|
|
if (size != configured)
|
|
{
|
|
logger.Warning(
|
|
"network.sendBufferSize {Configured} is not a power of two of at least {Minimum}; using {Adjusted}",
|
|
configured,
|
|
MinSendBufferSize,
|
|
size
|
|
);
|
|
}
|
|
|
|
return size;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Starts the network server on configured listening addresses.
|
|
/// </summary>
|
|
public static void Start()
|
|
{
|
|
HashSet<IPEndPoint> listeningAddresses = [];
|
|
List<nint> listeners = [];
|
|
|
|
var ring = _socketManager.Ring;
|
|
for (var i = 0; i < ServerConfiguration.Listeners.Count; i++)
|
|
{
|
|
var ipep = ServerConfiguration.Listeners[i];
|
|
var listener = ring.CreateListener(ipep.Address.ToString(), (ushort)ipep.Port, 256);
|
|
if (listener == -1)
|
|
{
|
|
logger.Warning("Failed to create listener for {Address}", ipep);
|
|
continue;
|
|
}
|
|
|
|
if (ipep.Address.Equals(IPAddress.Any) || ipep.Address.Equals(IPAddress.IPv6Any))
|
|
{
|
|
listeningAddresses.UnionWith(GetListeningAddresses(ipep));
|
|
}
|
|
else
|
|
{
|
|
listeningAddresses.Add(ipep);
|
|
}
|
|
|
|
listeners.Add(listener);
|
|
}
|
|
|
|
foreach (var ipep in listeningAddresses)
|
|
{
|
|
logger.Information("Listening: {Address}", ipep);
|
|
}
|
|
|
|
ListeningAddresses = listeningAddresses.ToArray();
|
|
|
|
// Register listeners to start accepting connections
|
|
RegisterListeners(listeners.ToArray());
|
|
}
|
|
|
|
/// <summary>
|
|
/// Shuts down the network server and closes all listeners.
|
|
/// </summary>
|
|
public static void Shutdown()
|
|
{
|
|
CloseListeners();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gets the actual listening addresses for a wildcard endpoint.
|
|
/// </summary>
|
|
public static IEnumerable<IPEndPoint> GetListeningAddresses(IPEndPoint ipep) =>
|
|
NetworkInterface.GetAllNetworkInterfaces().SelectMany(adapter =>
|
|
adapter.GetIPProperties().UnicastAddresses
|
|
.Where(uip => ipep.AddressFamily == uip.Address.AddressFamily)
|
|
.Select(uip => new IPEndPoint(uip.Address, ipep.Port))
|
|
);
|
|
|
|
/// <summary>
|
|
/// Registers listeners with the ring and starts accepting connections.
|
|
/// </summary>
|
|
private static void RegisterListeners(nint[] listeners)
|
|
{
|
|
_listeners = listeners;
|
|
|
|
var ring = _socketManager.Ring;
|
|
|
|
// Queue initial accept operations for each listener
|
|
for (var i = 0; i < _listeners.Length; i++)
|
|
{
|
|
var listener = _listeners[i];
|
|
for (var j = 0; j < PendingAcceptsPerListener; j++)
|
|
{
|
|
ring.PrepareAccept(listener, 0, 0, IORingUserData.EncodeAccept());
|
|
_pendingAcceptCount++;
|
|
}
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Closes all listeners.
|
|
/// </summary>
|
|
private static void CloseListeners()
|
|
{
|
|
var ring = _socketManager?.Ring;
|
|
if (ring == null)
|
|
{
|
|
return;
|
|
}
|
|
|
|
foreach (var listener in _listeners)
|
|
{
|
|
ring.CloseListener(listener);
|
|
}
|
|
|
|
_listeners = [];
|
|
}
|
|
|
|
private static void HandleAcceptCompletion(int result)
|
|
{
|
|
_pendingAcceptCount--;
|
|
|
|
var ring = _socketManager.Ring;
|
|
|
|
// EAGAIN (-11) means no connection pending - just re-queue
|
|
if (result == -11)
|
|
{
|
|
goto ReplenishAccepts;
|
|
}
|
|
|
|
if (result >= 0)
|
|
{
|
|
var clientSocket = (nint)result;
|
|
var remoteIP = SocketHelper.GetRemoteAddress(clientSocket);
|
|
|
|
if (remoteIP != null)
|
|
{
|
|
if (_ipRateLimiter != null && !_ipRateLimiter.Verify(remoteIP, out var totalAttempts))
|
|
{
|
|
logger.Debug("{Address} Past IP limit threshold ({TotalAttempts})", remoteIP, totalAttempts);
|
|
|
|
if (Bans.BanConfiguration.Settings.ReportRateLimitTrips)
|
|
{
|
|
// Enqueue-only contribution; NOT added to the local firewall set (the limiter already
|
|
// gates it here and the OS bouncer drops it at the kernel).
|
|
Bans.BanChannel.Report(remoteIP, Bans.BanConfiguration.Settings.AutoBanDuration, Bans.BanReasons.RateLimit);
|
|
}
|
|
}
|
|
else if (ConnectionFilters.ShouldDeny(remoteIP, out var deniedBy))
|
|
{
|
|
// Whatever a hit implies (persisting, promoting to an OS bouncer, contributing to the
|
|
// ban channel) is the filter's own business; the accept path just drops the socket.
|
|
logger.Debug("{Address} denied by connection filter '{Filter}'", remoteIP, deniedBy);
|
|
}
|
|
else
|
|
{
|
|
// Allow event handlers to reject the connection
|
|
var args = new SocketConnectEventArgs(remoteIP);
|
|
EventSink.InvokeSocketConnect(args);
|
|
|
|
if (args.AllowConnection)
|
|
{
|
|
ring.ConfigureSocket(clientSocket);
|
|
CreateFromSocket(clientSocket, remoteIP);
|
|
goto ReplenishAccepts;
|
|
}
|
|
|
|
logger.Debug("{Address} Rejected by socket handler", remoteIP);
|
|
}
|
|
}
|
|
|
|
ring.CloseSocket(clientSocket);
|
|
}
|
|
else if (result != -4) // EINTR
|
|
{
|
|
logger.Debug("Accept error: {Result}", result);
|
|
}
|
|
|
|
ReplenishAccepts:
|
|
var targetAccepts = _listeners.Length * PendingAcceptsPerListener;
|
|
while (_pendingAcceptCount < targetAccepts && _listeners.Length > 0)
|
|
{
|
|
var listenerIndex = _pendingAcceptCount % _listeners.Length;
|
|
ring.PrepareAccept(_listeners[listenerIndex], 0, 0, IORingUserData.EncodeAccept());
|
|
_pendingAcceptCount++;
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Creates a NetState from an accepted socket handle.
|
|
/// </summary>
|
|
internal static NetState CreateFromSocket(nint socketHandle, IPAddress address)
|
|
{
|
|
// Use socket manager to create managed socket (handles buffers, registration, recv posting)
|
|
var socket = _socketManager.CreateSocket(socketHandle);
|
|
if (socket == null)
|
|
{
|
|
logger.Debug("Failed to create socket (resources exhausted)");
|
|
_socketManager.Ring.CloseSocket(socketHandle);
|
|
return null;
|
|
}
|
|
|
|
// Create NetState and map by socket ID
|
|
var ns = new NetState(socket, address);
|
|
return _netStates[socket.Id] = ns;
|
|
}
|
|
|
|
private static void DisconnectUnattachedSockets()
|
|
{
|
|
var now = Core.Now;
|
|
|
|
// Process connecting queue with lazy removal - O(1) operations
|
|
while (_connectingQueue.TryPeek(out var ns))
|
|
{
|
|
// Lazy removal: skip already-authenticated or disconnected connections
|
|
if (!ns.Running || ns.Account != null)
|
|
{
|
|
_connectingQueue.Dequeue();
|
|
continue;
|
|
}
|
|
|
|
// If the socket has been connected for less than the limit, we can stop
|
|
// (queue is ordered by connection time, so remaining entries are newer)
|
|
if (now - ns.ConnectedOn < ConnectingSocketIdleLimit)
|
|
{
|
|
break;
|
|
}
|
|
|
|
_connectingQueue.Dequeue();
|
|
|
|
// Socket must have finished the entire authentication process or be forcibly disconnected
|
|
if (!ns.SentFirstPacket || !ns.Seeded)
|
|
{
|
|
// Only the totally silent ones are evidence. A connection that sent SOME data and ran out of
|
|
// time is far more likely a slow link, and banning those makes the player retry, trip the
|
|
// rate limiter, and compound it into an hours-long ban.
|
|
if (!ns._receivedData && Bans.BanConfiguration.Settings.ReportBadConnects)
|
|
{
|
|
Bans.BanChannel.Report(
|
|
ns.Address,
|
|
Bans.BanConfiguration.Settings.BadConnectDuration,
|
|
Bans.BanReasons.SilentConnect
|
|
);
|
|
}
|
|
|
|
ns.Disconnect(null);
|
|
|
|
// Force immediate cleanup - these are unauthenticated connections
|
|
// where graceful disconnect can get stuck with pending sends.
|
|
if (ns._socket is { DisconnectPending: true })
|
|
{
|
|
_socketManager.DisconnectImmediate(ns._socket);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
public static void FlushAll()
|
|
{
|
|
while (_flushPending.TryDequeue(out var ns))
|
|
{
|
|
if (ns == null)
|
|
{
|
|
continue;
|
|
}
|
|
|
|
// Reset flag to allow re-queueing if more data is added later
|
|
ns._flushQueued = false;
|
|
|
|
if (ns.Running)
|
|
{
|
|
ns._socket?.QueueSend();
|
|
}
|
|
}
|
|
|
|
// Submit any pending operations
|
|
_socketManager?.Submit();
|
|
}
|
|
|
|
public static void Slice()
|
|
{
|
|
var curTicks = Core.TickCount;
|
|
DisconnectUnattachedSockets();
|
|
|
|
// Process throttled states
|
|
while (_throttled.Count > 0)
|
|
{
|
|
var ns = _throttled.Dequeue();
|
|
if (ns.Running)
|
|
{
|
|
ns.HandleReceive(true);
|
|
}
|
|
}
|
|
|
|
// This is enqueued by HandleReceive if already throttled and still throttled
|
|
while (_throttledPending.Count > 0)
|
|
{
|
|
_throttled.Enqueue(_throttledPending.Dequeue());
|
|
}
|
|
|
|
// Process queued movements at proper intervals
|
|
MovementThrottle.ProcessAllQueues();
|
|
|
|
// Process all completions through the manager FIRST
|
|
// This ensures DataReceived events are processed and HandleReceive runs,
|
|
// which may call Send() and add to _flushPending
|
|
var eventCount = _socketManager.ProcessCompletions(_events);
|
|
|
|
for (var i = 0; i < eventCount; i++)
|
|
{
|
|
ref var evt = ref _events[i];
|
|
|
|
switch (evt.Type)
|
|
{
|
|
case RingSocketEventType.Accept:
|
|
{
|
|
// Handle accept - AcceptedSocketHandle contains the result
|
|
HandleAcceptCompletion((int)evt.AcceptedSocketHandle);
|
|
break;
|
|
}
|
|
|
|
case RingSocketEventType.DataReceived:
|
|
{
|
|
var nsRecv = _netStates[evt.Socket.Id];
|
|
// Verify generation via object identity to avoid stale completion issues
|
|
if (nsRecv != null && nsRecv._socket == evt.Socket)
|
|
{
|
|
nsRecv.NextActivityCheck = curTicks + 30000;
|
|
HandleDataReceived(nsRecv, evt.BytesTransferred);
|
|
}
|
|
break;
|
|
}
|
|
|
|
case RingSocketEventType.DataSent:
|
|
{
|
|
var nsSend = _netStates[evt.Socket.Id];
|
|
// Verify generation via object identity
|
|
if (nsSend != null && nsSend._socket == evt.Socket)
|
|
{
|
|
// Update activity check on successful send
|
|
nsSend.NextActivityCheck = curTicks + 30000;
|
|
}
|
|
break;
|
|
}
|
|
|
|
case RingSocketEventType.Disconnected:
|
|
{
|
|
var nsDisc = _netStates[evt.Socket.Id];
|
|
// Verify generation via object identity
|
|
if (nsDisc != null && nsDisc._socket == evt.Socket)
|
|
{
|
|
HandleDisconnected(nsDisc);
|
|
}
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
|
|
// Process flush queue AFTER event processing
|
|
// This ensures sends triggered by HandleReceive (via packet handlers like SendPlayServerAck)
|
|
// are queued in the SAME Slice, not the next one
|
|
while (_flushPending.TryDequeue(out var ns))
|
|
{
|
|
// Reset flag to allow re-queueing if more data is added later
|
|
ns._flushQueued = false;
|
|
|
|
if (ns.Running)
|
|
{
|
|
ns._socket?.QueueSend();
|
|
}
|
|
}
|
|
|
|
// CRITICAL: Process send queue NOW to post pending sends
|
|
// This ensures PostSend() runs and sets SendPending=true BEFORE disconnect checks
|
|
// Without this, Disconnect() would see SendPending=false even though data is queued
|
|
_socketManager.ProcessSendQueue();
|
|
|
|
// Process pending disconnects AFTER flush queue AND send queue processing
|
|
// This ensures the traditional order: Game Logic (Sends/Disconnects) → Receives → Flush → Disconnect
|
|
// Any Send() calls made after Disconnect() in the same tick are flushed before disconnect
|
|
while (_pendingDisconnects.TryDequeue(out var ns))
|
|
{
|
|
// Reset flag to allow re-queueing if reconnect happens
|
|
ns._disconnectQueued = false;
|
|
|
|
if (ns.Running && ns._socket != null)
|
|
{
|
|
// RingSocket.Disconnect() handles graceful disconnect:
|
|
// - Waits for pending sends to flush (if SendBuffer.ReadableBytes > 0)
|
|
// - Waits for in-flight I/O to complete
|
|
// - Ensures buffers aren't released while kernel is still using them
|
|
ns._socket.Disconnect();
|
|
}
|
|
}
|
|
|
|
// Submit any queued operations
|
|
_socketManager.Submit();
|
|
|
|
// Process disposes
|
|
while (_disposed.TryDequeue(out var ns))
|
|
{
|
|
ns.DisposeInternal();
|
|
}
|
|
|
|
// Check for dead connections AFTER processing all completions.
|
|
// Recv completions reset NextActivityCheck, so after a server stall,
|
|
// buffered client pings update timestamps before this check fires.
|
|
if (curTicks - _nextAliveCheck >= 0)
|
|
{
|
|
_nextAliveCheck = curTicks + AliveCheckIntervalMs;
|
|
CheckAllAlive();
|
|
}
|
|
}
|
|
|
|
private static void HandleDataReceived(NetState ns, int bytesReceived)
|
|
{
|
|
if (!ns._running)
|
|
{
|
|
return;
|
|
}
|
|
|
|
if (bytesReceived > 0)
|
|
{
|
|
ns._receivedData = true;
|
|
}
|
|
|
|
// Data is already committed to buffer by RingSocketManager
|
|
// Decode if encryption is enabled
|
|
ns.DecryptRecvBuffer(bytesReceived);
|
|
|
|
// Process packets
|
|
ns.HandleReceive();
|
|
}
|
|
|
|
private static void HandleDisconnected(NetState ns)
|
|
{
|
|
var slotId = ns._socket.Id;
|
|
|
|
// IMPORTANT: Check if the slot still points to this NetState
|
|
// During quick reconnect, the slot might have been reused for a new connection
|
|
var currentNs = _netStates[slotId];
|
|
if (currentNs != ns)
|
|
{
|
|
// Slot was already reused - don't clear it!
|
|
// Just mark this NetState as not running and queue for dispose
|
|
ns._running = false;
|
|
_disposed.Enqueue(ns);
|
|
return;
|
|
}
|
|
|
|
// Clear the NetState slot
|
|
_netStates[slotId] = null;
|
|
|
|
// Mark as not running and queue for dispose
|
|
ns._running = false;
|
|
_disposed.Enqueue(ns);
|
|
}
|
|
|
|
public static void CheckAllAlive()
|
|
{
|
|
try
|
|
{
|
|
var curTicks = Core.TickCount;
|
|
|
|
foreach (var ns in Instances)
|
|
{
|
|
ns.CheckAlive(curTicks);
|
|
}
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
TraceException(ex);
|
|
}
|
|
}
|
|
}
|