## Problem A connection's send buffer is a fixed 256 KB. A burst of world traffic (a crowded area, a mass spawn, a war) that outruns the client's acknowledgements fills it, `NetState.Send()` reports "send buffer exhausted", and the player is disconnected. Raising the size for everyone multiplies the per-connection footprint (4096 × 256 KB is already 1 GB at full occupancy, page-locked on Windows). ## What changes - **Growth.** When a packet does not fit (the write span is too small, the packet is larger than the span, or compression returns 0), `Send()` asks the transport to grow the buffer to the next power-of-two tier and retries, up to `network.sendBufferMaxSize` (2 MB). Compression retries once per tier since its output size is not known in advance, including when the buffer is completely full. Only when growth is refused does the existing exhaustion disconnect run. The success path is unchanged. - **Memory ceiling.** Growth is refused (with a once-a-minute warning) when the process working set exceeds `network.memoryCeilingPercent` (80) of the memory available to the process (container-aware; `0` turns the check off). The figure is sampled at startup and refreshed each maintenance tick. - **Shrink.** A grown socket returns to the base buffer once it is drained and 30 s have passed since its last growth, attempted from the `DataSent` handler and from the 5 s alive sweep. - **Retention.** Every minute a timer calls the transport's `Maintain()`, which trims idle tier slabs down to the peak concurrent usage of the last 15 minutes, so recurring bursts reuse buffers without allocation while rare ones give the memory back. The line logs at Debug, and only when capacity, usage, or the floor changed or a growth was refused (budget, at max, or ceiling), so an idle shard logs nothing. - **Budget.** `network.sendBufferGrowthBudget` (256 MB) caps the tier pools' capacity; a positive value below one tier slab is raised with a warning, a negative one is clamped to 0 (growth off). Worst case is base × connections plus the budget. - Settings are coerced with accurate warnings (power of two, minimum, 256 MB transport ceiling). `[dumpnetstates` gains the send buffer size. `dev-docs/server-requirements.md` describes the new memory story. ## Tests `NetStateSendBufferTests` (real loopback sockets): growth instead of disconnect with a byte-exact stream, compressed growth against the compressor's own output, the grow-then-copy path, growth with a send genuinely in flight, refusal past the maximum, refusal under the ceiling, refusal on a closing socket, shrink after the hold (direct and through the alive sweep), and the setting coercions. Server.Tests 891 passed, UOContent.Tests 1052 passed against the published 1.0.12. Reviewed per task, whole-branch, and adversarially by a second model (twice, the second time jointly with the transport branch); all findings addressed.
768 lines
28 KiB
C#
768 lines
28 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 DefaultMaxSendBufferSize = 1024 * 1024 * 2; // 2 MB
|
|
private const long DefaultSendBufferGrowthBudget = 1024L * 1024 * 256; // 256 MB
|
|
private const int DefaultMemoryCeilingPercent = 80;
|
|
|
|
// 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 MaxSendBufferSize { 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 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.
|
|
var 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;
|
|
|
|
const int maxBufferSlabs = 32;
|
|
var ring = IORingGroup.Create(
|
|
queueSize: MaxConnections * 2,
|
|
maxConnections: MaxConnections,
|
|
maxOutstandingSends: maxOutstandingSends,
|
|
maxRegisteredBuffers: RingSocketManager.RequiredRegisteredBuffers(MaxConnections, sendBufferSize, MaxSendBufferSize, _sendBufferGrowthBudget, maxBufferSlabs)
|
|
);
|
|
|
|
// Create socket manager which handles buffer pools and socket lifecycle
|
|
_socketManager = new RingSocketManager(
|
|
ring,
|
|
maxSockets: MaxConnections,
|
|
recvBufferSize: RecvBufferSize,
|
|
sendBufferSize: sendBufferSize,
|
|
initialBufferSlabs: 8,
|
|
maxBufferSlabs: maxBufferSlabs,
|
|
maxSendBufferSize: MaxSendBufferSize,
|
|
sendBufferGrowthBudget: _sendBufferGrowthBudget
|
|
);
|
|
|
|
_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;
|
|
_lastTierCapacityBytes = stats.TierCapacityBytes;
|
|
_lastTierInUse = stats.TierInUse;
|
|
_lastTierRetainFloor = stats.TierRetainFloor;
|
|
|
|
// Quiet unless something moved
|
|
if (changed || stats.BuffersReleased > 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}",
|
|
stats.TierCapacityBytes,
|
|
stats.TierInUse,
|
|
stats.TierRetainFloor,
|
|
stats.BuffersReleased,
|
|
stats.GrowthRefusals,
|
|
capRefusals,
|
|
ceilingRefusals
|
|
);
|
|
}
|
|
}
|
|
|
|
/// <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);
|
|
}
|
|
}
|
|
}
|