modernuo/Projects/Server/Network/NetState/NetState.Network.cs
Kamron Batman aae173a797
feat(network): allowlist false-positive IPs, escalate on behavior (#2556)
## Why

The shard owner, on a Starlink CGNAT address, was blocked by the imported reputation blocklist.

The cause was not CrowdSec. The address was a literal line in `ip-blocklist.txt`, so `BlocklistFilter` denied it at accept and then promoted it — and clearing the CrowdSec decision could not fix it either, because the file entry re-reports within `promoteSuppression` of every reconnect attempt.

This is structural, not a one-off. Reputation feeds list shared consumer address space constantly: on CGNAT one public address fronts many subscribers **at the same time**, so a single abusive customer gets the address listed and everyone else behind it is blocked with them. Where leases rotate, a listing says little about whoever holds the address now. Around 1,000 Starlink addresses sit in the current list.

So exemptions go where they cost nothing, and escalation is driven by what a connection actually does.

## Generator — `tools/Export-IpBlocklist.ps1`

`-AllowlistFile` takes multiple paths, subtracted from the merged set before the output is written. Defaults to every `ip-allowlist*.txt` beside the output, merged into one allow set:

- `ip-allowlist.txt` — operator exemptions, created once and **never rewritten**
- `ip-allowlist-<name>.txt` — a carve-out you built, regenerable and copyable between shards

**Subtraction is range-correct.** An allowlisted address inside a blocked CIDR splits that CIDR around the hole rather than being silently ignored. This also fixes `-ExcludeAnonymizers`, which parsed CIDR entries into `$anonCidr` and then only ever subtracted singles.

**No carve-out ships.** A carve-out names a real network, and which ones a shard should exempt depends on where its players actually are — so publishing one would make that policy call for every shard and put a specific provider's address space in the repo. The script builds them on request instead:

```powershell
.\Export-IpBlocklist.ps1 -AddCarveout starlink -Asn 14593
```

Carve-outs are **discovered, not configured**: every `ip-allowlist*.txt` beside the output is subtracted, by the generator and by the shard, so a file an admin adds needs no config edit and no code change. Each carries an `asn=` marker in its header, which is how `-RefreshCarveouts` rebuilds it without the script keeping a list of anyone's networks; a hand-written allowlist has no marker and is never rewritten.

Prefixes come from **announcements, not ownership records**, because registry data disagrees with what is actually routed and silently caps result sets: ARIN whois returns at most 256 rows and gives per-customer /24s, and `206.83.96.0/19` reads as APNIC in RDAP even though `206.83.96/21` is announced by Starlink.

Editing an allowlist bypasses `-MinInterval`, so a just-added exemption isn't indistinguishable from the allowlist not working. A Starlink carve-out, if you build one, costs **~4,300 IPs + ~144 CIDRs of 4.2M (0.10%)**.

## Allowlists

**`FileAllowlist`** reads the same files the generator subtracts, so an operator entry means "leave this address alone" for real. Subtraction alone only covers being *blocked*; behavioural detections never consult the blocklist, so without this a carve-out was quietly routed around — one scanner behind a shared address was enough to get everyone behind it contributed and firewalled, with nothing in the shard's own config explaining why. Reading the files also means an entry applies on the next reload rather than the next regeneration, which is what matters when someone is complaining now.

**`LoginAllowlist`** is earned by authenticating, with a 90-day TTL because an address that logged in years ago is a stranger. Its own store rather than `Account.LoginIPs`, which has no timestamps and cannot be backfilled. An entry is evidence rather than a licence: 10 suppressed contributions in an hour revokes it, and a fresh login forgives the tally.

Both are consulted **only after the blocklist has already matched**, so a normal accept pays nothing for them and the accept gate stays allowlist-free. `BanExemptions` combines them behind `BanChannel.IsExempt` and suppresses escalation only — every local defence still applies.

Two limits, both deliberate and documented in the class: `LoginAllowlist` **cannot bootstrap** (an entry is only earned by getting in, so it never repairs an existing false positive), and it is weakest on rotating CGNAT. That is why `FileAllowlist` is the fix for those, and why it is manual.

## Behavioural detection

| Reason | Trigger |
|---|---|
| `silent-connect` | Reaped after 5s having sent **zero bytes** |
| `invalid-seed` | Opened with a zero seed |
| `foreign-protocol` | Positively identified as HTTP, TLS or SSH |

**`ForeignProtocol` inverts the test.** Asking "is this a good UO client?" cannot work: `LoginEncryption.ClientDecrypt` is a byte-for-byte stream XOR, so a legitimate client with encryption enabled when the shard expects none sends a structurally perfect connection whose payload is noise. "Speaks HTTP" is safe where "unreadable" is not — however misconfigured a UO client is, it never sends `GET / HTTP/1.1`.

Nothing assumes arrival framing. TCP has no message boundaries, so a rule of the form "these bytes must arrive together" is broken by construction and drops real players on poor links. A prefix match with too few bytes to confirm waits for more. A four-byte seed can legitimately spell `GET ` (the address 71.69.84.32) or `0x16 0x03 0x0?` (22.3.x.x), so confirmation requires the request line to continue in printable ASCII or an actual ClientHello inside a plausible record — a real client's fifth byte is a packet id (`0x80`, `0x91`, `0xEF`), none of them printable, so those collisions fall through.

Everything is keyed on **bytes-received rather than elapsed time**. A connection that sent something and ran out of time is far more likely a slow link than an attack, and banning those produces the worst failure mode available: the player retries, trips the rate limiter, and compounds a bad connection into hours of being firewalled off.

## `AutoDenylist`

A short-lived local hold (15m) on behavioural detections, as `IConnectionFilter` + `IBanReporter` over one store so the engine detection sites never reach into content.

This closes the gap where a flood pays for a socket, buffer and `NetState` slot per connection while waiting for the OS bouncer — the verdicts that matter most are reachable only *after* reading bytes — and it is the entire defence on a shard running no bouncer, which is the default config. Not persisted: a holding pen that survives restarts is a ban without a ban's review.

Cost: one dictionary lookup on a usually-empty dict per accept.

## `BanReasons`

Centralises the reason slugs. `IsBehavioral` is an **opt-in** set, not "everything except manual", so a future reason escalates normally instead of silently inheriting an exemption or entering a local denylist.

This caught a real bug during review: the first cut of the exemption swallowed `manual` admin bans (`Commands.cs`, three sites in `AdminGump`) for any allowlisted address.

## Fixes found in review

- **`BanConfiguration.Settings` was null until `Configure()` ran**, while the reap path dereferences it every `Slice()`. A harness driving `NetState.Slice()` directly hit an NRE that presented as flaky because it depended on whether an earlier test had already called `Configure()` — which is why it failed on some CI platforms and not others. Now starts at the record's defaults, with idempotency tracked by a flag; this also removes the same latent NRE from the pre-existing rate-limit path.
- **`-AllowlistFile` was typed `[string]`** while documented and used as a list, so passing two paths would have collapsed them into one string.

## Layout and docs

Content network code moves out of `Misc/` into `UOContent/Network/`, one concern per folder — `AutoDenylist/`, `Blocklist/`, `CrowdSec/`, `Firewall/`, `LoginAllowlist/`, `Packets/`. **Namespaces are untouched**, so these are pure file moves (git tracks all 16 as renames).

`dev-docs/ip-bans-and-allowlists.md` documents the subsystem, leading with the operator process for unblocking a player — including the three things that look sufficient and are not: deleting the CrowdSec decision alone, editing `ip-blocklist.txt` by hand, and `cscli allowlists` alone. `.gitignore` covers the new config files.

## Testing

Build clean. **Server.Tests 810 passed**, **UOContent.Tests 637 passed**, zero warnings. This branch adds 38 tests; the rest of the delta is main's, since this is rebased on current `main`.

New coverage: TTL boundary and renewal, private-address exclusion, manual-ban-never-exempt, unopted-reason-never-exempt, strike revocation, quiet-window reset, login forgiveness, file-allowlist CIDR coverage, file-allowlist not spending the earned list's strikes, denylist expiry-on-read, cap enforcement, lapsed-entry reclaim, HTTP/TLS/SSH identification, seed-collision fall-through, and encrypted-login-is-not-foreign.

Generator verified end-to-end against live feeds: a clean run ships no carve-out, `-AddCarveout starlink -Asn 14593` fetches and collapses 213 prefixes to 115 ranges in 0.1s over 4.2M entries, `-RefreshCarveouts` rediscovers it by its `asn=` marker, a hand-written allowlist is left untouched, and deleting a carve-out drops it rather than having it rewritten. CIDR splitting verified exhaustively: a single-IP hole in a /24 leaves exactly 255 of 256 addresses blocked.

## Operator note

Existing installs are unaffected until the generator next runs, which creates `ip-allowlist.txt` and nothing else. To unblock someone: add the address to that file and delete any live CrowdSec decision — the existing ban outlives the config change. The shard picks the entry up on its next reload, so re-running the generator is optional.

A shard whose players are on CGNAT (satellite, mobile, or an ISP short on IPv4) will likely also want `-AddCarveout`; see `dev-docs/ip-bans-and-allowlists.md`.

## Also included: a latent CI failure this PR surfaced

`fix(tests): serialize test classes that rent through STArrayPool` touches a property-list test file that has nothing to do with this feature. It is here because it was failing macOS CI, and it is trivially cherry-pickable out if you would rather it went to `main` on its own — **which may be the better call, since it is failing `main` today.**

CI has since gone green with it applied.

`STArrayPool` is single-threaded by design and its bucket cache is a plain `static`, not `[ThreadStatic]`, with a check-then-act initialize in `Return()`:

```csharp
var cacheBuckets = _cacheBuckets ?? InitializeBuckets();
```

Two threads both see null, both initialize, and the loser trips `Debug.Assert(_cacheBuckets is null)`. Anything renting from it has to stay off parallel test threads — which is what the `DisableParallelization` collections are for.

- `ObjectPropertyListReentrancyTests` and `ObjectPropertyListNestedBuildTests` (added in #2555) build property lists, which rent the interpolation buffer, but were not in the sequential collection — unlike `PropertyListInvalidationDuringBuildTests` in the same file. This is a **latent failure already on `main`**; it is timing-dependent, so it shows on some platforms and not others.
- `AutoDenylistTests` (added here) has the same exposure: its cap tests reach `AutoDenylist.Sweep`, which rents a `PooledRefList` without `mt`. The blocklist tests need no marking because `BlocklistSnapshot.Build` asks for the `mt` pool explicitly.

No production change — `STArrayPool` is the right pool on the game loop, where both `Sweep` and the property list actually run.

## Deliberately not included

Waiting for a fragmented four-byte seed at `AwaitingSeed`. It looked like a bug but the disconnect is a deliberate defence: only pre-0xEF clients reach it (0xEF goes through `HandlePacket`, which already waits for its 21 bytes), and waiting converts an instant drop into a full 5s slot hold for a client sending one or two bytes, or a loris dribbling a byte every few seconds. Against a fixed 4096-entry `MaxConnections` table that trades capacity that matters for a fragmentation case a reconnect already fixes.
2026-07-30 23:12:17 -07:00

606 lines
21 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>
/// 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);
}
}
}