commit 7bd86e3331918290eca3caa15a57aeafcf29ab18 parent 09e8515d9c0aa5a8e72558567f2cd41cf68c1a03 Author: Mikolaj Lenczewski <mikolaj.lenczewski308@gmail.com> Date: Tue, 5 May 2020 21:03:16 +0100 Started implementing high-level connection. Diffstat:
34 files changed, 1911 insertions(+), 1485 deletions(-)
diff --git a/NetSharp/NetSharp/Connection.cs b/NetSharp/NetSharp/Connection.cs @@ -0,0 +1,251 @@ +using NetSharp.Raw; + +using System; +using System.Net; +using System.Net.Sockets; +using System.Threading.Tasks; + +namespace NetSharp +{ + // TODO complete class + public sealed class Connection : ConnectionBase + { + /// <inheritdoc /> + internal Connection(AddressFamily addressFamily, SocketType socketType, ProtocolType protocolType, + IRawNetworkTransportProvider<SocketAsyncEventArgs> transportProvider, EndPoint defaultEndPoint, int packetBufferSize, + int pooledBuffersPerBucket, uint preallocatedStateObjects) : base(addressFamily, socketType, protocolType, + transportProvider, defaultEndPoint, packetBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) + { + } + + /// <inheritdoc /> + /// TODO convert to use registered packet handlers + protected override bool HandleRawPacketReceived(in EndPoint remoteEndPoint, ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + Memory<byte> responseBuffer) + { + requestBuffer.CopyTo(responseBuffer); + + return true; + } + + /// <inheritdoc /> + public override void Connect(EndPoint remoteEndPoint) + { + Writer.Connect(remoteEndPoint); + } + + /// <inheritdoc /> + public override ValueTask ConnectAsync(EndPoint remoteEndPoint) + { + return Writer.ConnectAsync(remoteEndPoint); + } + + /// <inheritdoc /> + public override int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) + { + return Writer.Read(ref remoteEndPoint, readBuffer, flags); + } + + /// <inheritdoc /> + public override ValueTask<int> ReadAsync(EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) + { + return Writer.ReadAsync(remoteEndPoint, readBuffer, flags); + } + + /// <inheritdoc /> + public override void Start(ushort concurrentReadTasks) + { + Reader.Start(concurrentReadTasks); + } + + /// <inheritdoc /> + public override void Stop() + { + Reader.Stop(); + } + + /// <inheritdoc /> + public override int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, SocketFlags flags = SocketFlags.None) + { + return Writer.Write(remoteEndPoint, writeBuffer, flags); + } + + /// <inheritdoc /> + public override ValueTask<int> WriteAsync(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, SocketFlags flags = SocketFlags.None) + { + return Writer.ReadAsync(remoteEndPoint, Memory<byte>.Empty, flags); + } + } + + // TODO complete class + public abstract class ConnectionBase : IDisposable, INetworkReader, INetworkWriter + { + private readonly Socket connection; + + protected readonly RawNetworkReaderBase<SocketAsyncEventArgs> Reader; + + protected readonly RawNetworkWriterBase<SocketAsyncEventArgs> Writer; + + private protected ConnectionBase(AddressFamily addressFamily, SocketType socketType, ProtocolType protocolType, + IRawNetworkTransportProvider<SocketAsyncEventArgs> transportProvider, EndPoint defaultEndPoint, int packetBufferSize, + int pooledBuffersPerBucket, uint preallocatedStateObjects) + { + connection = new Socket(addressFamily, socketType, protocolType); + + Reader = transportProvider.GetReader(ref connection, defaultEndPoint, HandleRawPacketReceived, + packetBufferSize, pooledBuffersPerBucket, preallocatedStateObjects); + + Writer = transportProvider.GetWriter(ref connection, defaultEndPoint, packetBufferSize, + pooledBuffersPerBucket, preallocatedStateObjects); + } + + protected virtual void Dispose(bool disposing) + { + if (!disposing) return; + + Writer.Dispose(); + + Reader.Stop(); + Reader.Dispose(); + + connection.Dispose(); + } + + protected abstract bool HandleRawPacketReceived(in EndPoint remoteEndPoint, ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + Memory<byte> responseBuffer); + + /// <inheritdoc /> + public abstract void Connect(EndPoint remoteEndPoint); + + /// <inheritdoc /> + public abstract ValueTask ConnectAsync(EndPoint remoteEndPoint); + + /// <inheritdoc /> + public void Dispose() + { + Dispose(true); + GC.SuppressFinalize(this); + } + + /// <inheritdoc /> + public abstract int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None); + + /// <inheritdoc /> + public abstract ValueTask<int> ReadAsync(EndPoint remoteEndPoint, Memory<byte> readBuffer, + SocketFlags flags = SocketFlags.None); + + /// <inheritdoc /> + public abstract void Start(ushort concurrentReceiveTasks); + + /// <inheritdoc /> + public abstract void Stop(); + + /// <inheritdoc /> + public abstract int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, + SocketFlags flags = SocketFlags.None); + + /// <inheritdoc /> + public abstract ValueTask<int> WriteAsync(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, + SocketFlags flags = SocketFlags.None); + } + + public sealed class ConnectionBuilder + { + public static ConfiguredConnectionBuilder WithCustomTransport(IRawNetworkTransportProvider<SocketAsyncEventArgs> transportProvider, + ProtocolType transportProtocol) + { + return new ConfiguredConnectionBuilder(transportProvider, transportProtocol); + } + + public static ConfiguredConnectionBuilder WithDatagramTransport(ProtocolType transportProtocol = ProtocolType.Udp) => + WithCustomTransport(new DatagramRawNetworkTransportProvider(), transportProtocol); + + public static ConfiguredConnectionBuilder WithStreamTransport(ProtocolType transportProtocol = ProtocolType.Tcp) => + WithCustomTransport(new StreamRawNetworkTransportProvider(), transportProtocol); + + public sealed class ConfiguredConnectionBuilder + { + private readonly IRawNetworkTransportProvider<SocketAsyncEventArgs> transportProvider; + + internal ConfiguredConnectionBuilder(IRawNetworkTransportProvider<SocketAsyncEventArgs> transportProvider, ProtocolType transportProtocol) + { + this.transportProvider = transportProvider; + + SocketType = transportProvider.TransportProtocolType; + + ProtocolType = transportProtocol; + } + + public AddressFamily AddressFamily { get; private set; } = AddressFamily.InterNetwork; + public EndPoint DefaultEndPoint { get; private set; } = new IPEndPoint(IPAddress.Any, 0); + public ProtocolType ProtocolType { get; } + public SocketType SocketType { get; } + + public ConfiguredConnectionBuilder WithAddressFamily(AddressFamily addressFamily, EndPoint defaultEndPoint) + { + AddressFamily = addressFamily; + + DefaultEndPoint = defaultEndPoint; + + return this; + } + + public CompletedConnectionBuilder WithDefaultSettings(int packetBufferSize) => + WithSettings(packetBufferSize, 1000, 0); + + public ConfiguredConnectionBuilder WithInterNetwork(EndPoint defaultEndPoint) => + WithAddressFamily(AddressFamily.InterNetwork, defaultEndPoint); + + public ConfiguredConnectionBuilder WithInterNetworkV6(EndPoint defaultEndPoint) => + WithAddressFamily(AddressFamily.InterNetworkV6, defaultEndPoint); + + public CompletedConnectionBuilder WithSettings(int packetBufferSize, int pooledBuffersPerBucket, + uint preallocatedStateObjects) + { + return new CompletedConnectionBuilder(transportProvider, AddressFamily, DefaultEndPoint, SocketType, + ProtocolType, packetBufferSize, pooledBuffersPerBucket, preallocatedStateObjects); + } + + public sealed class CompletedConnectionBuilder + { + private readonly IRawNetworkTransportProvider<SocketAsyncEventArgs> transportProvider; + + internal CompletedConnectionBuilder( + IRawNetworkTransportProvider<SocketAsyncEventArgs> transportProvider, AddressFamily addressFamily, EndPoint defaultEndPoint, + SocketType transportProtocolType, ProtocolType transportProtocol, int maxBufferSize, int pooledBuffersPerBucket, + uint preallocatedStateObjects) + { + this.transportProvider = transportProvider; + + AddressFamily = addressFamily; + + DefaultEndPoint = defaultEndPoint; + + SocketType = transportProtocolType; + + ProtocolType = transportProtocol; + + BufferSize = maxBufferSize; + + PooledBuffersPerBucket = pooledBuffersPerBucket; + + PreallocatedStateObjects = preallocatedStateObjects; + } + + public AddressFamily AddressFamily { get; } + public int BufferSize { get; } + public EndPoint DefaultEndPoint { get; } + public int PooledBuffersPerBucket { get; } + public uint PreallocatedStateObjects { get; } + public ProtocolType ProtocolType { get; } + public SocketType SocketType { get; } + + public Connection BuildDefault() + { + return new Connection(AddressFamily, SocketType, ProtocolType, transportProvider, DefaultEndPoint, + BufferSize, PooledBuffersPerBucket, PreallocatedStateObjects); + } + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/INetworkReader.cs b/NetSharp/NetSharp/INetworkReader.cs @@ -0,0 +1,9 @@ +namespace NetSharp +{ + public interface INetworkReader + { + public void Start(ushort concurrentReadTasks); + + public void Stop(); + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/INetworkWriter.cs b/NetSharp/NetSharp/INetworkWriter.cs @@ -0,0 +1,26 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Threading.Tasks; + +namespace NetSharp +{ + public interface INetworkWriter + { + public void Connect(EndPoint remoteEndPoint); + + public ValueTask ConnectAsync(EndPoint remoteEndPoint); + + public int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, + SocketFlags flags = SocketFlags.None); + + public ValueTask<int> ReadAsync(EndPoint remoteEndPoint, Memory<byte> readBuffer, + SocketFlags flags = SocketFlags.None); + + public int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, + SocketFlags flags = SocketFlags.None); + + public ValueTask<int> WriteAsync(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, + SocketFlags flags = SocketFlags.None); + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/NetSharp.csproj b/NetSharp/NetSharp/NetSharp.csproj @@ -18,7 +18,6 @@ <ItemGroup> <PackageReference Include="Microsoft.CSharp" Version="4.7.0" /> - <PackageReference Include="Microsoft.Extensions.ObjectPool" Version="3.1.2" /> - <PackageReference Include="System.Threading.Channels" Version="4.7.0" /> + <PackageReference Include="Microsoft.Extensions.ObjectPool" Version="3.1.3" /> </ItemGroup> </Project> \ No newline at end of file diff --git a/NetSharp/NetSharp/NetSharp.xml b/NetSharp/NetSharp/NetSharp.xml @@ -4,6 +4,15 @@ <name>NetSharp</name> </assembly> <members> + <member name="M:NetSharp.Connection.#ctor(System.Net.Sockets.AddressFamily,System.Net.Sockets.SocketType,System.Net.Sockets.ProtocolType,NetSharp.Raw.IRawNetworkTransportProvider{System.Net.Sockets.SocketAsyncEventArgs},System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Connection.HandleRawPacketReceived(System.Net.EndPoint@,System.ReadOnlyMemory{System.Byte},System.Int32,System.Memory{System.Byte})"> + <inheritdoc /> + </member> + <member name="M:NetSharp.ConnectionBase.Dispose"> + <inheritdoc /> + </member> <member name="T:NetSharp.Packets.NetworkPacket"> <summary> Represents a raw packet sent across the network. @@ -156,123 +165,141 @@ The memory buffer to write the serialised packet header instance to. </param> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkReader.#ctor(System.Net.Sockets.Socket@,NetSharp.Raw.NetworkRequestHandler,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkReader.#ctor(System.Net.Sockets.Socket@,NetSharp.Raw.NetworkRequestHandler,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkReader.CanReuseStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkReader.CreateStateObject"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkReader.DestroyStateObject(System.Net.Sockets.SocketAsyncEventArgs)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkReader.ResetStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkReader.Start(System.UInt16)"> + <inheritdoc /> + </member> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkReader.CanReuseStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.CanReuseStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkReader.CreateStateObject"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.CreateStateObject"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkReader.DestroyStateObject(System.Net.Sockets.SocketAsyncEventArgs)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.DestroyStateObject(System.Net.Sockets.SocketAsyncEventArgs)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkReader.ResetStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.ResetStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkReader.Start(System.UInt16)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.Connect(System.Net.EndPoint)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.ConnectAsync(System.Net.EndPoint)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.CanReuseStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.Read(System.Net.EndPoint@,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.CreateStateObject"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.ReadAsync(System.Net.EndPoint,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.DestroyStateObject(System.Net.Sockets.SocketAsyncEventArgs)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.Write(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.ResetStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <member name="M:NetSharp.Raw.Datagram.RawDatagramNetworkWriter.WriteAsync(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.Connect(System.Net.EndPoint)"> + <member name="P:NetSharp.Raw.DatagramRawNetworkTransportProvider.TransportProtocolType"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.ConnectAsync(System.Net.EndPoint)"> + <member name="M:NetSharp.Raw.DatagramRawNetworkTransportProvider.GetReader(System.Net.Sockets.Socket@,System.Net.EndPoint,NetSharp.Raw.NetworkRequestHandler,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.Read(System.Net.EndPoint@,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.DatagramRawNetworkTransportProvider.GetWriter(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.ReadAsync(System.Net.EndPoint,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="P:NetSharp.Raw.StreamRawNetworkTransportProvider.TransportProtocolType"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.Write(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.StreamRawNetworkTransportProvider.GetReader(System.Net.Sockets.Socket@,System.Net.EndPoint,NetSharp.Raw.NetworkRequestHandler,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Datagram.DatagramNetworkWriter.WriteAsync(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.StreamRawNetworkTransportProvider.GetWriter(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.NetworkConnectionBase`1.Dispose(System.Boolean)"> + <member name="M:NetSharp.Raw.RawNetworkConnectionBase`1.Dispose(System.Boolean)"> <summary> Allows for inheritors to dispose of their own resources. </summary> </member> - <member name="M:NetSharp.Raw.NetworkConnectionBase`1.Dispose"> + <member name="M:NetSharp.Raw.RawNetworkConnectionBase`1.Dispose"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.NetworkReaderBase`1.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,NetSharp.Raw.NetworkRequestHandler,System.Int32,System.Int32,System.UInt32)"> + <member name="M:NetSharp.Raw.RawNetworkReaderBase`1.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,NetSharp.Raw.NetworkRequestHandler,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.NetworkReaderBase`1.Dispose(System.Boolean)"> + <member name="M:NetSharp.Raw.RawNetworkReaderBase`1.Dispose(System.Boolean)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.NetworkWriterBase`1.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <member name="M:NetSharp.Raw.RawNetworkWriterBase`1.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkReader.#ctor(System.Net.Sockets.Socket@,NetSharp.Raw.NetworkRequestHandler,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkReader.#ctor(System.Net.Sockets.Socket@,NetSharp.Raw.NetworkRequestHandler,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkReader.CanReuseStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkReader.CanReuseStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkReader.CreateStateObject"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkReader.CreateStateObject"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkReader.DestroyStateObject(System.Net.Sockets.SocketAsyncEventArgs)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkReader.DestroyStateObject(System.Net.Sockets.SocketAsyncEventArgs)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkReader.ResetStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkReader.ResetStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkReader.Start(System.UInt16)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkReader.Start(System.UInt16)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.#ctor(System.Net.Sockets.Socket@,System.Net.EndPoint,System.Int32,System.Int32,System.UInt32)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.CanReuseStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.CanReuseStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.CreateStateObject"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.CreateStateObject"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.DestroyStateObject(System.Net.Sockets.SocketAsyncEventArgs)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.DestroyStateObject(System.Net.Sockets.SocketAsyncEventArgs)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.ResetStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.ResetStateObject(System.Net.Sockets.SocketAsyncEventArgs@)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.Connect(System.Net.EndPoint)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.Connect(System.Net.EndPoint)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.ConnectAsync(System.Net.EndPoint)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.ConnectAsync(System.Net.EndPoint)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.Read(System.Net.EndPoint@,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.Read(System.Net.EndPoint@,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.ReadAsync(System.Net.EndPoint,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.ReadAsync(System.Net.EndPoint,System.Memory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.Write(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.Write(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> <inheritdoc /> </member> - <member name="M:NetSharp.Raw.Stream.StreamNetworkWriter.WriteAsync(System.Net.EndPoint,System.ReadOnlyMemory{System.Byte},System.Net.Sockets.SocketFlags)"> + <member name="M:NetSharp.Raw.Stream.RawStreamNetworkWriter.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/Datagram/DatagramNetworkReader.cs b/NetSharp/NetSharp/Raw/Datagram/DatagramNetworkReader.cs @@ -1,156 +0,0 @@ -using System.Net; -using System.Net.Sockets; - -namespace NetSharp.Raw.Datagram -{ - public sealed class DatagramNetworkReader : NetworkReaderBase<SocketAsyncEventArgs> - { - /// <inheritdoc /> - public DatagramNetworkReader(ref Socket rawConnection, NetworkRequestHandler? requestHandler, EndPoint defaultEndPoint, int maxPooledBufferSize, - int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, requestHandler, maxPooledBufferSize, - maxPooledBuffersPerBucket, preallocatedStateObjects) - { - } - - private void CompleteReceiveFrom(SocketAsyncEventArgs args) - { - byte[] receiveBuffer = args.Buffer; - - switch (args.SocketError) - { - case SocketError.Success: - byte[] responseBuffer = BufferPool.Rent(BufferSize); - - bool responseExists = - RequestHandler(args.RemoteEndPoint, receiveBuffer, args.BytesTransferred, responseBuffer); - BufferPool.Return(receiveBuffer, true); - - if (responseExists) - { - args.SetBuffer(responseBuffer, 0, BufferSize); - - StartSendTo(args); - - return; - } - - BufferPool.Return(responseBuffer, true); - break; - - default: - BufferPool.Return(receiveBuffer, true); - StateObjectPool.Return(args); - break; - } - } - - private void CompleteSendTo(SocketAsyncEventArgs args) - { - byte[] sendBuffer = args.Buffer; - - BufferPool.Return(sendBuffer, true); - StateObjectPool.Return(args); - } - - private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) - { - switch (args.LastOperation) - { - case SocketAsyncOperation.ReceiveFrom: - StartDefaultReceiveFrom(); - - CompleteReceiveFrom(args); - break; - - case SocketAsyncOperation.SendTo: - CompleteSendTo(args); - break; - } - } - - private void StartDefaultReceiveFrom() - { - if (ShutdownToken.IsCancellationRequested) - { - return; - } - - SocketAsyncEventArgs args = StateObjectPool.Rent(); - StartReceiveFrom(args); - } - - private void StartReceiveFrom(SocketAsyncEventArgs args) - { - byte[] receiveBuffer = BufferPool.Rent(BufferSize); - - if (ShutdownToken.IsCancellationRequested) - { - BufferPool.Return(receiveBuffer, true); - StateObjectPool.Return(args); - - return; - } - - args.SetBuffer(receiveBuffer, 0, BufferSize); - - if (Connection.ReceiveFromAsync(args)) return; - - StartDefaultReceiveFrom(); - CompleteReceiveFrom(args); - } - - private void StartSendTo(SocketAsyncEventArgs args) - { - byte[] sendBuffer = args.Buffer; - - if (ShutdownToken.IsCancellationRequested) - { - BufferPool.Return(sendBuffer, true); - StateObjectPool.Return(args); - - return; - } - - if (Connection.SendToAsync(args)) return; - - CompleteSendTo(args); - } - - /// <inheritdoc /> - protected override bool CanReuseStateObject(ref SocketAsyncEventArgs instance) - { - return true; - } - - /// <inheritdoc /> - protected override SocketAsyncEventArgs CreateStateObject() - { - SocketAsyncEventArgs instance = new SocketAsyncEventArgs { RemoteEndPoint = DefaultEndPoint }; - instance.Completed += HandleIoCompleted; - - return instance; - } - - /// <inheritdoc /> - protected override void DestroyStateObject(SocketAsyncEventArgs instance) - { - instance.Completed -= HandleIoCompleted; - instance.Dispose(); - } - - /// <inheritdoc /> - protected override void ResetStateObject(ref SocketAsyncEventArgs instance) - { - instance.RemoteEndPoint = DefaultEndPoint; - } - - /// <inheritdoc /> - public override void Start(ushort concurrentReadTasks) - { - for (ushort i = 0; i < concurrentReadTasks; i++) - { - StartDefaultReceiveFrom(); - } - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Datagram/DatagramNetworkWriter.cs b/NetSharp/NetSharp/Raw/Datagram/DatagramNetworkWriter.cs @@ -1,304 +0,0 @@ -using System; -using System.Net; -using System.Net.Sockets; -using System.Threading.Tasks; - -namespace NetSharp.Raw.Datagram -{ - public sealed class DatagramNetworkWriter : NetworkWriterBase<SocketAsyncEventArgs> - { - /// <inheritdoc /> - public DatagramNetworkWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, - uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxPooledBufferSize, maxPooledBuffersPerBucket, preallocatedStateObjects) - { - } - - private void CompleteConnect(SocketAsyncEventArgs args) - { - AsyncOperationToken token = (AsyncOperationToken)args.UserToken; - - switch (args.SocketError) - { - case SocketError.Success: - token.CompletionSource.SetResult(true); - break; - - case SocketError.OperationAborted: - token.CompletionSource.SetCanceled(); - break; - - default: - int errorCode = (int)args.SocketError; - token.CompletionSource.SetException(new SocketException(errorCode)); - break; - } - - StateObjectPool.Return(args); - } - - private void CompleteReceiveFrom(SocketAsyncEventArgs args) - { - AsyncDatagramReadToken token = (AsyncDatagramReadToken)args.UserToken; - - byte[] receiveBuffer = args.Buffer; - - switch (args.SocketError) - { - case SocketError.Success: - receiveBuffer.CopyTo(token.UserBuffer); - token.CompletionSource.SetResult(args.BytesTransferred); - 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); - StateObjectPool.Return(args); - } - - private void CompleteSendTo(SocketAsyncEventArgs args) - { - AsyncDatagramWriteToken token = (AsyncDatagramWriteToken)args.UserToken; - - byte[] sendBuffer = args.Buffer; - - switch (args.SocketError) - { - case SocketError.Success: - token.CompletionSource.SetResult(args.BytesTransferred); - 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); - StateObjectPool.Return(args); - } - - private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) - { - switch (args.LastOperation) - { - case SocketAsyncOperation.Connect: - CompleteConnect(args); - break; - - case SocketAsyncOperation.SendTo: - CompleteSendTo(args); - break; - - case SocketAsyncOperation.ReceiveFrom: - CompleteReceiveFrom(args); - break; - } - } - - /// <inheritdoc /> - protected override bool CanReuseStateObject(ref SocketAsyncEventArgs instance) - { - return true; - } - - /// <inheritdoc /> - protected override SocketAsyncEventArgs CreateStateObject() - { - SocketAsyncEventArgs instance = new SocketAsyncEventArgs(); - instance.Completed += HandleIoCompleted; - - return instance; - } - - /// <inheritdoc /> - protected override void DestroyStateObject(SocketAsyncEventArgs instance) - { - instance.Completed -= HandleIoCompleted; - instance.Dispose(); - } - - /// <inheritdoc /> - protected override void ResetStateObject(ref SocketAsyncEventArgs instance) - { - } - - /// <inheritdoc /> - public override void Connect(EndPoint remoteEndPoint) - { - Connection.Connect(remoteEndPoint); - } - - /// <inheritdoc /> - public override ValueTask ConnectAsync(EndPoint remoteEndPoint) - { - TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>(); - SocketAsyncEventArgs args = StateObjectPool.Rent(); - - args.RemoteEndPoint = remoteEndPoint; - - AsyncOperationToken token = new AsyncOperationToken(tcs); - args.UserToken = token; - - if (Connection.ConnectAsync(args)) return new ValueTask(tcs.Task); - - StateObjectPool.Return(args); - - return new ValueTask(); - } - - /// <inheritdoc /> - public override int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) - { - int totalBytes = readBuffer.Length; - if (totalBytes > BufferSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {BufferSize} bytes", - nameof(readBuffer.Length) - ); - } - - byte[] transmissionBuffer = BufferPool.Rent(BufferSize); - - int readBytes = Connection.ReceiveFrom(transmissionBuffer, flags, ref remoteEndPoint); - - 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 > BufferSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {BufferSize} bytes", - nameof(readBuffer.Length) - ); - } - - TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); - SocketAsyncEventArgs args = StateObjectPool.Rent(); - - byte[] transmissionBuffer = BufferPool.Rent(BufferSize); - - args.SetBuffer(transmissionBuffer, 0, BufferSize); - - args.RemoteEndPoint = remoteEndPoint; - args.SocketFlags = flags; - - AsyncDatagramReadToken token = new AsyncDatagramReadToken(tcs, in readBuffer); - args.UserToken = token; - - if (Connection.ReceiveFromAsync(args)) return new ValueTask<int>(tcs.Task); - - // inlining CompleteReceiveFrom(SocketAsyncEventArgs) for performance - int result = args.BytesTransferred; - - transmissionBuffer.CopyTo(readBuffer); - - BufferPool.Return(transmissionBuffer, true); - StateObjectPool.Return(args); - - return new ValueTask<int>(result); - } - - /// <inheritdoc /> - public override int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, - SocketFlags flags = SocketFlags.None) - { - int totalBytes = writeBuffer.Length; - if (totalBytes > BufferSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {BufferSize} bytes", - nameof(writeBuffer.Length) - ); - } - - byte[] transmissionBuffer = BufferPool.Rent(BufferSize); - writeBuffer.CopyTo(transmissionBuffer); - - int writtenBytes = Connection.SendTo(transmissionBuffer, flags, remoteEndPoint); - - 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 > BufferSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {BufferSize} bytes", - nameof(writeBuffer.Length) - ); - } - - TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); - SocketAsyncEventArgs args = StateObjectPool.Rent(); - - byte[] transmissionBuffer = BufferPool.Rent(BufferSize); - writeBuffer.CopyTo(transmissionBuffer); - - args.SetBuffer(transmissionBuffer, 0, BufferSize); - - args.RemoteEndPoint = remoteEndPoint; - args.SocketFlags = flags; - - AsyncDatagramWriteToken token = new AsyncDatagramWriteToken(tcs); - args.UserToken = token; - - if (Connection.SendToAsync(args)) return new ValueTask<int>(tcs.Task); - - // inlining CompleteSendTo(SocketAsyncEventArgs) for performance - int result = args.BytesTransferred; - - BufferPool.Return(transmissionBuffer, true); - StateObjectPool.Return(args); - - return new ValueTask<int>(result); - } - - private readonly struct AsyncDatagramReadToken - { - public readonly TaskCompletionSource<int> CompletionSource; - public readonly Memory<byte> UserBuffer; - - public AsyncDatagramReadToken(TaskCompletionSource<int> completionSource, in Memory<byte> userBuffer) - { - CompletionSource = completionSource; - - UserBuffer = userBuffer; - } - } - - private readonly struct AsyncDatagramWriteToken - { - public readonly TaskCompletionSource<int> CompletionSource; - - public AsyncDatagramWriteToken(TaskCompletionSource<int> completionSource) - { - CompletionSource = completionSource; - } - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Datagram/RawDatagramNetworkReader.cs b/NetSharp/NetSharp/Raw/Datagram/RawDatagramNetworkReader.cs @@ -0,0 +1,156 @@ +using System.Net; +using System.Net.Sockets; + +namespace NetSharp.Raw.Datagram +{ + public sealed class RawDatagramNetworkReader : RawNetworkReaderBase<SocketAsyncEventArgs> + { + /// <inheritdoc /> + public RawDatagramNetworkReader(ref Socket rawConnection, NetworkRequestHandler? requestHandler, EndPoint defaultEndPoint, int pooledPacketBufferSize, + int pooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, requestHandler, pooledPacketBufferSize, + pooledBuffersPerBucket, preallocatedStateObjects) + { + } + + private void CompleteReceiveFrom(SocketAsyncEventArgs args) + { + byte[] receiveBuffer = args.Buffer; + + switch (args.SocketError) + { + case SocketError.Success: + byte[] responseBuffer = BufferPool.Rent(PacketBufferSize); + + bool responseExists = + RequestHandler(args.RemoteEndPoint, receiveBuffer, args.BytesTransferred, responseBuffer); + BufferPool.Return(receiveBuffer, true); + + if (responseExists) + { + args.SetBuffer(responseBuffer, 0, PacketBufferSize); + + StartSendTo(args); + + return; + } + + BufferPool.Return(responseBuffer, true); + break; + + default: + BufferPool.Return(receiveBuffer, true); + StateObjectPool.Return(args); + break; + } + } + + private void CompleteSendTo(SocketAsyncEventArgs args) + { + byte[] sendBuffer = args.Buffer; + + BufferPool.Return(sendBuffer, true); + StateObjectPool.Return(args); + } + + private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) + { + switch (args.LastOperation) + { + case SocketAsyncOperation.ReceiveFrom: + StartDefaultReceiveFrom(); + + CompleteReceiveFrom(args); + break; + + case SocketAsyncOperation.SendTo: + CompleteSendTo(args); + break; + } + } + + private void StartDefaultReceiveFrom() + { + if (ShutdownToken.IsCancellationRequested) + { + return; + } + + SocketAsyncEventArgs args = StateObjectPool.Rent(); + StartReceiveFrom(args); + } + + private void StartReceiveFrom(SocketAsyncEventArgs args) + { + byte[] receiveBuffer = BufferPool.Rent(PacketBufferSize); + + if (ShutdownToken.IsCancellationRequested) + { + BufferPool.Return(receiveBuffer, true); + StateObjectPool.Return(args); + + return; + } + + args.SetBuffer(receiveBuffer, 0, PacketBufferSize); + + if (Connection.ReceiveFromAsync(args)) return; + + StartDefaultReceiveFrom(); + CompleteReceiveFrom(args); + } + + private void StartSendTo(SocketAsyncEventArgs args) + { + byte[] sendBuffer = args.Buffer; + + if (ShutdownToken.IsCancellationRequested) + { + BufferPool.Return(sendBuffer, true); + StateObjectPool.Return(args); + + return; + } + + if (Connection.SendToAsync(args)) return; + + CompleteSendTo(args); + } + + /// <inheritdoc /> + protected override bool CanReuseStateObject(ref SocketAsyncEventArgs instance) + { + return true; + } + + /// <inheritdoc /> + protected override SocketAsyncEventArgs CreateStateObject() + { + SocketAsyncEventArgs instance = new SocketAsyncEventArgs { RemoteEndPoint = DefaultEndPoint }; + instance.Completed += HandleIoCompleted; + + return instance; + } + + /// <inheritdoc /> + protected override void DestroyStateObject(SocketAsyncEventArgs instance) + { + instance.Completed -= HandleIoCompleted; + instance.Dispose(); + } + + /// <inheritdoc /> + protected override void ResetStateObject(ref SocketAsyncEventArgs instance) + { + instance.RemoteEndPoint = DefaultEndPoint; + } + + /// <inheritdoc /> + public override void Start(ushort concurrentReadTasks) + { + for (ushort i = 0; i < concurrentReadTasks; i++) + { + StartDefaultReceiveFrom(); + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Datagram/RawDatagramNetworkWriter.cs b/NetSharp/NetSharp/Raw/Datagram/RawDatagramNetworkWriter.cs @@ -0,0 +1,304 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Threading.Tasks; + +namespace NetSharp.Raw.Datagram +{ + public sealed class RawDatagramNetworkWriter : RawNetworkWriterBase<SocketAsyncEventArgs> + { + /// <inheritdoc /> + public RawDatagramNetworkWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int pooledPacketBufferSize, int pooledBuffersPerBucket = 1000, + uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, pooledPacketBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) + { + } + + private void CompleteConnect(SocketAsyncEventArgs args) + { + AsyncOperationToken token = (AsyncOperationToken)args.UserToken; + + switch (args.SocketError) + { + case SocketError.Success: + token.CompletionSource.SetResult(true); + break; + + case SocketError.OperationAborted: + token.CompletionSource.SetCanceled(); + break; + + default: + int errorCode = (int)args.SocketError; + token.CompletionSource.SetException(new SocketException(errorCode)); + break; + } + + StateObjectPool.Return(args); + } + + private void CompleteReceiveFrom(SocketAsyncEventArgs args) + { + AsyncDatagramReadToken token = (AsyncDatagramReadToken)args.UserToken; + + byte[] receiveBuffer = args.Buffer; + + switch (args.SocketError) + { + case SocketError.Success: + receiveBuffer.CopyTo(token.UserBuffer); + token.CompletionSource.SetResult(args.BytesTransferred); + 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); + StateObjectPool.Return(args); + } + + private void CompleteSendTo(SocketAsyncEventArgs args) + { + AsyncDatagramWriteToken token = (AsyncDatagramWriteToken)args.UserToken; + + byte[] sendBuffer = args.Buffer; + + switch (args.SocketError) + { + case SocketError.Success: + token.CompletionSource.SetResult(args.BytesTransferred); + 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); + StateObjectPool.Return(args); + } + + private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) + { + switch (args.LastOperation) + { + case SocketAsyncOperation.Connect: + CompleteConnect(args); + break; + + case SocketAsyncOperation.SendTo: + CompleteSendTo(args); + break; + + case SocketAsyncOperation.ReceiveFrom: + CompleteReceiveFrom(args); + break; + } + } + + /// <inheritdoc /> + protected override bool CanReuseStateObject(ref SocketAsyncEventArgs instance) + { + return true; + } + + /// <inheritdoc /> + protected override SocketAsyncEventArgs CreateStateObject() + { + SocketAsyncEventArgs instance = new SocketAsyncEventArgs(); + instance.Completed += HandleIoCompleted; + + return instance; + } + + /// <inheritdoc /> + protected override void DestroyStateObject(SocketAsyncEventArgs instance) + { + instance.Completed -= HandleIoCompleted; + instance.Dispose(); + } + + /// <inheritdoc /> + protected override void ResetStateObject(ref SocketAsyncEventArgs instance) + { + } + + /// <inheritdoc /> + public override void Connect(EndPoint remoteEndPoint) + { + Connection.Connect(remoteEndPoint); + } + + /// <inheritdoc /> + public override ValueTask ConnectAsync(EndPoint remoteEndPoint) + { + TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>(); + SocketAsyncEventArgs args = StateObjectPool.Rent(); + + args.RemoteEndPoint = remoteEndPoint; + + AsyncOperationToken token = new AsyncOperationToken(tcs); + args.UserToken = token; + + if (Connection.ConnectAsync(args)) return new ValueTask(tcs.Task); + + StateObjectPool.Return(args); + + return new ValueTask(); + } + + /// <inheritdoc /> + public override int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) + { + int totalBytes = readBuffer.Length; + if (totalBytes > PacketBufferSize) + { + throw new ArgumentException( + $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {PacketBufferSize} bytes", + nameof(readBuffer.Length) + ); + } + + byte[] transmissionBuffer = BufferPool.Rent(PacketBufferSize); + + int readBytes = Connection.ReceiveFrom(transmissionBuffer, flags, ref remoteEndPoint); + + 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 > PacketBufferSize) + { + throw new ArgumentException( + $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {PacketBufferSize} bytes", + nameof(readBuffer.Length) + ); + } + + TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); + SocketAsyncEventArgs args = StateObjectPool.Rent(); + + byte[] transmissionBuffer = BufferPool.Rent(PacketBufferSize); + + args.SetBuffer(transmissionBuffer, 0, PacketBufferSize); + + args.RemoteEndPoint = remoteEndPoint; + args.SocketFlags = flags; + + AsyncDatagramReadToken token = new AsyncDatagramReadToken(tcs, in readBuffer); + args.UserToken = token; + + if (Connection.ReceiveFromAsync(args)) return new ValueTask<int>(tcs.Task); + + // inlining CompleteReceiveFrom(SocketAsyncEventArgs) for performance + int result = args.BytesTransferred; + + transmissionBuffer.CopyTo(readBuffer); + + BufferPool.Return(transmissionBuffer, true); + StateObjectPool.Return(args); + + return new ValueTask<int>(result); + } + + /// <inheritdoc /> + public override int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, + SocketFlags flags = SocketFlags.None) + { + int totalBytes = writeBuffer.Length; + if (totalBytes > PacketBufferSize) + { + throw new ArgumentException( + $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {PacketBufferSize} bytes", + nameof(writeBuffer.Length) + ); + } + + byte[] transmissionBuffer = BufferPool.Rent(PacketBufferSize); + writeBuffer.CopyTo(transmissionBuffer); + + int writtenBytes = Connection.SendTo(transmissionBuffer, flags, remoteEndPoint); + + 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 > PacketBufferSize) + { + throw new ArgumentException( + $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {PacketBufferSize} bytes", + nameof(writeBuffer.Length) + ); + } + + TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); + SocketAsyncEventArgs args = StateObjectPool.Rent(); + + byte[] transmissionBuffer = BufferPool.Rent(PacketBufferSize); + writeBuffer.CopyTo(transmissionBuffer); + + args.SetBuffer(transmissionBuffer, 0, PacketBufferSize); + + args.RemoteEndPoint = remoteEndPoint; + args.SocketFlags = flags; + + AsyncDatagramWriteToken token = new AsyncDatagramWriteToken(tcs); + args.UserToken = token; + + if (Connection.SendToAsync(args)) return new ValueTask<int>(tcs.Task); + + // inlining CompleteSendTo(SocketAsyncEventArgs) for performance + int result = args.BytesTransferred; + + BufferPool.Return(transmissionBuffer, true); + StateObjectPool.Return(args); + + return new ValueTask<int>(result); + } + + private readonly struct AsyncDatagramReadToken + { + public readonly TaskCompletionSource<int> CompletionSource; + public readonly Memory<byte> UserBuffer; + + public AsyncDatagramReadToken(TaskCompletionSource<int> completionSource, in Memory<byte> userBuffer) + { + CompletionSource = completionSource; + + UserBuffer = userBuffer; + } + } + + private readonly struct AsyncDatagramWriteToken + { + public readonly TaskCompletionSource<int> CompletionSource; + + public AsyncDatagramWriteToken(TaskCompletionSource<int> completionSource) + { + CompletionSource = completionSource; + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/IRawNetworkTransportProvider.cs b/NetSharp/NetSharp/Raw/IRawNetworkTransportProvider.cs @@ -0,0 +1,64 @@ +using NetSharp.Raw.Datagram; +using NetSharp.Raw.Stream; + +using System.Net; +using System.Net.Sockets; + +namespace NetSharp.Raw +{ + public interface IRawNetworkTransportProvider<TState> where TState : class + { + SocketType TransportProtocolType { get; } + + RawNetworkReaderBase<TState> GetReader(ref Socket rawConnection, EndPoint defaultEndPoint, NetworkRequestHandler? requestHandler, int maxPooledBufferSize, + int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0); + + RawNetworkWriterBase<TState> GetWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, + uint preallocatedStateObjects = 0); + } + + public sealed class DatagramRawNetworkTransportProvider : IRawNetworkTransportProvider<SocketAsyncEventArgs> + { + /// <inheritdoc /> + public SocketType TransportProtocolType { get; } = SocketType.Dgram; + + /// <inheritdoc /> + public RawNetworkReaderBase<SocketAsyncEventArgs> GetReader(ref Socket rawConnection, EndPoint defaultEndPoint, + NetworkRequestHandler? requestHandler, int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, + uint preallocatedStateObjects = 0) + { + return new RawDatagramNetworkReader(ref rawConnection, requestHandler, defaultEndPoint, maxPooledBufferSize, + maxPooledBuffersPerBucket, preallocatedStateObjects); + } + + /// <inheritdoc /> + public RawNetworkWriterBase<SocketAsyncEventArgs> GetWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, + int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) + { + return new RawDatagramNetworkWriter(ref rawConnection, defaultEndPoint, maxPooledBufferSize, + maxPooledBuffersPerBucket, preallocatedStateObjects); + } + } + + public sealed class StreamRawNetworkTransportProvider : IRawNetworkTransportProvider<SocketAsyncEventArgs> + { + /// <inheritdoc /> + public SocketType TransportProtocolType { get; } = SocketType.Stream; + + /// <inheritdoc /> + public RawNetworkReaderBase<SocketAsyncEventArgs> GetReader(ref Socket rawConnection, EndPoint defaultEndPoint, NetworkRequestHandler? requestHandler, + int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) + { + return new RawStreamNetworkReader(ref rawConnection, requestHandler, defaultEndPoint, maxPooledBufferSize, + maxPooledBuffersPerBucket, preallocatedStateObjects); + } + + /// <inheritdoc /> + public RawNetworkWriterBase<SocketAsyncEventArgs> GetWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, + int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) + { + return new RawStreamNetworkWriter(ref rawConnection, defaultEndPoint, maxPooledBufferSize, + maxPooledBuffersPerBucket, preallocatedStateObjects); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/NetworkConnectionBase.cs b/NetSharp/NetSharp/Raw/NetworkConnectionBase.cs @@ -1,63 +0,0 @@ -using NetSharp.Utils; - -using System; -using System.Buffers; -using System.Net; -using System.Net.Sockets; - -namespace NetSharp.Raw -{ - public abstract class NetworkConnectionBase<TState> : IDisposable where TState : class - { - protected readonly ArrayPool<byte> BufferPool; - protected readonly int BufferSize; - protected readonly Socket Connection; - protected readonly EndPoint DefaultEndPoint; - protected readonly SlimObjectPool<TState> StateObjectPool; - - protected NetworkConnectionBase(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, - int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) - { - Connection = rawConnection; - - BufferSize = maxPooledBufferSize; - BufferPool = ArrayPool<byte>.Create(maxPooledBufferSize, maxPooledBuffersPerBucket); - - DefaultEndPoint = defaultEndPoint; - - StateObjectPool = - new SlimObjectPool<TState>(CreateStateObject, ResetStateObject, DestroyStateObject, CanReuseStateObject); - - // TODO implement pooling in better way - for (uint i = 0; i < preallocatedStateObjects; i++) - { - StateObjectPool.Return(CreateStateObject()); - } - } - - protected abstract bool CanReuseStateObject(ref TState instance); - - protected abstract TState CreateStateObject(); - - protected abstract void DestroyStateObject(TState instance); - - /// <summary> - /// Allows for inheritors to dispose of their own resources. - /// </summary> - protected virtual void Dispose(bool disposing) - { - if (!disposing) return; - - StateObjectPool.Dispose(); - } - - protected abstract void ResetStateObject(ref TState instance); - - /// <inheritdoc /> - public void Dispose() - { - Dispose(true); - GC.SuppressFinalize(this); - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/NetworkReaderBase.cs b/NetSharp/NetSharp/Raw/NetworkReaderBase.cs @@ -1,53 +0,0 @@ -using System; -using System.Net; -using System.Net.Sockets; -using System.Threading; - -namespace NetSharp.Raw -{ - public delegate bool NetworkRequestHandler(in EndPoint remoteEndPoint, ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, - Memory<byte> responseBuffer); - - public abstract class NetworkReaderBase<TState> : NetworkConnectionBase<TState> where TState : class - { - private readonly CancellationTokenSource shutdownTokenSource; - - protected readonly NetworkRequestHandler RequestHandler; - protected readonly CancellationToken ShutdownToken; - - /// <inheritdoc /> - protected NetworkReaderBase(ref Socket rawConnection, EndPoint defaultEndPoint, NetworkRequestHandler? requestHandler, int maxPooledBufferSize, - int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxPooledBufferSize, - maxPooledBuffersPerBucket, preallocatedStateObjects) - { - shutdownTokenSource = new CancellationTokenSource(); - ShutdownToken = shutdownTokenSource.Token; - - RequestHandler = requestHandler ?? DefaultRequestHandler; - } - - /// <inheritdoc /> - protected override void Dispose(bool disposing) - { - if (!disposing) return; - - shutdownTokenSource.Cancel(); - shutdownTokenSource.Dispose(); - - base.Dispose(disposing); - } - - public static bool DefaultRequestHandler(in EndPoint remoteEndPoint, ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, - Memory<byte> responseBuffer) - { - return requestBuffer.TryCopyTo(responseBuffer); - } - - public abstract void Start(ushort concurrentReadTasks); - - public void Stop() - { - shutdownTokenSource.Cancel(); - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/NetworkWriterBase.cs b/NetSharp/NetSharp/Raw/NetworkWriterBase.cs @@ -1,42 +0,0 @@ -using System; -using System.Net; -using System.Net.Sockets; -using System.Threading.Tasks; - -namespace NetSharp.Raw -{ - public abstract class NetworkWriterBase<TState> : NetworkConnectionBase<TState> where TState : class - { - /// <inheritdoc /> - protected NetworkWriterBase(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, - uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxPooledBufferSize, maxPooledBuffersPerBucket, preallocatedStateObjects) - { - } - - public abstract void Connect(EndPoint remoteEndPoint); - - public abstract ValueTask ConnectAsync(EndPoint remoteEndPoint); - - public abstract int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, - SocketFlags flags = SocketFlags.None); - - public abstract ValueTask<int> ReadAsync(EndPoint remoteEndPoint, Memory<byte> readBuffer, - SocketFlags flags = SocketFlags.None); - - public abstract int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, - SocketFlags flags = SocketFlags.None); - - public abstract ValueTask<int> WriteAsync(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, - SocketFlags flags = SocketFlags.None); - - protected readonly struct AsyncOperationToken - { - public readonly TaskCompletionSource<bool> CompletionSource; - - public AsyncOperationToken(TaskCompletionSource<bool> completionSource) - { - CompletionSource = completionSource; - } - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/RawNetworkConnectionBase.cs b/NetSharp/NetSharp/Raw/RawNetworkConnectionBase.cs @@ -0,0 +1,63 @@ +using NetSharp.Utils; + +using System; +using System.Buffers; +using System.Net; +using System.Net.Sockets; + +namespace NetSharp.Raw +{ + public abstract class RawNetworkConnectionBase<TState> : IDisposable where TState : class + { + protected readonly ArrayPool<byte> BufferPool; + protected readonly Socket Connection; + protected readonly EndPoint DefaultEndPoint; + protected readonly int PacketBufferSize; + protected readonly SlimObjectPool<TState> StateObjectPool; + + protected RawNetworkConnectionBase(ref Socket rawConnection, EndPoint defaultEndPoint, int pooledPacketBufferSize, + int pooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) + { + Connection = rawConnection; + + PacketBufferSize = pooledPacketBufferSize; + BufferPool = ArrayPool<byte>.Create(pooledPacketBufferSize, pooledBuffersPerBucket); + + DefaultEndPoint = defaultEndPoint; + + StateObjectPool = + new SlimObjectPool<TState>(CreateStateObject, ResetStateObject, DestroyStateObject, CanReuseStateObject); + + // TODO implement pooling in better way + for (uint i = 0; i < preallocatedStateObjects; i++) + { + StateObjectPool.Return(CreateStateObject()); + } + } + + protected abstract bool CanReuseStateObject(ref TState instance); + + protected abstract TState CreateStateObject(); + + protected abstract void DestroyStateObject(TState instance); + + /// <summary> + /// Allows for inheritors to dispose of their own resources. + /// </summary> + protected virtual void Dispose(bool disposing) + { + if (!disposing) return; + + StateObjectPool.Dispose(); + } + + protected abstract void ResetStateObject(ref TState instance); + + /// <inheritdoc /> + public void Dispose() + { + Dispose(true); + GC.SuppressFinalize(this); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/RawNetworkReaderBase.cs b/NetSharp/NetSharp/Raw/RawNetworkReaderBase.cs @@ -0,0 +1,55 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Threading; + +namespace NetSharp.Raw +{ + public delegate bool NetworkRequestHandler(in EndPoint remoteEndPoint, ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + Memory<byte> responseBuffer); + + public abstract class RawNetworkReaderBase<TState> : RawNetworkConnectionBase<TState>, INetworkReader where TState : class + { + private readonly CancellationTokenSource shutdownTokenSource; + + protected readonly NetworkRequestHandler RequestHandler; + protected readonly CancellationToken ShutdownToken; + + /// <inheritdoc /> + protected RawNetworkReaderBase(ref Socket rawConnection, EndPoint defaultEndPoint, NetworkRequestHandler? requestHandler, int pooledPacketBufferSize, + int pooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, pooledPacketBufferSize, + pooledBuffersPerBucket, preallocatedStateObjects) + { + shutdownTokenSource = new CancellationTokenSource(); + ShutdownToken = shutdownTokenSource.Token; + + RequestHandler = requestHandler ?? DefaultRequestHandler; + } + + /// <inheritdoc /> + protected override void Dispose(bool disposing) + { + if (!disposing) return; + + shutdownTokenSource.Cancel(); + shutdownTokenSource.Dispose(); + + base.Dispose(disposing); + } + + public static bool DefaultRequestHandler(in EndPoint remoteEndPoint, ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, + Memory<byte> responseBuffer) + { + return requestBuffer.TryCopyTo(responseBuffer); + } + + /// <inheritdoc /> + public abstract void Start(ushort concurrentReadTasks); + + /// <inheritdoc /> + public void Stop() + { + shutdownTokenSource.Cancel(); + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/RawNetworkWriterBase.cs b/NetSharp/NetSharp/Raw/RawNetworkWriterBase.cs @@ -0,0 +1,48 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Threading.Tasks; + +namespace NetSharp.Raw +{ + public abstract class RawNetworkWriterBase<TState> : RawNetworkConnectionBase<TState>, INetworkWriter where TState : class + { + /// <inheritdoc /> + protected RawNetworkWriterBase(ref Socket rawConnection, EndPoint defaultEndPoint, int pooledPacketBufferSize, int pooledBuffersPerBucket = 1000, + uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, pooledPacketBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) + { + } + + /// <inheritdoc /> + public abstract void Connect(EndPoint remoteEndPoint); + + /// <inheritdoc /> + public abstract ValueTask ConnectAsync(EndPoint remoteEndPoint); + + /// <inheritdoc /> + public abstract int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, + SocketFlags flags = SocketFlags.None); + + /// <inheritdoc /> + public abstract ValueTask<int> ReadAsync(EndPoint remoteEndPoint, Memory<byte> readBuffer, + SocketFlags flags = SocketFlags.None); + + /// <inheritdoc /> + public abstract int Write(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, + SocketFlags flags = SocketFlags.None); + + /// <inheritdoc /> + public abstract ValueTask<int> WriteAsync(EndPoint remoteEndPoint, ReadOnlyMemory<byte> writeBuffer, + SocketFlags flags = SocketFlags.None); + + protected readonly struct AsyncOperationToken + { + public readonly TaskCompletionSource<bool> CompletionSource; + + public AsyncOperationToken(TaskCompletionSource<bool> completionSource) + { + CompletionSource = completionSource; + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Stream/RawStreamNetworkReader.cs b/NetSharp/NetSharp/Raw/Stream/RawStreamNetworkReader.cs @@ -0,0 +1,314 @@ +using System.Net; +using System.Net.Sockets; +using System.Runtime.CompilerServices; + +namespace NetSharp.Raw.Stream +{ + public sealed class RawStreamNetworkReader : RawNetworkReaderBase<SocketAsyncEventArgs> + { + /// <inheritdoc /> + public RawStreamNetworkReader(ref Socket rawConnection, NetworkRequestHandler? requestHandler, EndPoint defaultEndPoint, int pooledPacketBufferSize, + int pooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, requestHandler, pooledPacketBufferSize, + pooledBuffersPerBucket, preallocatedStateObjects) + { + } + + private void CloseClientConnection(SocketAsyncEventArgs args) + { + byte[] rentedBuffer = args.Buffer; + BufferPool.Return(rentedBuffer, true); + + Socket clientSocket = args.AcceptSocket; + + clientSocket.Shutdown(SocketShutdown.Both); + clientSocket.Close(); + clientSocket.Dispose(); + + StateObjectPool.Return(args); + } + + private void CompleteAccept(SocketAsyncEventArgs args) + { + switch (args.SocketError) + { + case SocketError.Success: + StartReceive(args); + break; + + case SocketError.ConnectionReset: + /* + * The SocketAsyncEventArgs.Completed event can occur in some cases when no connection has been accepted and cause the SocketAsyncEventArgs.SocketError property to be set to ConnectionReset. + * This can occur as a result of port scanning using a half-open SYN type scan (a SYN -> SYN-ACK -> RST sequence). + * Applications using the AcceptAsync method should be prepared to handle this condition. + */ + StateObjectPool.Return(args); + break; + + default: + StateObjectPool.Return(args); + break; + } + + StartDefaultAccept(); + } + + private void CompleteReceive(SocketAsyncEventArgs args) + { + TransmissionToken token = (TransmissionToken)args.UserToken; + + byte[] receiveBuffer = args.Buffer; + int expectedBytes = receiveBuffer.Length; + + switch (args.SocketError) + { + case SocketError.Success: + int receivedBytes = args.BytesTransferred, totalReceivedBytes = token.BytesTransferred; + + if (totalReceivedBytes + receivedBytes == expectedBytes) // transmission complete + { + byte[] responseBufferHandle = BufferPool.Rent(expectedBytes); + + bool responseExists = + RequestHandler(args.AcceptSocket.RemoteEndPoint, receiveBuffer, totalReceivedBytes + receivedBytes, responseBufferHandle); + BufferPool.Return(receiveBuffer, true); + + if (responseExists) + { + args.SetBuffer(responseBufferHandle, 0, PacketBufferSize); + + TransmissionToken sendToken = new TransmissionToken(0); + args.UserToken = sendToken; + + StartSend(args); + return; + } + + BufferPool.Return(responseBufferHandle, true); + + StartReceive(args); + } + else if (0 < totalReceivedBytes + receivedBytes && totalReceivedBytes + receivedBytes < expectedBytes) // transmission not complete + { + token = new TransmissionToken(in token, args.BytesTransferred); + args.UserToken = token; + + args.SetBuffer(totalReceivedBytes, expectedBytes - receivedBytes); + + 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 = sendBuffer.Length; + + switch (args.SocketError) + { + case SocketError.Success: + int sentBytes = args.BytesTransferred, totalSentBytes = token.BytesTransferred; + + if (totalSentBytes + sentBytes == expectedBytes) // transmission complete + { + BufferPool.Return(sendBuffer, true); + + StartReceive(args); + } + else if (0 < totalSentBytes + sentBytes && totalSentBytes + sentBytes < expectedBytes) // transmission not complete + { + token = new TransmissionToken(in token, args.BytesTransferred); + args.UserToken = token; + + args.SetBuffer(totalSentBytes, expectedBytes - sentBytes); + + ContinueSend(args); + } + else if (sentBytes == 0) // connection is dead + { + CloseClientConnection(args); + } + break; + + default: + CloseClientConnection(args); + break; + } + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ContinueReceive(SocketAsyncEventArgs args) + { + if (ShutdownToken.IsCancellationRequested) + { + CloseClientConnection(args); + return; + } + + Socket clientSocket = args.AcceptSocket; + + if (clientSocket.ReceiveAsync(args)) return; + + CompleteReceive(args); + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ContinueSend(SocketAsyncEventArgs args) + { + if (ShutdownToken.IsCancellationRequested) + { + CloseClientConnection(args); + return; + } + + Socket clientSocket = args.AcceptSocket; + + if (clientSocket.SendAsync(args)) return; + + CompleteSend(args); + } + + private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) + { + switch (args.LastOperation) + { + case SocketAsyncOperation.Accept: + StartDefaultAccept(); + CompleteAccept(args); + break; + + case SocketAsyncOperation.Send: + CompleteSend(args); + break; + + case SocketAsyncOperation.Receive: + CompleteReceive(args); + break; + } + } + + private void StartAccept(SocketAsyncEventArgs args) + { + if (ShutdownToken.IsCancellationRequested) + { + return; + } + + if (Connection.AcceptAsync(args)) return; + + StartDefaultAccept(); + CompleteAccept(args); + } + + private void StartDefaultAccept() + { + if (ShutdownToken.IsCancellationRequested) + { + return; + } + + SocketAsyncEventArgs args = StateObjectPool.Rent(); + StartAccept(args); + } + + private void StartReceive(SocketAsyncEventArgs args) + { + if (ShutdownToken.IsCancellationRequested) + { + CloseClientConnection(args); + return; + } + + Socket clientSocket = args.AcceptSocket; + + byte[] receiveBuffer = BufferPool.Rent(PacketBufferSize); + + args.SetBuffer(receiveBuffer, 0, PacketBufferSize); + + TransmissionToken token = new TransmissionToken(0); + args.UserToken = token; + + if (clientSocket.ReceiveAsync(args)) return; + + CompleteReceive(args); + } + + private void StartSend(SocketAsyncEventArgs args) + { + if (ShutdownToken.IsCancellationRequested) + { + CloseClientConnection(args); + return; + } + + Socket clientSocket = args.AcceptSocket; + + if (clientSocket.SendAsync(args)) return; + + CompleteSend(args); + } + + /// <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) + { + for (ushort i = 0; i < concurrentReadTasks; i++) + { + StartDefaultAccept(); + } + } + + private readonly struct TransmissionToken + { + public readonly int BytesTransferred; + + public TransmissionToken(int bytesTransferred) + { + BytesTransferred = bytesTransferred; + } + + public TransmissionToken(in TransmissionToken token, int newlyTransferredBytes) + { + BytesTransferred = token.BytesTransferred + newlyTransferredBytes; + } + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Stream/RawStreamNetworkWriter.cs b/NetSharp/NetSharp/Raw/Stream/RawStreamNetworkWriter.cs @@ -0,0 +1,481 @@ +using System; +using System.Net; +using System.Net.Sockets; +using System.Runtime.CompilerServices; +using System.Threading.Tasks; + +namespace NetSharp.Raw.Stream +{ + public sealed class RawStreamNetworkWriter : RawNetworkWriterBase<SocketAsyncEventArgs> + { + /// <inheritdoc /> + public RawStreamNetworkWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int pooledPacketBufferSize, int pooledBuffersPerBucket = 1000, + uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, pooledPacketBufferSize, pooledBuffersPerBucket, preallocatedStateObjects) + { + } + + private void CompleteConnect(SocketAsyncEventArgs args) + { + AsyncOperationToken token = (AsyncOperationToken)args.UserToken; + + switch (args.SocketError) + { + case SocketError.Success: + token.CompletionSource.SetResult(true); + break; + + case SocketError.OperationAborted: + token.CompletionSource.SetCanceled(); + break; + + default: + int errorCode = (int)args.SocketError; + token.CompletionSource.SetException(new SocketException(errorCode)); + break; + } + + StateObjectPool.Return(args); + } + + private void CompleteDisconnect(SocketAsyncEventArgs args) + { + AsyncOperationToken token = (AsyncOperationToken)args.UserToken; + + switch (args.SocketError) + { + case SocketError.Success: + token.CompletionSource.SetResult(true); + break; + + case SocketError.OperationAborted: + token.CompletionSource.SetCanceled(); + break; + + default: + int errorCode = (int)args.SocketError; + token.CompletionSource.SetException(new SocketException(errorCode)); + break; + } + + StateObjectPool.Return(args); + } + + private 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); + StateObjectPool.Return(args); + } + + private 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); + StateObjectPool.Return(args); + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void ContinueReceive(SocketAsyncEventArgs args) + { + if (Connection.ReceiveAsync(args)) return; + + CompleteReceive(args); + } + + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private 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() + { + 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) + { + } + + /// <inheritdoc /> + public override void Connect(EndPoint remoteEndPoint) + { + Connection.Connect(remoteEndPoint); + } + + /// <inheritdoc /> + public override ValueTask ConnectAsync(EndPoint remoteEndPoint) + { + TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>(); + SocketAsyncEventArgs args = StateObjectPool.Rent(); + + args.RemoteEndPoint = remoteEndPoint; + + AsyncOperationToken token = new AsyncOperationToken(tcs); + args.UserToken = token; + + if (Connection.ConnectAsync(args)) return new ValueTask(tcs.Task); + + StateObjectPool.Return(args); + + return new ValueTask(); + } + + public void Disconnect(bool reuseSocket) + { + Connection.Disconnect(reuseSocket); + } + + public ValueTask DisconnectAsync(bool reuseSocket) + { + TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>(); + SocketAsyncEventArgs args = StateObjectPool.Rent(); + + args.DisconnectReuseSocket = reuseSocket; + + AsyncOperationToken token = new AsyncOperationToken(tcs); + args.UserToken = token; + + if (Connection.DisconnectAsync(args)) return new ValueTask(tcs.Task); + + StateObjectPool.Return(args); + + return new ValueTask(); + } + + /// <inheritdoc /> + public override int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) + { + int totalBytes = readBuffer.Length; + if (totalBytes > PacketBufferSize) + { + throw new ArgumentException( + $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {PacketBufferSize} 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 > PacketBufferSize) + { + throw new ArgumentException( + $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {PacketBufferSize} bytes", + nameof(readBuffer.Length) + ); + } + + TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); + SocketAsyncEventArgs args = StateObjectPool.Rent(); + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + + args.SetBuffer(transmissionBuffer, 0, PacketBufferSize); + + 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 == PacketBufferSize) // transmission complete + { + transmissionBuffer.CopyTo(readBuffer); + + BufferPool.Return(transmissionBuffer, true); + StateObjectPool.Return(args); + + return new ValueTask<int>(totalReceivedBytes + receivedBytes); + } + else if (0 < totalReceivedBytes + receivedBytes && totalReceivedBytes + receivedBytes < PacketBufferSize) // 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, PacketBufferSize - 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 > PacketBufferSize) + { + throw new ArgumentException( + $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {PacketBufferSize} 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 > PacketBufferSize) + { + throw new ArgumentException( + $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {PacketBufferSize} bytes", + nameof(writeBuffer.Length) + ); + } + + TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); + SocketAsyncEventArgs args = StateObjectPool.Rent(); + + byte[] transmissionBuffer = BufferPool.Rent(totalBytes); + writeBuffer.CopyTo(transmissionBuffer); + + args.SetBuffer(transmissionBuffer, 0, PacketBufferSize); + + 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 == PacketBufferSize) // transmission complete + { + BufferPool.Return(transmissionBuffer, true); + StateObjectPool.Return(args); + + return new ValueTask<int>(totalSentBytes + sentBytes); + } + else if (0 < totalSentBytes + sentBytes && totalSentBytes + sentBytes < PacketBufferSize) // 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, PacketBufferSize - 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/StreamNetworkReader.cs b/NetSharp/NetSharp/Raw/Stream/StreamNetworkReader.cs @@ -1,314 +0,0 @@ -using System.Net; -using System.Net.Sockets; -using System.Runtime.CompilerServices; - -namespace NetSharp.Raw.Stream -{ - public sealed class StreamNetworkReader : NetworkReaderBase<SocketAsyncEventArgs> - { - /// <inheritdoc /> - public StreamNetworkReader(ref Socket rawConnection, NetworkRequestHandler? requestHandler, EndPoint defaultEndPoint, int maxPooledBufferSize, - int maxPooledBuffersPerBucket = 1000, uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, requestHandler, maxPooledBufferSize, - maxPooledBuffersPerBucket, preallocatedStateObjects) - { - } - - private void CloseClientConnection(SocketAsyncEventArgs args) - { - byte[] rentedBuffer = args.Buffer; - BufferPool.Return(rentedBuffer, true); - - Socket clientSocket = args.AcceptSocket; - - clientSocket.Shutdown(SocketShutdown.Both); - clientSocket.Close(); - clientSocket.Dispose(); - - StateObjectPool.Return(args); - } - - private void CompleteAccept(SocketAsyncEventArgs args) - { - switch (args.SocketError) - { - case SocketError.Success: - StartReceive(args); - break; - - case SocketError.ConnectionReset: - /* - * The SocketAsyncEventArgs.Completed event can occur in some cases when no connection has been accepted and cause the SocketAsyncEventArgs.SocketError property to be set to ConnectionReset. - * This can occur as a result of port scanning using a half-open SYN type scan (a SYN -> SYN-ACK -> RST sequence). - * Applications using the AcceptAsync method should be prepared to handle this condition. - */ - StateObjectPool.Return(args); - break; - - default: - StateObjectPool.Return(args); - break; - } - - StartDefaultAccept(); - } - - private void CompleteReceive(SocketAsyncEventArgs args) - { - TransmissionToken token = (TransmissionToken)args.UserToken; - - byte[] receiveBuffer = args.Buffer; - int expectedBytes = receiveBuffer.Length; - - switch (args.SocketError) - { - case SocketError.Success: - int receivedBytes = args.BytesTransferred, totalReceivedBytes = token.BytesTransferred; - - if (totalReceivedBytes + receivedBytes == expectedBytes) // transmission complete - { - byte[] responseBufferHandle = BufferPool.Rent(expectedBytes); - - bool responseExists = - RequestHandler(args.AcceptSocket.RemoteEndPoint, receiveBuffer, totalReceivedBytes + receivedBytes, responseBufferHandle); - BufferPool.Return(receiveBuffer, true); - - if (responseExists) - { - args.SetBuffer(responseBufferHandle, 0, BufferSize); - - TransmissionToken sendToken = new TransmissionToken(0); - args.UserToken = sendToken; - - StartSend(args); - return; - } - - BufferPool.Return(responseBufferHandle, true); - - StartReceive(args); - } - else if (0 < totalReceivedBytes + receivedBytes && totalReceivedBytes + receivedBytes < expectedBytes) // transmission not complete - { - token = new TransmissionToken(in token, args.BytesTransferred); - args.UserToken = token; - - args.SetBuffer(totalReceivedBytes, expectedBytes - receivedBytes); - - 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 = sendBuffer.Length; - - switch (args.SocketError) - { - case SocketError.Success: - int sentBytes = args.BytesTransferred, totalSentBytes = token.BytesTransferred; - - if (totalSentBytes + sentBytes == expectedBytes) // transmission complete - { - BufferPool.Return(sendBuffer, true); - - StartReceive(args); - } - else if (0 < totalSentBytes + sentBytes && totalSentBytes + sentBytes < expectedBytes) // transmission not complete - { - token = new TransmissionToken(in token, args.BytesTransferred); - args.UserToken = token; - - args.SetBuffer(totalSentBytes, expectedBytes - sentBytes); - - ContinueSend(args); - } - else if (sentBytes == 0) // connection is dead - { - CloseClientConnection(args); - } - break; - - default: - CloseClientConnection(args); - break; - } - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ContinueReceive(SocketAsyncEventArgs args) - { - if (ShutdownToken.IsCancellationRequested) - { - CloseClientConnection(args); - return; - } - - Socket clientSocket = args.AcceptSocket; - - if (clientSocket.ReceiveAsync(args)) return; - - CompleteReceive(args); - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ContinueSend(SocketAsyncEventArgs args) - { - if (ShutdownToken.IsCancellationRequested) - { - CloseClientConnection(args); - return; - } - - Socket clientSocket = args.AcceptSocket; - - if (clientSocket.SendAsync(args)) return; - - CompleteSend(args); - } - - private void HandleIoCompleted(object sender, SocketAsyncEventArgs args) - { - switch (args.LastOperation) - { - case SocketAsyncOperation.Accept: - StartDefaultAccept(); - CompleteAccept(args); - break; - - case SocketAsyncOperation.Send: - CompleteSend(args); - break; - - case SocketAsyncOperation.Receive: - CompleteReceive(args); - break; - } - } - - private void StartAccept(SocketAsyncEventArgs args) - { - if (ShutdownToken.IsCancellationRequested) - { - return; - } - - if (Connection.AcceptAsync(args)) return; - - StartDefaultAccept(); - CompleteAccept(args); - } - - private void StartDefaultAccept() - { - if (ShutdownToken.IsCancellationRequested) - { - return; - } - - SocketAsyncEventArgs args = StateObjectPool.Rent(); - StartAccept(args); - } - - private void StartReceive(SocketAsyncEventArgs args) - { - if (ShutdownToken.IsCancellationRequested) - { - CloseClientConnection(args); - return; - } - - Socket clientSocket = args.AcceptSocket; - - byte[] receiveBuffer = BufferPool.Rent(BufferSize); - - args.SetBuffer(receiveBuffer, 0, BufferSize); - - TransmissionToken token = new TransmissionToken(0); - args.UserToken = token; - - if (clientSocket.ReceiveAsync(args)) return; - - CompleteReceive(args); - } - - private void StartSend(SocketAsyncEventArgs args) - { - if (ShutdownToken.IsCancellationRequested) - { - CloseClientConnection(args); - return; - } - - Socket clientSocket = args.AcceptSocket; - - if (clientSocket.SendAsync(args)) return; - - CompleteSend(args); - } - - /// <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) - { - for (ushort i = 0; i < concurrentReadTasks; i++) - { - StartDefaultAccept(); - } - } - - private readonly struct TransmissionToken - { - public readonly int BytesTransferred; - - public TransmissionToken(int bytesTransferred) - { - BytesTransferred = bytesTransferred; - } - - public TransmissionToken(in TransmissionToken token, int newlyTransferredBytes) - { - BytesTransferred = token.BytesTransferred + newlyTransferredBytes; - } - } - } -} -\ No newline at end of file diff --git a/NetSharp/NetSharp/Raw/Stream/StreamNetworkWriter.cs b/NetSharp/NetSharp/Raw/Stream/StreamNetworkWriter.cs @@ -1,481 +0,0 @@ -using System; -using System.Net; -using System.Net.Sockets; -using System.Runtime.CompilerServices; -using System.Threading.Tasks; - -namespace NetSharp.Raw.Stream -{ - public sealed class StreamNetworkWriter : NetworkWriterBase<SocketAsyncEventArgs> - { - /// <inheritdoc /> - public StreamNetworkWriter(ref Socket rawConnection, EndPoint defaultEndPoint, int maxPooledBufferSize, int maxPooledBuffersPerBucket = 1000, - uint preallocatedStateObjects = 0) : base(ref rawConnection, defaultEndPoint, maxPooledBufferSize, maxPooledBuffersPerBucket, preallocatedStateObjects) - { - } - - private void CompleteConnect(SocketAsyncEventArgs args) - { - AsyncOperationToken token = (AsyncOperationToken)args.UserToken; - - switch (args.SocketError) - { - case SocketError.Success: - token.CompletionSource.SetResult(true); - break; - - case SocketError.OperationAborted: - token.CompletionSource.SetCanceled(); - break; - - default: - int errorCode = (int)args.SocketError; - token.CompletionSource.SetException(new SocketException(errorCode)); - break; - } - - StateObjectPool.Return(args); - } - - private void CompleteDisconnect(SocketAsyncEventArgs args) - { - AsyncOperationToken token = (AsyncOperationToken)args.UserToken; - - switch (args.SocketError) - { - case SocketError.Success: - token.CompletionSource.SetResult(true); - break; - - case SocketError.OperationAborted: - token.CompletionSource.SetCanceled(); - break; - - default: - int errorCode = (int)args.SocketError; - token.CompletionSource.SetException(new SocketException(errorCode)); - break; - } - - StateObjectPool.Return(args); - } - - private 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); - StateObjectPool.Return(args); - } - - private 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); - StateObjectPool.Return(args); - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - private void ContinueReceive(SocketAsyncEventArgs args) - { - if (Connection.ReceiveAsync(args)) return; - - CompleteReceive(args); - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - private 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() - { - 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) - { - } - - /// <inheritdoc /> - public override void Connect(EndPoint remoteEndPoint) - { - Connection.Connect(remoteEndPoint); - } - - /// <inheritdoc /> - public override ValueTask ConnectAsync(EndPoint remoteEndPoint) - { - TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>(); - SocketAsyncEventArgs args = StateObjectPool.Rent(); - - args.RemoteEndPoint = remoteEndPoint; - - AsyncOperationToken token = new AsyncOperationToken(tcs); - args.UserToken = token; - - if (Connection.ConnectAsync(args)) return new ValueTask(tcs.Task); - - StateObjectPool.Return(args); - - return new ValueTask(); - } - - public void Disconnect(bool reuseSocket) - { - Connection.Disconnect(reuseSocket); - } - - public ValueTask DisconnectAsync(bool reuseSocket) - { - TaskCompletionSource<bool> tcs = new TaskCompletionSource<bool>(); - SocketAsyncEventArgs args = StateObjectPool.Rent(); - - args.DisconnectReuseSocket = reuseSocket; - - AsyncOperationToken token = new AsyncOperationToken(tcs); - args.UserToken = token; - - if (Connection.DisconnectAsync(args)) return new ValueTask(tcs.Task); - - StateObjectPool.Return(args); - - return new ValueTask(); - } - - /// <inheritdoc /> - public override int Read(ref EndPoint remoteEndPoint, Memory<byte> readBuffer, SocketFlags flags = SocketFlags.None) - { - int totalBytes = readBuffer.Length; - if (totalBytes > BufferSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {BufferSize} 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 > BufferSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {BufferSize} bytes", - nameof(readBuffer.Length) - ); - } - - TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); - SocketAsyncEventArgs args = StateObjectPool.Rent(); - - byte[] transmissionBuffer = BufferPool.Rent(totalBytes); - - args.SetBuffer(transmissionBuffer, 0, BufferSize); - - 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 == BufferSize) // transmission complete - { - transmissionBuffer.CopyTo(readBuffer); - - BufferPool.Return(transmissionBuffer, true); - StateObjectPool.Return(args); - - return new ValueTask<int>(totalReceivedBytes + receivedBytes); - } - else if (0 < totalReceivedBytes + receivedBytes && totalReceivedBytes + receivedBytes < BufferSize) // 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, BufferSize - 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 > BufferSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {BufferSize} 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 > BufferSize) - { - throw new ArgumentException( - $"Cannot rent a temporary buffer of size: {totalBytes} bytes; maximum temporary buffer size: {BufferSize} bytes", - nameof(writeBuffer.Length) - ); - } - - TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); - SocketAsyncEventArgs args = StateObjectPool.Rent(); - - byte[] transmissionBuffer = BufferPool.Rent(totalBytes); - writeBuffer.CopyTo(transmissionBuffer); - - args.SetBuffer(transmissionBuffer, 0, BufferSize); - - 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 == BufferSize) // transmission complete - { - BufferPool.Return(transmissionBuffer, true); - StateObjectPool.Return(args); - - return new ValueTask<int>(totalSentBytes + sentBytes); - } - else if (0 < totalSentBytes + sentBytes && totalSentBytes + sentBytes < BufferSize) // 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, BufferSize - 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/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkReaderBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkReaderBenchmark.cs @@ -10,7 +10,7 @@ using System.Threading.Tasks; namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks { - public class DatagramNetworkReaderBenchmark : INetSharpExample, INetSharpBenchmark + public class DatagramNetworkReaderBenchmark : INetSharpBenchmark { private const int PacketSize = 8192, PacketCount = 1_000_000, ClientCount = 12; @@ -23,7 +23,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Datagram Network Reader Benchmark"; + public string Name { get; } = "Datagram Raw Network Reader Benchmark"; private static bool RequestHandler(in EndPoint remoteEndPoint, ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, Memory<byte> responseBuffer) { @@ -102,7 +102,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp); rawSocket.Bind(ServerEndPoint); - using DatagramNetworkReader reader = new DatagramNetworkReader(ref rawSocket, RequestHandler, defaultRemoteEndPoint, PacketSize); + using RawDatagramNetworkReader reader = new RawDatagramNetworkReader(ref rawSocket, RequestHandler, defaultRemoteEndPoint, PacketSize); reader.Start(ClientCount); ServerReadyEvent.Set(); diff --git a/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterAsyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterAsyncBenchmark.cs @@ -9,7 +9,7 @@ using System.Threading.Tasks; namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks { - public class DatagramNetworkWriterAsyncBenchmark : INetSharpExample, INetSharpBenchmark + public class DatagramNetworkWriterAsyncBenchmark : INetSharpBenchmark { private const int PacketSize = 8192, PacketCount = 1_000_000; @@ -21,7 +21,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Datagram Network Writer Benchmark (Asynchronous)"; + public string Name { get; } = "Datagram Raw Network Writer Benchmark (Asynchronous)"; private static Task ServerTask(CancellationToken cancellationToken) { @@ -59,7 +59,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp); rawSocket.Bind(ClientEndPoint); - using DatagramNetworkWriter writer = new DatagramNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); + using RawDatagramNetworkWriter writer = new RawDatagramNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); using CancellationTokenSource serverCts = new CancellationTokenSource(); Task serverTask = Task.Factory.StartNew(state => ServerTask((CancellationToken)state), serverCts.Token, TaskCreationOptions.LongRunning); diff --git a/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterSyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Datagram Network Connection Benchmarks/DatagramNetworkWriterSyncBenchmark.cs @@ -9,7 +9,7 @@ using System.Threading.Tasks; namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks { - public class DatagramNetworkWriterSyncBenchmark : INetSharpExample, INetSharpBenchmark + public class DatagramNetworkWriterSyncBenchmark : INetSharpBenchmark { private const int PacketSize = 8192, PacketCount = 1_000_000; @@ -21,7 +21,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Datagram Network Writer Benchmark (Synchronous)"; + public string Name { get; } = "Datagram Raw Network Writer Benchmark (Synchronous)"; private static Task ServerTask(CancellationToken cancellationToken) { @@ -59,7 +59,7 @@ namespace NetSharpExamples.Benchmarks.Datagram_Network_Connection_Benchmarks Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp); rawSocket.Bind(ClientEndPoint); - using DatagramNetworkWriter writer = new DatagramNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); + using RawDatagramNetworkWriter writer = new RawDatagramNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); using CancellationTokenSource serverCts = new CancellationTokenSource(); Task serverTask = Task.Factory.StartNew(state => ServerTask((CancellationToken)state), serverCts.Token, TaskCreationOptions.LongRunning); diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkReaderBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkReaderBenchmark.cs @@ -10,7 +10,7 @@ using System.Threading.Tasks; namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks { - public class StreamNetworkReaderBenchmark : INetSharpExample, INetSharpBenchmark + public class StreamNetworkReaderBenchmark : INetSharpBenchmark { private const int PacketSize = 8192, PacketCount = 1_000_000, ClientCount = 12; @@ -23,7 +23,7 @@ namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Stream Network Reader Benchmark"; + public string Name { get; } = "Stream Raw Network Reader Benchmark"; private static bool RequestHandler(in EndPoint remoteEndPoint, ReadOnlyMemory<byte> requestBuffer, int receivedRequestBytes, Memory<byte> responseBuffer) { @@ -122,7 +122,7 @@ namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks rawSocket.Bind(ServerEndPoint); rawSocket.Listen(ClientCount); - using StreamNetworkReader reader = new StreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize); + using RawStreamNetworkReader reader = new RawStreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize); reader.Start(ClientCount); ServerReadyEvent.Set(); diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterAsyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterAsyncBenchmark.cs @@ -9,7 +9,7 @@ using System.Threading.Tasks; namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks { - public class StreamNetworkWriterAsyncBenchmark : INetSharpExample, INetSharpBenchmark + public class StreamNetworkWriterAsyncBenchmark : INetSharpBenchmark { private const int PacketSize = 8192, PacketCount = 1_000_000; @@ -21,7 +21,7 @@ namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Stream Network Writer Benchmark (Asynchronous)"; + public string Name { get; } = "Stream Raw Network Writer Benchmark (Asynchronous)"; private static Task ServerTask(CancellationToken cancellationToken) { @@ -81,7 +81,7 @@ namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); rawSocket.Bind(ClientEndPoint); - using StreamNetworkWriter writer = new StreamNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); + 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); diff --git a/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterSyncBenchmark.cs b/NetSharp/NetSharpExamples/Benchmarks/Stream Network Connection Benchmarks/StreamNetworkWriterSyncBenchmark.cs @@ -9,7 +9,7 @@ using System.Threading.Tasks; namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks { - public class StreamNetworkWriterSyncBenchmark : INetSharpExample, INetSharpBenchmark + public class StreamNetworkWriterSyncBenchmark : INetSharpBenchmark { private const int PacketSize = 8192, PacketCount = 1_000_000; @@ -21,7 +21,7 @@ namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks public static readonly ManualResetEventSlim ServerReadyEvent = new ManualResetEventSlim(); /// <inheritdoc /> - public string Name { get; } = "Stream Network Writer Benchmark (Synchronous)"; + public string Name { get; } = "Stream Raw Network Writer Benchmark (Synchronous)"; private static Task ServerTask(CancellationToken cancellationToken) { @@ -81,7 +81,7 @@ namespace NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); rawSocket.Bind(ClientEndPoint); - using StreamNetworkWriter writer = new StreamNetworkWriter(ref rawSocket, defaultRemoteEndPoint, PacketSize); + 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); diff --git a/NetSharp/NetSharpExamples/Examples/Connection Examples/DefaultConnectionExample.cs b/NetSharp/NetSharpExamples/Examples/Connection Examples/DefaultConnectionExample.cs @@ -0,0 +1,31 @@ +using System.Net; +using System.Threading; +using System.Threading.Tasks; +using NetSharp; +using NetSharp.Packets; + +namespace NetSharpExamples.Examples.Connection_Examples +{ + public class DefaultConnectionExample : INetSharpExample + { + /// <inheritdoc /> + public string Name { get; } = "Connection Instantiation Code Example"; + + /// <inheritdoc /> + public Task RunAsync() + { + using Connection streamConnection = ConnectionBuilder + .WithStreamTransport() + .WithInterNetwork(new IPEndPoint(IPAddress.Any, 0)) + .WithSettings(NetworkPacket.TotalSize, 1000, 0) + .BuildDefault(); + + using Connection datagramConnection = ConnectionBuilder + .WithDatagramTransport() + .WithDefaultSettings(1024) + .BuildDefault(); + + return Task.CompletedTask; + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Examples/Datagram Network Connection Examples/DatagramNetworkReaderExample.cs b/NetSharp/NetSharpExamples/Examples/Datagram Network Connection Examples/DatagramNetworkReaderExample.cs @@ -37,7 +37,7 @@ namespace NetSharpExamples.Examples.Datagram_Network_Connection_Examples Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp); rawSocket.Bind(ServerEndPoint); - using DatagramNetworkReader reader = new DatagramNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize, 100); + using RawDatagramNetworkReader reader = new RawDatagramNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize, 100); reader.Start(ExpectedClientCount); Console.WriteLine($"Started datagram server at {ServerEndPoint}! Enter any key to stop the server..."); diff --git a/NetSharp/NetSharpExamples/Examples/Datagram Network Connection Examples/DatagramNetworkWriterAsyncExample.cs b/NetSharp/NetSharpExamples/Examples/Datagram Network Connection Examples/DatagramNetworkWriterAsyncExample.cs @@ -28,7 +28,7 @@ namespace NetSharpExamples.Examples.Datagram_Network_Connection_Examples Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp); rawSocket.Bind(ClientEndPoint); - using DatagramNetworkWriter writer = new DatagramNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); + using RawDatagramNetworkWriter writer = new RawDatagramNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); byte[] transmissionBuffer = new byte[PacketSize]; diff --git a/NetSharp/NetSharpExamples/Examples/Datagram Network Connection Examples/DatagramNetworkWriterSyncExample.cs b/NetSharp/NetSharpExamples/Examples/Datagram Network Connection Examples/DatagramNetworkWriterSyncExample.cs @@ -28,7 +28,7 @@ namespace NetSharpExamples.Examples.Datagram_Network_Connection_Examples Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp); rawSocket.Bind(ClientEndPoint); - using DatagramNetworkWriter writer = new DatagramNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); + using RawDatagramNetworkWriter writer = new RawDatagramNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); byte[] transmissionBuffer = new byte[PacketSize]; diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkReaderExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkReaderExample.cs @@ -38,7 +38,7 @@ namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples rawSocket.Bind(ServerEndPoint); rawSocket.Listen(ExpectedClientCount); - using StreamNetworkReader reader = new StreamNetworkReader(ref rawSocket, RequestHandler, defaultEndPoint, PacketSize, 100); + 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..."); diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterAsyncExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterAsyncExample.cs @@ -28,7 +28,7 @@ namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); rawSocket.Bind(ClientEndPoint); - using StreamNetworkWriter writer = new StreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); + using RawStreamNetworkWriter writer = new RawStreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); await writer.ConnectAsync(ServerEndPoint); byte[] transmissionBuffer = new byte[PacketSize]; diff --git a/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterSyncExample.cs b/NetSharp/NetSharpExamples/Examples/Stream Network Connection Examples/StreamNetworkWriterSyncExample.cs @@ -28,7 +28,7 @@ namespace NetSharpExamples.Examples.Stream_Network_Connection_Examples Socket rawSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); rawSocket.Bind(ClientEndPoint); - using StreamNetworkWriter writer = new StreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); + using RawStreamNetworkWriter writer = new RawStreamNetworkWriter(ref rawSocket, defaultEndPoint, PacketSize); writer.Connect(ServerEndPoint); byte[] transmissionBuffer = new byte[PacketSize]; diff --git a/NetSharp/NetSharpExamples/NetSharpExamples.xml b/NetSharp/NetSharpExamples/NetSharpExamples.xml @@ -40,6 +40,12 @@ <member name="M:NetSharpExamples.Benchmarks.Stream_Network_Connection_Benchmarks.StreamNetworkWriterSyncBenchmark.RunAsync"> <inheritdoc /> </member> + <member name="P:NetSharpExamples.Examples.Connection_Examples.DefaultConnectionExample.Name"> + <inheritdoc /> + </member> + <member name="M:NetSharpExamples.Examples.Connection_Examples.DefaultConnectionExample.RunAsync"> + <inheritdoc /> + </member> <member name="P:NetSharpExamples.Examples.Datagram_Network_Connection_Examples.DatagramNetworkReaderExample.Name"> <inheritdoc /> </member>