feat: Upgrades networking to use io_uring. (#2315)

> [!IMPORTANT]
> **Breaking Changes**
> - DecodePacket and EncodePacket delegates replaced with IClientEncryption interface
> - NetState.Connection (Socket) replaced with internal RingSocket management
> - NetState.RecvPipe and NetState.SendPipe removed (buffers managed internally)

## Summary

Upgrades the networking stack from PollGroup-based I/O to io_uring, significantly improving I/O performance on Linux.
This also adds native client encryption support for encrypted UO clients.

## Major Changes

io_uring Networking Architecture
- Replaced PollGroup with IORingGroup for async socket I/O operations
- Removed Pipe.cs (mirrored ring buffer) and TcpServer.cs in favor of RingSocketManager
- Added NetState.Network.cs - centralized network infrastructure handling accept, recv, send, and disconnect
completions
- Added SocketHelper.cs - platform-specific socket utilities for raw socket handle operations (getpeername,
getsockname)
- Buffer management now handled by RingSocketManager with configurable slab allocation

### Client Encryption Support
- Added full encryption stack in Network/Encryption/:
  - EncryptionConfig.cs - configurable encryption modes (None, Unencrypted, Encrypted, Both)
  - EncryptionManager.cs - encryption detection and initialization for login/game packets
  - LoginEncryption.cs - handles login packet encryption with version-derived keys
  - GameEncryption.cs - handles game server encryption using Twofish
  - TwofishEngine.cs - optimized Twofish block cipher implementation
  - LoginKeys.cs - encryption key table for client versions
  - IClientEncryption.cs - interface for client encryption implementations

### NetState Improvements
- Replaced Socket Connection with RingSocket _socket for managed socket lifecycle
- Changed from GCHandle polling to event-based completion processing
- Disconnect handling now properly waits for pending sends to flush
- Simplified connecting socket management using lazy queue removal

### Configuration
- New settings: network.encryptionMode and network.encryptionDebug
- Encryption mode flags: Unencrypted, Encrypted, or Both

### Dependencies
- Replaced PollGroup NuGet package with IORingGroup
- Linux requires liburing-dev / liburing-devel package

### Test plan

- Verify server starts and accepts connections on Linux with io_uring
- Verify server starts and accepts connections on Windows (fallback to IOCP)
- Test unencrypted client connections (ClassicUO with encryption disabled)
- Test encrypted client connections if available
- Verify graceful disconnect flushes pending data
- Confirm CI builds pass on all target platforms
This commit is contained in:
Kamron Batman 2026-02-01 16:02:32 -08:00 committed by GitHub
parent 91a553b0bc
commit 3c0d6cb9d6
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
70 changed files with 2227 additions and 1468 deletions

View file

@ -32,7 +32,7 @@ public static class DumpNetStates
foreach (var ns in NetState.Instances)
{
file.WriteLine($"{ns}, {ns.ConnectedOn}, {ns.NextActivityCheck}, {ns.Connection.Connected}, {ns._protocolState}, {ns._parserState}");
file.WriteLine($"{ns}, {ns.ConnectedOn}, {ns.NextActivityCheck}, {ns.IsConnected}, {ns._protocolState}, {ns._parserState}");
}
}
}

View file

