From: Forest Date: Wed, 29 Apr 2020 22:46:21 +0000 (-0700) Subject: Fix a severe issue with redundant acks X-Git-Tag: 1.0.0~27^2~7 X-Git-Url: https://git.deb.at/?a=commitdiff_plain;h=a8b72ec7a989aee3fd360c71b9f63b6111709452;p=rhonda%2Fimpostor.hazel.git Fix a severe issue with redundant acks Fix a potential issue in ObjectPool on Unity IL2CPP Improve naming of thread limited classes. --- diff --git a/Hazel/FewerThreads/ThreadLimitedUdpConnectionListener.cs b/Hazel/FewerThreads/ThreadLimitedUdpConnectionListener.cs new file mode 100644 index 0000000..df5e000 --- /dev/null +++ b/Hazel/FewerThreads/ThreadLimitedUdpConnectionListener.cs @@ -0,0 +1,322 @@ +using System; +using System.Collections.Concurrent; +using System.Collections.Generic; +using System.Net; +using System.Net.Sockets; +using System.Threading; + +namespace Hazel.Udp.FewerThreads +{ + /// + /// Listens for new UDP connections and creates UdpConnections for them. + /// + /// + public class ThreadLimitedUdpConnectionListener : IDisposable + { + private struct SendMessageInfo + { + public byte[] Buffer; + public EndPoint Recipient; + } + + private struct ReceiveMessageInfo + { + public MessageReader Message; + public EndPoint Sender; + } + + private const int SendReceiveBufferSize = 1024 * 1024; + private const int BufferSize = ushort.MaxValue; + + public event Action NewConnection; + + /// + /// A callback for early connection rejection. + /// * Return false to reject connection. + /// * A null response is ok, we just won't send anything. + /// + public AcceptConnectionCheck AcceptConnection; + public delegate bool AcceptConnectionCheck(IPEndPoint endPoint, byte[] input, out byte[] response); + + private Socket socket; + private ILogger Logger; + + public IPEndPoint EndPoint { get; } + public IPMode IPMode { get; } + + private Thread reliablePacketThread; + private Thread receiveThread; + private Thread sendThread; + private HazelThreadPool processThreads; + + private ConcurrentDictionary allConnections = new ConcurrentDictionary(); + + private Queue receiveQueue = new Queue(); + private Queue sendQueue = new Queue(); + + public int ConnectionCount { get { return this.allConnections.Count; } } + public int SendQueueLength { get { lock(this.sendQueue) return this.sendQueue.Count; } } + public int ReceiveQueueLength { get { lock (this.receiveQueue) return this.receiveQueue.Count; } } + + private bool isActive; + + public ThreadLimitedUdpConnectionListener(int numWorkers, IPEndPoint endPoint, ILogger logger, IPMode ipMode = IPMode.IPv4) + { + this.Logger = logger; + this.EndPoint = endPoint; + this.IPMode = ipMode; + + this.socket = UdpConnection.CreateSocket(this.IPMode); + this.socket.Blocking = false; + + this.socket.ReceiveBufferSize = SendReceiveBufferSize; + this.socket.SendBufferSize = SendReceiveBufferSize; + + this.reliablePacketThread = new Thread(ManageReliablePackets); + this.sendThread = new Thread(SendLoop); + this.receiveThread = new Thread(ReceiveLoop); + this.processThreads = new HazelThreadPool(numWorkers, ProcessingLoop); + } + + ~ThreadLimitedUdpConnectionListener() + { + this.Dispose(false); + } + + private void ManageReliablePackets() + { + while (this.isActive) + { + foreach (var kvp in this.allConnections) + { + var sock = kvp.Value; + sock.ManageReliablePackets(); + } + + Thread.Sleep(100); + } + } + + public void Start() + { + try + { + socket.Bind(EndPoint); + } + catch (SocketException e) + { + throw new HazelException("Could not start listening as a SocketException occurred", e); + } + + this.isActive = true; + this.reliablePacketThread.Start(); + this.sendThread.Start(); + this.receiveThread.Start(); + this.processThreads.Start(); + } + + private void ReceiveLoop() + { + while (this.isActive) + { + if (this.socket.Poll(Timeout.Infinite, SelectMode.SelectRead)) + { + EndPoint remoteEP = new IPEndPoint(this.EndPoint.Address, this.EndPoint.Port); + MessageReader message = MessageReader.GetSized(BufferSize); + try + { + message.Length = socket.ReceiveFrom(message.Buffer, 0, message.Buffer.Length, SocketFlags.None, ref remoteEP); + } + catch (SocketException sx) + { + message.Recycle(); + this.Logger.WriteError("Socket Ex in StartListening: " + sx.Message); + continue; + } + catch (Exception ex) + { + message.Recycle(); + this.Logger.WriteError("Stopped due to: " + ex.Message); + return; + } + + lock (this.receiveQueue) + { + this.receiveQueue.Enqueue(new ReceiveMessageInfo() { Message = message, Sender = remoteEP }); + Monitor.Pulse(this.receiveQueue); + } + } + } + } + + private void ProcessingLoop() + { + while (this.isActive) + { + ReceiveMessageInfo msg; + lock (this.receiveQueue) + { + if (this.receiveQueue.Count == 0) + { + Monitor.Wait(this.receiveQueue); + + if (this.receiveQueue.Count == 0) + { + continue; + } + } + + msg = this.receiveQueue.Dequeue(); + } + + try + { + this.ReadCallback(msg.Message, msg.Sender); + } + catch + { + } + } + } + + private void SendLoop() + { + while (this.isActive) + { + SendMessageInfo msg; + lock (this.sendQueue) + { + if (this.sendQueue.Count == 0) + { + Monitor.Wait(this.sendQueue); + + if (this.sendQueue.Count == 0) + { + continue; + } + } + + msg = this.sendQueue.Dequeue(); + } + + try + { + this.socket.SendTo(msg.Buffer, 0, msg.Buffer.Length, SocketFlags.None, msg.Recipient); + } + catch { } + } + } + + void ReadCallback(MessageReader message, EndPoint remoteEndPoint) + { + int bytesReceived = message.Length; + bool aware = true; + bool isHello = message.Buffer[0] == (byte)UdpSendOption.Hello; + + // If we're aware of this connection use the one already + // If this is a new client then connect with them! + ThreadLimitedUdpServerConnection connection; + if (!this.allConnections.TryGetValue(remoteEndPoint, out connection)) + { + lock (this.allConnections) + { + if (!this.allConnections.TryGetValue(remoteEndPoint, out connection)) + { + // Check for malformed connection attempts + if (!isHello) + { + message.Recycle(); + return; + } + + if (AcceptConnection != null) + { + if (!AcceptConnection((IPEndPoint)remoteEndPoint, message.Buffer, out var response)) + { + message.Recycle(); + if (response != null) + { + SendDataRaw(response, remoteEndPoint); + } + + return; + } + } + + aware = false; + connection = new ThreadLimitedUdpServerConnection(this, (IPEndPoint)remoteEndPoint, this.IPMode); + if (!this.allConnections.TryAdd(remoteEndPoint, connection)) + { + throw new HazelException("Failed to add a connection. This should never happen."); + } + } + } + } + + // If it's a new connection invoke the NewConnection event. + // This needs to happen before handling the message because in localhost scenarios, the ACK and + // subsequent messages can happen before the NewConnection event sets up OnDataRecieved handlers + if (!aware) + { + // Skip header and hello byte; + message.Offset = 4; + message.Length = bytesReceived - 4; + message.Position = 0; + this.NewConnection?.Invoke(new NewConnectionEventArgs(message, connection)); + } + + // Inform the connection of the buffer (new connections need to send an ack back to client) + connection.HandleReceive(message, bytesReceived); + + if (isHello) + { + message.Recycle(); + } + } + + internal void SendDataRaw(byte[] response, EndPoint remoteEndPoint) + { + lock (this.sendQueue) + { + this.sendQueue.Enqueue(new SendMessageInfo() { Buffer = response, Recipient = remoteEndPoint }); + Monitor.Pulse(this.sendQueue); + } + } + + /// + /// Removes a virtual connection from the list. + /// + /// The endpoint of the virtual connection. + internal bool RemoveConnectionTo(EndPoint endPoint) + { + return this.allConnections.TryRemove(endPoint, out var conn); + } + + protected virtual void Dispose(bool disposing) + { + foreach (var kvp in this.allConnections) + { + kvp.Value.Dispose(); + } + + try { this.socket.Shutdown(SocketShutdown.Both); } catch { } + try { this.socket.Close(); } catch { } + try { this.socket.Dispose(); } catch { } + + this.isActive = false; + + lock (this.sendQueue) Monitor.PulseAll(this.sendQueue); + lock (this.receiveQueue) Monitor.PulseAll(this.receiveQueue); + + this.reliablePacketThread.Join(); + this.sendThread.Join(); + this.receiveThread.Join(); + this.processThreads.Join(); + } + + public void Dispose() + { + this.Dispose(true); + } + } +} diff --git a/Hazel/FewerThreads/ThreadLimitedUdpServerConnection.cs b/Hazel/FewerThreads/ThreadLimitedUdpServerConnection.cs new file mode 100644 index 0000000..ee5d2cd --- /dev/null +++ b/Hazel/FewerThreads/ThreadLimitedUdpServerConnection.cs @@ -0,0 +1,101 @@ +using System; +using System.Net; + +namespace Hazel.Udp.FewerThreads +{ + /// + /// Represents a servers's connection to a client that uses the UDP protocol. + /// + /// + internal sealed class ThreadLimitedUdpServerConnection : UdpConnection + { + /// + /// The connection listener that we use the socket of. + /// + /// + /// Udp server connections utilize the same socket in the listener for sends/receives, this is the listener that + /// created this connection and is hence the listener this conenction sends and receives via. + /// + public ThreadLimitedUdpConnectionListener Listener { get; private set; } + + /// + /// Creates a UdpConnection for the virtual connection to the endpoint. + /// + /// The listener that created this connection. + /// The endpoint that we are connected to. + /// The IPMode we are connected using. + internal ThreadLimitedUdpServerConnection(ThreadLimitedUdpConnectionListener listener, IPEndPoint endPoint, IPMode IPMode) + : base() + { + this.Listener = listener; + this.RemoteEndPoint = endPoint; + this.EndPoint = endPoint; + this.IPMode = IPMode; + + State = ConnectionState.Connected; + this.InitializeKeepAliveTimer(); + } + + /// + protected override void WriteBytesToConnection(byte[] bytes, int length) + { + if (bytes.Length != length) throw new ArgumentException("I made an assumption here. I hope you see this error."); + + Listener.SendDataRaw(bytes, RemoteEndPoint); + } + + /// + /// + /// This will always throw a HazelException. + /// + public override void Connect(byte[] bytes = null, int timeout = 5000) + { + throw new InvalidOperationException("Cannot manually connect a UdpServerConnection, did you mean to use UdpClientConnection?"); + } + + /// + /// + /// This will always throw a HazelException. + /// + public override void ConnectAsync(byte[] bytes = null) + { + throw new InvalidOperationException("Cannot manually connect a UdpServerConnection, did you mean to use UdpClientConnection?"); + } + + /// + /// Sends a disconnect message to the end point. + /// + protected override bool SendDisconnect(MessageWriter data = null) + { + if (!Listener.RemoveConnectionTo(RemoteEndPoint)) return false; + this._state = ConnectionState.NotConnected; + + var bytes = EmptyDisconnectBytes; + if (data != null && data.Length > 0) + { + if (data.SendOption != SendOption.None) throw new ArgumentException("Disconnect messages can only be unreliable."); + + bytes = data.ToByteArray(true); + bytes[0] = (byte)UdpSendOption.Disconnect; + } + + try + { + Listener.SendDataRaw(bytes, RemoteEndPoint); + } + catch { } + + return true; + } + + protected override void Dispose(bool disposing) + { + if (disposing) + { + SendDisconnect(); + } + + base.Dispose(disposing); + } + } +} diff --git a/Hazel/FewerThreads/UdpConnectionListener2.cs b/Hazel/FewerThreads/UdpConnectionListener2.cs deleted file mode 100644 index d4e8876..0000000 --- a/Hazel/FewerThreads/UdpConnectionListener2.cs +++ /dev/null @@ -1,319 +0,0 @@ -using System; -using System.Collections.Concurrent; -using System.Collections.Generic; -using System.Net; -using System.Net.Sockets; -using System.Threading; - -namespace Hazel.Udp.FewerThreads -{ - /// - /// Listens for new UDP connections and creates UdpConnections for them. - /// - /// - public class UdpConnectionListener2 : IDisposable - { - private struct SendMessageInfo - { - public byte[] Buffer; - public EndPoint Recipient; - } - - private struct ReceiveMessageInfo - { - public MessageReader Message; - public EndPoint Sender; - } - - private const int SendReceiveBufferSize = 1024 * 1024; - private const int BufferSize = ushort.MaxValue; - - public event Action NewConnection; - - /// - /// A callback for early connection rejection. - /// * Return false to reject connection. - /// * A null response is ok, we just won't send anything. - /// - public AcceptConnectionCheck AcceptConnection; - public delegate bool AcceptConnectionCheck(IPEndPoint endPoint, byte[] input, out byte[] response); - - private Socket socket; - private ILogger Logger; - - public IPEndPoint EndPoint { get; } - public IPMode IPMode { get; } - - private Thread reliablePacketThread; - private Thread receiveThread; - private Thread sendThread; - private HazelThreadPool processThreads; - - private ConcurrentDictionary allConnections = new ConcurrentDictionary(); - - private Queue receiveQueue = new Queue(); - private Queue sendQueue = new Queue(); - - public int ConnectionCount { get { return this.allConnections.Count; } } - public int SendQueueLength { get { lock(this.sendQueue) return this.sendQueue.Count; } } - public int ReceiveQueueLength { get { lock (this.receiveQueue) return this.receiveQueue.Count; } } - - private bool isActive; - - public UdpConnectionListener2(IPEndPoint endPoint, ILogger logger, IPMode ipMode = IPMode.IPv4) - { - this.Logger = logger; - this.EndPoint = endPoint; - this.IPMode = ipMode; - - this.socket = UdpConnection.CreateSocket(this.IPMode); - this.socket.Blocking = false; - - this.socket.ReceiveBufferSize = SendReceiveBufferSize; - this.socket.SendBufferSize = SendReceiveBufferSize; - - this.reliablePacketThread = new Thread(ManageReliablePackets); - this.sendThread = new Thread(SendLoop); - this.receiveThread = new Thread(ReceiveLoop); - this.processThreads = new HazelThreadPool(4, ProcessingLoop); - } - - ~UdpConnectionListener2() - { - this.Dispose(false); - } - - private void ManageReliablePackets() - { - while (this.isActive) - { - foreach (var kvp in this.allConnections) - { - var sock = kvp.Value; - sock.ManageReliablePackets(); - } - - Thread.Sleep(100); - } - } - - public void Start() - { - try - { - socket.Bind(EndPoint); - } - catch (SocketException e) - { - throw new HazelException("Could not start listening as a SocketException occurred", e); - } - - this.isActive = true; - this.reliablePacketThread.Start(); - this.sendThread.Start(); - this.receiveThread.Start(); - this.processThreads.Start(); - } - - private void ReceiveLoop() - { - while (this.isActive) - { - if (this.socket.Poll(Timeout.Infinite, SelectMode.SelectRead)) - { - EndPoint remoteEP = new IPEndPoint(IPMode == IPMode.IPv4 ? IPAddress.Any : IPAddress.IPv6Any, this.EndPoint.Port); - MessageReader message = MessageReader.GetSized(BufferSize); - try - { - message.Length = socket.ReceiveFrom(message.Buffer, 0, message.Buffer.Length, SocketFlags.None, ref remoteEP); - } - catch (SocketException sx) - { - message.Recycle(); - this.Logger.WriteError("Socket Ex in StartListening: " + sx.Message); - continue; - } - catch (Exception ex) - { - message.Recycle(); - this.Logger.WriteError("Stopped due to: " + ex.Message); - return; - } - - lock (this.receiveQueue) - { - this.receiveQueue.Enqueue(new ReceiveMessageInfo() { Message = message, Sender = remoteEP }); - Monitor.Pulse(this.receiveQueue); - } - } - } - } - - private void ProcessingLoop() - { - while (this.isActive) - { - ReceiveMessageInfo msg; - lock (this.receiveQueue) - { - if (this.receiveQueue.Count == 0) - { - Monitor.Wait(this.receiveQueue); - - if (this.receiveQueue.Count == 0) - { - continue; - } - } - - msg = this.receiveQueue.Dequeue(); - } - - try - { - this.ReadCallback(msg.Message, msg.Sender); - } - catch - { - } - } - } - - private void SendLoop() - { - while (this.isActive) - { - SendMessageInfo msg; - lock (this.sendQueue) - { - if (this.sendQueue.Count == 0) - { - Monitor.Wait(this.sendQueue); - - if (this.sendQueue.Count == 0) - { - continue; - } - } - - msg = this.sendQueue.Dequeue(); - } - - try - { - this.socket.SendTo(msg.Buffer, 0, msg.Buffer.Length, SocketFlags.None, msg.Recipient); - } - catch { } - } - } - - void ReadCallback(MessageReader message, EndPoint remoteEndPoint) - { - int bytesReceived = message.Length; - bool aware = true; - bool isHello = message.Buffer[0] == (byte)UdpSendOption.Hello; - - // If we're aware of this connection use the one already - // If this is a new client then connect with them! - UdpServerConnection2 connection; - if (!this.allConnections.TryGetValue(remoteEndPoint, out connection)) - { - lock (this.allConnections) - { - if (!this.allConnections.TryGetValue(remoteEndPoint, out connection)) - { - // Check for malformed connection attempts - if (!isHello) - { - message.Recycle(); - return; - } - - if (AcceptConnection != null) - { - if (!AcceptConnection((IPEndPoint)remoteEndPoint, message.Buffer, out var response)) - { - message.Recycle(); - if (response != null) - { - SendDataRaw(response, remoteEndPoint); - } - - return; - } - } - - aware = false; - connection = new UdpServerConnection2(this, (IPEndPoint)remoteEndPoint, this.IPMode); - if (!this.allConnections.TryAdd(remoteEndPoint, connection)) - { - throw new HazelException("Failed to add a connection. This should never happen."); - } - } - } - } - - //Inform the connection of the buffer (new connections need to send an ack back to client) - connection.HandleReceive(message, bytesReceived); - - //If it's a new connection invoke the NewConnection event. - if (!aware) - { - // Skip header and hello byte; - message.Offset = 4; - message.Length = bytesReceived - 4; - message.Position = 0; - this.NewConnection?.Invoke(new NewConnectionEventArgs(message, connection)); - } - else if (isHello) - { - message.Recycle(); - } - } - - internal void SendDataRaw(byte[] response, EndPoint remoteEndPoint) - { - lock (this.sendQueue) - { - this.sendQueue.Enqueue(new SendMessageInfo() { Buffer = response, Recipient = remoteEndPoint }); - Monitor.Pulse(this.sendQueue); - } - } - - /// - /// Removes a virtual connection from the list. - /// - /// The endpoint of the virtual connection. - internal bool RemoveConnectionTo(EndPoint endPoint) - { - return this.allConnections.TryRemove(endPoint, out var conn); - } - - protected virtual void Dispose(bool disposing) - { - foreach (var kvp in this.allConnections) - { - kvp.Value.Dispose(); - } - - try { this.socket.Shutdown(SocketShutdown.Both); } catch { } - try { this.socket.Close(); } catch { } - try { this.socket.Dispose(); } catch { } - - this.isActive = false; - - lock (this.sendQueue) Monitor.PulseAll(this.sendQueue); - lock (this.receiveQueue) Monitor.PulseAll(this.receiveQueue); - - this.reliablePacketThread.Join(); - this.sendThread.Join(); - this.receiveThread.Join(); - this.processThreads.Join(); - } - - public void Dispose() - { - this.Dispose(true); - } - } -} diff --git a/Hazel/FewerThreads/UdpServerConnection2.cs b/Hazel/FewerThreads/UdpServerConnection2.cs deleted file mode 100644 index c5e6795..0000000 --- a/Hazel/FewerThreads/UdpServerConnection2.cs +++ /dev/null @@ -1,101 +0,0 @@ -using System; -using System.Net; - -namespace Hazel.Udp.FewerThreads -{ - /// - /// Represents a servers's connection to a client that uses the UDP protocol. - /// - /// - internal sealed class UdpServerConnection2 : UdpConnection - { - /// - /// The connection listener that we use the socket of. - /// - /// - /// Udp server connections utilize the same socket in the listener for sends/receives, this is the listener that - /// created this connection and is hence the listener this conenction sends and receives via. - /// - public UdpConnectionListener2 Listener { get; private set; } - - /// - /// Creates a UdpConnection for the virtual connection to the endpoint. - /// - /// The listener that created this connection. - /// The endpoint that we are connected to. - /// The IPMode we are connected using. - internal UdpServerConnection2(UdpConnectionListener2 listener, IPEndPoint endPoint, IPMode IPMode) - : base() - { - this.Listener = listener; - this.RemoteEndPoint = endPoint; - this.EndPoint = endPoint; - this.IPMode = IPMode; - - State = ConnectionState.Connected; - this.InitializeKeepAliveTimer(); - } - - /// - protected override void WriteBytesToConnection(byte[] bytes, int length) - { - if (bytes.Length != length) throw new ArgumentException("I made an assumption here. I hope you see this error."); - - Listener.SendDataRaw(bytes, RemoteEndPoint); - } - - /// - /// - /// This will always throw a HazelException. - /// - public override void Connect(byte[] bytes = null, int timeout = 5000) - { - throw new InvalidOperationException("Cannot manually connect a UdpServerConnection, did you mean to use UdpClientConnection?"); - } - - /// - /// - /// This will always throw a HazelException. - /// - public override void ConnectAsync(byte[] bytes = null) - { - throw new InvalidOperationException("Cannot manually connect a UdpServerConnection, did you mean to use UdpClientConnection?"); - } - - /// - /// Sends a disconnect message to the end point. - /// - protected override bool SendDisconnect(MessageWriter data = null) - { - if (!Listener.RemoveConnectionTo(RemoteEndPoint)) return false; - this._state = ConnectionState.NotConnected; - - var bytes = EmptyDisconnectBytes; - if (data != null && data.Length > 0) - { - if (data.SendOption != SendOption.None) throw new ArgumentException("Disconnect messages can only be unreliable."); - - bytes = data.ToByteArray(true); - bytes[0] = (byte)UdpSendOption.Disconnect; - } - - try - { - Listener.SendDataRaw(bytes, RemoteEndPoint); - } - catch { } - - return true; - } - - protected override void Dispose(bool disposing) - { - if (disposing) - { - SendDisconnect(); - } - - base.Dispose(disposing); - } - } -} diff --git a/Hazel/Hazel.csproj b/Hazel/Hazel.csproj index e249199..f6b10e6 100644 --- a/Hazel/Hazel.csproj +++ b/Hazel/Hazel.csproj @@ -28,7 +28,7 @@ 7.3 - pdbonly + portable true bin\Release\ TRACE @@ -39,6 +39,7 @@ true false 7.3 + true true @@ -73,8 +74,8 @@ - - + + @@ -103,7 +104,6 @@ -