Some checks are pending
Build / Build (MacOS 15) (push) Waiting to run
Build / Build (MacOS 26) (push) Waiting to run
Build / Build (AlmaLinux 10) (push) Waiting to run
Build / Build (Debian 12) (push) Waiting to run
Build / Build (Debian 13) (push) Waiting to run
Build / Build (Fedora 44) (push) Waiting to run
Build / Build (CentOS 10 Stream) (push) Waiting to run
Build / Build (CentOS 9 Stream) (push) Waiting to run
Build / Build (Ubuntu 26) (push) Waiting to run
Build / Build (Ubuntu 22) (push) Waiting to run
Build / Build (Ubuntu 24) (push) Waiting to run
**Follow-on to #2639 (merged). References a local IORingGroup `1.0.13-preview.11` pack until 1.0.13 (modernuo/IORingGroup#15) is published; do not merge before that switch.** ## Summary Consumes IORingGroup's lean base pools (modernuo/IORingGroup#15): both network pools now start with one slab, grow a slab at a time with the population, and trim idle slabs back after quiet periods. - Fixes the transport's send-pool cap: previously only 1024 of the 4096 connections could get a send buffer; connection 1025 was closed at accept. - Network memory at boot drops from about 96 MB to about 10 MB at the defaults; a full 4096 logged-in connections is about 1.25 GB of base buffers plus the growth budget. - New settings: `network.initialBufferSlabs` (default 1; slabs of each pool held from boot and the trim floor) and `network.maxBufferSlabs` (default 128; divides the connection maximum into slabs, 32 connections per slab). Both are coerced with a warning; the same value feeds the ring table and the manager so they cannot drift. - The Debug-only maintenance line includes base-pool capacity and releases. - `dev-docs/server-requirements.md` rewrites the network memory story and adds the two settings. ## Pre-auth buffers Every connection starts on the transport's platform-minimum buffers (4 KB receive, 4 KB send) instead of the base pools. It is promoted to full-size buffers (64 KB receive, `network.sendBufferSize` send) when the game server verifies its account — the point where `NetState.Account` is assigned in the `GameServer_AwaitingGameServerLogin` or `GameServer_LoggedIn` state, so a verdict that lands after the parser has moved on still promotes. Nothing ever moves back. The login-server pass stays on the small buffers for its whole lifetime. Before credentials verify, nothing promotes: the 4 KB send ring is the entire pre-auth send budget, and a connection that overruns it is dropped as exhausted, exactly as the receive side drops a packet header declaring more bytes than the receive buffer can hold (new guard in `HandlePacket`; it also closes the old 65535-byte edge on 64 KB buffers). The stock login sequence sends under 2 KB. After credentials verify, the send path promotes on demand if it ever needs to (unbudgeted, outside the memory ceiling and the shrink bookkeeping — promotion is not growth), and the oversize-packet guard waits on a pending receive promotion or retries a stalled one once for a verified account before disconnecting. A completion that fills the receive buffer arms no receive, so `HandleReceive` now calls `RingSocket.ResumeReceive()` after the parse loop — at 4 KB a burst of small packets fills the buffer in one completion. Net effect: a flood of unauthenticated connections tops out at about 32 MB across the full 4096-connection cap where the platform minimum is 4 KB (the transport's retained slabs and the base pools used by logged-in players are separate), and never allocates a base-pool slab. Platform note: on Windows Server 2012 R2 / 2016 the transport's legacy mapping path floors at 64 KB: the pre-auth receive pool is off there (its base is 64 KB), while the pre-auth send buffer starts at 64 KB under the 256 KB base. The server logs the effective sizes at startup. ## Testing Server.Tests (905) and UOContent.Tests green on the preview pack (one pre-existing `FamiliarAITests` failure from #2644 reproduces on `main`, tracked separately). New tests cover both coercions, that the ring's registration table equals `RequiredRegisteredBuffers` for the configured values, promotion on game-server auth and on a late account, no promotion on the login server, pre-credential overrun ending in exhaustion, post-credential on-demand promotion, the oversize-packet guard through loopback (error, wait, retry with an account, and a promotion made pending mid-parse), and receiving again after a burst fills the initial buffer.
867 lines
32 KiB
C#
867 lines
32 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
|
|
internal 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 DefaultMaxSendBufferSize = 1024 * 1024 * 2; // 2 MB
|
|
private const long DefaultSendBufferGrowthBudget = 1024L * 1024 * 256; // 256 MB
|
|
private const int DefaultMemoryCeilingPercent = 80;
|
|
private const int DefaultInitialBufferSlabs = 1; // one slab of each base pool warm at boot
|
|
internal const int DefaultMaxBufferSlabs = 128; // 32 connections per slab at MaxConnections
|
|
|
|
// Transport ceiling; larger values overflow its tier enumeration
|
|
private const int TransportMaxSendBufferSize = 1024 * 1024 * 256; // 256 MB
|
|
private const int MaxConnections = 4096; // Max concurrent connections
|
|
|
|
internal static int SendBufferSize { get; private set; }
|
|
internal static int MaxSendBufferSize { get; private set; }
|
|
|
|
// Pre-auth buffers are the smallest the platform can map; a platform whose floor is not below
|
|
// the base size (the Windows legacy mapping path) starts sockets on base
|
|
internal static int InitialRecvBufferSize { get; private set; }
|
|
internal static int InitialSendBufferSize { get; private set; }
|
|
private static long _sendBufferGrowthBudget;
|
|
private static int _memoryCeilingPercent;
|
|
|
|
// Refreshed by the maintenance sweep; internal for tests
|
|
internal static long _availableMemoryBytes;
|
|
|
|
private static Timer.DelayCallTimer _maintenanceTimer;
|
|
private static long _lastTierCapacityBytes;
|
|
private static int _lastTierInUse;
|
|
private static int _lastTierRetainFloor;
|
|
private static long _lastBaseCapacityBytes;
|
|
|
|
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;
|
|
|
|
// Test hook
|
|
internal static RingSocketManager SocketManager => _socketManager;
|
|
|
|
/// <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 && _pendingDisconnects.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;
|
|
}
|
|
|
|
// Seed from a real tick; a zero default suppresses the sweep when ticks start negative
|
|
_nextAliveCheck = Core.TickCount;
|
|
|
|
// 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);
|
|
|
|
// Per-connection send buffer: the lever for "send buffer exhausted" disconnects, and the
|
|
// per-connection memory ceiling.
|
|
SendBufferSize = GetSendBufferSize();
|
|
MaxSendBufferSize = GetPowerOfTwoSetting("network.sendBufferMaxSize", DefaultMaxSendBufferSize, SendBufferSize);
|
|
_sendBufferGrowthBudget = CoerceSendBufferGrowthBudget(
|
|
ServerConfiguration.GetOrUpdateSetting("network.sendBufferGrowthBudget", DefaultSendBufferGrowthBudget),
|
|
SendBufferSize,
|
|
MaxSendBufferSize
|
|
);
|
|
// 0 disables the ceiling
|
|
_memoryCeilingPercent = Math.Clamp(ServerConfiguration.GetOrUpdateSetting("network.memoryCeilingPercent", DefaultMemoryCeilingPercent), 0, 100);
|
|
_availableMemoryBytes = GC.GetGCMemoryInfo().TotalAvailableMemoryBytes;
|
|
|
|
// Divides MaxConnections into base-pool slabs; both pools still reach MaxConnections, so
|
|
// this sets slab granularity, not a connection or memory ceiling.
|
|
var maxBufferSlabs = CoerceMaxBufferSlabs(
|
|
ServerConfiguration.GetOrUpdateSetting("network.maxBufferSlabs", DefaultMaxBufferSlabs)
|
|
);
|
|
|
|
// Slabs of each base pool held from boot; the pools grow and trim from here
|
|
var initialBufferSlabs = CoerceInitialBufferSlabs(
|
|
ServerConfiguration.GetOrUpdateSetting("network.initialBufferSlabs", DefaultInitialBufferSlabs),
|
|
BasePoolSlabCount(maxBufferSlabs)
|
|
);
|
|
|
|
// Until the game server verifies credentials a connection holds the smallest buffers the
|
|
// platform can map; a flood can fill these pools but never reaches the base pools. This is a
|
|
// request: the manager raises it to the platform floor and turns the pool off if that leaves
|
|
// no room below the base size (the Windows legacy mapping path floors at 64 KiB, which can
|
|
// equal RecvBufferSize). The effective sizes are read back below once the manager knows them.
|
|
var initialBufferSize = IORingBuffer.MinimumSize;
|
|
|
|
// Both calls take the same slab count; the manager throws if the table is smaller. Sized by
|
|
// the request, not the effective size the manager may coerce down to - conservative, never small.
|
|
var ring = IORingGroup.Create(
|
|
queueSize: MaxConnections * 2,
|
|
maxConnections: MaxConnections,
|
|
maxOutstandingSends: maxOutstandingSends,
|
|
maxRegisteredBuffers: RingSocketManager.RequiredRegisteredBuffers(
|
|
MaxConnections, SendBufferSize, MaxSendBufferSize, _sendBufferGrowthBudget, maxBufferSlabs,
|
|
initialBufferSize, initialBufferSize
|
|
)
|
|
);
|
|
|
|
// Create socket manager which handles buffer pools and socket lifecycle
|
|
_socketManager = new RingSocketManager(
|
|
ring,
|
|
maxSockets: MaxConnections,
|
|
recvBufferSize: RecvBufferSize,
|
|
sendBufferSize: SendBufferSize,
|
|
initialBufferSlabs: initialBufferSlabs,
|
|
maxBufferSlabs: maxBufferSlabs,
|
|
maxSendBufferSize: MaxSendBufferSize,
|
|
sendBufferGrowthBudget: _sendBufferGrowthBudget,
|
|
initialRecvBufferSize: initialBufferSize,
|
|
initialSendBufferSize: initialBufferSize
|
|
);
|
|
|
|
InitialRecvBufferSize = _socketManager.InitialRecvBufferSize;
|
|
InitialSendBufferSize = _socketManager.InitialSendBufferSize;
|
|
|
|
if (InitialRecvBufferSize == 0 || InitialSendBufferSize == 0)
|
|
{
|
|
logger.Information(
|
|
"Pre-auth buffers off where the platform floor ({Minimum} bytes) leaves no room below the base size (recv {Recv}, send {Send})",
|
|
IORingBuffer.MinimumSize,
|
|
RecvBufferSize,
|
|
SendBufferSize
|
|
);
|
|
}
|
|
|
|
_maintenanceTimer = Timer.DelayCall(TimeSpan.FromMinutes(1), TimeSpan.FromMinutes(1), MaintainSendBuffers);
|
|
}
|
|
|
|
internal static void MaintainSendBuffers()
|
|
{
|
|
if (_socketManager == null)
|
|
{
|
|
return;
|
|
}
|
|
|
|
// Container limits can change
|
|
_availableMemoryBytes = GC.GetGCMemoryInfo().TotalAvailableMemoryBytes;
|
|
|
|
var ceilingRefusals = _ceilingRefusals;
|
|
var capRefusals = _capRefusals;
|
|
_ceilingRefusals = 0;
|
|
_capRefusals = 0;
|
|
|
|
var stats = _socketManager.Maintain();
|
|
var changed = stats.TierCapacityBytes != _lastTierCapacityBytes ||
|
|
stats.TierInUse != _lastTierInUse ||
|
|
stats.TierRetainFloor != _lastTierRetainFloor ||
|
|
stats.BaseCapacityBytes != _lastBaseCapacityBytes;
|
|
_lastTierCapacityBytes = stats.TierCapacityBytes;
|
|
_lastTierInUse = stats.TierInUse;
|
|
_lastTierRetainFloor = stats.TierRetainFloor;
|
|
_lastBaseCapacityBytes = stats.BaseCapacityBytes;
|
|
|
|
// Quiet unless something moved
|
|
if (changed || stats.BuffersReleased > 0 || stats.BaseBuffersReleased > 0 || stats.GrowthRefusals > 0 ||
|
|
capRefusals > 0 || ceilingRefusals > 0)
|
|
{
|
|
logger.Debug(
|
|
"Send buffer tiers: {Capacity} bytes of tier capacity, {InUse} buffers in use, floor {Floor} buffers, released {Released}, refused: budget {BudgetRefusals}, at max {CapRefusals}, ceiling {CeilingRefusals}, base {BaseCapacity} bytes, released {BaseReleased}",
|
|
stats.TierCapacityBytes,
|
|
stats.TierInUse,
|
|
stats.TierRetainFloor,
|
|
stats.BuffersReleased,
|
|
stats.GrowthRefusals,
|
|
capRefusals,
|
|
ceilingRefusals,
|
|
stats.BaseCapacityBytes,
|
|
stats.BaseBuffersReleased
|
|
);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Slabs each base pool is divided into, after the transport applies its minimum slab size.
|
|
/// </summary>
|
|
internal static int BasePoolSlabCount(int maxBufferSlabs) =>
|
|
RingSocketManager.BasePoolSlabCount(MaxConnections, maxBufferSlabs);
|
|
|
|
/// <summary>
|
|
/// Clamps the base-pool slab divisor. Fewer than one slab is meaningless, and more slabs than
|
|
/// connections cannot make a slab any smaller.
|
|
/// </summary>
|
|
internal static int CoerceMaxBufferSlabs(int configured)
|
|
{
|
|
var slabs = Math.Clamp(configured, 1, MaxConnections);
|
|
|
|
if (slabs != configured)
|
|
{
|
|
logger.Warning(
|
|
"network.maxBufferSlabs {Configured} is outside 1..{Maximum}; using {Adjusted}",
|
|
configured,
|
|
MaxConnections,
|
|
slabs
|
|
);
|
|
}
|
|
|
|
return slabs;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Clamps the slabs of each base pool held from boot; a pool never holds more slabs than it has.
|
|
/// </summary>
|
|
internal static int CoerceInitialBufferSlabs(int configured, int slabCount)
|
|
{
|
|
var slabs = Math.Clamp(configured, 1, slabCount);
|
|
|
|
if (slabs != configured)
|
|
{
|
|
logger.Warning(
|
|
"network.initialBufferSlabs {Configured} is outside 1..{Maximum}; using {Adjusted}",
|
|
configured,
|
|
slabCount,
|
|
slabs
|
|
);
|
|
}
|
|
|
|
return slabs;
|
|
}
|
|
|
|
/// <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() =>
|
|
GetPowerOfTwoSetting("network.sendBufferSize", DefaultSendBufferSize, MinSendBufferSize);
|
|
|
|
private static int GetPowerOfTwoSetting(string key, int defaultValue, int minimum) =>
|
|
CoercePowerOfTwoSetting(key, ServerConfiguration.GetOrUpdateSetting(key, defaultValue), minimum);
|
|
|
|
/// <summary>
|
|
/// Clamps to a power of two between <paramref name="minimum"/> and the transport ceiling.
|
|
/// </summary>
|
|
internal static int CoercePowerOfTwoSetting(string key, int configured, int minimum)
|
|
{
|
|
var size = configured;
|
|
|
|
if (size > TransportMaxSendBufferSize)
|
|
{
|
|
logger.Warning(
|
|
"{Key} {Configured} is above the transport maximum {Maximum} (capped); using {Adjusted}",
|
|
key,
|
|
configured,
|
|
TransportMaxSendBufferSize,
|
|
TransportMaxSendBufferSize
|
|
);
|
|
|
|
size = TransportMaxSendBufferSize;
|
|
}
|
|
|
|
if (size < minimum)
|
|
{
|
|
logger.Warning(
|
|
"{Key} {Configured} is below the minimum {Minimum} (raised); using {Adjusted}",
|
|
key,
|
|
configured,
|
|
minimum,
|
|
minimum
|
|
);
|
|
|
|
size = minimum;
|
|
}
|
|
|
|
if (!BitOperations.IsPow2(size))
|
|
{
|
|
// Rounding up cannot cross the ceiling, itself a power of two
|
|
var rounded = (int)BitOperations.RoundUpToPowerOf2((uint)size);
|
|
|
|
logger.Warning("{Key} {Configured} is not a power of two; using {Adjusted}", key, configured, rounded);
|
|
size = rounded;
|
|
}
|
|
|
|
return size;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Negative disables growth; below one tier slab is raised to the minimum.
|
|
/// </summary>
|
|
internal static long CoerceSendBufferGrowthBudget(long configured, int sendBufferSize, int maxSendBufferSize)
|
|
{
|
|
if (configured < 0)
|
|
{
|
|
logger.Warning("network.sendBufferGrowthBudget {Configured} is negative; using 0 (growth disabled)", configured);
|
|
return 0;
|
|
}
|
|
|
|
if (configured == 0 || maxSendBufferSize <= sendBufferSize)
|
|
{
|
|
return configured;
|
|
}
|
|
|
|
// Tier buffers are allocated a slab at a time.
|
|
var minimumBudget = RingSocketManager.MinimumSendBufferGrowthBudget(sendBufferSize);
|
|
if (configured < minimumBudget)
|
|
{
|
|
logger.Warning(
|
|
"network.sendBufferGrowthBudget {Configured} is below one tier slab; using {Minimum}",
|
|
configured,
|
|
minimumBudget
|
|
);
|
|
|
|
return minimumBudget;
|
|
}
|
|
|
|
return configured;
|
|
}
|
|
|
|
/// <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;
|
|
nsSend.TryShrinkSendBuffer(curTicks);
|
|
}
|
|
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();
|
|
|
|
if (ns._socket.DisconnectPending)
|
|
{
|
|
ns.ArmDrainDeadline(curTicks);
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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);
|
|
ns.TryShrinkSendBuffer(curTicks);
|
|
}
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
TraceException(ex);
|
|
}
|
|
}
|
|
}
|