From 263b21174b33b2006d2d69b6eaac4719a1a8e7d7 Mon Sep 17 00:00:00 2001 From: Forest Date: Sun, 16 Dec 2018 23:04:21 -0800 Subject: [PATCH] Remove unused fragmented code, reduce more allocations, fix tests that wouldn't fail --- Hazel.UnitTests/Hazel.UnitTests.csproj | 1 + Hazel.UnitTests/TestHelper.cs | 49 +++--- Hazel.UnitTests/UdpConnectionTests.cs | 30 +--- Hazel.UnitTests/UnitTest1.cs | 29 ++++ Hazel/Connection.cs | 24 ++- Hazel/ConnectionListener.cs | 12 +- Hazel/Hazel.csproj | 1 - Hazel/MessageReader.cs | 18 ++- Hazel/MessageWriter.cs | 28 +++- Hazel/SendOption.cs | 12 -- Hazel/Udp/UdpClientConnection.cs | 3 +- Hazel/Udp/UdpConnection.Fragmented.cs | 209 ------------------------- Hazel/Udp/UdpConnection.Reliable.cs | 10 +- Hazel/Udp/UdpConnection.cs | 63 +++----- Hazel/Udp/UdpConnectionListener.cs | 46 +++--- 15 files changed, 175 insertions(+), 360 deletions(-) create mode 100644 Hazel.UnitTests/UnitTest1.cs delete mode 100644 Hazel/Udp/UdpConnection.Fragmented.cs diff --git a/Hazel.UnitTests/Hazel.UnitTests.csproj b/Hazel.UnitTests/Hazel.UnitTests.csproj index 0cd7986..d29fb81 100644 --- a/Hazel.UnitTests/Hazel.UnitTests.csproj +++ b/Hazel.UnitTests/Hazel.UnitTests.csproj @@ -63,6 +63,7 @@ + diff --git a/Hazel.UnitTests/TestHelper.cs b/Hazel.UnitTests/TestHelper.cs index 08c44b7..4fb4374 100644 --- a/Hazel.UnitTests/TestHelper.cs +++ b/Hazel.UnitTests/TestHelper.cs @@ -23,34 +23,42 @@ namespace Hazel.UnitTests ManualResetEvent mutex = new ManualResetEvent(false); //Setup listener - listener.NewConnection += delegate(object sender, NewConnectionEventArgs args) + listener.NewConnection += delegate(object sender, NewConnectionEventArgs ncArgs) { - args.Connection.SendBytes(data, sendOption); + ncArgs.Connection.SendBytes(data, sendOption); }; listener.Start(); + DataReceivedEventArgs args = null; //Setup conneciton - connection.DataReceived += delegate(object sender, DataReceivedEventArgs args) + connection.DataReceived += delegate(object sender, DataReceivedEventArgs a) { Trace.WriteLine("Data was received correctly."); - Assert.AreEqual(data.Length, args.Message.Length); - - for (int i = 0; i < data.Length; i++) + try { - Assert.AreEqual(data[i], args.Message.ReadByte()); + args = a; + } + finally + { + mutex.Set(); } - - Assert.AreEqual(sendOption, args.SendOption); - - mutex.Set(); }; connection.Connect(); //Wait until data is received mutex.WaitOne(); + + Assert.AreEqual(data.Length, args.Message.Length); + + for (int i = 0; i < data.Length; i++) + { + Assert.AreEqual(data[i], args.Message.ReadByte()); + } + + Assert.AreEqual(sendOption, args.SendOption); } /// @@ -66,20 +74,14 @@ namespace Hazel.UnitTests ManualResetEvent mutex2 = new ManualResetEvent(false); //Setup listener + DataReceivedEventArgs result = null; listener.NewConnection += delegate(object sender, NewConnectionEventArgs args) { args.Connection.DataReceived += delegate(object innerSender, DataReceivedEventArgs innerArgs) { Trace.WriteLine("Data was received correctly."); - Assert.AreEqual(data.Length, innerArgs.Message.Length); - - for (int i = 0; i < data.Length; i++) - { - Assert.AreEqual(data[i], innerArgs.Message.ReadByte()); - } - - Assert.AreEqual(sendOption, innerArgs.SendOption); + result = innerArgs; mutex2.Set(); }; @@ -98,6 +100,15 @@ namespace Hazel.UnitTests //Wait until data is received mutex2.WaitOne(); + + Assert.AreEqual(data.Length, result.Message.Length); + + for (int i = 0; i < data.Length; i++) + { + Assert.AreEqual(data[i], result.Message.ReadByte()); + } + + Assert.AreEqual(sendOption, result.SendOption); } /// diff --git a/Hazel.UnitTests/UdpConnectionTests.cs b/Hazel.UnitTests/UdpConnectionTests.cs index 575cdbf..e5bc74a 100644 --- a/Hazel.UnitTests/UdpConnectionTests.cs +++ b/Hazel.UnitTests/UdpConnectionTests.cs @@ -202,20 +202,7 @@ namespace Hazel.UnitTests TestHelper.RunServerToClientTest(listener, connection, 10, SendOption.Reliable); } } - - /// - /// Tests server to client reliable communication on the UdpConnection. - /// - [TestMethod] - public void UdpFragmentedServerToClientTest() - { - using (UdpConnectionListener listener = new UdpConnectionListener(new NetworkEndPoint(IPAddress.Any, 4296))) - using (UdpConnection connection = new UdpClientConnection(new NetworkEndPoint(IPAddress.Loopback, 4296))) - { - TestHelper.RunServerToClientTest(listener, connection, (int)(UdpConnection.FragmentSize * 9.5), SendOption.FragmentedReliable); - } - } - + /// /// Tests server to client unreliable communication on the UdpConnection. /// @@ -241,20 +228,7 @@ namespace Hazel.UnitTests TestHelper.RunClientToServerTest(listener, connection, 10, SendOption.Reliable); } } - - /// - /// Tests server to client reliable communication on the UdpConnection. - /// - [TestMethod] - public void UdpFragmentedClientToServerTest() - { - using (UdpConnectionListener listener = new UdpConnectionListener(new NetworkEndPoint(IPAddress.Any, 4296))) - using (UdpConnection connection = new UdpClientConnection(new NetworkEndPoint(IPAddress.Loopback, 4296))) - { - TestHelper.RunClientToServerTest(listener, connection, (int)(UdpConnection.FragmentSize * 9.5), SendOption.FragmentedReliable); - } - } - + /// /// Tests the keepalive functionality from the client, /// diff --git a/Hazel.UnitTests/UnitTest1.cs b/Hazel.UnitTests/UnitTest1.cs new file mode 100644 index 0000000..0abb779 --- /dev/null +++ b/Hazel.UnitTests/UnitTest1.cs @@ -0,0 +1,29 @@ +using System; +using System.Net; +using System.Threading; +using System.Threading.Tasks; +using Hazel.Udp; +using Microsoft.VisualStudio.TestTools.UnitTesting; + +namespace Hazel.UnitTests +{ + [TestClass] + public class UnitTest1 + { + // [TestMethod] + public void StressTest() + { + Parallel.For(0, 10000, + new ParallelOptions { MaxDegreeOfParallelism = 16 }, + (i) => + { + var ep = new NetworkEndPoint(IPAddress.Loopback, 22023); + using (var connection = new UdpClientConnection(ep)) + { + connection.Connect(); + Thread.Sleep(100); + } + }); + } + } +} diff --git a/Hazel/Connection.cs b/Hazel/Connection.cs index 5a02b40..b328530 100644 --- a/Hazel/Connection.cs +++ b/Hazel/Connection.cs @@ -245,12 +245,18 @@ namespace Hazel /// protected void InvokeDataReceived(MessageReader msg, SendOption sendOption, ushort reliableId) { - DataReceivedEventArgs args = DataReceivedEventArgs.GetObject(); - args.Set(msg, sendOption, reliableId); - //Make a copy to avoid race condition between null check and invocation EventHandler handler = DataReceived; - if (handler != null) handler.Invoke(this, args); + if (handler != null) + { + DataReceivedEventArgs args = DataReceivedEventArgs.GetObject(); + args.Set(msg, sendOption, reliableId); + handler.Invoke(this, args); + } + else + { + msg.Recycle(); + } } /// @@ -264,12 +270,14 @@ namespace Hazel /// protected void InvokeDisconnected(Exception e = null) { - DisconnectedEventArgs args = DisconnectedEventArgs.GetObject(); - args.Set(e); - //Make a copy to avoid race condition between null check and invocation EventHandler handler = Disconnected; - if (handler != null) handler.Invoke(this, args); + if (handler != null) + { + DisconnectedEventArgs args = DisconnectedEventArgs.GetObject(); + args.Set(e); + handler.Invoke(this, args); + } } /// diff --git a/Hazel/ConnectionListener.cs b/Hazel/ConnectionListener.cs index eadf2b1..51fe187 100644 --- a/Hazel/ConnectionListener.cs +++ b/Hazel/ConnectionListener.cs @@ -76,14 +76,18 @@ namespace Hazel /// protected void InvokeNewConnection(MessageReader msg, Connection connection) { - //Get new args - NewConnectionEventArgs args = NewConnectionEventArgs.GetObject(); - args.Set(msg, connection); - //Make a copy to avoid race condition between null check and invocation EventHandler handler = NewConnection; if (handler != null) + { + NewConnectionEventArgs args = NewConnectionEventArgs.GetObject(); + args.Set(msg, connection); handler(this, args); + } + else + { + msg.Recycle(); + } } /// diff --git a/Hazel/Hazel.csproj b/Hazel/Hazel.csproj index 7f33975..4f1be81 100644 --- a/Hazel/Hazel.csproj +++ b/Hazel/Hazel.csproj @@ -80,7 +80,6 @@ Code - diff --git a/Hazel/MessageReader.cs b/Hazel/MessageReader.cs index f6951b3..07677e2 100644 --- a/Hazel/MessageReader.cs +++ b/Hazel/MessageReader.cs @@ -13,8 +13,8 @@ namespace Hazel public byte Tag; public int Length; + public int Offset; - public int Offset { get; private set; } public int Position { get { return this._position; } @@ -28,6 +28,21 @@ namespace Hazel private int _position; private int readHead; + + public static MessageReader GetSized(int minSize) + { + var output = ReaderPool.GetObject(); + if (output.Buffer == null || output.Buffer.Length < minSize) + { + output.Buffer = new byte[minSize]; + } + + output.Offset = 0; + output.Position = 0; + output.Length = minSize; + output.Tag = byte.MaxValue; + return output; + } public static MessageReader GetRaw(byte[] bytes, int offset, int length) { @@ -36,6 +51,7 @@ namespace Hazel output.Offset = offset; output.Position = 0; output.Length = length; + output.Tag = byte.MaxValue; return output; } diff --git a/Hazel/MessageWriter.cs b/Hazel/MessageWriter.cs index 165c41c..01c3d9c 100644 --- a/Hazel/MessageWriter.cs +++ b/Hazel/MessageWriter.cs @@ -19,6 +19,12 @@ namespace Hazel private Stack messageStarts = new Stack(); + public MessageWriter(byte[] buffer) + { + this.Buffer = buffer; + this.Length = this.Buffer.Length; + } + /// public MessageWriter(int bufferSize) { @@ -114,8 +120,6 @@ namespace Hazel case SendOption.Reliable: this.Length = this.Position = 3; break; - case SendOption.FragmentedReliable: - throw new NotImplementedException("Sry bruh"); } } @@ -127,6 +131,26 @@ namespace Hazel } #region WriteMethods + + public void CopyFrom(MessageReader target) + { + int offset, length; + if (target.Tag == byte.MaxValue) + { + offset = target.Offset; + length = target.Length; + } + else + { + offset = target.Offset - 3; + length = target.Length + 3; + } + + System.Buffer.BlockCopy(target.Buffer, offset, this.Buffer, this.Position, length); + this.Position += length; + if (this.Position > this.Length) this.Length = this.Position; + } + public void Write(bool value) { this.Buffer[this.Position++] = (byte)(value ? 1 : 0); diff --git a/Hazel/SendOption.cs b/Hazel/SendOption.cs index e7ebec2..c2ffb22 100644 --- a/Hazel/SendOption.cs +++ b/Hazel/SendOption.cs @@ -31,17 +31,5 @@ namespace Hazel /// a larger number of protocol bytes and can be slower than unreliable delivery. /// Reliable = 1, - - /// - /// Requests data be sent so that large messages are fragmented into smaller chunks of - /// data and reassembled when received. - /// - /// - /// Fragmented messages allow large amounts of data to be transmitted in smaller chunks when using connections - /// that do not support the transmission of large messages. By specifying reliable delivery messages are - /// guaranteed to arrive and to arrive only once but the sending process may require more memory, processing, - /// a larger number protocol bytes and may be slower than sending unreliably. - /// - FragmentedReliable = 2 } } diff --git a/Hazel/Udp/UdpClientConnection.cs b/Hazel/Udp/UdpClientConnection.cs index 3976af6..0616299 100644 --- a/Hazel/Udp/UdpClientConnection.cs +++ b/Hazel/Udp/UdpClientConnection.cs @@ -285,7 +285,8 @@ namespace Hazel.Udp Thread.Sleep(this.TestLagMs); } - HandleReceive(bytes); + MessageReader msg = MessageReader.GetRaw(bytes, 0, bytesReceived); + HandleReceive(msg); } /// diff --git a/Hazel/Udp/UdpConnection.Fragmented.cs b/Hazel/Udp/UdpConnection.Fragmented.cs deleted file mode 100644 index 9e525b2..0000000 --- a/Hazel/Udp/UdpConnection.Fragmented.cs +++ /dev/null @@ -1,209 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Text; - -namespace Hazel.Udp -{ - partial class UdpConnection - { - /// - /// The amount of data that can be put into a fragment. - /// - public const int FragmentSize = 65507 - 1 - 2 - 2 - 2; - - /// - /// The last fragmented message ID that was written. - /// - volatile ushort lastFragmentIDAllocated; - - Dictionary fragmentedMessagesReceived = new Dictionary(); - - /// - /// Sends a message fragmenting it as needed to pass over the network. - /// - /// The send option the message was sent with. - /// The data of the message to send. - void FragmentedSend(byte[] data) - { - //Get an ID not used yet. - ushort id = ++lastFragmentIDAllocated; //TODO is extra code needed to manage loop around? - - for (ushort i = 0; i < Math.Ceiling(data.Length / (double)FragmentSize); i++) - { - byte[] buffer = new byte[Math.Min(data.Length - (FragmentSize * i), FragmentSize) + 7]; - - //Add send option - buffer[0] = i == 0 ? (byte)SendOption.FragmentedReliable : (byte)UdpSendOption.Fragment; - - //Add fragment message ID - buffer[1] = (byte)((id >> 8) & 0xFF); - buffer[2] = (byte)id; - - //Add length or fragment id - if (i == 0) - { - ushort fragments = (ushort)Math.Ceiling(data.Length / (double)FragmentSize); - buffer[3] = (byte)((fragments >> 8) & 0xFF); - buffer[4] = (byte)fragments; - } - else - { - buffer[3] = (byte)((i >> 8) & 0xFF); - buffer[4] = (byte)i; - } - - //Pass fragment to reliable send code to ensure it will arrive - AttachReliableID(buffer, 5, buffer.Length); - - //Copy data into fragment - Buffer.BlockCopy(data, FragmentSize * i, buffer, 7, buffer.Length - 7); - - //Send - WriteBytesToConnection(buffer, buffer.Length); - } - } - - /// - /// Gets a message from those we've begun receiving or adds a new one. - /// - /// The Id of the message to find. - /// - FragmentedMessage GetFragmentedMessage(ushort messageId) - { - lock (fragmentedMessagesReceived) - { - FragmentedMessage message; - if (fragmentedMessagesReceived.ContainsKey(messageId)) - { - message = fragmentedMessagesReceived[messageId]; - } - else - { - message = new FragmentedMessage(); - - fragmentedMessagesReceived.Add(messageId, message); - } - - return message; - } - } - - /// - /// Handles a the start message of a fragmented message. - /// - /// The buffer received. - void FragmentedStartMessageReceive(byte[] buffer) - { - //Send to reliable code to send the acknowledgement - ushort reliableId; - if (!ProcessReliableReceive(buffer, 5, out reliableId)) - return; - - ushort id = (ushort)((buffer[1] << 8) + buffer[2]); - - ushort length = (ushort)((buffer[3] << 8) + buffer[4]); - - FragmentedMessage message; - bool messageComplete; - lock (fragmentedMessagesReceived) - { - message = GetFragmentedMessage(id); - message.received.Add(new FragmentedMessage.Fragment(0, buffer, 7)); - message.noFragments = length; - - messageComplete = message.noFragments == message.received.Count; - } - - if (messageComplete) - FinalizeFragmentedMessage(message); - } - - /// - /// Handles a fragment message of a fragmented message. - /// - /// The buffer received. - void FragmentedMessageReceive(byte[] buffer) - { - //Send to reliable code to send the acknowledgement - ushort reliableId; - if (!ProcessReliableReceive(buffer, 5, out reliableId)) - return; - - ushort id = (ushort)((buffer[1] << 8) + buffer[2]); - - ushort fragmentID = (ushort)((buffer[3] << 8) + buffer[4]); - - FragmentedMessage message; - bool messageComplete; - lock (fragmentedMessagesReceived) - { - message = GetFragmentedMessage(id); - message.received.Add(new FragmentedMessage.Fragment(fragmentID, buffer, 7)); - - messageComplete = message.noFragments == message.received.Count; - } - - if (messageComplete) - FinalizeFragmentedMessage(message); - } - - /// - /// Finalizes a completed fragmented message and invokes message received events. - /// - /// The message received. - void FinalizeFragmentedMessage(FragmentedMessage message) - { - IEnumerable orderedFragments = message.received.OrderBy((x) => x.fragmentID); - FragmentedMessage.Fragment last = orderedFragments.Last(); - - byte[] completeData = new byte[(orderedFragments.Count() - 1) * FragmentSize + last.data.Length - last.offset]; - int ptr = 0; - foreach (FragmentedMessage.Fragment fragment in orderedFragments) - { - Buffer.BlockCopy(fragment.data, fragment.offset, completeData, ptr, fragment.data.Length - fragment.offset); - ptr += fragment.data.Length - fragment.offset; - } - - var reader = MessageReader.GetRaw(completeData, 0, completeData.Length); - try - { - InvokeDataReceived(reader, SendOption.FragmentedReliable, 0); - } - finally - { - reader.Recycle(); - } - } - - /// - /// Holding class for the parts of a fragmented message so far received. - /// - private class FragmentedMessage - { - /// - /// The total number of fragments expected. - /// - public int noFragments = -1; - - /// - /// The fragments received so far. - /// - public List received = new List(); - - public struct Fragment - { - public int fragmentID; - public byte[] data; - public int offset; - - public Fragment(int fragmentID, byte[] data, int offset) - { - this.fragmentID = fragmentID; - this.data = data; - this.offset = offset; - } - } - } - } -} diff --git a/Hazel/Udp/UdpConnection.Reliable.cs b/Hazel/Udp/UdpConnection.Reliable.cs index 6314f5f..74e2824 100644 --- a/Hazel/Udp/UdpConnection.Reliable.cs +++ b/Hazel/Udp/UdpConnection.Reliable.cs @@ -286,14 +286,14 @@ namespace Hazel.Udp /// /// Handles a reliable message being received and invokes the data event. /// - /// The buffer received. - void ReliableMessageReceive(byte[] buffer) + /// The buffer received. + void ReliableMessageReceive(MessageReader message) { ushort id; - if (ProcessReliableReceive(buffer, 1, out id)) - InvokeDataReceived(SendOption.Reliable, buffer, 3, id); + if (ProcessReliableReceive(message.Buffer, 1, out id)) + InvokeDataReceived(SendOption.Reliable, message, 3, id); - Statistics.LogReliableReceive(buffer.Length - 3, buffer.Length); + Statistics.LogReliableReceive(message.Length - 3, message.Length); } /// diff --git a/Hazel/Udp/UdpConnection.cs b/Hazel/Udp/UdpConnection.cs index 926a616..2e82f84 100644 --- a/Hazel/Udp/UdpConnection.cs +++ b/Hazel/Udp/UdpConnection.cs @@ -55,9 +55,6 @@ namespace Hazel.Udp Statistics.LogReliableSend(buffer.Length - 3, buffer.Length); break; - case SendOption.FragmentedReliable: - throw new NotImplementedException("Not yet"); - default: WriteBytesToConnection(buffer, buffer.Length); Statistics.LogUnreliableSend(buffer.Length - 1, buffer.Length);; @@ -111,12 +108,7 @@ namespace Hazel.Udp case SendOption.Reliable: ReliableSend((byte)sendOption, bytes, offset, length); break; - - case SendOption.FragmentedReliable: - throw new NotImplementedException(); - // FragmentedSend(data); - // break; - + //Treat all else as unreliable default: UnreliableSend((byte)sendOption, bytes, offset, length); @@ -140,11 +132,7 @@ namespace Hazel.Udp case (byte)UdpSendOption.Hello: ReliableSend(sendOption, data, ackCallback); break; - - case (byte)SendOption.FragmentedReliable: - FragmentedSend(data); - break; - + //Treat all else as unreliable default: UnreliableSend(sendOption, data); @@ -155,48 +143,39 @@ namespace Hazel.Udp /// /// Handles the receiving of data. /// - /// The buffer containing the bytes received. - protected internal void HandleReceive(byte[] buffer) + /// The buffer containing the bytes received. + protected internal void HandleReceive(MessageReader message) { - InvokeDataReceivedRaw(buffer); + InvokeDataReceivedRaw(message.Buffer); - switch (buffer[0]) + switch (message.Buffer[0]) { //Handle reliable receives case (byte)SendOption.Reliable: - ReliableMessageReceive(buffer); + ReliableMessageReceive(message); break; //Handle acknowledgments case (byte)UdpSendOption.Acknowledgement: - AcknowledgementMessageReceive(buffer); + AcknowledgementMessageReceive(message.Buffer); break; //We need to acknowledge hello and ping messages but dont want to invoke any events! case (byte)UdpSendOption.Ping: case (byte)UdpSendOption.Hello: ushort id; - ProcessReliableReceive(buffer, 1, out id); - Statistics.LogHelloReceive(buffer.Length); + ProcessReliableReceive(message.Buffer, 1, out id); + Statistics.LogHelloReceive(message.Length); break; case (byte)UdpSendOption.Disconnect: HandleDisconnect(new HazelException("The remote sent a disconnect request")); break; - - //Handle fragmented messages - case (byte)SendOption.FragmentedReliable: - FragmentedStartMessageReceive(buffer); - break; - - case (byte)UdpSendOption.Fragment: - FragmentedMessageReceive(buffer); - break; - + //Treat everything else as unreliable default: - InvokeDataReceived(SendOption.None, buffer, 1, 0); - Statistics.LogUnreliableReceive(buffer.Length - 1, buffer.Length); + InvokeDataReceived(SendOption.None, message, 1, 0); + Statistics.LogUnreliableReceive(message.Length - 1, message.Length); break; } } @@ -240,17 +219,13 @@ namespace Hazel.Udp /// The send option the message was received with. /// The buffer received. /// The offset of data in the buffer. - void InvokeDataReceived(SendOption sendOption, byte[] buffer, int dataOffset, ushort reliableId) + void InvokeDataReceived(SendOption sendOption, MessageReader buffer, int dataOffset, ushort reliableId) { - var reader = MessageReader.GetRaw(buffer, dataOffset, buffer.Length - dataOffset); - try - { - InvokeDataReceived(reader, sendOption, reliableId); - } - finally - { - reader.Recycle(); - } + buffer.Offset = dataOffset; + buffer.Length = buffer.Length - dataOffset; + buffer.Position = 0; + + InvokeDataReceived(buffer, sendOption, reliableId); } /// diff --git a/Hazel/Udp/UdpConnectionListener.cs b/Hazel/Udp/UdpConnectionListener.cs index c8f5333..acb47e6 100644 --- a/Hazel/Udp/UdpConnectionListener.cs +++ b/Hazel/Udp/UdpConnectionListener.cs @@ -4,7 +4,7 @@ using System.Linq; using System.Net; using System.Net.Sockets; using System.Text; - +using System.Threading; namespace Hazel.Udp { @@ -18,12 +18,7 @@ namespace Hazel.Udp /// The socket listening for connections. /// Socket listener; - - /// - /// Buffer to store incoming data in. - /// - byte[] dataBuffer = new byte[ushort.MaxValue]; - + /// /// The connections we currently hold /// @@ -92,7 +87,9 @@ namespace Hazel.Udp try { - listener.BeginReceiveFrom(dataBuffer, 0, dataBuffer.Length, SocketFlags.None, ref remoteEP, ReadCallback, dataBuffer); + var message = MessageReader.GetSized(ushort.MaxValue); + listener.BeginReceiveFrom(message.Buffer, 0, message.Buffer.Length, SocketFlags.None, ref remoteEP, ReadCallback, message); + Interlocked.Increment(ref ActiveThreads); } catch (ObjectDisposedException) { @@ -111,6 +108,8 @@ namespace Hazel.Udp /// Called when data has been received by the listener. /// /// The asyncronous operation's result. + + private int ActiveThreads; void ReadCallback(IAsyncResult result) { int bytesReceived; @@ -119,6 +118,7 @@ namespace Hazel.Udp //End the receive operation try { + Interlocked.Decrement(ref ActiveThreads); bytesReceived = listener.EndReceiveFrom(result, ref remoteEndPoint); } catch (ObjectDisposedException) @@ -133,7 +133,6 @@ namespace Hazel.Udp //This thread suggests the IP is not passed out from WinSoc so maybe not possible //http://stackoverflow.com/questions/2576926/python-socket-error-on-udp-data-receive-10054 - StartListeningForData(); return; } @@ -141,14 +140,13 @@ namespace Hazel.Udp //Exit if no bytes read, we've closed. if (bytesReceived == 0) return; - - //Copy to new buffer - byte[] buffer = new byte[bytesReceived]; - Buffer.BlockCopy((byte[])result.AsyncState, 0, buffer, 0, bytesReceived); - + //Begin receiving again StartListeningForData(); + var message = (MessageReader)result.AsyncState; + message.Length = bytesReceived; + bool aware; UdpServerConnection connection; lock (connections) @@ -158,7 +156,7 @@ namespace Hazel.Udp if (!(aware = connections.TryGetValue(remoteEndPoint, out connection))) { //Check for malformed connection attempts - if (buffer[0] != (byte)UdpSendOption.Hello) + if (message.Buffer[0] != (byte)UdpSendOption.Hello) return; connection = new UdpServerConnection(this, remoteEndPoint, IPMode); @@ -167,20 +165,16 @@ namespace Hazel.Udp } //Inform the connection of the buffer (new connections need to send an ack back to client) - connection.HandleReceive(buffer); - + connection.HandleReceive(message); + //If it's a new connection invoke the NewConnection event. if (!aware) { - var reader = MessageReader.GetRaw(buffer, 4, buffer.Length - 4); - try - { - InvokeNewConnection(reader, connection); - } - finally - { - reader.Recycle(); - } + // Skip header and hello byte; + message.Offset = 4; + message.Length = bytesReceived - 4; + message.Position = 0; + InvokeNewConnection(message, connection); } } -- 2.39.5