@ -0,0 +1,498 @@
/*************************************************************************
* ModernUO *
* Copyright 2019-2025 - 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
// 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;
/// <summary>
/// Gets the IORingGroup instance for socket operations.
/// </summary>
public static IIORingGroup Ring => _socketManager?.Ring;
/// <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);
}
else if (Firewall.IsBlocked(remoteIP))
{
logger.Debug("{Address} Firewalled", remoteIP);
}
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);
}
}
}
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()
{
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 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)
{
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 = Core.TickCount + 90000;
}
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.Dispose();
}
}
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);
}
}
}

View file

@ -1,6 +1,6 @@
/*************************************************************************
* ModernUO *
* Copyright 2019-2023 - ModernUO Development Team *
* Copyright 2019-2025 - ModernUO Development Team *
* Email: hi@modernuo.com *
* File: NetState.cs *
* *
@ -25,55 +25,46 @@ using System.Collections.Concurrent;
using System.Collections.Generic;
using System.IO;
using System.Net;
using System.Net.Sockets;
using System.Network;
using System.Runtime.CompilerServices;
using System.Runtime.InteropServices;
namespace Server.Network;
public delegate void DecodePacket(Span<byte> buffer, ref int length);
public delegate int EncodePacket(ReadOnlySpan<byte> inputBuffer, Span<byte> outputBuffer);
public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetState>, IDisposable
{
private static readonly ILogger logger = LogFactory.GetLogger(typeof(NetState));
private static readonly TimeSpan ConnectingSocketIdleLimit = TimeSpan.FromMilliseconds(5000); // 5 seconds
private const int RecvPipeSize = 1024 * 64;
private const int SendPipeSize = 1024 * 256;
private const int HuePickerCap = 512;
private const int MenuCap = 512;
private const int PacketPerSecondThreshold = 3000;
private static readonly GCHandle[] _polledStates = new GCHandle[2048];
private static readonly IPollGroup _pollGroup = PollGroup.Create();
private static readonly Queue<NetState> _flushPending = new(2048);
private static readonly Queue<NetState> _flushedPartials = new(256);
private static readonly Queue<NetState> _pendingDisconnects = new(256); // Processed AFTER flush
private static readonly ConcurrentQueue<NetState> _disposed = new();
private static readonly Queue<NetState> _throttled = new(256);
private static readonly Queue<NetState> _throttledPending = new(256);
private static readonly SortedSet<NetState> _connecting = new(NetStateConnectingComparer.Instance);
private static readonly Queue<NetState> _connectingQueue = new(2048);
private static readonly HashSet<NetState> _instances = new(2048);
public static IReadOnlySet<NetState> Instances => _instances;
private readonly string _toString;
private ClientVersion _version;
private bool _running = true;
private volatile DecodePacket _packetDecoder;
private volatile EncodePacket _packetEncoder;
private IClientEncryption _encryption;
private bool _flushQueued;
private bool _disconnectQueued; // Queued for disconnect processing (after flush)
private long[] _packetThrottles;
private long[] _packetCounts;
private string _disconnectReason = string.Empty;
internal ParserState _parserState = ParserState.AwaitingNextPacket;
internal ProtocolState _protocolState = ProtocolState.AwaitingSeed;
internal GCHandle _handle;
private bool _packetLogging;
public GCHandle Handle => _handle;
// Managed socket with buffers (handles lifecycle automatically)
internal RingSocket _socket;
// Speed Hack Prevention
internal long _movementCredit;
@ -108,6 +99,9 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
public static void Configure()
{
_packetLoggingPath = ServerConfiguration.GetSetting("netstate.packetLoggingPath", Path.Combine(Core.BaseDirectory, "Packets"));
// Initialize IORingGroup and buffer pools
ConfigureNetwork();
}
public static void Initialize()
@ -115,45 +109,24 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
Timer.DelayCall(TimeSpan.FromMinutes(1), TimeSpan.FromMinutes(1.5), CheckAllAlive);
}
public NetState(Socket connection)
// Internal constructor for accepted sockets
private NetState(RingSocket socket, IPAddress address)
{
Connection = connection;
_socket = socket;
Address = address;
Seeded = false;
HuePickers = [];
Menus = [];
Trades = [];
RecvPipe = new Pipe(RecvPipeSize);
SendPipe = new Pipe(SendPipeSize);
NextActivityCheck = Core.TickCount + 30000;
ConnectedOn = Core.Now;
try
{
Address = Utility.Intern((Connection?.RemoteEndPoint as IPEndPoint)?.Address);
_toString = Address?.ToString() ?? "(error)";
}
catch (Exception ex)
{
TraceException(ex);
Address = IPAddress.None;
_toString = "(error)";
}
_toString = address?.ToString() ?? "(error)";
_instances.Add(this);
_connecting.Add(this);
_handle = GCHandle.Alloc(this);
_connectingQueue.Enqueue(this);
LogInfo($"Connected. [{_instances.Count} Online]");
try
{
_pollGroup.Add(connection, _handle);
}
catch (Exception ex)
{
TraceException(ex);
Disconnect("Unable to add socket to poll group");
}
}
// Sectors
@ -188,16 +161,10 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
public IPAddress Address { get; }
public DecodePacket PacketDecoder
public IClientEncryption Encryption
{
get => _packetDecoder;
set => _packetDecoder = value;
}
public EncodePacket PacketEncoder
{
get => _packetEncoder;
set => _packetEncoder = value;
get => _encryption;
set => _encryption = value;
}
public int CurrentPacket { get; internal set; }
@ -210,13 +177,27 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
public bool Seeded { get; set; }
public Pipe RecvPipe { get; }
public Pipe SendPipe { get; }
public bool Running => _running;
public Socket Connection { get; private set; }
/// <summary>
/// Gets whether the socket is connected.
/// </summary>
public bool IsConnected => _running && _socket != null;
/// <summary>
/// Gets the socket handle.
/// </summary>
public nint SocketHandle => _socket?.Handle ?? 0;
/// <summary>
/// Gets the local endpoint (address/port) the client connected to.
/// </summary>
public IPEndPoint LocalEndPoint => _socket != null ? SocketHelper.GetLocalEndPoint(_socket.Handle) : null;
/// <summary>
/// Gets the send buffer for this connection.
/// </summary>
internal IORingBuffer SendBuffer => _socket?.SendBuffer;
public bool CompressionEnabled { get; set; }
@ -235,15 +216,7 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
public IAccount Account
{
get => _account;
set
{
if (_account != null)
{
_connecting.Remove(this);
}
_account = value;
}
set => _account = value;
}
public string Assistant { get; set; }
@ -464,8 +437,14 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
return false;
}
#endif
buffer = SendPipe.Writer.AvailableToWrite();
return !(SendPipe.Writer.IsClosed || buffer.Length <= 0);
if (!_running || _socket == null)
{
buffer = Span<byte>.Empty;
return false;
}
buffer = _socket.SendBuffer.GetWriteSpan();
return buffer.Length > 0;
}
public void Send(ReadOnlySpan<byte> span)
@ -483,21 +462,25 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
try
{
if (_packetEncoder != null)
// Apply encoding first (e.g., compression from UOContent)
if (CompressionEnabled)
{
length = _packetEncoder(span, buffer);
length = NetworkCompression.Compress(span, buffer);
}
else
{
span.CopyTo(buffer);
}
// Then encrypt (if encryption is enabled)
_encryption?.ServerEncrypt(buffer[..length]);
if (PacketLogging)
{
LogPacket(span, false);
}
SendPipe.Writer.Advance((uint)length);
_socket.SendBuffer.CommitWrite(length);
if (!_flushQueued)
{
@ -554,26 +537,34 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
}
}
public void HandleReceive(bool throttled = false)
private void DecryptRecvBuffer(int bytesReceived)
{
if (!_running)
if (_socket == null || _encryption == null)
{
return;
}
if (!throttled)
// Get the portion of the buffer that was just written (the new data)
var readSpan = _socket.RecvBuffer.GetReadSpan();
var newDataStart = Math.Max(0, readSpan.Length - bytesReceived);
_encryption?.ClientDecrypt(readSpan.Slice(newDataStart, bytesReceived));
}
public void HandleReceive(bool throttled = false)
{
if (!_running || _socket == null)
{
ReceiveData();
return;
}
var reader = RecvPipe.Reader;
// Data already in recv buffer from recv completion - no need to call ReceiveData
try
{
// Process as many packets as we can synchronously
while (_running && _parserState != ParserState.Error && _protocolState != ProtocolState.Error)
{
var buffer = reader.AvailableToRead();
var buffer = _socket.RecvBuffer.GetReadSpan();
var length = buffer.Length;
if (length <= 0)
@ -632,13 +623,63 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
case ProtocolState.LoginServer_AwaitingLogin:
{
if (packetId != 0x80)
// Check for unencrypted login packet
if (packetId == 0x80)
{
// Unencrypted - check if allowed
if (EncryptionManager.Enabled && !EncryptionManager.Mode.HasFlag(EncryptionMode.Unencrypted))
{
LogInfo("Unencrypted client rejected by encryption policy.");
HandleError(packetId, packetLength);
return;
}
_parserState = ParserState.ProcessingPacket;
_parserState = HandlePacket(packetReader, packetId, out packetLength);
if (_parserState == ParserState.AwaitingNextPacket)
{
_protocolState = ProtocolState.LoginServer_AwaitingServerSelect;
}
break;
}
// First byte isn't 0x80 - might be encrypted
if (!EncryptionManager.Enabled)
{
LogInfo("Possible encrypted client detected, disconnecting...");
HandleError(packetId, packetLength);
return;
}
// Need 62 bytes for login packet to attempt decryption
if (length < 62)
{
_parserState = ParserState.AwaitingPartialPacket;
break;
}
// Try to detect and decrypt encrypted login
if (!this.DetectLoginEncryption(buffer[..62], out var loginEncryption))
{
LogInfo("Encrypted client detection failed, disconnecting...");
HandleError(packetId, packetLength);
return;
}
// Decryption succeeded - set up encryption and process
if (loginEncryption != null)
{
_encryption = loginEncryption;
// Decrypt the buffer in place for processing
var mutableBuffer = _socket.RecvBuffer.GetReadSpan();
loginEncryption.ClientDecrypt(mutableBuffer[..62]);
}
// Now process as normal (first byte should now be 0x80)
packetReader = new SpanReader(buffer);
packetId = packetReader.ReadByte();
_parserState = ParserState.ProcessingPacket;
_parserState = HandlePacket(packetReader, packetId, out packetLength);
if (_parserState == ParserState.AwaitingNextPacket)
@ -680,17 +721,68 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
case ProtocolState.GameServer_AwaitingGameServerLogin:
{
// Some clients send 0x80 on game server connection
if (packetId == 0x80)
{
goto case ProtocolState.LoginServer_AwaitingLogin;
}
if (packetId != 0x91)
// Check for unencrypted game login packet
if (packetId == 0x91)
{
// Unencrypted - check if allowed
if (EncryptionManager.Enabled && !EncryptionManager.Mode.HasFlag(EncryptionMode.Unencrypted))
{
LogInfo("Unencrypted game client rejected by encryption policy.");
HandleError(packetId, packetLength);
return;
}
_parserState = ParserState.ProcessingPacket;
_parserState = HandlePacket(packetReader, packetId, out packetLength);
if (_parserState == ParserState.AwaitingNextPacket)
{
_protocolState = ProtocolState.GameServer_LoggedIn;
}
break;
}
// First byte isn't 0x91 - might be encrypted
if (!EncryptionManager.Enabled)
{
HandleError(packetId, packetLength);
return;
}
// Need 65 bytes for game login packet to attempt decryption
if (length < 65)
{
_parserState = ParserState.AwaitingPartialPacket;
break;
}
// Try to detect and decrypt encrypted game login
if (!this.DetectGameEncryption(buffer[..65], out var gameEncryption))
{
LogInfo("Encrypted game client detection failed, disconnecting...");
HandleError(packetId, packetLength);
return;
}
// Decryption succeeded - set up encryption and process
if (gameEncryption != null)
{
_encryption = gameEncryption;
// Decrypt the buffer in place for processing
var mutableBuffer = _socket.RecvBuffer.GetReadSpan();
gameEncryption.ClientDecrypt(mutableBuffer[..65]);
}
// Now process as normal (first byte should now be 0x91)
packetReader = new SpanReader(buffer);
packetId = packetReader.ReadByte();
_parserState = ParserState.ProcessingPacket;
_parserState = HandlePacket(packetReader, packetId, out packetLength);
if (_parserState == ParserState.AwaitingNextPacket)
@ -711,7 +803,7 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
if (_parserState is ParserState.AwaitingNextPacket)
{
reader.Advance((uint)packetLength);
_socket.RecvBuffer.CommitRead(packetLength);
}
else if (_parserState is ParserState.Throttled)
{
@ -844,225 +936,15 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
return ParserState.AwaitingNextPacket;
}
private bool Flush()
{
_flushQueued = false;
// We don't have a running check since we need to send the last bits of data even after a disconnect, but before a dispose.
if (Connection == null)
{
return true;
}
var reader = SendPipe.Reader;
var buffer = reader.AvailableToRead();
if (reader.IsClosed || buffer.Length == 0)
{
return true;
}
var bytesWritten = 0;
try
{
bytesWritten = Connection.Send(buffer, SocketFlags.None);
}
catch (SocketException ex)
{
if (ex.SocketErrorCode != SocketError.WouldBlock)
{
logger.Debug(ex, "Disconnected due to a socket exception");
Disconnect(string.Empty);
return true;
}
}
catch (Exception ex)
{
Disconnect($"Disconnected with error: {ex}");
TraceException(ex);
return true;
}
if (bytesWritten > 0)
{
NextActivityCheck = Core.TickCount + 90000;
reader.Advance((uint)bytesWritten);
}
return bytesWritten == buffer.Length;
}
private void DecodePacket(Span<byte> buffer, ref int length)
{
_packetDecoder?.Invoke(buffer, ref length);
}
private void ReceiveData()
{
var writer = RecvPipe.Writer;
var buffer = writer.AvailableToWrite();
if (writer.IsClosed || buffer.Length == 0)
{
return;
}
var bytesWritten = 0;
try
{
bytesWritten = Connection.Receive(buffer, SocketFlags.None);
}
catch (SocketException ex)
{
if (ex.ErrorCode is not 54 and not 89 and not 995)
{
logger.Debug(ex, "Disconnected due to a socket exception");
}
Disconnect(string.Empty);
}
catch (Exception ex)
{
Disconnect($"Disconnected with error: {ex}");
TraceException(ex);
}
if (bytesWritten <= 0)
{
Disconnect(string.Empty);
return;
}
DecodePacket(buffer, ref bytesWritten);
writer.Advance((uint)bytesWritten);
NextActivityCheck = Core.TickCount + 90000;
}
private static void DisconnectUnattachedSockets()
{
var now = Core.Now;
// Clear out any sockets that have been connecting for too long
while (_connecting.Count > 0)
{
var ns = _connecting.Min;
var socketTime = ns.ConnectedOn;
// If the socket has been connected for less than the limit, we can stop checking
if (now - socketTime < ConnectingSocketIdleLimit)
{
break;
}
// Socket must have finished the entire authentication process or be forcibly disconnected.
if (!ns.Running || !ns.SentFirstPacket || !ns.Seeded || ns.Account == null)
{
// Not sending a message because it will fill up the logs.
ns.Disconnect(null);
}
_connecting.Remove(ns);
}
}
public static void FlushAll()
{
while (_flushPending.Count != 0)
{
_flushPending.Dequeue()?.Flush();
}
}
public static void Slice()
{
DisconnectUnattachedSockets();
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());
}
var count = _pollGroup.Poll(_polledStates);
if (count > 0)
{
for (var i = 0; i < count; i++)
{
(_polledStates[i].Target as NetState)?.HandleReceive();
_polledStates[i] = default;
}
}
while (_flushPending.TryDequeue(out var ns))
{
if (!ns.Flush())
{
// Incomplete data, so we need to requeue
_flushedPartials.Enqueue(ns);
}
}
var hasDisposes = false;
while (_disposed.TryDequeue(out var ns))
{
hasDisposes = true;
ns.Dispose();
}
// If they weren't disconnected, requeue them
while (_flushedPartials.TryDequeue(out var ns))
{
if (ns.Running)
{
_flushPending.Enqueue(ns);
}
}
if (hasDisposes)
{
_pollGroup.Poll(_polledStates.Length);
}
}
public void CheckAlive(long curTicks)
{
if (Connection != null && NextActivityCheck - curTicks < 0)
if (_socket != null && NextActivityCheck - curTicks < 0)
{
LogInfo("Disconnecting due to inactivity...");
Disconnect("Disconnecting due to inactivity.");
}
}
public static void CheckAllAlive()
{
try
{
var curTicks = Core.TickCount;
foreach (var ns in Instances)
{
ns.CheckAlive(curTicks);
}
}
catch (Exception ex)
{
TraceException(ex);
}
}
public void Trace(ReadOnlySpan<byte> buffer)
{
// We don't have data, so nothing to trace
@ -1105,17 +987,24 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
Console.WriteLine(ex);
}
/// <summary>
/// Requests a graceful disconnect. The disconnect is queued and processed after the flush
/// queue in Slice(), ensuring Send() calls made in the same tick are processed first.
/// </summary>
public void Disconnect(string reason)
{
if (!_running)
if (!_running || _socket == null)
{
return;
}
_running = false;
_disconnectReason = reason;
_disposed.Enqueue(this);
if (!_disconnectQueued)
{
_disconnectQueued = true;
_pendingDisconnects.Enqueue(this);
}
}
public static void TraceDisconnect(string reason, string ip)
@ -1148,16 +1037,18 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
{
_running = false;
// It's possible we could queue for dispose multiple times
if (Connection == null)
if (_socket == null)
{
return;
}
TraceDisconnect(_disconnectReason, _toString);
// If still running, force immediate disconnect
if (_running)
{
throw new Exception("Disconnected a NetState that is still running.");
_running = false;
_socketManager?.DisconnectImmediate(_socket);
}
var m = Mobile;
@ -1167,21 +1058,17 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
}
_instances.Remove(this);
_connecting.Remove(this);
try
// Clear the NetState slot
var slotId = _socket.Id;
if (slotId >= 0 && slotId < _netStates.Length && _netStates[slotId] == this)
{
_pollGroup.Remove(Connection, _handle);
}
catch (Exception ex)
{
TraceException(ex);
_netStates[slotId] = null;
}
Connection.Close();
_handle.Free();
RecvPipe.Dispose();
SendPipe.Dispose();
// Note: RingSocketManager handles cleanup of ring resources (unregister, close, buffer release)
// when it processes the disconnect event. We just clear our reference.
_socket = null;
Mobile = null;
@ -1192,46 +1079,9 @@ public partial class NetState : IComparable<NetState>, IValueLinkListNode<NetSta
Account = null;
ServerInfo = null;
CityInfo = null;
Connection = null;
var count = _instances.Count;
LogInfo(a != null ? $"Disconnected. [{count} Online] [{a}]" : $"Disconnected. [{count} Online]");
}
private class NetStateConnectingComparer : IComparer<NetState>
{
public static readonly IComparer<NetState> Instance = new NetStateConnectingComparer();
public int Compare(NetState x, NetState y)
{
if (x == null && y == null)
{
return 0;
}
if (x == null)
{
return -1;
}
if (y == null)
{
return 1;
}
if (ReferenceEquals(x, y))
{
return 0;
}
var connectedOn = x.ConnectedOn.CompareTo(y.ConnectedOn);
if (connectedOn != 0)
{
return connectedOn;
}
return x.CompareTo(y);
}
}
}