namespace Supercell.Laser.Server.Message { using Supercell.Laser.Server.Networking; using Supercell.Laser.Server.Message; using System.Collections.Concurrent; using Supercell.Laser.Logic.Message; using Supercell.Laser.Logic.Util; public static class Processor { private static ConcurrentQueue IncomingQueue; private static ConcurrentQueue OutgoingQueue; private static ManualResetEvent ReceiveEvent; private static ManualResetEvent SendEvent; private static Thread ReceiveThread; private static Thread SendThread; private struct QueueItem { public readonly Connection Connection; public readonly GameMessage Message; public QueueItem(Connection connection, GameMessage message) { Connection = connection; Message = message; } } public static void Init() { IncomingQueue = new ConcurrentQueue(); OutgoingQueue = new ConcurrentQueue(); ReceiveEvent = new ManualResetEvent(false); SendEvent = new ManualResetEvent(false); ReceiveThread = new Thread(UpdateReceive); SendThread = new Thread(UpdateSend); ReceiveThread.Start(); SendThread.Start(); } public static bool Receive(Connection connection, GameMessage message) { if (message == null) return false; ServerDiagnostics.RecordIncomingPacket(message); if (ServerDiagnostics.ShouldLogPacket(message.GetMessageType())) { Logger.Packet($"RECV {(connection == null ? "conn#?" : connection.Describe())} <= {ServerDiagnostics.DescribeMessage(message)}"); } if (IncomingQueue.Count >= 1024) { Logger.Warning($"Processor: Incoming message queue full. {ServerDiagnostics.DescribeMessage(message)} discarded."); return false; } IncomingQueue.Enqueue(new QueueItem(connection, message)); ReceiveEvent.Set(); return true; } public static void Send(Connection connection, GameMessage message) { if (message == null) return; ServerDiagnostics.RecordOutgoingPacket(message); if (ServerDiagnostics.ShouldLogPacket(message.GetMessageType())) { Logger.Packet($"SEND {(connection == null ? "conn#?" : connection.Describe())} => {ServerDiagnostics.DescribeMessage(message)}"); } if (OutgoingQueue.Count >= 1024) { Logger.Warning($"Processor: Outgoing message queue full. {ServerDiagnostics.DescribeMessage(message)} discarded."); return; } OutgoingQueue.Enqueue(new QueueItem(connection, message)); SendEvent.Set(); } public static void UpdateReceive() { while (true) { ReceiveEvent.WaitOne(); while (IncomingQueue.TryDequeue(out QueueItem item)) { try { item.Connection.MessageManager.ReceiveMessage(item.Message); } catch (Exception exception) { Logger.Exception($"Incoming message handling failed for {(item.Connection == null ? "conn#?" : item.Connection.Describe())} {ServerDiagnostics.DescribeMessage(item.Message)}", exception); } } ReceiveEvent.Reset(); } } public static void UpdateSend() { while (true) { SendEvent.WaitOne(); while (OutgoingQueue.TryDequeue(out QueueItem item)) { try { item.Connection.Messaging.EncryptAndWrite(item.Message); } catch (Exception exception) { Logger.Exception($"Outgoing message handling failed for {(item.Connection == null ? "conn#?" : item.Connection.Describe())} {ServerDiagnostics.DescribeMessage(item.Message)}", exception); } } SendEvent.Reset(); } } } }