NetSharp

NetSharp.git
git clone git://git.lenczewski.org/NetSharp.git
Log | Files | Refs | README | LICENSE

commit e616c6787ffcb5570713f5fa5780eee8dfe8dab9
parent f5ccc28b2f4a061fd7fe0d37b35f700ff42321a5
Author: Mikolaj Lenczewski <mikolaj.lenczewski308@gmail.com>
Date:   Thu, 21 May 2020 16:25:34 +0100

In progress of implemening variable-length stream packet sending. We dont build

Diffstat:
MNetSharp/NetSharp/NetSharp.xml | 77++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
MNetSharp/NetSharp/Raw/IRawNetworkTransportProvider.cs | 32+++++++++++++++++++++++++++-----
MNetSharp/NetSharp/Raw/RawNetworkConnectionBase.cs | 6+++---
MNetSharp/NetSharp/Raw/RawNetworkReaderBase.cs | 4++--
MNetSharp/NetSharp/Raw/RawNetworkWriterBase.cs | 4++--
ANetSharp/NetSharp/Raw/Stream/FixedPacketRawStreamNetworkReader.cs | 149+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharp/Raw/Stream/FixedPacketRawStreamNetworkWriter.cs | 281+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
MNetSharp/NetSharp/Raw/Stream/RawStreamNetworkReader.cs | 352+++++++++++++------------------------------------------------------------------
MNetSharp/NetSharp/Raw/Stream/RawStreamNetworkWriter.cs | 368+++++++------------------------------------------------------------------------
MNetSharp/NetSharp/Raw/Stream/RawStreamPacket.cs | 61+++++++++++++++++++++++++++++++++++++++++++++++--------------
ANetSharp/NetSharp/Raw/Stream/VariablePacketRawStreamNetworkReader.cs | 183+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharp/Raw/Stream/VariablePacketRawStreamNetworkWriter.cs | 340+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
MNetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkReaderBenchmark.cs | 2+-
MNetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterAsyncBenchmark.cs | 2+-
MNetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterSyncBenchmark.cs | 2+-
ANetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkReaderBenchmark.cs | 141+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkWriterAsyncBenchmark.cs | 131+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkWriterSyncBenchmark.cs | 135+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
DNetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkReaderBenchmark.cs | 145-------------------------------------------------------------------------------
DNetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterAsyncBenchmark.cs | 131-------------------------------------------------------------------------------
DNetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterSyncBenchmark.cs | 135-------------------------------------------------------------------------------
ANetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkReaderBenchmark.cs | 145+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkWriterAsyncBenchmark.cs | 131+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkWriterSyncBenchmark.cs | 135+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkReaderExample.cs | 57+++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkWriterAsyncExample.cs | 62++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkWriterSyncExample.cs | 64++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
DNetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkReaderExample.cs | 57---------------------------------------------------------
DNetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterAsyncExample.cs | 62--------------------------------------------------------------
DNetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterSyncExample.cs | 64----------------------------------------------------------------
ANetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkReaderExample.cs | 57+++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkWriterAsyncExample.cs | 62++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
ANetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkWriterSyncExample.cs | 64++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
MNetSharp/NetSharpExamples/NetSharpExamples.xml | 60++++++++++++++++++++++++++++++++++++++++++++++++------------
34 files changed, 2430 insertions(+), 1271 deletions(-)

