Reshapes IP banning around one idea: **core owns the question, content owns every answer.**
Core gains a single accept-path seam — `IConnectionFilter` — and loses everything that used to implement one. The firewall moves to UOContent, a new file-backed blocklist joins it there, and CrowdSec is repositioned from an in-app enforcer to a contribute-first reporter.
## The seam
```csharp
public interface IConnectionFilter
{
string Name { get; }
void Configure();
void Start(CancellationToken token);
void Stop();
bool ShouldDeny(IPAddress address);
}
```
The accept path went from hardcoded branches to one question:
```csharp
else if (ConnectionFilters.ShouldDeny(remoteIP, out var deniedBy))
{
logger.Debug("{Address} denied by connection filter '{Filter}'", remoteIP, deniedBy);
}
```
Filters register during the Configure sweep. The registry is a plain array walked by an indexed loop — no enumerator, no closure, no allocation — and the first denial short-circuits. An interface dispatch is noise next to the `accept()` syscall, so pluggability costs nothing measurable on the path that has to survive a DDoS.
Whatever a hit implies — persisting, promoting to an OS bouncer, contributing to the ban channel — is the filter's business, not the accept path's.
A filter that throws is **unregistered and the connection fails open**. A filter that faults once faults for every subsequent connection, so leaving it registered means an exception and a log line per accept — exactly the amplification an attacker wants — and a broken filter must not be able to deny everyone either.
This deliberately does **not** reuse `EventSink.InvokeSocketConnect`: that fires later and allocates a `SocketConnectEventArgs` per connection, which is what the accept path avoids for rejected traffic.
## What ships behind it
**`firewall`** (UOContent) — the existing admin-curated set. Collapsed from `Firewall` + `AdminFirewall` + a threaded enforcer into one single-threaded store with **zero concurrency primitives**: the accept path, admin gump, TTL expiry and boot load all run on the game loop. Persists to `Configuration/firewall.json` with automatic migration from the legacy `firewall.cfg`. No behavior change for operators — same namespace, same gump, same commands.
**`blocklist`** (UOContent) — new. Holds a millions-strong list in-app and **demand-pages** hits up to CrowdSec, which promotes them to the OS firewall.
The motivation is concrete: CrowdSec's Windows bouncer cannot load the ~3.9M IPs that 91 community feeds produce, but it handles ~100k fine. So the millions live in-process behind a binary search, and only addresses that *actually connect* get promoted. A `PromotedGuard` suppresses re-reporting an address until the bouncer picks it up.
The list is parsed straight from UTF-8 file bytes with no per-line string allocation, off the game loop, and published as an immutable snapshot swapped through a single `volatile` reference. Reloads yield to world saves.
**`tools/Export-IpBlocklist.ps1`** — the producer. Requires PowerShell 7 and runs on Windows, Linux and macOS; Windows PowerShell 5.1 is refused up front via `#requires`. Merges a thin, non-overlapping feed set into one de-duplicated, bogon-filtered file. Parsing runs in a compiled `Add-Type` hot loop (~1s for ~4M lines instead of minutes). Written to a `.tmp` sibling and swapped with `File.Replace`, so the shard never reads a half-written list, and a total feed outage refuses to overwrite a good list with an empty one. Re-running is idempotent — it exits without downloading anything while the list on disk is younger than `-MinInterval` (default 2h, the anchor feed's own refresh period), so a misconfigured scheduler can't hammer upstream.
## CrowdSec: contribute-first
`IBanReporter` + `BanChannel` fan locally-decided bans out to external systems. `CrowdSecReporter` (UOContent) posts to LAPI `POST /v1/alerts` and retracts via `DELETE /v1/decisions`.
Reporting is **enqueue-only** on the accept path: a bounded, coalescing channel drained off-loop with bounded retry, counted drops on overflow, and a flush on shutdown. Under a DDoS the accept path never does synchronous or lock-contending per-IP work.
### Why not pull decisions from CrowdSec?
The original design streamed decisions into an in-app snapshot and enforced them at the accept gate. That's the wrong layer: by the time the shard sees the connection, the TCP handshake and socket setup are already paid for. `cs-firewall-bouncer` drops the same traffic **at the kernel**, and it's what CrowdSec is built to do. So the shard now contributes what it uniquely knows (rate-limit trips, blocklist hits from real connection attempts) and lets the OS enforce.
The one thing the OS can't do — hold millions of entries on Windows — is exactly what the in-app blocklist covers, and it feeds the same pipeline.
## Threading policy
`CLAUDE.md` rule #3 is rewritten as an explicit three-part policy, with rule #10 restated in tandem:
- Anything touching game state runs **only** on the main loop.
- Heavy work that *needs* game state must be **chunked** across ticks, never threaded.
- Heavy work that does *not* need game state (large-file parse, external I/O) **must** run off-loop **and must yield to world saves**.
Results come back via an immutable snapshot swapped through a single `volatile` reference, or `Core.LoopContext.Post` — never by letting the scheduler decide where heavy work runs. Both new subsystems follow it.
## Shared primitives
`SortedRangeIndex<T> where T : IBinaryInteger<T>` — coalesced disjoint interval arrays plus a binary search. The firewall, the blocklist, and (as of this PR) core's reserved-network tables all use it.
Coalescing is a correctness requirement, not an optimization: multi-feed lists nest CIDRs (`/24` containing a `/32`), and a search that inspects only the rightmost run whose minimum is ≤ the value is sound **only** over disjoint runs. That bug was caught in review and is covered by regression tests.
`IPAddressUtility` collects the allocation-free `IPAddress` ↔ `UInt128` conversions and CIDR parsing that were previously scattered or duplicated.
## Config
| File | Owner | Keys |
|---|---|---|
| `Configuration/bans.json` | core | `reportRateLimitTrips`, `autoBanDuration` |
| `Configuration/blocklist.json` | content | `file`, `reloadInterval`, `reportHits`, `banDuration`, `promoteSuppression` |
| `Configuration/crowdsec.json` | content | `lapiUrl`, `machineId`, `password`, `origin`, `manualBanDuration`, `flushInterval`, `maxQueue` |
| `Configuration/firewall.json` | content | persisted firewall entries (migrated from `firewall.cfg`) |
Everything is inert by default. CrowdSec self-disables without credentials; the blocklist self-disables until its file exists. A shard that changes nothing sees no behavior change.
## Notes for review
- **Core no longer references `Firewall` or `IFirewallEntry` anywhere.** `NetworkUtilities` used to build its reserved-network tables out of `CidrFirewallEntry`, which coupled core to the firewall for something unrelated to banning; those are now a `SortedRangeIndex<UInt128>`, same semantics and public API.
- **`BanChannel.Stop()` no longer persists the firewall** — a contribution coordinator has no business saving an enforcement store. That's the firewall filter's `Stop()`.
- **A dead `whitelisted` parameter was dropped** from the blocklist gate: it was hardcoded `false` at its only call site, and no whitelist concept exists in core.
- **The blocklist filter is an instance, not a static.** The static version forced its tests onto the sequential collection with a reset hook; they now run in parallel.
- `dev-docs/networking-packets.md` documents the seam for content authors, plus a known wart in the `IPAddress` ↔ `UInt128` normalization flagged for a follow-up PR.
- The generator was verified on Linux, macOS and Windows under a temporary CI matrix (since removed). It caught two portability bugs — a Windows-only path separator, and a culture-sensitive duration parse that read `2.5` as `25` on comma-decimal locales and *silently* turned a 2.5h cooldown into 25h — plus a third that made the script unparseable on Windows PowerShell 5.1. The source is ASCII-only for that last reason: `#requires` is only honored once a file parses, so non-ASCII in a BOM-less script produces parse errors instead of the version message.
## Tests
**1344 pass** (782 `Server.Tests`, 562 `UOContent.Tests`). New coverage: filter registry (registration, short-circuit, fault-disable), blocklist parsing/CIDR/coalescing, snapshot reload markers, promote-guard TTL, ban-channel fan-out, CrowdSec alert building/dedup/flush-on-stop, and the generator's output-format contract pinned against the reader.
544 lines
19 KiB
C#
544 lines
19 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;
|
|
|
|
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 SendBufferSize = 1024 * 256; // 256KB send buffers
|
|
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
|
|
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>
|
|
/// 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);
|
|
|
|
// Initialize IORingGroup
|
|
var ring = IORingGroup.Create(queueSize: MaxConnections * 2, maxConnections: MaxConnections);
|
|
|
|
// 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>
|
|
/// 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, "rate-limit");
|
|
}
|
|
}
|
|
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)
|
|
{
|
|
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;
|
|
}
|
|
|
|
// 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);
|
|
}
|
|
}
|
|
}
|