diff --git a/NetSharp/NetSharp/NetSharp.xml b/NetSharp/NetSharp/NetSharp.xml @@ -64,13 +64,22 @@ <member name="M:NetSharp.Raw.DatagramRawNetworkTransportProvider.GetWriter(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="P:NetSharp.Raw.StreamRawNetworkTransportProvider.TransportProtocolType"> + <member name="P:NetSharp.Raw.FixedPacketRawStreamNetworkTransportProvider.TransportProtocolType"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.StreamRawNetworkTransportProvider.GetReader(System.Net.Sockets.Socket@,System.Net.EndPoint,NetSharp.Raw.Stream.RawStreamRequestHandler,System.Int32,System.Int32,System.UInt32)"> + <member name="M:NetSharp.Raw.FixedPacketRawStreamNetworkTransportProvider.GetReader(System.Net.Sockets.Socket@,System.Net.EndPoint,NetSharp.Raw.Stream.RawStreamRequestHandler,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.StreamRawNetworkTransportProvider.GetWriter(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <member name="M:NetSharp.Raw.FixedPacketRawStreamNetworkTransportProvider.GetWriter(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <inheritdoc /> + </member> + <member name="P:NetSharp.Raw.VariablePacketRawStreamNetworkTransportProvider.TransportProtocolType"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.VariablePacketRawStreamNetworkTransportProvider.GetReader(System.Net.Sockets.Socket@,System.Net.EndPoint,NetSharp.Raw.Stream.RawStreamRequestHandler,System.Int32,System.Int32,System.UInt32)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.VariablePacketRawStreamNetworkTransportProvider.GetWriter(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> <member name="M:NetSharp.Raw.RawNetworkConnectionBase.Dispose(System.Boolean)"> @@ -114,6 +123,39 @@ <member name="M:NetSharp.Raw.RawNetworkWriterBase.WriteAsync(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkReader.#ctor(System.Net.Sockets.Socket@,NetSharp.Raw.Stream.RawStreamRequestHandler,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkReader.CompleteAccept(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkReader.CompleteReceive(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkReader.CompleteSend(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkWriter.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkWriter.CompleteReceive(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkWriter.CompleteSend(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkWriter.Read(System.Net.EndPoint@,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkWriter.ReadAsync(System.Net.EndPoint,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkWriter.Write(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.FixedPacketRawStreamNetworkWriter.WriteAsync(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> + <inheritdoc /> + </member> <member name="M:NetSharp.Raw.Stream.RawStreamNetworkReader.#ctor(System.Net.Sockets.Socket@,NetSharp.Raw.Stream.RawStreamRequestHandler,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> @@ -153,16 +195,37 @@ <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.ConnectAsync(System.Net.EndPoint)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.Read(System.Net.EndPoint@,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkReader.#ctor(System.Net.Sockets.Socket@,NetSharp.Raw.Stream.RawStreamRequestHandler,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkReader.CompleteAccept(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkReader.CompleteReceive(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkReader.CompleteSend(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkWriter.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkWriter.CompleteReceive(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkWriter.CompleteSend(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkWriter.Read(System.Net.EndPoint@,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.ReadAsync(System.Net.EndPoint,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkWriter.ReadAsync(System.Net.EndPoint,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.Write(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkWriter.Write(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.WriteAsync(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.Stream.VariablePacketRawStreamNetworkWriter.WriteAsync(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> <member name="T:NetSharp.Utils.BiDictionary`2"> diff --git a/NetSharp/NetSharp/Raw/IRawNetworkTransportProvider.cs b/NetSharp/NetSharp/Raw/IRawNetworkTransportProvider.cs @@ -1,7 +1,7 @@ -using System; -using NetSharp.Raw.Datagram; +using NetSharp.Raw.Datagram; using NetSharp.Raw.Stream; +using System; using System.Net; using System.Net.Sockets; @@ -47,7 +47,7 @@ namespace NetSharp.Raw } } - public sealed class StreamRawNetworkTransportProvider : IRawNetworkTransportProvider<RawStreamRequestHandler> + public sealed class FixedPacketRawStreamNetworkTransportProvider : IRawNetworkTransportProvider<RawStreamRequestHandler> { /// <inheritdoc /> public SocketType TransportProtocolType { get; } = SocketType.Stream; @@ -56,7 +56,7 @@ namespace NetSharp.Raw public RawNetworkReaderBase GetReader(ref Socket rawConnection, EndPoint defaultEndPoint, RawStreamRequestHandler? requestHandler, int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) { - return new RawStreamNetworkReader(ref rawConnection, requestHandler, defaultEndPoint, maxPooledBufferSize, + return new FixedPacketRawStreamNetworkReader(ref rawConnection, requestHandler, defaultEndPoint, maxPooledBufferSize, maxPooledBuffersPerBucket, preallocatedStateObjects); } @@ -64,7 +64,29 @@ namespace NetSharp.Raw public RawNetworkWriterBase GetWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) { - return new RawStreamNetworkWriter(ref rawConnection, defaultEndPoint, maxPooledBufferSize, + return new FixedPacketRawStreamNetworkWriter(ref rawConnection, defaultEndPoint, maxPooledBufferSize, + maxPooledBuffersPerBucket, preallocatedStateObjects); + } + } + + public sealed class VariablePacketRawStreamNetworkTransportProvider : IRawNetworkTransportProvider<RawStreamRequestHandler> + { + /// <inheritdoc /> + public SocketType TransportProtocolType { get; } = SocketType.Stream; + + /// <inheritdoc /> + public RawNetworkReaderBase GetReader(ref Socket rawConnection, EndPoint defaultEndPoint, RawStreamRequestHandler? requestHandler, + int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) + { + return new VariablePacketRawStreamNetworkReader(ref rawConnection, requestHandler, defaultEndPoint, maxPooledBufferSize, + maxPooledBuffersPerBucket, preallocatedStateObjects); + } + + /// <inheritdoc /> + public RawNetworkWriterBase GetWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, + int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) + { + return new VariablePacketRawStreamNetworkWriter(ref rawConnection, defaultEndPoint, maxPooledBufferSize, maxPooledBuffersPerBucket, preallocatedStateObjects); } } diff --git a/NetSharp/NetSharp/Raw/RawNetworkConnectionBase.cs b/NetSharp/NetSharp/Raw/RawNetworkConnectionBase.cs @@ -17,14 +17,14 @@ namespace NetSharp.Raw protected readonly Socket Connection; protected readonly EndPoint DefaultEndPoint; - protected RawNetworkConnectionBase(ref Socket rawConnection, EndPoint defaultEndPoint, int pooledPacketBufferSize, + protected RawNetworkConnectionBase(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, int pooledBuffersPerBucket = 50, uint preallocatedStateObjects = 0) { Connection = rawConnection; - BufferPool = pooledPacketBufferSize <= DefaultMaxPooledBufferSize && pooledBuffersPerBucket <= DefaultMaxPooledBuffersPerBucket + BufferPool = maxPooledBufferSize <= DefaultMaxPooledBufferSize && pooledBuffersPerBucket <= DefaultMaxPooledBuffersPerBucket ? ArrayPool<byte>.Shared - : BufferPool = ArrayPool<byte>.Create(pooledPacketBufferSize, pooledBuffersPerBucket); + : BufferPool = ArrayPool<byte>.Create(maxPooledBufferSize, pooledBuffersPerBucket); DefaultEndPoint = defaultEndPoint; diff --git a/NetSharp/NetSharp/Raw/RawNetworkReaderBase.cs b/NetSharp/NetSharp/Raw/RawNetworkReaderBase.cs @@ -12,8 +12,8 @@ namespace NetSharp.Raw protected readonly CancellationToken ShutdownToken; /// <inheritdoc /> - private protected RawNetworkReaderBase(ref Socket rawConnection, EndPoint defaultEndPoint, int pooledPacketBufferSize, int pooledBuffersPerBucket = 50, - uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, pooledPacketBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) + private protected RawNetworkReaderBase(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, int pooledBuffersPerBucket = 50, + uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxPooledBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) { shutdownTokenSource = new CancellationTokenSource(); ShutdownToken = shutdownTokenSource.Token; diff --git a/NetSharp/NetSharp/Raw/RawNetworkWriterBase.cs b/NetSharp/NetSharp/Raw/RawNetworkWriterBase.cs @@ -8,8 +8,8 @@ namespace NetSharp.Raw public abstract class RawNetworkWriterBase : RawNetworkConnectionBase, INetworkWriter { /// <inheritdoc /> - protected RawNetworkWriterBase(ref Socket rawConnection, EndPoint defaultEndPoint, int pooledPacketBufferSize, int pooledBuffersPerBucket = 50, - uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, pooledPacketBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) + protected RawNetworkWriterBase(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, int pooledBuffersPerBucket = 50, + uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxPooledBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) { } diff --git a/NetSharp/NetSharp/Raw/Stream/FixedPacketRawStreamNetworkReader.cs b/NetSharp/NetSharp/Raw/Stream/FixedPacketRawStreamNetworkReader.cs @@ -0,0 +1,148 @@ +using System.Net; +using System.Net.Sockets; +using System.Runtime.CompilerServices; + +namespace NetSharp.Raw.Stream +{ + public sealed class FixedPacketRawStreamNetworkReader : RawStreamNetworkReader + { + private readonly int messageSize; + + /// <inheritdoc /> + public FixedPacketRawStreamNetworkReader(ref Socket rawConnection, RawStreamRequestHandler? requestHandler, EndPoint defaultEndPoint, int messageSize, + int pooledBuffersPerBucket = 50, uint preallocatedStateObjects = 0) : base(ref rawConnection, requestHandler, defaultEndPoint, messageSize, + pooledBuffersPerBucket, preallocatedStateObjects) + { + this.messageSize = messageSize; + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ConfigureReceive(SocketAsyncEventArgs args, int dataSize) + { + byte[] receiveBuffer = BufferPool.Rent(dataSize); + args.SetBuffer(receiveBuffer, 0, dataSize); + + TransmissionToken token = new TransmissionToken(dataSize, 0); + args.UserToken = token; + } + + /// <inheritdoc /> + protected override void CompleteAccept(SocketAsyncEventArgs args) + { + switch (args.SocketError) + { + case SocketError.Success: + ConfigureReceive(args, messageSize); + StartReceive(args); + break; + + default: + ArgsPool.Return(args); + break; + } + + StartDefaultAccept(); + } + + /// <inheritdoc /> + protected override void CompleteReceive(SocketAsyncEventArgs args) + { + TransmissionToken token = (TransmissionToken)args.UserToken; + + byte[] receiveBuffer = args.Buffer; + + int expectedBytes = token.ExpectedBytes; + + switch (args.SocketError) + { + case SocketError.Success: + int receivedBytes = args.BytesTransferred, previousReceivedBytes = token.BytesTransferred, totalReceivedBytes = previousReceivedBytes + receivedBytes; + + if (totalReceivedBytes == expectedBytes) // transmission complete + { + EndPoint clientEndPoint = args.AcceptSocket.RemoteEndPoint; + + byte[] responseBuffer = BufferPool.Rent(messageSize); + + bool responseExists = RequestHandler(clientEndPoint, receiveBuffer, totalReceivedBytes, responseBuffer); + BufferPool.Return(receiveBuffer, true); + + if (responseExists) + { + args.SetBuffer(responseBuffer, 0, messageSize); + + TransmissionToken sendToken = new TransmissionToken(messageSize, 0); + args.UserToken = sendToken; + + StartSend(args); + return; + } + + BufferPool.Return(responseBuffer, true); + + ConfigureReceive(args, messageSize); + StartReceive(args); + } + else if (0 < totalReceivedBytes && totalReceivedBytes < expectedBytes) // transmission not complete + { + token = new TransmissionToken(in token, receivedBytes); + args.UserToken = token; + + args.SetBuffer(totalReceivedBytes, expectedBytes - totalReceivedBytes); + + ContinueReceive(args); + } + else if (receivedBytes == 0) // connection is dead + { + CloseClientConnection(args); + } + break; + + default: + CloseClientConnection(args); + break; + } + } + + /// <inheritdoc /> + protected override void CompleteSend(SocketAsyncEventArgs args) + { + TransmissionToken token = (TransmissionToken)args.UserToken; + + byte[] sendBuffer = args.Buffer; + int expectedBytes = token.ExpectedBytes; + + switch (args.SocketError) + { + case SocketError.Success: + int sentBytes = args.BytesTransferred, previousSentBytes = token.BytesTransferred, totalSentBytes = previousSentBytes + sentBytes; + + if (totalSentBytes == expectedBytes) // transmission complete + { + BufferPool.Return(sendBuffer, true); + + ConfigureReceive(args, messageSize); + StartReceive(args); + } + else if (0 < totalSentBytes && totalSentBytes < expectedBytes) // transmission not complete + { + token = new TransmissionToken(in token, sentBytes); + args.UserToken = token; + + args.SetBuffer(totalSentBytes, expectedBytes - totalSentBytes); + + ContinueSend(args); + } + else if (sentBytes == 0) // connection is dead + { + CloseClientConnection(args); + } + break; + + default: + CloseClientConnection(args); + break; + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Stream/FixedPacketRawStreamNetworkWriter.cs b/NetSharp/NetSharp/Raw/Stream/FixedPacketRawStreamNetworkWriter.cs @@ -0,0 +1,280 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Threading.Tasks; + +namespace NetSharp.Raw.Stream +{ + public sealed class FixedPacketRawStreamNetworkWriter : RawStreamNetworkWriter + { + private readonly int messageSize; + + /// <inheritdoc /> + public FixedPacketRawStreamNetworkWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int messageSize, int pooledBuffersPerBucket = 50, + uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, messageSize, pooledBuffersPerBucket, preallocatedStateObjects) + { + this.messageSize = messageSize; + } + + /// <inheritdoc /> + protected override void CompleteReceive(SocketAsyncEventArgs args) + { + AsyncStreamReadToken token = (AsyncStreamReadToken)args.UserToken; + + byte[] receiveBuffer = args.Buffer; + int expectedBytes = receiveBuffer.Length; + + switch (args.SocketError) + { + case SocketError.Success: + int receivedBytes = args.BytesTransferred, totalReceivedBytes = token.TotalReadBytes; + + if (totalReceivedBytes + receivedBytes == expectedBytes) // transmission complete + { + receiveBuffer.CopyTo(token.UserBuffer); + token.CompletionSource.SetResult(totalReceivedBytes + receivedBytes); + } + else if (0 < totalReceivedBytes + receivedBytes && totalReceivedBytes + receivedBytes < expectedBytes) // transmission not complete + { + // update user token to take account of newly read bytes + token = new AsyncStreamReadToken(in token, receivedBytes); + args.UserToken = token; + + args.SetBuffer(totalReceivedBytes, expectedBytes - receivedBytes); + + ContinueReceive(args); + return; + } + else if (receivedBytes == 0) // connection is dead + { + token.CompletionSource.SetException(new SocketException((int)SocketError.HostDown)); + } + break; + + case SocketError.OperationAborted: + token.CompletionSource.SetCanceled(); + break; + + default: + int errorCode = (int)args.SocketError; + token.CompletionSource.SetException(new SocketException(errorCode)); + break; + } + + BufferPool.Return(receiveBuffer, true); + ArgsPool.Return(args); + } + + /// <inheritdoc /> + protected override void CompleteSend(SocketAsyncEventArgs args) + { + AsyncStreamWriteToken token = (AsyncStreamWriteToken)args.UserToken; + + byte[] sendBuffer = args.Buffer; + int expectedBytes = sendBuffer.Length; + + switch (args.SocketError) + { + case SocketError.Success: + int sentBytes = args.BytesTransferred, totalSentBytes = token.TotalWrittenBytes; + + if (totalSentBytes + sentBytes == expectedBytes) // transmission complete + { + token.CompletionSource.SetResult(totalSentBytes + sentBytes); + } + else if (0 < totalSentBytes + sentBytes && totalSentBytes + sentBytes < expectedBytes) // transmission not complete + { + // update user token to take account of newly written bytes + token = new AsyncStreamWriteToken(in token, sentBytes); + args.UserToken = token; + + args.SetBuffer(totalSentBytes, expectedBytes - sentBytes); + + ContinueSend(args); + return; + } + else if (sentBytes == 0) // connection is dead + { + token.CompletionSource.SetException(new SocketException((int)SocketError.HostDown)); + } + break; + + case SocketError.OperationAborted: + token.CompletionSource.SetCanceled(); + break; + + default: + int errorCode = (int)args.SocketError; + token.CompletionSource.SetException(new SocketException(errorCode)); + break; + } + + BufferPool.Return(sendBuffer, true); + ArgsPool.Return(args); + } + + /// <inheritdoc /> + public override int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) + { + int totalBytes = readBuffer.Length; + if (totalBytes > messageSize) + { + throw new ArgumentException( + $"Cannot receive a message of size: {totalBytes} bytes; maximum message size: {messageSize} bytes", + nameof(readBuffer.Length) + ); + } + + int readBytes = 0; + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + + do + { + readBytes += Connection.Receive(transmissionBuffer, readBytes, totalBytes - readBytes, flags); + } while (readBytes < totalBytes && readBytes != 0); + + transmissionBuffer.CopyTo(readBuffer); + BufferPool.Return(transmissionBuffer, true); + + return readBytes; + } + + /// <inheritdoc /> + public override ValueTask<int> ReadAsync(EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) + { + int totalBytes = readBuffer.Length; + if (totalBytes > messageSize) + { + throw new ArgumentException( + $"Cannot receive a message of size: {totalBytes} bytes; maximum message size: {messageSize} bytes", + nameof(readBuffer.Length) + ); + } + + TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); + SocketAsyncEventArgs args = ArgsPool.Rent(); + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + + args.SetBuffer(transmissionBuffer, 0, messageSize); + + args.RemoteEndPoint = remoteEndPoint; + args.SocketFlags = flags; + + AsyncStreamReadToken token = new AsyncStreamReadToken(tcs, 0, in readBuffer); + args.UserToken = token; + + if (Connection.ReceiveAsync(args)) return new ValueTask<int>(tcs.Task); + + CompleteReceive(args); + + return new ValueTask<int>(tcs.Task); + } + + /// <inheritdoc /> + public override int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, SocketFlags flags = SocketFlags.None) + { + int totalBytes = writeBuffer.Length; + if (totalBytes > messageSize) + { + throw new ArgumentException( + $"Cannot send a message of size: {totalBytes} bytes; maximum message size: {messageSize} bytes", + nameof(writeBuffer.Length) + ); + } + + int writtenBytes = 0; + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + writeBuffer.CopyTo(transmissionBuffer); + + do + { + writtenBytes += Connection.Send(transmissionBuffer, writtenBytes, totalBytes - writtenBytes, flags); + } while (writtenBytes < totalBytes && writtenBytes != 0); + + BufferPool.Return(transmissionBuffer); + + return writtenBytes; + } + + /// <inheritdoc /> + public override ValueTask<int> WriteAsync(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, SocketFlags flags = SocketFlags.None) + { + int totalBytes = writeBuffer.Length; + if (totalBytes > messageSize) + { + throw new ArgumentException( + $"Cannot send a message of size: {totalBytes} bytes; maximum message size: {messageSize} bytes", + nameof(writeBuffer.Length) + ); + } + + TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); + SocketAsyncEventArgs args = ArgsPool.Rent(); + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + writeBuffer.CopyTo(transmissionBuffer); + + args.SetBuffer(transmissionBuffer, 0, messageSize); + + args.RemoteEndPoint = remoteEndPoint; + args.SocketFlags = flags; + + AsyncStreamWriteToken token = new AsyncStreamWriteToken(tcs, 0); + args.UserToken = token; + + if (Connection.SendAsync(args)) return new ValueTask<int>(tcs.Task); + + CompleteSend(args); + + return new ValueTask<int>(tcs.Task); + } + + private readonly struct AsyncStreamReadToken + { + public readonly TaskCompletionSource<int> CompletionSource; + public readonly int TotalReadBytes; + public readonly Memory<byte> UserBuffer; + + public AsyncStreamReadToken(TaskCompletionSource<int> completionSource, int totalReadBytes, in Memory<byte> userBuffer) + { + CompletionSource = completionSource; + + TotalReadBytes = totalReadBytes; + + UserBuffer = userBuffer; + } + + public AsyncStreamReadToken(in AsyncStreamReadToken previousToken, int newlyReadBytes) + { + CompletionSource = previousToken.CompletionSource; + + TotalReadBytes = previousToken.TotalReadBytes + newlyReadBytes; + + UserBuffer = previousToken.UserBuffer; + } + } + + private readonly struct AsyncStreamWriteToken + { + public readonly TaskCompletionSource<int> CompletionSource; + public readonly int TotalWrittenBytes; + + public AsyncStreamWriteToken(TaskCompletionSource<int> completionSource, int totalWrittenBytes) + { + CompletionSource = completionSource; + + TotalWrittenBytes = totalWrittenBytes; + } + + public AsyncStreamWriteToken(in AsyncStreamWriteToken previousToken, int newlyWrittenBytes) + { + CompletionSource = previousToken.CompletionSource; + + TotalWrittenBytes = previousToken.TotalWrittenBytes + newlyWrittenBytes; + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Stream/RawStreamNetworkReader.cs b/NetSharp/NetSharp/Raw/Stream/RawStreamNetworkReader.cs @@ -2,95 +2,62 @@ using System.Net; using System.Net.Sockets; using System.Runtime.CompilerServices; -using NetSharp.Utils.Conversion; namespace NetSharp.Raw.Stream { public delegate bool RawStreamRequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, in Memory<byte> responseBuffer); - public readonly struct RawMessage + public abstract class RawStreamNetworkReader : RawNetworkReaderBase { - public readonly Memory<byte> MsgData; - public readonly Header MsgHeader; + protected readonly RawStreamRequestHandler RequestHandler; - private RawMessage(in Header header, in Memory<byte> data) - { - MsgHeader = header; - - MsgData = data; - } - - public RawMessage(in Memory<byte> data) - { - MsgHeader = new Header(data.Length); - - MsgData = data; - } - - public static RawMessage Deserialise(in Memory<byte> buffer) + /// <inheritdoc /> + protected RawStreamNetworkReader(ref Socket rawConnection, RawStreamRequestHandler? requestHandler, EndPoint defaultEndPoint, int maxMessageSize, + int pooledBuffersPerBucket = 50, uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxMessageSize, + pooledBuffersPerBucket, preallocatedStateObjects) { - Memory<byte> serialisedHeader = buffer.Slice(0, Header.TotalHeaderSize); - Header header = Header.Deserialise(in serialisedHeader); - - Memory<byte> serialisedData = buffer.Slice(Header.TotalHeaderSize); + if (maxMessageSize <= 0) + { + throw new ArgumentOutOfRangeException(nameof(maxMessageSize), maxMessageSize, + $"The message size must be greater than 0"); + } - return new RawMessage(in header, in serialisedData); + RequestHandler = requestHandler ?? DefaultRequestHandler; } - public void Serialise(in Memory<byte> buffer) + private static bool DefaultRequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + in Memory<byte> responseBuffer) { - MsgHeader.Serialise(buffer.Slice(0, Header.TotalHeaderSize)); - - MsgData.CopyTo(buffer.Slice(Header.TotalHeaderSize, MsgData.Length)); + return requestBuffer.TryCopyTo(responseBuffer); } - public readonly struct Header + private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) { - public const int TotalHeaderSize = sizeof(int); - - public readonly int DataSize; - - internal Header(int dataSize) - { - DataSize = dataSize; - } - - public static Header Deserialise(in Memory<byte> buffer) + switch (args.LastOperation) { - Span<byte> serialisedDataSize = buffer.Slice(0, sizeof(int)).Span; - int dataSize = EndianAwareBitConverter.ToInt32(serialisedDataSize); + case SocketAsyncOperation.Accept: + StartDefaultAccept(); + CompleteAccept(args); + break; - return new Header(dataSize); - } + case SocketAsyncOperation.Send: + CompleteSend(args); + break; - public void Serialise(in Memory<byte> buffer) - { - Span<byte> serialisedDataSize = EndianAwareBitConverter.GetBytes(DataSize); - serialisedDataSize.CopyTo(buffer.Slice(0, sizeof(int)).Span); + case SocketAsyncOperation.Receive: + CompleteReceive(args); + break; } } - } - - public sealed class RawStreamNetworkReader : RawNetworkReaderBase - { - private readonly RawStreamRequestHandler requestHandler; /// <inheritdoc /> - public RawStreamNetworkReader(ref Socket rawConnection, RawStreamRequestHandler? requestHandler, EndPoint defaultEndPoint, int maxMessageSize, - int pooledBuffersPerBucket = 50, uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxMessageSize, - pooledBuffersPerBucket, preallocatedStateObjects) - { - this.requestHandler = requestHandler ?? DefaultRequestHandler; - } - - private static bool DefaultRequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, - in Memory<byte> responseBuffer) + protected sealed override bool CanReuseStateObject(ref SocketAsyncEventArgs instance) { - return requestBuffer.TryCopyTo(responseBuffer); + return true; } - private void CloseClientConnection(SocketAsyncEventArgs args) + protected void CloseClientConnection(SocketAsyncEventArgs args) { args.BufferList = null; @@ -106,192 +73,14 @@ namespace NetSharp.Raw.Stream ArgsPool.Return(args); } - private void CompleteAccept(SocketAsyncEventArgs args) - { - switch (args.SocketError) - { - case SocketError.Success: - //ConfigureReceiveHeader(args); // inlined for performance - byte[] receiveBuffer = BufferPool.Rent(RawMessage.Header.TotalHeaderSize); - args.SetBuffer(receiveBuffer, 0, RawMessage.Header.TotalHeaderSize); - - args.UserToken = new TransmissionToken(RawMessage.Header.TotalHeaderSize, 0); - - StartReceive(args); - break; - - default: - ArgsPool.Return(args); - break; - } - - StartDefaultAccept(); - } - - private void CompleteReceive(SocketAsyncEventArgs args) - { - TransmissionToken token = (TransmissionToken)args.UserToken; + protected abstract void CompleteAccept(SocketAsyncEventArgs args); - byte[] receiveBuffer = args.Buffer; - Memory<byte> receiveBufferMemory = new Memory<byte>(receiveBuffer); + protected abstract void CompleteReceive(SocketAsyncEventArgs args); - int expectedBytes = token.ExpectedBytes; - - bool readHeader = expectedBytes == RawMessage.Header.TotalHeaderSize; - - switch (args.SocketError) - { - case SocketError.Success: - int receivedBytes = args.BytesTransferred, previousReceivedBytes = token.BytesTransferred, totalReceivedBytes = previousReceivedBytes + receivedBytes; - - if (totalReceivedBytes == expectedBytes) // transmission complete - { - if (readHeader) // handle a received message header - { - Memory<byte> headerBuffer = receiveBufferMemory.Slice(0, RawMessage.Header.TotalHeaderSize); - RawMessage.Header header = RawMessage.Header.Deserialise(in headerBuffer); - - //ConfigureReceiveData(args, in header); // inlined for performance - byte[] newReceiveBuffer = BufferPool.Rent(header.DataSize); - args.SetBuffer(newReceiveBuffer, 0, header.DataSize); - - token = new TransmissionToken(header.DataSize, 0); - args.UserToken = token; - - StartReceive(args); - } - else // handle a received message header - { - EndPoint clientEndPoint = args.AcceptSocket.RemoteEndPoint; - - int responseBufferSize = RawMessage.Header.TotalHeaderSize + expectedBytes; - byte[] responseBuffer = BufferPool.Rent(responseBufferSize); - Memory<byte> responseBufferMemory = new Memory<byte>(responseBuffer); - - Memory<byte> headerMemory = responseBufferMemory.Slice(0, RawMessage.Header.TotalHeaderSize); - Memory<byte> responseMemory = responseBufferMemory.Slice(RawMessage.Header.TotalHeaderSize, expectedBytes); - - bool responseExists = requestHandler(clientEndPoint, receiveBuffer[..expectedBytes], totalReceivedBytes, responseMemory); - BufferPool.Return(receiveBuffer, true); - - if (responseExists) - { - RawMessage.Header responseHeader = new RawMessage.Header(responseMemory.Length); - responseHeader.Serialise(in headerMemory); - - args.SetBuffer(responseBuffer, 0, responseBufferSize); - - TransmissionToken sendToken = new TransmissionToken(responseBufferSize, 0); - args.UserToken = sendToken; - - StartSend(args); - return; - } - - BufferPool.Return(responseBuffer, true); - - //ConfigureReceiveHeader(args); // inlined for performance - byte[] newReceiveBuffer = BufferPool.Rent(RawMessage.Header.TotalHeaderSize); - args.SetBuffer(newReceiveBuffer, 0, RawMessage.Header.TotalHeaderSize); - - args.UserToken = new TransmissionToken(RawMessage.Header.TotalHeaderSize, 0); - - StartReceive(args); - } - } - else if (0 < totalReceivedBytes && totalReceivedBytes < expectedBytes) // transmission not complete - { - token = new TransmissionToken(in token, receivedBytes); - args.UserToken = token; - - args.SetBuffer(totalReceivedBytes, expectedBytes - totalReceivedBytes); - - ContinueReceive(args); - } - else if (receivedBytes == 0) // connection is dead - { - CloseClientConnection(args); - } - break; - - default: - CloseClientConnection(args); - break; - } - } - - private void CompleteSend(SocketAsyncEventArgs args) - { - TransmissionToken token = (TransmissionToken)args.UserToken; - - byte[] sendBuffer = args.Buffer; - int expectedBytes = token.ExpectedBytes; - - switch (args.SocketError) - { - case SocketError.Success: - int sentBytes = args.BytesTransferred, previousSentBytes = token.BytesTransferred, totalSentBytes = previousSentBytes + sentBytes; - - if (totalSentBytes == expectedBytes) // transmission complete - { - BufferPool.Return(sendBuffer, true); - - //ConfigureReceiveHeader(args); // inlined for performance - byte[] newReceiveBuffer = BufferPool.Rent(RawMessage.Header.TotalHeaderSize); - args.SetBuffer(newReceiveBuffer, 0, RawMessage.Header.TotalHeaderSize); - - args.UserToken = new TransmissionToken(RawMessage.Header.TotalHeaderSize, 0); - - StartReceive(args); - } - else if (0 < totalSentBytes && totalSentBytes < expectedBytes) // transmission not complete - { - token = new TransmissionToken(in token, sentBytes); - args.UserToken = token; - - args.SetBuffer(totalSentBytes, expectedBytes - totalSentBytes); - - ContinueSend(args); - } - else if (sentBytes == 0) // connection is dead - { - CloseClientConnection(args); - } - break; - - default: - CloseClientConnection(args); - break; - } - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ConfigureReceiveData(SocketAsyncEventArgs args, in RawMessage.Header header) - { - byte[] receiveBuffer = BufferPool.Rent(header.DataSize); - args.SetBuffer(receiveBuffer, 0, header.DataSize); - - TransmissionToken token = new TransmissionToken(header.DataSize, 0); - args.UserToken = token; - } + protected abstract void CompleteSend(SocketAsyncEventArgs args); [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ConfigureReceiveHeader(SocketAsyncEventArgs args) - { - byte[] receiveBuffer = BufferPool.Rent(RawMessage.Header.TotalHeaderSize); - args.SetBuffer(receiveBuffer, 0, RawMessage.Header.TotalHeaderSize); - - TransmissionToken token = new TransmissionToken(RawMessage.Header.TotalHeaderSize, 0); - args.UserToken = token; - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ConfigureSend(SocketAsyncEventArgs args) - { - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ContinueReceive(SocketAsyncEventArgs args) + protected void ContinueReceive(SocketAsyncEventArgs args) { if (ShutdownToken.IsCancellationRequested) { @@ -307,7 +96,7 @@ namespace NetSharp.Raw.Stream } [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ContinueSend(SocketAsyncEventArgs args) + protected void ContinueSend(SocketAsyncEventArgs args) { if (ShutdownToken.IsCancellationRequested) { @@ -322,26 +111,29 @@ namespace NetSharp.Raw.Stream CompleteSend(args); } - private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) + /// <inheritdoc /> + protected sealed override SocketAsyncEventArgs CreateStateObject() { - switch (args.LastOperation) - { - case SocketAsyncOperation.Accept: - StartDefaultAccept(); - CompleteAccept(args); - break; + SocketAsyncEventArgs args = new SocketAsyncEventArgs(); + args.Completed += HandleIoCompleted; - case SocketAsyncOperation.Send: - CompleteSend(args); - break; + return args; + } - case SocketAsyncOperation.Receive: - CompleteReceive(args); - break; - } + /// <inheritdoc /> + protected sealed override void DestroyStateObject(SocketAsyncEventArgs instance) + { + instance.Completed -= HandleIoCompleted; + instance.Dispose(); + } + + /// <inheritdoc /> + protected sealed override void ResetStateObject(ref SocketAsyncEventArgs instance) + { + instance.AcceptSocket = null; } - private void StartAccept(SocketAsyncEventArgs args) + protected void StartAccept(SocketAsyncEventArgs args) { if (ShutdownToken.IsCancellationRequested) { @@ -354,7 +146,7 @@ namespace NetSharp.Raw.Stream CompleteAccept(args); } - private void StartDefaultAccept() + protected void StartDefaultAccept() { if (ShutdownToken.IsCancellationRequested) { @@ -365,7 +157,7 @@ namespace NetSharp.Raw.Stream StartAccept(args); } - private void StartReceive(SocketAsyncEventArgs args) + protected void StartReceive(SocketAsyncEventArgs args) { if (ShutdownToken.IsCancellationRequested) { @@ -380,7 +172,7 @@ namespace NetSharp.Raw.Stream CompleteReceive(args); } - private void StartSend(SocketAsyncEventArgs args) + protected void StartSend(SocketAsyncEventArgs args) { if (ShutdownToken.IsCancellationRequested) { @@ -396,35 +188,7 @@ namespace NetSharp.Raw.Stream } /// <inheritdoc /> - protected override bool CanReuseStateObject(ref SocketAsyncEventArgs instance) - { - return true; - } - - /// <inheritdoc /> - protected override SocketAsyncEventArgs CreateStateObject() - { - SocketAsyncEventArgs args = new SocketAsyncEventArgs(); - args.Completed += HandleIoCompleted; - - return args; - } - - /// <inheritdoc /> - protected override void DestroyStateObject(SocketAsyncEventArgs instance) - { - instance.Completed -= HandleIoCompleted; - instance.Dispose(); - } - - /// <inheritdoc /> - protected override void ResetStateObject(ref SocketAsyncEventArgs instance) - { - instance.AcceptSocket = null; - } - - /// <inheritdoc /> - public override void Start(ushort concurrentReadTasks) + public sealed override void Start(ushort concurrentReadTasks) { for (ushort i = 0; i < concurrentReadTasks; i++) { @@ -432,7 +196,7 @@ namespace NetSharp.Raw.Stream } } - private readonly struct TransmissionToken + protected readonly struct TransmissionToken { public readonly int BytesTransferred; public readonly int ExpectedBytes; diff --git a/NetSharp/NetSharp/Raw/Stream/RawStreamNetworkWriter.cs b/NetSharp/NetSharp/Raw/Stream/RawStreamNetworkWriter.cs @@ -6,16 +6,17 @@ using System.Threading.Tasks; namespace NetSharp.Raw.Stream { - public sealed class RawStreamNetworkWriter : RawNetworkWriterBase + public abstract class RawStreamNetworkWriter : RawNetworkWriterBase { - // TODO replace with proper packet size - private readonly int datagramSize; - /// <inheritdoc /> - public RawStreamNetworkWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int pooledPacketBufferSize, int pooledBuffersPerBucket = 50, - uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, pooledPacketBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) + protected RawStreamNetworkWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxMessageSize, int pooledBuffersPerBucket = 50, + uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxMessageSize, pooledBuffersPerBucket, preallocatedStateObjects) { - datagramSize = pooledPacketBufferSize; + if (maxMessageSize <= 0) + { + throw new ArgumentOutOfRangeException(nameof(maxMessageSize), maxMessageSize, + $"The message size must be greater than 0"); + } } private void CompleteConnect(SocketAsyncEventArgs args) @@ -64,103 +65,40 @@ namespace NetSharp.Raw.Stream ArgsPool.Return(args); } - private void CompleteReceive(SocketAsyncEventArgs args) + private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) { - AsyncStreamReadToken token = (AsyncStreamReadToken)args.UserToken; - - byte[] receiveBuffer = args.Buffer; - int expectedBytes = receiveBuffer.Length; - - switch (args.SocketError) + switch (args.LastOperation) { - case SocketError.Success: - int receivedBytes = args.BytesTransferred, totalReceivedBytes = token.TotalReadBytes; - - if (totalReceivedBytes + receivedBytes == expectedBytes) // transmission complete - { - receiveBuffer.CopyTo(token.UserBuffer); - token.CompletionSource.SetResult(totalReceivedBytes + receivedBytes); - } - else if (0 < totalReceivedBytes + receivedBytes && totalReceivedBytes + receivedBytes < expectedBytes) // transmission not complete - { - // update user token to take account of newly read bytes - token = new AsyncStreamReadToken(in token, receivedBytes); - args.UserToken = token; - - args.SetBuffer(totalReceivedBytes, expectedBytes - receivedBytes); - - ContinueReceive(args); - return; - } - else if (receivedBytes == 0) // connection is dead - { - token.CompletionSource.SetException(new SocketException((int)SocketError.HostDown)); - } + case SocketAsyncOperation.Connect: + CompleteConnect(args); break; - case SocketError.OperationAborted: - token.CompletionSource.SetCanceled(); + case SocketAsyncOperation.Disconnect: + CompleteDisconnect(args); break; - default: - int errorCode = (int)args.SocketError; - token.CompletionSource.SetException(new SocketException(errorCode)); + case SocketAsyncOperation.Receive: + CompleteReceive(args); break; - } - BufferPool.Return(receiveBuffer, true); - ArgsPool.Return(args); + case SocketAsyncOperation.Send: + CompleteSend(args); + break; + } } - private void CompleteSend(SocketAsyncEventArgs args) + /// <inheritdoc /> + protected sealed override bool CanReuseStateObject(ref SocketAsyncEventArgs instance) { - AsyncStreamWriteToken token = (AsyncStreamWriteToken)args.UserToken; - - byte[] sendBuffer = args.Buffer; - int expectedBytes = sendBuffer.Length; - - switch (args.SocketError) - { - case SocketError.Success: - int sentBytes = args.BytesTransferred, totalSentBytes = token.TotalWrittenBytes; - - if (totalSentBytes + sentBytes == expectedBytes) // transmission complete - { - token.CompletionSource.SetResult(totalSentBytes + sentBytes); - } - else if (0 < totalSentBytes + sentBytes && totalSentBytes + sentBytes < expectedBytes) // transmission not complete - { - // update user token to take account of newly written bytes - token = new AsyncStreamWriteToken(in token, sentBytes); - args.UserToken = token; - - args.SetBuffer(totalSentBytes, expectedBytes - sentBytes); - - ContinueSend(args); - return; - } - else if (sentBytes == 0) // connection is dead - { - token.CompletionSource.SetException(new SocketException((int)SocketError.HostDown)); - } - break; - - case SocketError.OperationAborted: - token.CompletionSource.SetCanceled(); - break; + return true; + } - default: - int errorCode = (int)args.SocketError; - token.CompletionSource.SetException(new SocketException(errorCode)); - break; - } + protected abstract void CompleteReceive(SocketAsyncEventArgs args); - BufferPool.Return(sendBuffer, true); - ArgsPool.Return(args); - } + protected abstract void CompleteSend(SocketAsyncEventArgs args); [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ContinueReceive(SocketAsyncEventArgs args) + protected void ContinueReceive(SocketAsyncEventArgs args) { if (Connection.ReceiveAsync(args)) return; @@ -168,43 +106,15 @@ namespace NetSharp.Raw.Stream } [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ContinueSend(SocketAsyncEventArgs args) + protected void ContinueSend(SocketAsyncEventArgs args) { if (Connection.SendAsync(args)) return; CompleteSend(args); } - private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) - { - switch (args.LastOperation) - { - case SocketAsyncOperation.Connect: - CompleteConnect(args); - break; - - case SocketAsyncOperation.Disconnect: - CompleteDisconnect(args); - break; - - case SocketAsyncOperation.Receive: - CompleteReceive(args); - break; - - case SocketAsyncOperation.Send: - CompleteSend(args); - break; - } - } - /// <inheritdoc /> - protected override bool CanReuseStateObject(ref SocketAsyncEventArgs instance) - { - return true; - } - - /// <inheritdoc /> - protected override SocketAsyncEventArgs CreateStateObject() + protected sealed override SocketAsyncEventArgs CreateStateObject() { SocketAsyncEventArgs args = new SocketAsyncEventArgs(); args.Completed += HandleIoCompleted; @@ -213,25 +123,25 @@ namespace NetSharp.Raw.Stream } /// <inheritdoc /> - protected override void DestroyStateObject(SocketAsyncEventArgs instance) + protected sealed override void DestroyStateObject(SocketAsyncEventArgs instance) { instance.Completed -= HandleIoCompleted; instance.Dispose(); } /// <inheritdoc /> - protected override void ResetStateObject(ref SocketAsyncEventArgs instance) + protected sealed override void ResetStateObject(ref SocketAsyncEventArgs instance) { } /// <inheritdoc /> - public override void Connect(EndPoint remoteEndPoint) + public sealed override void Connect(EndPoint remoteEndPoint) { Connection.Connect(remoteEndPoint); } /// <inheritdoc /> - public override ValueTask ConnectAsync(EndPoint remoteEndPoint) + public sealed override ValueTask ConnectAsync(EndPoint remoteEndPoint) { TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>(); SocketAsyncEventArgs args = ArgsPool.Rent(); @@ -269,217 +179,5 @@ namespace NetSharp.Raw.Stream return new ValueTask(); } - - /// <inheritdoc /> - public override int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) - { - int totalBytes = readBuffer.Length; - if (totalBytes > datagramSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {datagramSize} bytes", - nameof(readBuffer.Length) - ); - } - - int readBytes = 0; - - byte[] transmissionBuffer = BufferPool.Rent(totalBytes); - - do - { - readBytes += Connection.Receive(transmissionBuffer, readBytes, totalBytes - readBytes, flags); - } while (readBytes < totalBytes && readBytes != 0); - - transmissionBuffer.CopyTo(readBuffer); - BufferPool.Return(transmissionBuffer, true); - - return readBytes; - } - - /// <inheritdoc /> - public override ValueTask<int> ReadAsync(EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) - { - int totalBytes = readBuffer.Length; - if (totalBytes > datagramSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {datagramSize} bytes", - nameof(readBuffer.Length) - ); - } - - TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); - SocketAsyncEventArgs args = ArgsPool.Rent(); - - byte[] transmissionBuffer = BufferPool.Rent(totalBytes); - - args.SetBuffer(transmissionBuffer, 0, datagramSize); - - args.RemoteEndPoint = remoteEndPoint; - args.SocketFlags = flags; - - AsyncStreamReadToken token = new AsyncStreamReadToken(tcs, 0, in readBuffer); - args.UserToken = token; - - if (Connection.ReceiveAsync(args)) return new ValueTask<int>(tcs.Task); - - // inlining CompleteReceive(SocketAsyncEventArgs) for performance - int receivedBytes = args.BytesTransferred, totalReceivedBytes = token.TotalReadBytes; - - if (totalReceivedBytes + receivedBytes == datagramSize) // transmission complete - { - transmissionBuffer.CopyTo(readBuffer); - - BufferPool.Return(transmissionBuffer, true); - ArgsPool.Return(args); - - return new ValueTask<int>(totalReceivedBytes + receivedBytes); - } - else if (0 < totalReceivedBytes + receivedBytes && totalReceivedBytes + receivedBytes < datagramSize) // transmission not complete - { - // update user token to take account of newly read bytes - token = new AsyncStreamReadToken(in token, receivedBytes); - args.UserToken = token; - - args.SetBuffer(totalReceivedBytes, datagramSize - receivedBytes); - - ContinueReceive(args); - } - else if (receivedBytes == 0) // connection is dead - { - token.CompletionSource.SetException(new SocketException((int)SocketError.HostDown)); - } - - return new ValueTask<int>(tcs.Task); - } - - /// <inheritdoc /> - public override int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, SocketFlags flags = SocketFlags.None) - { - int totalBytes = writeBuffer.Length; - if (totalBytes > datagramSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {datagramSize} bytes", - nameof(writeBuffer.Length) - ); - } - - int writtenBytes = 0; - - byte[] transmissionBuffer = BufferPool.Rent(totalBytes); - writeBuffer.CopyTo(transmissionBuffer); - - do - { - writtenBytes += Connection.Send(transmissionBuffer, writtenBytes, totalBytes - writtenBytes, flags); - } while (writtenBytes < totalBytes && writtenBytes != 0); - - BufferPool.Return(transmissionBuffer); - - return writtenBytes; - } - - /// <inheritdoc /> - public override ValueTask<int> WriteAsync(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, SocketFlags flags = SocketFlags.None) - { - int totalBytes = writeBuffer.Length; - if (totalBytes > datagramSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {datagramSize} bytes", - nameof(writeBuffer.Length) - ); - } - - TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); - SocketAsyncEventArgs args = ArgsPool.Rent(); - - byte[] transmissionBuffer = BufferPool.Rent(totalBytes); - writeBuffer.CopyTo(transmissionBuffer); - - args.SetBuffer(transmissionBuffer, 0, datagramSize); - - args.RemoteEndPoint = remoteEndPoint; - args.SocketFlags = flags; - - AsyncStreamWriteToken token = new AsyncStreamWriteToken(tcs, 0); - args.UserToken = token; - - if (Connection.SendAsync(args)) return new ValueTask<int>(tcs.Task); - - // inlining CompleteSend(SocketAsyncEventArgs) for performance - int sentBytes = args.BytesTransferred, totalSentBytes = token.TotalWrittenBytes; - - if (totalSentBytes + sentBytes == datagramSize) // transmission complete - { - BufferPool.Return(transmissionBuffer, true); - ArgsPool.Return(args); - - return new ValueTask<int>(totalSentBytes + sentBytes); - } - else if (0 < totalSentBytes + sentBytes && totalSentBytes + sentBytes < datagramSize) // transmission not complete - { - // update user token to take account of newly written bytes - token = new AsyncStreamWriteToken(in token, sentBytes); - args.UserToken = token; - - args.SetBuffer(totalSentBytes, datagramSize - sentBytes); - - ContinueSend(args); - } - else if (sentBytes == 0) // connection is dead - { - token.CompletionSource.SetException(new SocketException((int)SocketError.HostDown)); - } - - return new ValueTask<int>(tcs.Task); - } - - private readonly struct AsyncStreamReadToken - { - public readonly TaskCompletionSource<int> CompletionSource; - public readonly int TotalReadBytes; - public readonly Memory<byte> UserBuffer; - - public AsyncStreamReadToken(TaskCompletionSource<int> completionSource, int totalReadBytes, in Memory<byte> userBuffer) - { - CompletionSource = completionSource; - - TotalReadBytes = totalReadBytes; - - UserBuffer = userBuffer; - } - - public AsyncStreamReadToken(in AsyncStreamReadToken previousToken, int newlyReadBytes) - { - CompletionSource = previousToken.CompletionSource; - - TotalReadBytes = previousToken.TotalReadBytes + newlyReadBytes; - - UserBuffer = previousToken.UserBuffer; - } - } - - private readonly struct AsyncStreamWriteToken - { - public readonly TaskCompletionSource<int> CompletionSource; - public readonly int TotalWrittenBytes; - - public AsyncStreamWriteToken(TaskCompletionSource<int> completionSource, int totalWrittenBytes) - { - CompletionSource = completionSource; - - TotalWrittenBytes = totalWrittenBytes; - } - - public AsyncStreamWriteToken(in AsyncStreamWriteToken previousToken, int newlyWrittenBytes) - { - CompletionSource = previousToken.CompletionSource; - - TotalWrittenBytes = previousToken.TotalWrittenBytes + newlyWrittenBytes; - } - } } } \ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Stream/RawStreamPacket.cs b/NetSharp/NetSharp/Raw/Stream/RawStreamPacket.cs @@ -5,31 +5,64 @@ namespace NetSharp.Raw.Stream { public readonly struct RawStreamPacket { - } + public readonly Memory<byte> Data; + public readonly Header PacketHeader; - public readonly struct RawStreamPacketHeader - { - public const int HeaderSize = sizeof(int); - public readonly int PacketSize; + private RawStreamPacket(in Header header, in Memory<byte> data) + { + PacketHeader = header; + + Data = data; + } - private RawStreamPacketHeader(int packetSize) + public RawStreamPacket(in Memory<byte> data) { - PacketSize = packetSize; + PacketHeader = new Header(data.Length); + + Data = data; } - public static RawStreamPacketHeader Deserialise(Memory<byte> buffer) + public static RawStreamPacket Deserialise(in Memory<byte> buffer) { - Span<byte> packetSizeSpan = buffer.Slice(0, sizeof(int)).Span; + Memory<byte> serialisedHeader = buffer.Slice(0, Header.TotalHeaderSize); + Header header = Header.Deserialise(in serialisedHeader); - int packetSize = EndianAwareBitConverter.ToInt32(packetSizeSpan); + Memory<byte> serialisedData = buffer.Slice(Header.TotalHeaderSize); - return new RawStreamPacketHeader(packetSize); + return new RawStreamPacket(in header, in serialisedData); } - public static void Serialise(RawStreamPacketHeader instance, Memory<byte> buffer) + public void Serialise(in Memory<byte> buffer) { - Span<byte> packetSizeSpan = buffer.Slice(0, sizeof(int)).Span; - EndianAwareBitConverter.GetBytes(instance.PacketSize).CopyTo(packetSizeSpan); + PacketHeader.Serialise(buffer.Slice(0, Header.TotalHeaderSize)); + + Data.CopyTo(buffer.Slice(Header.TotalHeaderSize, Data.Length)); + } + + public readonly struct Header + { + public const int TotalHeaderSize = sizeof(int); + + public readonly int DataSize; + + internal Header(int dataSize) + { + DataSize = dataSize; + } + + public static Header Deserialise(in Memory<byte> buffer) + { + Span<byte> serialisedDataSize = buffer.Slice(0, sizeof(int)).Span; + int dataSize = EndianAwareBitConverter.ToInt32(serialisedDataSize); + + return new Header(dataSize); + } + + public void Serialise(in Memory<byte> buffer) + { + Span<byte> serialisedDataSize = EndianAwareBitConverter.GetBytes(DataSize); + serialisedDataSize.CopyTo(buffer.Slice(0, sizeof(int)).Span); + } } } } \ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Stream/VariablePacketRawStreamNetworkReader.cs b/NetSharp/NetSharp/Raw/Stream/VariablePacketRawStreamNetworkReader.cs @@ -0,0 +1,182 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Runtime.CompilerServices; + +namespace NetSharp.Raw.Stream +{ + public sealed class VariablePacketRawStreamNetworkReader : RawStreamNetworkReader + { + /// <inheritdoc /> + public VariablePacketRawStreamNetworkReader(ref Socket rawConnection, RawStreamRequestHandler? requestHandler, EndPoint defaultEndPoint, int maxMessageSize, + int pooledBuffersPerBucket = 50, uint preallocatedStateObjects = 0) : base(ref rawConnection, requestHandler, defaultEndPoint, maxMessageSize, + pooledBuffersPerBucket, preallocatedStateObjects) + { + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ConfigureReceiveData(SocketAsyncEventArgs args, in RawStreamPacket.Header header) + { + byte[] receiveBuffer = BufferPool.Rent(header.DataSize); + args.SetBuffer(receiveBuffer, 0, header.DataSize); + + TransmissionToken token = new TransmissionToken(header.DataSize, 0); + args.UserToken = token; + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ConfigureReceiveHeader(SocketAsyncEventArgs args) + { + byte[] receiveBuffer = BufferPool.Rent(RawStreamPacket.Header.TotalHeaderSize); + args.SetBuffer(receiveBuffer, 0, RawStreamPacket.Header.TotalHeaderSize); + + TransmissionToken token = new TransmissionToken(RawStreamPacket.Header.TotalHeaderSize, 0); + args.UserToken = token; + } + + /// <inheritdoc /> + protected override void CompleteAccept(SocketAsyncEventArgs args) + { + switch (args.SocketError) + { + case SocketError.Success: + ConfigureReceiveHeader(args); + + StartReceive(args); + break; + + default: + ArgsPool.Return(args); + break; + } + + StartDefaultAccept(); + } + + /// <inheritdoc /> + protected override void CompleteReceive(SocketAsyncEventArgs args) + { + TransmissionToken token = (TransmissionToken)args.UserToken; + + byte[] receiveBuffer = args.Buffer; + Memory<byte> receiveBufferMemory = new Memory<byte>(receiveBuffer); + + int expectedBytes = token.ExpectedBytes; + + bool readHeader = expectedBytes == RawStreamPacket.Header.TotalHeaderSize; + + switch (args.SocketError) + { + case SocketError.Success: + int receivedBytes = args.BytesTransferred, previousReceivedBytes = token.BytesTransferred, totalReceivedBytes = previousReceivedBytes + receivedBytes; + + if (totalReceivedBytes == expectedBytes) // transmission complete + { + if (readHeader) // handle a received message header + { + Memory<byte> headerBuffer = receiveBufferMemory.Slice(0, RawStreamPacket.Header.TotalHeaderSize); + RawStreamPacket.Header header = RawStreamPacket.Header.Deserialise(in headerBuffer); + + ConfigureReceiveData(args, in header); + + StartReceive(args); + } + else // handle a received message header + { + EndPoint clientEndPoint = args.AcceptSocket.RemoteEndPoint; + + int responseBufferSize = RawStreamPacket.Header.TotalHeaderSize + expectedBytes; + byte[] responseBuffer = BufferPool.Rent(responseBufferSize); + Memory<byte> responseBufferMemory = new Memory<byte>(responseBuffer); + + Memory<byte> headerMemory = responseBufferMemory.Slice(0, RawStreamPacket.Header.TotalHeaderSize); + Memory<byte> responseMemory = responseBufferMemory.Slice(RawStreamPacket.Header.TotalHeaderSize, expectedBytes); + + bool responseExists = RequestHandler(clientEndPoint, receiveBuffer[..expectedBytes], totalReceivedBytes, responseMemory); + BufferPool.Return(receiveBuffer, true); + + if (responseExists) + { + RawStreamPacket.Header responseHeader = new RawStreamPacket.Header(responseMemory.Length); + responseHeader.Serialise(in headerMemory); + + args.SetBuffer(responseBuffer, 0, responseBufferSize); + + TransmissionToken sendToken = new TransmissionToken(responseBufferSize, 0); + args.UserToken = sendToken; + + StartSend(args); + return; + } + + BufferPool.Return(responseBuffer, true); + + ConfigureReceiveHeader(args); + + StartReceive(args); + } + } + else if (0 < totalReceivedBytes && totalReceivedBytes < expectedBytes) // transmission not complete + { + token = new TransmissionToken(in token, receivedBytes); + args.UserToken = token; + + args.SetBuffer(totalReceivedBytes, expectedBytes - totalReceivedBytes); + + ContinueReceive(args); + } + else if (receivedBytes == 0) // connection is dead + { + CloseClientConnection(args); + } + break; + + default: + CloseClientConnection(args); + break; + } + } + + /// <inheritdoc /> + protected override void CompleteSend(SocketAsyncEventArgs args) + { + TransmissionToken token = (TransmissionToken)args.UserToken; + + byte[] sendBuffer = args.Buffer; + int expectedBytes = token.ExpectedBytes; + + switch (args.SocketError) + { + case SocketError.Success: + int sentBytes = args.BytesTransferred, previousSentBytes = token.BytesTransferred, totalSentBytes = previousSentBytes + sentBytes; + + if (totalSentBytes == expectedBytes) // transmission complete + { + BufferPool.Return(sendBuffer, true); + + ConfigureReceiveHeader(args); + + StartReceive(args); + } + else if (0 < totalSentBytes && totalSentBytes < expectedBytes) // transmission not complete + { + token = new TransmissionToken(in token, sentBytes); + args.UserToken = token; + + args.SetBuffer(totalSentBytes, expectedBytes - totalSentBytes); + + ContinueSend(args); + } + else if (sentBytes == 0) // connection is dead + { + CloseClientConnection(args); + } + break; + + default: + CloseClientConnection(args); + break; + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Stream/VariablePacketRawStreamNetworkWriter.cs b/NetSharp/NetSharp/Raw/Stream/VariablePacketRawStreamNetworkWriter.cs @@ -0,0 +1,339 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Runtime.CompilerServices; +using System.Threading.Tasks; + +namespace NetSharp.Raw.Stream +{ + public sealed class VariablePacketRawStreamNetworkWriter : RawStreamNetworkWriter + { + //TODO replace as soon as possible + private const int MESSAGE_SIZE = 8192; + + /// <inheritdoc /> + public VariablePacketRawStreamNetworkWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxMessageSize, int pooledBuffersPerBucket = 50, + uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxMessageSize, pooledBuffersPerBucket, preallocatedStateObjects) + { + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ConfigureReceiveData(SocketAsyncEventArgs args, in RawStreamPacket.Header header) + { + byte[] receiveBuffer = BufferPool.Rent(header.DataSize); + args.SetBuffer(receiveBuffer, 0, header.DataSize); + + AsyncStreamReadToken token = new AsyncStreamReadToken(header.DataSize, 0); + args.UserToken = token; + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ConfigureReceiveHeader(SocketAsyncEventArgs args) + { + byte[] receiveBuffer = BufferPool.Rent(RawStreamPacket.Header.TotalHeaderSize); + args.SetBuffer(receiveBuffer, 0, RawStreamPacket.Header.TotalHeaderSize); + + AsyncStreamReadToken token = new AsyncStreamReadToken(RawStreamPacket.Header.TotalHeaderSize, 0); + args.UserToken = token; + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ConfigureSendData(SocketAsyncEventArgs args, in RawStreamPacket.Header header) + { + byte[] receiveBuffer = BufferPool.Rent(header.DataSize); + args.SetBuffer(receiveBuffer, 0, header.DataSize); + + AsyncStreamWriteToken token = new AsyncStreamWriteToken(header.DataSize, 0); + args.UserToken = token; + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ConfigureSendHeader(SocketAsyncEventArgs args) + { + byte[] receiveBuffer = BufferPool.Rent(RawStreamPacket.Header.TotalHeaderSize); + args.SetBuffer(receiveBuffer, 0, RawStreamPacket.Header.TotalHeaderSize); + + AsyncStreamWriteToken token = new AsyncStreamWriteToken(RawStreamPacket.Header.TotalHeaderSize, 0); + args.UserToken = token; + } + + /// <inheritdoc /> + protected override void CompleteReceive(SocketAsyncEventArgs args) + { + AsyncStreamReadToken token = (AsyncStreamReadToken)args.UserToken; + + byte[] receiveBuffer = args.Buffer; + int expectedBytes = receiveBuffer.Length; + + switch (args.SocketError) + { + case SocketError.Success: + int receivedBytes = args.BytesTransferred, totalReceivedBytes = token.TotalReadBytes; + + if (totalReceivedBytes + receivedBytes == expectedBytes) // transmission complete + { + receiveBuffer.CopyTo(token.UserBuffer); + token.CompletionSource.SetResult(totalReceivedBytes + receivedBytes); + } + else if (0 < totalReceivedBytes + receivedBytes && totalReceivedBytes + receivedBytes < expectedBytes) // transmission not complete + { + // update user token to take account of newly read bytes + token = new AsyncStreamReadToken(in token, receivedBytes); + args.UserToken = token; + + args.SetBuffer(totalReceivedBytes, expectedBytes - receivedBytes); + + ContinueReceive(args); + return; + } + else if (receivedBytes == 0) // connection is dead + { + token.CompletionSource.SetException(new SocketException((int)SocketError.HostDown)); + } + break; + + case SocketError.OperationAborted: + token.CompletionSource.SetCanceled(); + break; + + default: + int errorCode = (int)args.SocketError; + token.CompletionSource.SetException(new SocketException(errorCode)); + break; + } + + BufferPool.Return(receiveBuffer, true); + ArgsPool.Return(args); + } + + /// <inheritdoc /> + protected override void CompleteSend(SocketAsyncEventArgs args) + { + AsyncStreamWriteToken token = (AsyncStreamWriteToken)args.UserToken; + + byte[] sendBuffer = args.Buffer; + int expectedBytes = sendBuffer.Length; + + switch (args.SocketError) + { + case SocketError.Success: + int sentBytes = args.BytesTransferred, totalSentBytes = token.TotalWrittenBytes; + + if (totalSentBytes + sentBytes == expectedBytes) // transmission complete + { + token.CompletionSource.SetResult(totalSentBytes + sentBytes); + } + else if (0 < totalSentBytes + sentBytes && totalSentBytes + sentBytes < expectedBytes) // transmission not complete + { + // update user token to take account of newly written bytes + token = new AsyncStreamWriteToken(in token, sentBytes); + args.UserToken = token; + + args.SetBuffer(totalSentBytes, expectedBytes - sentBytes); + + ContinueSend(args); + return; + } + else if (sentBytes == 0) // connection is dead + { + token.CompletionSource.SetException(new SocketException((int)SocketError.HostDown)); + } + break; + + case SocketError.OperationAborted: + token.CompletionSource.SetCanceled(); + break; + + default: + int errorCode = (int)args.SocketError; + token.CompletionSource.SetException(new SocketException(errorCode)); + break; + } + + BufferPool.Return(sendBuffer, true); + ArgsPool.Return(args); + } + + /// <inheritdoc /> + public override int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) + { + byte[] transmissionHeaderBuffer = BufferPool.Rent(RawStreamPacket.Header.TotalHeaderSize); + int readHeaderBytes = 0; + + do + { + readHeaderBytes += Connection.Receive(transmissionHeaderBuffer, readHeaderBytes, + RawStreamPacket.Header.TotalHeaderSize - readHeaderBytes, flags); + } while (readHeaderBytes < RawStreamPacket.Header.TotalHeaderSize && readHeaderBytes != 0); + + RawStreamPacket.Header header = RawStreamPacket.Header.Deserialise(transmissionHeaderBuffer); + int totalBytes = header.DataSize; + BufferPool.Return(transmissionHeaderBuffer, true); + + if (totalBytes > readBuffer.Length) + throw new ArgumentException("Read buffer is too small to fully contain message.", nameof(readBuffer)); + + int readBytes = 0; + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + + do + { + readBytes += Connection.Receive(transmissionBuffer, readBytes, totalBytes - readBytes, flags); + } while (readBytes < totalBytes && readBytes != 0); + + transmissionBuffer.CopyTo(readBuffer); + BufferPool.Return(transmissionBuffer, true); + + return readBytes; + } + + /// <inheritdoc /> + public override ValueTask<int> ReadAsync(EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) + { + int totalBytes = readBuffer.Length; + if (totalBytes > MESSAGE_SIZE) + { + throw new ArgumentException( + $"Cannot receive a message of size: {totalBytes} bytes; maximum message size: {MESSAGE_SIZE} bytes", + nameof(readBuffer.Length) + ); + } + + TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); + SocketAsyncEventArgs args = ArgsPool.Rent(); + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + + args.SetBuffer(transmissionBuffer, 0, MESSAGE_SIZE); + + args.RemoteEndPoint = remoteEndPoint; + args.SocketFlags = flags; + + AsyncStreamReadToken token = new AsyncStreamReadToken(tcs, 0, in readBuffer); + args.UserToken = token; + + if (Connection.ReceiveAsync(args)) return new ValueTask<int>(tcs.Task); + + CompleteReceive(args); + + return new ValueTask<int>(tcs.Task); + } + + /// <inheritdoc /> + public override int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, SocketFlags flags = SocketFlags.None) + { + byte[] transmissionHeaderBuffer = BufferPool.Rent(RawStreamPacket.Header.TotalHeaderSize); + new RawStreamPacket.Header(writeBuffer.Length).Serialise(transmissionHeaderBuffer); + + int writtenHeaderBytes = 0; + do + { + writtenHeaderBytes += Connection.Send(transmissionHeaderBuffer, writtenHeaderBytes, + RawStreamPacket.Header.TotalHeaderSize - writtenHeaderBytes, flags); + } while (writtenHeaderBytes < RawStreamPacket.Header.TotalHeaderSize && writtenHeaderBytes != 0); + + int totalBytes = writeBuffer.Length; + int writtenBytes = 0; + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + writeBuffer.CopyTo(transmissionBuffer); + + do + { + writtenBytes += Connection.Send(transmissionBuffer, writtenBytes, totalBytes - writtenBytes, flags); + } while (writtenBytes < totalBytes && writtenBytes != 0); + + BufferPool.Return(transmissionBuffer); + + return writtenBytes; + } + + /// <inheritdoc /> + public override ValueTask<int> WriteAsync(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, SocketFlags flags = SocketFlags.None) + { + int totalBytes = writeBuffer.Length; + if (totalBytes > MESSAGE_SIZE) + { + throw new ArgumentException( + $"Cannot send a message of size: {totalBytes} bytes; maximum message size: {MESSAGE_SIZE} bytes", + nameof(writeBuffer.Length) + ); + } + + TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); + SocketAsyncEventArgs args = ArgsPool.Rent(); + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + writeBuffer.CopyTo(transmissionBuffer); + + args.SetBuffer(transmissionBuffer, 0, MESSAGE_SIZE); + + args.RemoteEndPoint = remoteEndPoint; + args.SocketFlags = flags; + + AsyncStreamWriteToken token = new AsyncStreamWriteToken(tcs, 0); + args.UserToken = token; + + if (Connection.SendAsync(args)) return new ValueTask<int>(tcs.Task); + + CompleteSend(args); + + return new ValueTask<int>(tcs.Task); + } + + private readonly struct AsyncStreamReadToken + { + public readonly TaskCompletionSource<int> CompletionSource; + public readonly int ExpectedBytes; + public readonly int TotalReadBytes; + public readonly Memory<byte> UserBuffer; + + public AsyncStreamReadToken(TaskCompletionSource<int> completionSource, int expectedBytes, int totalReadBytes, in Memory<byte> userBuffer) + { + CompletionSource = completionSource; + + ExpectedBytes = expectedBytes; + + TotalReadBytes = totalReadBytes; + + UserBuffer = userBuffer; + } + + public AsyncStreamReadToken(in AsyncStreamReadToken previousToken, int newlyReadBytes) + { + CompletionSource = previousToken.CompletionSource; + + ExpectedBytes = previousToken.ExpectedBytes; + + TotalReadBytes = previousToken.TotalReadBytes + newlyReadBytes; + + UserBuffer = previousToken.UserBuffer; + } + } + + private readonly struct AsyncStreamWriteToken + { + public readonly TaskCompletionSource<int> CompletionSource; + public readonly int ExpectedBytes; + public readonly int TotalWrittenBytes; + + public AsyncStreamWriteToken(TaskCompletionSource<int> completionSource, int expectedBytes, int totalWrittenBytes) + { + CompletionSource = completionSource; + + ExpectedBytes = expectedBytes; + + TotalWrittenBytes = totalWrittenBytes; + } + + public AsyncStreamWriteToken(in AsyncStreamWriteToken previousToken, int newlyWrittenBytes) + { + CompletionSource = previousToken.CompletionSource; + + ExpectedBytes = previousToken.ExpectedBytes; + + TotalWrittenBytes = previousToken.TotalWrittenBytes + newlyWrittenBytes; + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkReaderBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkReaderBenchmark.cs @@ -23,7 +23,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Datagram Raw Network Reader Benchmark"; + public string Name { get; } = "Raw Datagram Network Reader Benchmark"; private static bool RequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, in Memory<byte> responseBuffer) diff --git a/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterAsyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterAsyncBenchmark.cs @@ -21,7 +21,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Datagram Raw Network Writer Benchmark (Asynchronous)"; + public string Name { get; } = "Raw Datagram Network Writer Benchmark (Asynchronous)"; private static Task ServerTask(CancellationToken cancellationToken) { diff --git a/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterSyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterSyncBenchmark.cs @@ -21,7 +21,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Datagram Raw Network Writer Benchmark (Synchronous)"; + public string Name { get; } = "Raw Datagram Network Writer Benchmark (Synchronous)"; private static Task ServerTask(CancellationToken cancellationToken) { diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkReaderBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkReaderBenchmark.cs @@ -0,0 +1,140 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Linq; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading; +using System.Threading.Tasks; + +namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks +{ + public class FixedPacketStreamNetworkReaderBenchmark : INetSharpBenchmark + { + private const int PacketSize = 8192, PacketCount = 1_000_000, ClientCount = 12; + + private double[] ClientBandwidths; + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = Encoding.UTF8; + public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12373); + + public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); + + /// <inheritdoc /> + public string Name { get; } = "Raw Fixed Packet-size Stream Network Reader Benchmark"; + + private static bool RequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + in Memory<byte> responseBuffer) + { + requestBuffer.CopyTo(responseBuffer); + + return true; + } + + private Task BenchmarkClientTask(object idObj) + { + int id = (int)idObj; + + BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); + + Socket clientSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + clientSocket.Bind(ClientEndPoint); + + ServerReadyEvent.Wait(); + clientSocket.Connect(ServerEndPoint); + + byte[] sendBuffer = new byte[PacketSize]; + byte[] receiveBuffer = new byte[PacketSize]; + + EndPoint remoteEndPoint = ServerEndPoint; + + lock (typeof(Console)) + { + Console.WriteLine($"[Client {id}] Starting client; sending messages to {remoteEndPoint}"); + } + + for (int i = 0; i < PacketCount; i++) + { + ServerEncoding.GetBytes($"[Client {id}] Hello World! (Packet {i})").CopyTo(sendBuffer, 0); + + benchmarkHelper.StartStopwatch(); + + int totalSent = 0; + do + { + totalSent += clientSocket.Send(sendBuffer, totalSent, sendBuffer.Length - totalSent, + SocketFlags.None); + } while (totalSent != 0 && totalSent < sendBuffer.Length); + + if (totalSent == 0) + { + break; + } + + int totalReceived = 0; + do + { + totalReceived += clientSocket.Receive(receiveBuffer, totalReceived, receiveBuffer.Length - totalReceived, + SocketFlags.None); + } while (totalReceived != 0 && totalReceived < receiveBuffer.Length); + + if (totalReceived == 0) + { + break; + } + + benchmarkHelper.StopStopwatch(); + + benchmarkHelper.SnapshotRttStats(); + } + + clientSocket.Disconnect(true); + clientSocket.Close(); + + benchmarkHelper.PrintBandwidthStats(id, PacketCount, PacketSize); + benchmarkHelper.PrintRttStats(id); + + ClientBandwidths[id] = benchmarkHelper.CalcBandwidth(PacketCount, PacketSize); + + return Task.CompletedTask; + } + + /// <inheritdoc /> + public async Task RunAsync() + { + if (PacketCount > 10_000) + { + Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); + } + + ClientBandwidths = new double[ClientCount]; + Task[] clientTasks = new Task[ClientCount]; + for (int i = 0; i < clientTasks.Length; i++) + { + clientTasks[i] = Task.Factory.StartNew(BenchmarkClientTask, i, TaskCreationOptions.LongRunning); + } + + EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ServerEndPoint); + rawSocket.Listen(ClientCount); + + using RawStreamNetworkReader reader = new FixedPacketRawStreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize); + reader.Start(ClientCount); + + ServerReadyEvent.Set(); + + await Task.WhenAll(clientTasks); + + Console.WriteLine($"Total estimated bandwidth: {ClientBandwidths.Sum():F3}"); + + reader.Stop(); + + rawSocket.Close(); + rawSocket.Dispose(); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkWriterAsyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkWriterAsyncBenchmark.cs @@ -0,0 +1,130 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading; +using System.Threading.Tasks; + +namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks +{ + public class FixedPacketStreamNetworkWriterAsyncBenchmark : INetSharpBenchmark + { + private const int PacketSize = 8192, PacketCount = 1_000_000; + + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = Encoding.UTF8; + public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12374); + + public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); + + /// <inheritdoc /> + public string Name { get; } = "Raw Fixed Packet-size Stream Network Writer Benchmark (Asynchronous)"; + + private static Task ServerTask(CancellationToken cancellationToken) + { + using Socket server = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + + server.Bind(ServerEndPoint); + ServerReadyEvent.Set(); + + byte[] transmissionBuffer = new byte[PacketSize]; + + server.Listen(1); + Socket clientSocket = server.Accept(); + + while (!cancellationToken.IsCancellationRequested) + { + int expectedBytes = transmissionBuffer.Length; + + int receivedBytes = 0; + do + { + receivedBytes += clientSocket.Receive(transmissionBuffer, receivedBytes, expectedBytes - receivedBytes, SocketFlags.None); + } while (receivedBytes != 0 && receivedBytes < expectedBytes); + + if (receivedBytes == 0) + { + break; + } + + int sentBytes = 0; + do + { + sentBytes += clientSocket.Send(transmissionBuffer, sentBytes, expectedBytes - sentBytes, SocketFlags.None); + } while (sentBytes != 0 && sentBytes < expectedBytes); + + if (sentBytes == 0) + { + break; + } + } + + server.Shutdown(SocketShutdown.Both); + server.Close(); + + return Task.CompletedTask; + } + + /// <inheritdoc /> + public async Task RunAsync() + { + if (PacketCount > 10_000) + { + Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); + } + + EndPoint defaultRemoteEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ClientEndPoint); + + using RawStreamNetworkWriter writer = new FixedPacketRawStreamNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); + + using CancellationTokenSource serverCts = new CancellationTokenSource(); + Task serverTask = Task.Factory.StartNew(state => ServerTask((CancellationToken)state), serverCts.Token, TaskCreationOptions.LongRunning); + + ServerReadyEvent.Wait(); + rawSocket.Connect(ServerEndPoint); + + BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); + + byte[] sendBuffer = new byte[PacketSize]; + byte[] receiveBuffer = new byte[PacketSize]; + + for (int i = 0; i < PacketCount; i++) + { + byte[] packetBuffer = ServerEncoding.GetBytes($"[Client 0] Hello World! (Packet {i})"); + packetBuffer.CopyTo(sendBuffer, 0); + + benchmarkHelper.StartStopwatch(); + int sendResult = await writer.WriteAsync(ServerEndPoint, sendBuffer); + + int receiveResult = await writer.ReadAsync(ServerEndPoint, receiveBuffer); + benchmarkHelper.StopStopwatch(); + + benchmarkHelper.SnapshotRttStats(); + } + + benchmarkHelper.PrintBandwidthStats(0, PacketCount, PacketSize); + benchmarkHelper.PrintRttStats(0); + + serverCts.Cancel(); + try + { + serverTask.Dispose(); + } + catch (Exception) + { + // ignored + } + + rawSocket.Disconnect(false); + rawSocket.Shutdown(SocketShutdown.Both); + rawSocket.Close(); + rawSocket.Dispose(); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkWriterSyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/FixedPacketStreamNetworkWriterSyncBenchmark.cs @@ -0,0 +1,134 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading; +using System.Threading.Tasks; + +namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks +{ + public class FixedPacketStreamNetworkWriterSyncBenchmark : INetSharpBenchmark + { + private const int PacketSize = 8192, PacketCount = 1_000_000; + + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = FixedPacketStreamNetworkReaderBenchmark.ServerEncoding; + public static readonly EndPoint ServerEndPoint = FixedPacketStreamNetworkReaderBenchmark.ServerEndPoint; + + public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); + + /// <inheritdoc /> + public string Name { get; } = "Raw Fixed Packet-size Stream Network Writer Benchmark (Synchronous)"; + + private static Task ServerTask(CancellationToken cancellationToken) + { + using Socket server = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + + server.Bind(ServerEndPoint); + ServerReadyEvent.Set(); + + byte[] transmissionBuffer = new byte[PacketSize]; + + server.Listen(1); + Socket clientSocket = server.Accept(); + + while (!cancellationToken.IsCancellationRequested) + { + int expectedBytes = transmissionBuffer.Length; + + int receivedBytes = 0; + do + { + receivedBytes += clientSocket.Receive(transmissionBuffer, receivedBytes, expectedBytes - receivedBytes, SocketFlags.None); + } while (receivedBytes != 0 && receivedBytes < expectedBytes); + + if (receivedBytes == 0) + { + break; + } + + int sentBytes = 0; + do + { + sentBytes += clientSocket.Send(transmissionBuffer, sentBytes, expectedBytes - sentBytes, SocketFlags.None); + } while (sentBytes != 0 && sentBytes < expectedBytes); + + if (sentBytes == 0) + { + break; + } + } + + server.Shutdown(SocketShutdown.Both); + server.Close(); + + return Task.CompletedTask; + } + + /// <inheritdoc /> + public Task RunAsync() + { + if (PacketCount > 10_000) + { + Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); + } + + EndPoint defaultRemoteEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ClientEndPoint); + + using RawStreamNetworkWriter writer = new FixedPacketRawStreamNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); + + using CancellationTokenSource serverCts = new CancellationTokenSource(); + Task serverTask = Task.Factory.StartNew(state => ServerTask((CancellationToken)state), serverCts.Token, TaskCreationOptions.LongRunning); + + ServerReadyEvent.Wait(); + rawSocket.Connect(ServerEndPoint); + + BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); + + byte[] sendBuffer = new byte[PacketSize]; + byte[] receiveBuffer = new byte[PacketSize]; + + EndPoint remoteEndPoint = ServerEndPoint; + + for (int i = 0; i < PacketCount; i++) + { + byte[] packetBuffer = ServerEncoding.GetBytes($"[Client 0] Hello World! (Packet {i})"); + packetBuffer.CopyTo(sendBuffer, 0); + + benchmarkHelper.StartStopwatch(); + int sendResult = writer.Write(ServerEndPoint, sendBuffer); + + int receiveResult = writer.Read(ref remoteEndPoint, receiveBuffer); + benchmarkHelper.StopStopwatch(); + + benchmarkHelper.SnapshotRttStats(); + } + + benchmarkHelper.PrintBandwidthStats(0, PacketCount, PacketSize); + benchmarkHelper.PrintRttStats(0); + + serverCts.Cancel(); + try + { + serverTask.Dispose(); + } + catch (Exception) + { + // ignored + } + + rawSocket.Disconnect(false); + rawSocket.Shutdown(SocketShutdown.Both); + rawSocket.Close(); + rawSocket.Dispose(); + + return Task.CompletedTask; + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkReaderBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkReaderBenchmark.cs @@ -1,144 +0,0 @@ -using NetSharp.Raw.Stream; - -using System; -using System.Linq; -using System.Net; -using System.Net.Sockets; -using System.Text; -using System.Threading; -using System.Threading.Tasks; - -namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks -{ - public class StreamNetworkReaderBenchmark : INetSharpBenchmark - { - private const int PacketSize = 8192, PacketCount = 1_000_000, ClientCount = 12; - - private double[] ClientBandwidths; - public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); - - public static readonly Encoding ServerEncoding = Encoding.UTF8; - public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12373); - - public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); - - /// <inheritdoc /> - public string Name { get; } = "Stream Raw Network Reader Benchmark"; - - private static bool RequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, - in Memory<byte> responseBuffer) - { - requestBuffer.CopyTo(responseBuffer); - - return true; - } - - private Task BenchmarkClientTask(object idObj) - { - int id = (int)idObj; - - BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); - - Socket clientSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - clientSocket.Bind(ClientEndPoint); - - ServerReadyEvent.Wait(); - clientSocket.Connect(ServerEndPoint); - - byte[] sendBuffer = new byte[PacketSize + RawMessage.Header.TotalHeaderSize]; - byte[] receiveBuffer = new byte[PacketSize + RawMessage.Header.TotalHeaderSize]; - byte[] packetBuffer = new byte[PacketSize]; - - EndPoint remoteEndPoint = ServerEndPoint; - - lock (typeof(Console)) - { - Console.WriteLine($"[Client {id}] Starting client; sending messages to {remoteEndPoint}"); - } - - for (int i = 0; i < PacketCount; i++) - { - ServerEncoding.GetBytes($"[Client {id}] Hello World! (Packet {i})").CopyTo(packetBuffer, 0); - - RawMessage message = new RawMessage(packetBuffer); - message.Serialise(sendBuffer); - - benchmarkHelper.StartStopwatch(); - - int totalSent = 0; - do - { - totalSent += clientSocket.Send(sendBuffer, totalSent, sendBuffer.Length - totalSent, - SocketFlags.None); - } while (totalSent != 0 && totalSent < sendBuffer.Length); - - if (totalSent == 0) - { - break; - } - - int totalReceived = 0; - do - { - totalReceived += clientSocket.Receive(receiveBuffer, totalReceived, receiveBuffer.Length - totalReceived, - SocketFlags.None); - } while (totalReceived != 0 && totalReceived < receiveBuffer.Length); - - if (totalReceived == 0) - { - break; - } - - benchmarkHelper.StopStopwatch(); - - benchmarkHelper.SnapshotRttStats(); - } - - clientSocket.Disconnect(true); - clientSocket.Close(); - - benchmarkHelper.PrintBandwidthStats(id, PacketCount, PacketSize); - benchmarkHelper.PrintRttStats(id); - - ClientBandwidths[id] = benchmarkHelper.CalcBandwidth(PacketCount, PacketSize); - - return Task.CompletedTask; - } - - /// <inheritdoc /> - public async Task RunAsync() - { - if (PacketCount > 10_000) - { - Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); - } - - ClientBandwidths = new double[ClientCount]; - Task[] clientTasks = new Task[ClientCount]; - for (int i = 0; i < clientTasks.Length; i++) - { - clientTasks[i] = Task.Factory.StartNew(BenchmarkClientTask, i, TaskCreationOptions.LongRunning); - } - - EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); - - Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - rawSocket.Bind(ServerEndPoint); - rawSocket.Listen(ClientCount); - - using RawStreamNetworkReader reader = new RawStreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize); - reader.Start(ClientCount); - - ServerReadyEvent.Set(); - - await Task.WhenAll(clientTasks); - - Console.WriteLine($"Total estimated bandwidth: {ClientBandwidths.Sum():F3}"); - - reader.Stop(); - - rawSocket.Close(); - rawSocket.Dispose(); - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterAsyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterAsyncBenchmark.cs @@ -1,130 +0,0 @@ -using NetSharp.Raw.Stream; - -using System; -using System.Net; -using System.Net.Sockets; -using System.Text; -using System.Threading; -using System.Threading.Tasks; - -namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks -{ - public class StreamNetworkWriterAsyncBenchmark : INetSharpBenchmark - { - private const int PacketSize = 8192, PacketCount = 1_000_000; - - public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); - - public static readonly Encoding ServerEncoding = Encoding.UTF8; - public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12374); - - public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); - - /// <inheritdoc /> - public string Name { get; } = "Stream Raw Network Writer Benchmark (Asynchronous)"; - - private static Task ServerTask(CancellationToken cancellationToken) - { - using Socket server = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - - server.Bind(ServerEndPoint); - ServerReadyEvent.Set(); - - byte[] transmissionBuffer = new byte[PacketSize]; - - server.Listen(1); - Socket clientSocket = server.Accept(); - - while (!cancellationToken.IsCancellationRequested) - { - int expectedBytes = transmissionBuffer.Length; - - int receivedBytes = 0; - do - { - receivedBytes += clientSocket.Receive(transmissionBuffer, receivedBytes, expectedBytes - receivedBytes, SocketFlags.None); - } while (receivedBytes != 0 && receivedBytes < expectedBytes); - - if (receivedBytes == 0) - { - break; - } - - int sentBytes = 0; - do - { - sentBytes += clientSocket.Send(transmissionBuffer, sentBytes, expectedBytes - sentBytes, SocketFlags.None); - } while (sentBytes != 0 && sentBytes < expectedBytes); - - if (sentBytes == 0) - { - break; - } - } - - server.Shutdown(SocketShutdown.Both); - server.Close(); - - return Task.CompletedTask; - } - - /// <inheritdoc /> - public async Task RunAsync() - { - if (PacketCount > 10_000) - { - Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); - } - - EndPoint defaultRemoteEndPoint = new IPEndPoint(IPAddress.Any, 0); - - Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - rawSocket.Bind(ClientEndPoint); - - using RawStreamNetworkWriter writer = new RawStreamNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); - - using CancellationTokenSource serverCts = new CancellationTokenSource(); - Task serverTask = Task.Factory.StartNew(state => ServerTask((CancellationToken)state), serverCts.Token, TaskCreationOptions.LongRunning); - - ServerReadyEvent.Wait(); - rawSocket.Connect(ServerEndPoint); - - BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); - - byte[] sendBuffer = new byte[PacketSize]; - byte[] receiveBuffer = new byte[PacketSize]; - - for (int i = 0; i < PacketCount; i++) - { - byte[] packetBuffer = ServerEncoding.GetBytes($"[Client 0] Hello World! (Packet {i})"); - packetBuffer.CopyTo(sendBuffer, 0); - - benchmarkHelper.StartStopwatch(); - int sendResult = await writer.WriteAsync(ServerEndPoint, sendBuffer); - - int receiveResult = await writer.ReadAsync(ServerEndPoint, receiveBuffer); - benchmarkHelper.StopStopwatch(); - - benchmarkHelper.SnapshotRttStats(); - } - - benchmarkHelper.PrintBandwidthStats(0, PacketCount, PacketSize); - benchmarkHelper.PrintRttStats(0); - - serverCts.Cancel(); - try - { - serverTask.Dispose(); - } - catch (Exception) - { - // ignored - } - - rawSocket.Disconnect(false); - rawSocket.Shutdown(SocketShutdown.Both); - rawSocket.Close(); - rawSocket.Dispose(); - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterSyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterSyncBenchmark.cs @@ -1,134 +0,0 @@ -using NetSharp.Raw.Stream; - -using System; -using System.Net; -using System.Net.Sockets; -using System.Text; -using System.Threading; -using System.Threading.Tasks; - -namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks -{ - public class StreamNetworkWriterSyncBenchmark : INetSharpBenchmark - { - private const int PacketSize = 8192, PacketCount = 1_000_000; - - public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); - - public static readonly Encoding ServerEncoding = Encoding.UTF8; - public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12375); - - public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); - - /// <inheritdoc /> - public string Name { get; } = "Stream Raw Network Writer Benchmark (Synchronous)"; - - private static Task ServerTask(CancellationToken cancellationToken) - { - using Socket server = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - - server.Bind(ServerEndPoint); - ServerReadyEvent.Set(); - - byte[] transmissionBuffer = new byte[PacketSize]; - - server.Listen(1); - Socket clientSocket = server.Accept(); - - while (!cancellationToken.IsCancellationRequested) - { - int expectedBytes = transmissionBuffer.Length; - - int receivedBytes = 0; - do - { - receivedBytes += clientSocket.Receive(transmissionBuffer, receivedBytes, expectedBytes - receivedBytes, SocketFlags.None); - } while (receivedBytes != 0 && receivedBytes < expectedBytes); - - if (receivedBytes == 0) - { - break; - } - - int sentBytes = 0; - do - { - sentBytes += clientSocket.Send(transmissionBuffer, sentBytes, expectedBytes - sentBytes, SocketFlags.None); - } while (sentBytes != 0 && sentBytes < expectedBytes); - - if (sentBytes == 0) - { - break; - } - } - - server.Shutdown(SocketShutdown.Both); - server.Close(); - - return Task.CompletedTask; - } - - /// <inheritdoc /> - public Task RunAsync() - { - if (PacketCount > 10_000) - { - Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); - } - - EndPoint defaultRemoteEndPoint = new IPEndPoint(IPAddress.Any, 0); - - Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - rawSocket.Bind(ClientEndPoint); - - using RawStreamNetworkWriter writer = new RawStreamNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); - - using CancellationTokenSource serverCts = new CancellationTokenSource(); - Task serverTask = Task.Factory.StartNew(state => ServerTask((CancellationToken)state), serverCts.Token, TaskCreationOptions.LongRunning); - - ServerReadyEvent.Wait(); - rawSocket.Connect(ServerEndPoint); - - BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); - - byte[] sendBuffer = new byte[PacketSize]; - byte[] receiveBuffer = new byte[PacketSize]; - - EndPoint remoteEndPoint = ServerEndPoint; - - for (int i = 0; i < PacketCount; i++) - { - byte[] packetBuffer = ServerEncoding.GetBytes($"[Client 0] Hello World! (Packet {i})"); - packetBuffer.CopyTo(sendBuffer, 0); - - benchmarkHelper.StartStopwatch(); - int sendResult = writer.Write(ServerEndPoint, sendBuffer); - - int receiveResult = writer.Read(ref remoteEndPoint, receiveBuffer); - benchmarkHelper.StopStopwatch(); - - benchmarkHelper.SnapshotRttStats(); - } - - benchmarkHelper.PrintBandwidthStats(0, PacketCount, PacketSize); - benchmarkHelper.PrintRttStats(0); - - serverCts.Cancel(); - try - { - serverTask.Dispose(); - } - catch (Exception) - { - // ignored - } - - rawSocket.Disconnect(false); - rawSocket.Shutdown(SocketShutdown.Both); - rawSocket.Close(); - rawSocket.Dispose(); - - return Task.CompletedTask; - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkReaderBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkReaderBenchmark.cs @@ -0,0 +1,144 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Linq; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading; +using System.Threading.Tasks; + +namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks +{ + public class VariablePacketStreamNetworkReaderBenchmark : INetSharpBenchmark + { + private const int PacketSize = 8192, PacketCount = 1_000_000, ClientCount = 12; + + private double[] ClientBandwidths; + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = Encoding.UTF8; + public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12373); + + public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); + + /// <inheritdoc /> + public string Name { get; } = "Raw Variable Packet-size Stream Network Reader Benchmark"; + + private static bool RequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + in Memory<byte> responseBuffer) + { + requestBuffer.CopyTo(responseBuffer); + + return true; + } + + private Task BenchmarkClientTask(object idObj) + { + int id = (int)idObj; + + BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); + + Socket clientSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + clientSocket.Bind(ClientEndPoint); + + ServerReadyEvent.Wait(); + clientSocket.Connect(ServerEndPoint); + + byte[] sendBuffer = new byte[PacketSize + RawStreamPacket.Header.TotalHeaderSize]; + byte[] receiveBuffer = new byte[PacketSize + RawStreamPacket.Header.TotalHeaderSize]; + byte[] packetBuffer = new byte[PacketSize]; + + EndPoint remoteEndPoint = ServerEndPoint; + + lock (typeof(Console)) + { + Console.WriteLine($"[Client {id}] Starting client; sending messages to {remoteEndPoint}"); + } + + for (int i = 0; i < PacketCount; i++) + { + ServerEncoding.GetBytes($"[Client {id}] Hello World! (Packet {i})").CopyTo(packetBuffer, 0); + + RawStreamPacket streamPacket = new RawStreamPacket(packetBuffer); + streamPacket.Serialise(sendBuffer); + + benchmarkHelper.StartStopwatch(); + + int totalSent = 0; + do + { + totalSent += clientSocket.Send(sendBuffer, totalSent, sendBuffer.Length - totalSent, + SocketFlags.None); + } while (totalSent != 0 && totalSent < sendBuffer.Length); + + if (totalSent == 0) + { + break; + } + + int totalReceived = 0; + do + { + totalReceived += clientSocket.Receive(receiveBuffer, totalReceived, receiveBuffer.Length - totalReceived, + SocketFlags.None); + } while (totalReceived != 0 && totalReceived < receiveBuffer.Length); + + if (totalReceived == 0) + { + break; + } + + benchmarkHelper.StopStopwatch(); + + benchmarkHelper.SnapshotRttStats(); + } + + clientSocket.Disconnect(true); + clientSocket.Close(); + + benchmarkHelper.PrintBandwidthStats(id, PacketCount, PacketSize); + benchmarkHelper.PrintRttStats(id); + + ClientBandwidths[id] = benchmarkHelper.CalcBandwidth(PacketCount, PacketSize); + + return Task.CompletedTask; + } + + /// <inheritdoc /> + public async Task RunAsync() + { + if (PacketCount > 10_000) + { + Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); + } + + ClientBandwidths = new double[ClientCount]; + Task[] clientTasks = new Task[ClientCount]; + for (int i = 0; i < clientTasks.Length; i++) + { + clientTasks[i] = Task.Factory.StartNew(BenchmarkClientTask, i, TaskCreationOptions.LongRunning); + } + + EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ServerEndPoint); + rawSocket.Listen(ClientCount); + + using RawStreamNetworkReader reader = new VariablePacketRawStreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize); + reader.Start(ClientCount); + + ServerReadyEvent.Set(); + + await Task.WhenAll(clientTasks); + + Console.WriteLine($"Total estimated bandwidth: {ClientBandwidths.Sum():F3}"); + + reader.Stop(); + + rawSocket.Close(); + rawSocket.Dispose(); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkWriterAsyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkWriterAsyncBenchmark.cs @@ -0,0 +1,130 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading; +using System.Threading.Tasks; + +namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks +{ + public class VariablePacketStreamNetworkWriterAsyncBenchmark : INetSharpBenchmark + { + private const int PacketSize = 8192, PacketCount = 1_000_000; + + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = VariablePacketStreamNetworkReaderBenchmark.ServerEncoding; + public static readonly EndPoint ServerEndPoint = VariablePacketStreamNetworkReaderBenchmark.ServerEndPoint; + + public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); + + /// <inheritdoc /> + public string Name { get; } = "Raw Variable Packet-size Stream Network Writer Benchmark (Asynchronous)"; + + private static Task ServerTask(CancellationToken cancellationToken) + { + using Socket server = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + + server.Bind(ServerEndPoint); + ServerReadyEvent.Set(); + + byte[] transmissionBuffer = new byte[PacketSize]; + + server.Listen(1); + Socket clientSocket = server.Accept(); + + while (!cancellationToken.IsCancellationRequested) + { + int expectedBytes = transmissionBuffer.Length; + + int receivedBytes = 0; + do + { + receivedBytes += clientSocket.Receive(transmissionBuffer, receivedBytes, expectedBytes - receivedBytes, SocketFlags.None); + } while (receivedBytes != 0 && receivedBytes < expectedBytes); + + if (receivedBytes == 0) + { + break; + } + + int sentBytes = 0; + do + { + sentBytes += clientSocket.Send(transmissionBuffer, sentBytes, expectedBytes - sentBytes, SocketFlags.None); + } while (sentBytes != 0 && sentBytes < expectedBytes); + + if (sentBytes == 0) + { + break; + } + } + + server.Shutdown(SocketShutdown.Both); + server.Close(); + + return Task.CompletedTask; + } + + /// <inheritdoc /> + public async Task RunAsync() + { + if (PacketCount > 10_000) + { + Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); + } + + EndPoint defaultRemoteEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ClientEndPoint); + + using RawStreamNetworkWriter writer = new VariablePacketRawStreamNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); + + using CancellationTokenSource serverCts = new CancellationTokenSource(); + Task serverTask = Task.Factory.StartNew(state => ServerTask((CancellationToken)state), serverCts.Token, TaskCreationOptions.LongRunning); + + ServerReadyEvent.Wait(); + rawSocket.Connect(ServerEndPoint); + + BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); + + byte[] sendBuffer = new byte[PacketSize]; + byte[] receiveBuffer = new byte[PacketSize]; + + for (int i = 0; i < PacketCount; i++) + { + byte[] packetBuffer = ServerEncoding.GetBytes($"[Client 0] Hello World! (Packet {i})"); + packetBuffer.CopyTo(sendBuffer, 0); + + benchmarkHelper.StartStopwatch(); + int sendResult = await writer.WriteAsync(ServerEndPoint, sendBuffer); + + int receiveResult = await writer.ReadAsync(ServerEndPoint, receiveBuffer); + benchmarkHelper.StopStopwatch(); + + benchmarkHelper.SnapshotRttStats(); + } + + benchmarkHelper.PrintBandwidthStats(0, PacketCount, PacketSize); + benchmarkHelper.PrintRttStats(0); + + serverCts.Cancel(); + try + { + serverTask.Dispose(); + } + catch (Exception) + { + // ignored + } + + rawSocket.Disconnect(false); + rawSocket.Shutdown(SocketShutdown.Both); + rawSocket.Close(); + rawSocket.Dispose(); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkWriterSyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/VariablePacketStreamNetworkWriterSyncBenchmark.cs @@ -0,0 +1,134 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading; +using System.Threading.Tasks; + +namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks +{ + public class VariablePacketStreamNetworkWriterSyncBenchmark : INetSharpBenchmark + { + private const int PacketSize = 8192, PacketCount = 1_000_000; + + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = VariablePacketStreamNetworkReaderBenchmark.ServerEncoding; + public static readonly EndPoint ServerEndPoint = VariablePacketStreamNetworkReaderBenchmark.ServerEndPoint; + + public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); + + /// <inheritdoc /> + public string Name { get; } = "Raw Variable Packet-size Stream Network Writer Benchmark (Synchronous)"; + + private static Task ServerTask(CancellationToken cancellationToken) + { + using Socket server = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + + server.Bind(ServerEndPoint); + ServerReadyEvent.Set(); + + byte[] transmissionBuffer = new byte[PacketSize]; + + server.Listen(1); + Socket clientSocket = server.Accept(); + + while (!cancellationToken.IsCancellationRequested) + { + int expectedBytes = transmissionBuffer.Length; + + int receivedBytes = 0; + do + { + receivedBytes += clientSocket.Receive(transmissionBuffer, receivedBytes, expectedBytes - receivedBytes, SocketFlags.None); + } while (receivedBytes != 0 && receivedBytes < expectedBytes); + + if (receivedBytes == 0) + { + break; + } + + int sentBytes = 0; + do + { + sentBytes += clientSocket.Send(transmissionBuffer, sentBytes, expectedBytes - sentBytes, SocketFlags.None); + } while (sentBytes != 0 && sentBytes < expectedBytes); + + if (sentBytes == 0) + { + break; + } + } + + server.Shutdown(SocketShutdown.Both); + server.Close(); + + return Task.CompletedTask; + } + + /// <inheritdoc /> + public Task RunAsync() + { + if (PacketCount > 10_000) + { + Console.WriteLine($"{PacketCount} packets will be sent per client. This could take a long time (maybe more than a minute)!"); + } + + EndPoint defaultRemoteEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ClientEndPoint); + + using RawStreamNetworkWriter writer = new VariablePacketRawStreamNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); + + using CancellationTokenSource serverCts = new CancellationTokenSource(); + Task serverTask = Task.Factory.StartNew(state => ServerTask((CancellationToken)state), serverCts.Token, TaskCreationOptions.LongRunning); + + ServerReadyEvent.Wait(); + rawSocket.Connect(ServerEndPoint); + + BenchmarkHelper benchmarkHelper = new BenchmarkHelper(); + + byte[] sendBuffer = new byte[PacketSize]; + byte[] receiveBuffer = new byte[PacketSize]; + + EndPoint remoteEndPoint = ServerEndPoint; + + for (int i = 0; i < PacketCount; i++) + { + byte[] packetBuffer = ServerEncoding.GetBytes($"[Client 0] Hello World! (Packet {i})"); + packetBuffer.CopyTo(sendBuffer, 0); + + benchmarkHelper.StartStopwatch(); + int sendResult = writer.Write(ServerEndPoint, sendBuffer); + + int receiveResult = writer.Read(ref remoteEndPoint, receiveBuffer); + benchmarkHelper.StopStopwatch(); + + benchmarkHelper.SnapshotRttStats(); + } + + benchmarkHelper.PrintBandwidthStats(0, PacketCount, PacketSize); + benchmarkHelper.PrintRttStats(0); + + serverCts.Cancel(); + try + { + serverTask.Dispose(); + } + catch (Exception) + { + // ignored + } + + rawSocket.Disconnect(false); + rawSocket.Shutdown(SocketShutdown.Both); + rawSocket.Close(); + rawSocket.Dispose(); + + return Task.CompletedTask; + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkReaderExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkReaderExample.cs @@ -0,0 +1,56 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading.Tasks; + +namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples +{ + public class FixedPacketStreamNetworkReaderExample : INetSharpExample + { + private const int PacketSize = 8192, ExpectedClientCount = 8; + public static readonly Encoding ServerEncoding = Encoding.UTF8; + public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12377); + + /// <inheritdoc /> + public string Name { get; } = "Raw Fixed Packet-size Stream Network Reader Example"; + + private static bool RequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + in Memory<byte> responseBuffer) + { + requestBuffer.CopyTo(responseBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Received {receivedRequestBytes} bytes from {remoteEndPoint}! Echoing back..."); + } + + return true; + } + + /// <inheritdoc /> + public Task RunAsync() + { + EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ServerEndPoint); + rawSocket.Listen(ExpectedClientCount); + + using RawStreamNetworkReader reader = new FixedPacketRawStreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize, 100); + reader.Start(ExpectedClientCount); + + Console.WriteLine($"Started stream server at {ServerEndPoint}! Enter any key to stop the server..."); + Console.ReadLine(); + + reader.Stop(); + + rawSocket.Close(); + rawSocket.Dispose(); + + return Task.CompletedTask; + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkWriterAsyncExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkWriterAsyncExample.cs @@ -0,0 +1,61 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading.Tasks; + +namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples +{ + public class FixedPacketStreamNetworkWriterAsyncExample : INetSharpExample + { + private const int PacketSize = 8192; + + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = FixedPacketStreamNetworkReaderExample.ServerEncoding; + public static readonly EndPoint ServerEndPoint = FixedPacketStreamNetworkReaderExample.ServerEndPoint; + + /// <inheritdoc /> + public string Name { get; } = "Raw Fixed Packet-size Stream Network Writer Example (Asynchronous)"; + + /// <inheritdoc /> + public async Task RunAsync() + { + EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ClientEndPoint); + + using RawStreamNetworkWriter writer = new FixedPacketRawStreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); + await writer.ConnectAsync(ServerEndPoint); + + byte[] transmissionBuffer = new byte[PacketSize]; + + EndPoint remoteEndPoint = ServerEndPoint; + + while (true) + { + int sent = await writer.WriteAsync(remoteEndPoint, transmissionBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Sent {sent} bytes to {remoteEndPoint}!"); + } + + int received = await writer.ReadAsync(remoteEndPoint, transmissionBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Received {received} bytes from {remoteEndPoint}!"); + } + } + + await writer.DisconnectAsync(false); + + rawSocket.Close(); + rawSocket.Dispose(); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkWriterSyncExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/FixedPacketStreamNetworkWriterSyncExample.cs @@ -0,0 +1,63 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading.Tasks; + +namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples +{ + public class FixedPacketStreamNetworkWriterSyncExample : INetSharpExample + { + private const int PacketSize = 8192; + + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = FixedPacketStreamNetworkReaderExample.ServerEncoding; + public static readonly EndPoint ServerEndPoint = FixedPacketStreamNetworkReaderExample.ServerEndPoint; + + /// <inheritdoc /> + public string Name { get; } = "Raw Fixed Packet-size Stream Network Writer Example (Synchronous)"; + + /// <inheritdoc /> + public Task RunAsync() + { + EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ClientEndPoint); + + using RawStreamNetworkWriter writer = new FixedPacketRawStreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); + writer.Connect(ServerEndPoint); + + byte[] transmissionBuffer = new byte[PacketSize]; + + EndPoint remoteEndPoint = ServerEndPoint; + + while (true) + { + int sent = writer.Write(remoteEndPoint, transmissionBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Sent {sent} bytes to {remoteEndPoint}!"); + } + + int received = writer.Read(ref remoteEndPoint, transmissionBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Received {received} bytes from {remoteEndPoint}!"); + } + } + + writer.Disconnect(false); + + rawSocket.Close(); + rawSocket.Dispose(); + + return Task.CompletedTask; + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkReaderExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkReaderExample.cs @@ -1,56 +0,0 @@ -using NetSharp.Raw.Stream; - -using System; -using System.Net; -using System.Net.Sockets; -using System.Text; -using System.Threading.Tasks; - -namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples -{ - public class StreamNetworkReaderExample : INetSharpExample - { - private const int PacketSize = 8192, ExpectedClientCount = 8; - public static readonly Encoding ServerEncoding = Encoding.UTF8; - public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12377); - - /// <inheritdoc /> - public string Name { get; } = "Stream Network Reader Example"; - - private static bool RequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, - in Memory<byte> responseBuffer) - { - requestBuffer.CopyTo(responseBuffer); - - lock (typeof(Console)) - { - Console.WriteLine($"Received {receivedRequestBytes} bytes from {remoteEndPoint}! Echoing back..."); - } - - return true; - } - - /// <inheritdoc /> - public Task RunAsync() - { - EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); - - Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - rawSocket.Bind(ServerEndPoint); - rawSocket.Listen(ExpectedClientCount); - - using RawStreamNetworkReader reader = new RawStreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize, 100); - reader.Start(ExpectedClientCount); - - Console.WriteLine($"Started stream server at {ServerEndPoint}! Enter any key to stop the server..."); - Console.ReadLine(); - - reader.Stop(); - - rawSocket.Close(); - rawSocket.Dispose(); - - return Task.CompletedTask; - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterAsyncExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterAsyncExample.cs @@ -1,61 +0,0 @@ -using NetSharp.Raw.Stream; - -using System; -using System.Net; -using System.Net.Sockets; -using System.Text; -using System.Threading.Tasks; - -namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples -{ - public class StreamNetworkWriterAsyncExample : INetSharpExample - { - private const int PacketSize = 8192; - - public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); - - public static readonly Encoding ServerEncoding = StreamNetworkReaderExample.ServerEncoding; - public static readonly EndPoint ServerEndPoint = StreamNetworkReaderExample.ServerEndPoint; - - /// <inheritdoc /> - public string Name { get; } = "Stream Network Writer Example (Asynchronous)"; - - /// <inheritdoc /> - public async Task RunAsync() - { - EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); - - Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - rawSocket.Bind(ClientEndPoint); - - using RawStreamNetworkWriter writer = new RawStreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); - await writer.ConnectAsync(ServerEndPoint); - - byte[] transmissionBuffer = new byte[PacketSize]; - - EndPoint remoteEndPoint = ServerEndPoint; - - while (true) - { - int sent = await writer.WriteAsync(remoteEndPoint, transmissionBuffer); - - lock (typeof(Console)) - { - Console.WriteLine($"Sent {sent} bytes to {remoteEndPoint}!"); - } - - int received = await writer.ReadAsync(remoteEndPoint, transmissionBuffer); - - lock (typeof(Console)) - { - Console.WriteLine($"Received {received} bytes from {remoteEndPoint}!"); - } - } - - await writer.DisconnectAsync(false); - - rawSocket.Close(); - rawSocket.Dispose(); - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterSyncExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterSyncExample.cs @@ -1,63 +0,0 @@ -using NetSharp.Raw.Stream; - -using System; -using System.Net; -using System.Net.Sockets; -using System.Text; -using System.Threading.Tasks; - -namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples -{ - public class StreamNetworkWriterSyncExample : INetSharpExample - { - private const int PacketSize = 8192; - - public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); - - public static readonly Encoding ServerEncoding = StreamNetworkReaderExample.ServerEncoding; - public static readonly EndPoint ServerEndPoint = StreamNetworkReaderExample.ServerEndPoint; - - /// <inheritdoc /> - public string Name { get; } = "Stream Network Writer Example (Synchronous)"; - - /// <inheritdoc /> - public Task RunAsync() - { - EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); - - Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - rawSocket.Bind(ClientEndPoint); - - using RawStreamNetworkWriter writer = new RawStreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); - writer.Connect(ServerEndPoint); - - byte[] transmissionBuffer = new byte[PacketSize]; - - EndPoint remoteEndPoint = ServerEndPoint; - - while (true) - { - int sent = writer.Write(remoteEndPoint, transmissionBuffer); - - lock (typeof(Console)) - { - Console.WriteLine($"Sent {sent} bytes to {remoteEndPoint}!"); - } - - int received = writer.Read(ref remoteEndPoint, transmissionBuffer); - - lock (typeof(Console)) - { - Console.WriteLine($"Received {received} bytes from {remoteEndPoint}!"); - } - } - - writer.Disconnect(false); - - rawSocket.Close(); - rawSocket.Dispose(); - - return Task.CompletedTask; - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkReaderExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkReaderExample.cs @@ -0,0 +1,56 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading.Tasks; + +namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples +{ + public class VariablePacketStreamNetworkReaderExample : INetSharpExample + { + private const int PacketSize = 8192, ExpectedClientCount = 8; + public static readonly Encoding ServerEncoding = Encoding.UTF8; + public static readonly EndPoint ServerEndPoint = new IPEndPoint(IPAddress.Loopback, 12377); + + /// <inheritdoc /> + public string Name { get; } = "Raw Variable Packet-size Stream Network Reader Example"; + + private static bool RequestHandler(EndPoint remoteEndPoint, in ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + in Memory<byte> responseBuffer) + { + requestBuffer.CopyTo(responseBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Received {receivedRequestBytes} bytes from {remoteEndPoint}! Echoing back..."); + } + + return true; + } + + /// <inheritdoc /> + public Task RunAsync() + { + EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ServerEndPoint); + rawSocket.Listen(ExpectedClientCount); + + using RawStreamNetworkReader reader = new VariablePacketRawStreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize, 100); + reader.Start(ExpectedClientCount); + + Console.WriteLine($"Started stream server at {ServerEndPoint}! Enter any key to stop the server..."); + Console.ReadLine(); + + reader.Stop(); + + rawSocket.Close(); + rawSocket.Dispose(); + + return Task.CompletedTask; + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkWriterAsyncExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkWriterAsyncExample.cs @@ -0,0 +1,61 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading.Tasks; + +namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples +{ + public class VariablePacketStreamNetworkWriterAsyncExample : INetSharpExample + { + private const int PacketSize = 8192; + + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = VariablePacketStreamNetworkReaderExample.ServerEncoding; + public static readonly EndPoint ServerEndPoint = VariablePacketStreamNetworkReaderExample.ServerEndPoint; + + /// <inheritdoc /> + public string Name { get; } = "Raw Variable Packet-size Stream Network Writer Example (Asynchronous)"; + + /// <inheritdoc /> + public async Task RunAsync() + { + EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ClientEndPoint); + + using RawStreamNetworkWriter writer = new VariablePacketRawStreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); + await writer.ConnectAsync(ServerEndPoint); + + byte[] transmissionBuffer = new byte[PacketSize]; + + EndPoint remoteEndPoint = ServerEndPoint; + + while (true) + { + int sent = await writer.WriteAsync(remoteEndPoint, transmissionBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Sent {sent} bytes to {remoteEndPoint}!"); + } + + int received = await writer.ReadAsync(remoteEndPoint, transmissionBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Received {received} bytes from {remoteEndPoint}!"); + } + } + + await writer.DisconnectAsync(false); + + rawSocket.Close(); + rawSocket.Dispose(); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkWriterSyncExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/VariablePacketStreamNetworkWriterSyncExample.cs @@ -0,0 +1,63 @@ +using NetSharp.Raw.Stream; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Threading.Tasks; + +namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples +{ + public class VariablePacketStreamNetworkWriterSyncExample : INetSharpExample + { + private const int PacketSize = 8192; + + public static readonly EndPoint ClientEndPoint = new IPEndPoint(IPAddress.Loopback, 0); + + public static readonly Encoding ServerEncoding = VariablePacketStreamNetworkReaderExample.ServerEncoding; + public static readonly EndPoint ServerEndPoint = VariablePacketStreamNetworkReaderExample.ServerEndPoint; + + /// <inheritdoc /> + public string Name { get; } = "Raw Variable Packet-size Stream Network Writer Example (Synchronous)"; + + /// <inheritdoc /> + public Task RunAsync() + { + EndPoint defaultEndPoint = new IPEndPoint(IPAddress.Any, 0); + + Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + rawSocket.Bind(ClientEndPoint); + + using RawStreamNetworkWriter writer = new VariablePacketRawStreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); + writer.Connect(ServerEndPoint); + + byte[] transmissionBuffer = new byte[PacketSize]; + + EndPoint remoteEndPoint = ServerEndPoint; + + while (true) + { + int sent = writer.Write(remoteEndPoint, transmissionBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Sent {sent} bytes to {remoteEndPoint}!"); + } + + int received = writer.Read(ref remoteEndPoint, transmissionBuffer); + + lock (typeof(Console)) + { + Console.WriteLine($"Received {received} bytes from {remoteEndPoint}!"); + } + } + + writer.Disconnect(false); + + rawSocket.Close(); + rawSocket.Dispose(); + + return Task.CompletedTask; + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/NetSharpExamples.xml b/NetSharp/NetSharpExamples/NetSharpExamples.xml @@ -22,22 +22,40 @@ <member name="M:NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks.DatagramNetworkWriterSyncBenchmark.RunAsync"> <inheritdoc /> </member> - <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.StreamNetworkReaderBenchmark.Name"> + <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.FixedPacketStreamNetworkReaderBenchmark.Name"> <inheritdoc /> </member> - <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.StreamNetworkReaderBenchmark.RunAsync"> + <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.FixedPacketStreamNetworkReaderBenchmark.RunAsync"> <inheritdoc /> </member> - <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.StreamNetworkWriterAsyncBenchmark.Name"> + <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.FixedPacketStreamNetworkWriterAsyncBenchmark.Name"> <inheritdoc /> </member> - <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.StreamNetworkWriterAsyncBenchmark.RunAsync"> + <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.FixedPacketStreamNetworkWriterAsyncBenchmark.RunAsync"> <inheritdoc /> </member> - <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.StreamNetworkWriterSyncBenchmark.Name"> + <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.FixedPacketStreamNetworkWriterSyncBenchmark.Name"> <inheritdoc /> </member> - <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.StreamNetworkWriterSyncBenchmark.RunAsync"> + <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.FixedPacketStreamNetworkWriterSyncBenchmark.RunAsync"> + <inheritdoc /> + </member> + <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.VariablePacketStreamNetworkReaderBenchmark.Name"> + <inheritdoc /> + </member> + <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.VariablePacketStreamNetworkReaderBenchmark.RunAsync"> + <inheritdoc /> + </member> + <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.VariablePacketStreamNetworkWriterAsyncBenchmark.Name"> + <inheritdoc /> + </member> + <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.VariablePacketStreamNetworkWriterAsyncBenchmark.RunAsync"> + <inheritdoc /> + </member> + <member name="P:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.VariablePacketStreamNetworkWriterSyncBenchmark.Name"> + <inheritdoc /> + </member> + <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.VariablePacketStreamNetworkWriterSyncBenchmark.RunAsync"> <inheritdoc /> </member> <member name="P:NetSharpExamples.Examples.Datagram_Network_Connection_Examples.DatagramNetworkReaderExample.Name"> @@ -58,22 +76,40 @@ <member name="M:NetSharpExamples.Examples.Datagram_Network_Connection_Examples.DatagramNetworkWriterSyncExample.RunAsync"> <inheritdoc /> </member> - <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.StreamNetworkReaderExample.Name"> + <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.FixedPacketStreamNetworkReaderExample.Name"> + <inheritdoc /> + </member> + <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.FixedPacketStreamNetworkReaderExample.RunAsync"> + <inheritdoc /> + </member> + <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.FixedPacketStreamNetworkWriterAsyncExample.Name"> + <inheritdoc /> + </member> + <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.FixedPacketStreamNetworkWriterAsyncExample.RunAsync"> + <inheritdoc /> + </member> + <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.FixedPacketStreamNetworkWriterSyncExample.Name"> + <inheritdoc /> + </member> + <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.FixedPacketStreamNetworkWriterSyncExample.RunAsync"> + <inheritdoc /> + </member> + <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.VariablePacketStreamNetworkReaderExample.Name"> <inheritdoc /> </member> - <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.StreamNetworkReaderExample.RunAsync"> + <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.VariablePacketStreamNetworkReaderExample.RunAsync"> <inheritdoc /> </member> - <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.StreamNetworkWriterAsyncExample.Name"> + <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.VariablePacketStreamNetworkWriterAsyncExample.Name"> <inheritdoc /> </member> - <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.StreamNetworkWriterAsyncExample.RunAsync"> + <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.VariablePacketStreamNetworkWriterAsyncExample.RunAsync"> <inheritdoc /> </member> - <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.StreamNetworkWriterSyncExample.Name"> + <member name="P:NetSharpExamples.Examples.Stream_Network_Connection_Examples.VariablePacketStreamNetworkWriterSyncExample.Name"> <inheritdoc /> </member> - <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.StreamNetworkWriterSyncExample.RunAsync"> + <member name="M:NetSharpExamples.Examples.Stream_Network_Connection_Examples.VariablePacketStreamNetworkWriterSyncExample.RunAsync"> <inheritdoc /> </member> <member name="T:NetSharpExamples.INetSharpExample">