NetSharp

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

commit a8d2d1c315a0f0643d0316e9161e7a6dccc653a2
parent 8c0479a3b8fa784ba1988e8948a359b927664bf3
Author: Mikolaj Lenczewski <33129490+EnderRifter@users.noreply.github.com>
Date:   Fri, 28 Feb 2020 22:11:03 +0000

Temporary UDP echoing and basice packet pipeline functionality included in Connection!

Diffstat:
MNetSharp/NetSharp/Connection.cs | 296++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
MNetSharp/NetSharp/ConnectionBuilder.cs | 65++++++-----------------------------------------------------------
MNetSharp/NetSharp/Extensions/ConnectionBuilderExtensions.cs | 10----------
MNetSharp/NetSharp/Extensions/ConnectionExtensions.cs | 4++--
MNetSharp/NetSharp/NetSharp.xml | 88+++++++++++++++++++++++++++++++++++++++----------------------------------------
MNetSharp/NetSharp/Packets/NetworkPacket.cs | 33++++++++++++++++++++++++---------
MNetSharp/NetSharp/Sockets/SocketReader.cs | 30+++++++++++++++++++++---------
MNetSharp/NetSharp/Sockets/SocketWriter.cs | 42+++++++++++++++++++++++++++---------------
ANetSharp/NetSharp/Utils/CryptographyHelpers.cs | 97+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
MNetSharp/NetSharpExamples/Program.cs | 44+++++++++++++++++++++++++-------------------
10 files changed, 500 insertions(+), 209 deletions(-)

diff --git a/NetSharp/NetSharp/Connection.cs b/NetSharp/NetSharp/Connection.cs @@ -4,7 +4,9 @@ using System.IO; using System.Net; using System.Net.Sockets; using System.Threading; +using System.Threading.Channels; using System.Threading.Tasks; +using NetSharp.Deprecated; using NetSharp.Logging; using NetSharp.Packets; using NetSharp.Pipelines; @@ -13,15 +15,41 @@ using NetSharp.Utils; namespace NetSharp { + /// <summary> + /// Encapsulates a connection capable of receiving packets and responding to them with registered packet handlers. + /// </summary> public class Connection : IDisposable { + /// <summary> + /// Represents any remote endpoint for datagram operations. + /// </summary> + private static readonly EndPoint AnyRemoteEndPoint = new IPEndPoint(IPAddress.Any, 0); + private readonly SocketAcceptor acceptor; - private readonly ConcurrentDictionary<EndPoint, Connection> connections; - private readonly CancellationTokenSource connectionShutdownTokenSource; + + private readonly ConcurrentDictionary<EndPoint, DatagramClientArgs> datagramConnections; + private readonly Socket datagramSocket; + private readonly Channel<(EndPoint origin, Memory<byte> packet)> incomingPacketChannel; + + /// <summary> + /// Pipeline to convert incoming byte buffers to <see cref="NetworkPacket"/> instances. + /// </summary> private readonly PacketPipeline<Memory<byte>, Memory<byte>, NetworkPacket> incomingPacketPipeline; + private readonly SocketReader listener; + private readonly Channel<(EndPoint destination, NetworkPacket packet)> outgoingPacketChannel; + + /// <summary> + /// Pipeline to convert outgoing <see cref="NetworkPacket"/> instances to a byte buffer for sending. + /// </summary> private readonly PacketPipeline<NetworkPacket, Memory<byte>, Memory<byte>> outgoingPacketPipeline; - private readonly Socket socket; + + private readonly Channel<(EndPoint origin, IRequestPacket request)> requestChannel; + private readonly CancellationTokenSource serverShutdownTokenSource; + private readonly ConcurrentDictionary<EndPoint, StreamClientArgs> streamConnections; + + private readonly Socket streamSocket; + private readonly SocketWriter transmitter; /// <summary> @@ -32,30 +60,169 @@ namespace NetSharp Dispose(false); } - private void AcceptorWork(object shutdownToken) + private async Task AcceptorWork(object shutdownToken) { CancellationToken cancellationToken = (CancellationToken)shutdownToken; + logger.LogMessage("Started stream acceptor task."); + + async Task StreamListenerWork(object clientSocketObj) + { + Socket clientSocket = (Socket)clientSocketObj; + logger.LogMessage("Started stream listener task."); + + while (!cancellationToken.IsCancellationRequested) + { + // TODO: implement receive buffer pooling + byte[] receiveBuffer = new byte[NetworkPacket.PacketSize]; + Memory<byte> receiveBufferMemory = new Memory<byte>(receiveBuffer); + + TransmissionResult result = + await listener.ReceiveAsync(clientSocket, SocketFlags.None, receiveBufferMemory, cancellationToken); + + await incomingPacketChannel.Writer.WriteAsync((result.RemoteEndPoint, receiveBufferMemory), + cancellationToken); + } + + logger.LogMessage("Stopped stream listener task."); + } + while (!cancellationToken.IsCancellationRequested) { + Socket clientSocket = await acceptor.AcceptAsync(streamSocket, cancellationToken); + + if (!streamConnections.ContainsKey(clientSocket.RemoteEndPoint)) + { + streamConnections[clientSocket.RemoteEndPoint] = new StreamClientArgs(); + } + + await Task.Factory.StartNew(StreamListenerWork, clientSocket, ServerShutdownToken, + TaskCreationOptions.LongRunning, TaskScheduler.Default); } + + logger.LogMessage("Stopped stream acceptor task."); } - private void ListenerWork(object shutdownToken) + private async Task DatagramListenerWork(object shutdownToken) { CancellationToken cancellationToken = (CancellationToken)shutdownToken; + logger.LogMessage("Started datagram listener task."); + while (!cancellationToken.IsCancellationRequested) { + // TODO: implement receive buffer pooling + byte[] receiveBuffer = new byte[NetworkPacket.PacketSize]; + Memory<byte> receiveBufferMemory = new Memory<byte>(receiveBuffer); + + TransmissionResult result = + await listener.ReceiveFromAsync(datagramSocket, AnyRemoteEndPoint, SocketFlags.None, + receiveBufferMemory, cancellationToken); + + if (!datagramConnections.ContainsKey(result.RemoteEndPoint)) + { + datagramConnections[result.RemoteEndPoint] = new DatagramClientArgs(); + } + + await incomingPacketChannel.Writer.WriteAsync((result.RemoteEndPoint, receiveBufferMemory), cancellationToken); } + + logger.LogMessage("Stopped datagram listener task."); } - private void PacketPipelineWork(object shutdownToken) + private async Task IncomingPacketHandlerWork(object shutdownToken) { CancellationToken cancellationToken = (CancellationToken)shutdownToken; + logger.LogMessage("Started incoming packet handler task."); + while (!cancellationToken.IsCancellationRequested) { + (EndPoint origin, Memory<byte> packet) = + await incomingPacketChannel.Reader.ReadAsync(cancellationToken); + + NetworkPacket deserialisedRequest = incomingPacketPipeline.ProcessPacket(packet); + + // TODO implement deserialisation according to registered packet deserialisers + + // TODO: write to requestChannel, not to outgoingPacketChannel + await outgoingPacketChannel.Writer.WriteAsync((origin, deserialisedRequest), cancellationToken); + } + + logger.LogMessage("Stopped incoming packet handler task."); + } + + private async Task OutgoingPacketHandlerWork(object shutdownToken) + { + CancellationToken cancellationToken = (CancellationToken)shutdownToken; + + logger.LogMessage("Started outgoing packet handler task."); + + while (!cancellationToken.IsCancellationRequested) + { + (EndPoint destination, NetworkPacket packet) = + await outgoingPacketChannel.Reader.ReadAsync(cancellationToken); + + Memory<byte> serialisedResponse = outgoingPacketPipeline.ProcessPacket(packet); + + if (datagramConnections.ContainsKey(destination)) + { + await transmitter.SendToAsync(datagramSocket, destination, + SocketFlags.None, serialisedResponse, cancellationToken); + } + else if (streamConnections.ContainsKey(destination)) + { + StreamClientArgs streamClientArgs = streamConnections[destination]; + + await transmitter.SendAsync(streamClientArgs.ClientSocket, + SocketFlags.None, serialisedResponse, cancellationToken); + } + else + { + logger.LogWarning($"Dropping packet destined for unknown destination; {destination}"); + } + } + + logger.LogMessage("Stopped outgoing packet handler task."); + } + + private async Task RequestHandlerInvocationWork(object shutdownToken) + { + CancellationToken cancellationToken = (CancellationToken)shutdownToken; + + logger.LogMessage("Started request handler invocation task."); + + while (!cancellationToken.IsCancellationRequested) + { + (EndPoint origin, IRequestPacket request) = await requestChannel.Reader.ReadAsync(cancellationToken); + + // TODO: implement proper request handling, and conversion to IResponsePacket<IRequestPacke> + + NetworkPacket serialisedResponsePacket = new NetworkPacket(); + + await outgoingPacketChannel.Writer.WriteAsync((origin, serialisedResponsePacket), cancellationToken); + } + + logger.LogMessage("Stopped request handler invocation task."); + } + + private readonly struct DatagramClientArgs + { + public readonly EndPoint ClientEndPoint; + + public DatagramClientArgs(EndPoint clientEndPoint) + { + ClientEndPoint = clientEndPoint; + } + } + + private readonly struct StreamClientArgs + { + public readonly Socket ClientSocket; + + public StreamClientArgs(Socket clientSocket) + { + ClientSocket = clientSocket; } } @@ -65,9 +232,9 @@ namespace NetSharp protected readonly object loggerLockObject = new object(); /// <summary> - /// Cancellation token which allows observing the shutdown of the server. It is set when <see cref="Shutdown()"/> is called. + /// Cancellation token which allows observing the shutdown of the server. It is set when <see cref="ShutdownServer"/> is called. /// </summary> - protected readonly CancellationToken ShutdownToken; + protected readonly CancellationToken ServerShutdownToken; /// <summary> /// A logger object allowing for writing debug messages to an output stream. @@ -82,33 +249,64 @@ namespace NetSharp { if (disposing) { - socket.Dispose(); + streamSocket.Dispose(); } } - internal Connection(AddressFamily addressFamily, SocketType socketType, ProtocolType protocolType, - PacketPipeline<Memory<byte>, Memory<byte>, NetworkPacket> incomingPacketPipeline, - PacketPipeline<NetworkPacket, Memory<byte>, Memory<byte>> outgoingPacketPipeline, - int objectPoolSize = 10, bool preallocateBuffers = false, Stream? loggingStream = default, - LogLevel minimumLoggedSeverity = LogLevel.Info) + internal Connection( + PacketPipeline<Memory<byte>, Memory<byte>, NetworkPacket> incomingPacketPipeline, + PacketPipeline<NetworkPacket, Memory<byte>, Memory<byte>> outgoingPacketPipeline, + int objectPoolSize = 10, bool preallocateBuffers = false, Stream? loggingStream = default, + LogLevel minimumLoggedSeverity = LogLevel.Info) { - connectionShutdownTokenSource = new CancellationTokenSource(); - ShutdownToken = connectionShutdownTokenSource.Token; + serverShutdownTokenSource = new CancellationTokenSource(); + ServerShutdownToken = serverShutdownTokenSource.Token; - socket = new Socket(addressFamily, socketType, protocolType); + streamSocket = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + datagramSocket = new Socket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp); acceptor = new SocketAcceptor(objectPoolSize); listener = new SocketReader(NetworkPacket.PacketSize, objectPoolSize, preallocateBuffers); transmitter = new SocketWriter(NetworkPacket.PacketSize, objectPoolSize, preallocateBuffers); - connections = new ConcurrentDictionary<EndPoint, Connection>(); + streamConnections = new ConcurrentDictionary<EndPoint, StreamClientArgs>(); + datagramConnections = new ConcurrentDictionary<EndPoint, DatagramClientArgs>(); this.incomingPacketPipeline = incomingPacketPipeline; + BoundedChannelOptions incomingChannelOptions = new BoundedChannelOptions(MaximumPacketBacklog) + { + FullMode = BoundedChannelFullMode.DropOldest, + SingleReader = true, + SingleWriter = false, + }; + incomingPacketChannel = Channel.CreateBounded<(EndPoint origin, Memory<byte> packet)>(incomingChannelOptions); + + BoundedChannelOptions requestChannelOptions = new BoundedChannelOptions(MaximumPacketBacklog) + { + FullMode = BoundedChannelFullMode.DropOldest, + SingleReader = true, + SingleWriter = true, + }; + requestChannel = Channel.CreateBounded<(EndPoint origin, IRequestPacket request)>(requestChannelOptions); + this.outgoingPacketPipeline = outgoingPacketPipeline; + BoundedChannelOptions outgoingChannelOptions = new BoundedChannelOptions(MaximumPacketBacklog) + { + FullMode = BoundedChannelFullMode.DropOldest, + SingleReader = true, + SingleWriter = true, + }; + outgoingPacketChannel = Channel.CreateBounded<(EndPoint destination, NetworkPacket packet)>(outgoingChannelOptions); logger = new Logger(loggingStream ?? Stream.Null, minimumLoggedSeverity); } + /// <summary> + /// The maximum number of packets that will be stored before older packets start to be dropped. + /// </summary> + /// TODO change this to a configurable builder option + public const int MaximumPacketBacklog = 64; + /// <inheritdoc /> public void Dispose() { @@ -120,58 +318,72 @@ namespace NetSharp { using CancellationTokenSource timeoutCancellationTokenSource = new CancellationTokenSource(timeout); using CancellationTokenSource cts = - CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ShutdownToken); + CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ServerShutdownToken); - return listener.ReceiveAsync(socket, flags, inputBuffer, cts.Token); + return listener.ReceiveAsync(streamSocket, flags, inputBuffer, cts.Token); } public Task<TransmissionResult> ReceiveFromAsync(EndPoint remoteEndPoint, Memory<byte> inputBuffer, SocketFlags flags, TimeSpan timeout) { using CancellationTokenSource timeoutCancellationTokenSource = new CancellationTokenSource(timeout); using CancellationTokenSource cts = - CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ShutdownToken); + CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ServerShutdownToken); - return listener.ReceiveFromAsync(socket, remoteEndPoint, flags, inputBuffer, cts.Token); + return listener.ReceiveFromAsync(datagramSocket, remoteEndPoint, flags, inputBuffer, cts.Token); } /// <summary> - /// Makes the connection listen for incoming client request packets, and handle them according to registered packet handler delegates. - /// This work can be cancelled by calling <see cref="Shutdown()"/>. + /// Makes the connection listen for incoming request packets, and handle them according to registered packet handler delegates. + /// This work can be cancelled by calling <see cref="ShutdownServer"/>. /// </summary> /// <returns>The task representing the connection work.</returns> - public Task RunAsync() + public Task RunServerAsync() { + logger.LogMessage("Starting stream acceptor task..."); Task acceptorThread = - Task.Factory.StartNew(AcceptorWork, ShutdownToken, ShutdownToken, TaskCreationOptions.LongRunning, + Task.Factory.StartNew(AcceptorWork, ServerShutdownToken, ServerShutdownToken, TaskCreationOptions.LongRunning, + TaskScheduler.Default).Result; + + logger.LogMessage("Starting datagram listener task..."); + Task datagramListenerThread = + Task.Factory.StartNew(DatagramListenerWork, ServerShutdownToken, ServerShutdownToken, TaskCreationOptions.LongRunning, TaskScheduler.Default); - Task listenerThread = - Task.Factory.StartNew(ListenerWork, ShutdownToken, ShutdownToken, TaskCreationOptions.LongRunning, + logger.LogMessage("Starting incoming packet handler task..."); + Task incomingPacketHandlerThread = + Task.Factory.StartNew(IncomingPacketHandlerWork, ServerShutdownToken, ServerShutdownToken, TaskCreationOptions.LongRunning, TaskScheduler.Default); - Task packetPipelineThread = - Task.Factory.StartNew(PacketPipelineWork, ShutdownToken, ShutdownToken, TaskCreationOptions.LongRunning, + logger.LogMessage("Starting request packet invocation task..."); + Task requestHandlerInvocationThread = + Task.Factory.StartNew(RequestHandlerInvocationWork, ServerShutdownToken, ServerShutdownToken, TaskCreationOptions.LongRunning, + TaskScheduler.Default).Result; + + logger.LogMessage("Starting outgoing packet handler task..."); + Task outgoingPacketHandlerThread = + Task.Factory.StartNew(OutgoingPacketHandlerWork, ServerShutdownToken, ServerShutdownToken, TaskCreationOptions.LongRunning, TaskScheduler.Default); - return Task.WhenAll(acceptorThread, listenerThread, packetPipelineThread); + return Task.WhenAll(acceptorThread, datagramListenerThread, + incomingPacketHandlerThread, requestHandlerInvocationThread, outgoingPacketHandlerThread); } - public Task<int> SendAsync(Memory<byte> outputBuffer, SocketFlags flags, TimeSpan timeout) + public ValueTask<int> SendAsync(Memory<byte> outputBuffer, SocketFlags flags, TimeSpan timeout) { using CancellationTokenSource timeoutCancellationTokenSource = new CancellationTokenSource(timeout); using CancellationTokenSource cts = - CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ShutdownToken); + CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ServerShutdownToken); - return transmitter.SendAsync(socket, flags, outputBuffer, cts.Token); + return transmitter.SendAsync(streamSocket, flags, outputBuffer, cts.Token); } - public Task<int> SendToAsync(EndPoint remoteEndPoint, Memory<byte> outputBuffer, SocketFlags flags, TimeSpan timeout) + public ValueTask<int> SendToAsync(EndPoint remoteEndPoint, Memory<byte> outputBuffer, SocketFlags flags, TimeSpan timeout) { using CancellationTokenSource timeoutCancellationTokenSource = new CancellationTokenSource(timeout); using CancellationTokenSource cts = - CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ShutdownToken); + CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ServerShutdownToken); - return transmitter.SendToAsync(socket, remoteEndPoint, flags, outputBuffer, cts.Token); + return transmitter.SendToAsync(datagramSocket, remoteEndPoint, flags, outputBuffer, cts.Token); } /// <summary> @@ -191,9 +403,10 @@ namespace NetSharp /// <summary> /// Shuts down the connection, and releases managed and unmanaged resources. /// </summary> - public void Shutdown() + public void ShutdownServer() { - connectionShutdownTokenSource.Cancel(); + logger.LogMessage("Signalling shutdown to all client connection handlers..."); + serverShutdownTokenSource.Cancel(); } /// <summary> @@ -217,13 +430,14 @@ namespace NetSharp { using CancellationTokenSource timeoutCancellationTokenSource = new CancellationTokenSource(timeout); using CancellationTokenSource cts = - CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ShutdownToken); + CancellationTokenSource.CreateLinkedTokenSource(timeoutCancellationTokenSource.Token, ServerShutdownToken); try { return await Task.Run(() => { - socket.Bind(localEndPoint); + streamSocket.Bind(localEndPoint); + datagramSocket.Bind(localEndPoint); return true; }, cts.Token); diff --git a/NetSharp/NetSharp/ConnectionBuilder.cs b/NetSharp/NetSharp/ConnectionBuilder.cs @@ -2,6 +2,7 @@ using System.Collections.Generic; using System.IO; using System.Net.Sockets; +using System.Security.Cryptography; using NetSharp.Logging; using NetSharp.Packets; using NetSharp.Pipelines; @@ -13,8 +14,11 @@ namespace NetSharp /// </summary> public sealed class ConnectionBuilder { - private static readonly LoggingSettings DefaultLoggingSettings = new LoggingSettings(Stream.Null, LogLevel.Warn); - private static readonly PoolingSettings DefaultPoolingSettings = new PoolingSettings(10, false); + private static readonly LoggingSettings DefaultLoggingSettings = + new LoggingSettings(Stream.Null, LogLevel.Warn); + + private static readonly PoolingSettings DefaultPoolingSettings = + new PoolingSettings(10, false); private readonly List<Func<Memory<byte>, Memory<byte>>> incomingPipelineStages = new List<Func<Memory<byte>, Memory<byte>>>(); @@ -24,7 +28,6 @@ namespace NetSharp private LoggingSettings? loggingSettings; private PoolingSettings? poolingSettings; - private SocketSettings? socketSettings; /// <summary> /// The number of stages in the currently configured incoming packet pipeline. @@ -46,16 +49,8 @@ namespace NetSharp /// Returns a new <see cref="Connection"/> instance with the current configuration. /// </summary> /// <returns>The configured <see cref="Connection"/> instance.</returns> - /// <exception cref="ArgumentNullException"> - /// Thrown when <see cref="WithSocket"/> has not been called. - /// </exception> public Connection Build() { - if (socketSettings == null) - { - throw new ArgumentNullException(nameof(socketSettings), $"{nameof(WithSocket)} has not been called."); - } - PacketPipelineBuilder<Memory<byte>, Memory<byte>, NetworkPacket> incomingPipelineBuilder = new PacketPipelineBuilder<Memory<byte>, Memory<byte>, NetworkPacket>(); @@ -79,9 +74,6 @@ namespace NetSharp outgoingPipelineBuilder.WithOutputStage(memory => memory); Connection connection = new Connection( - socketSettings.Value.AddressFamily, - socketSettings.Value.SocketType, - socketSettings.Value.ProtocolType, incomingPipelineBuilder.Build(), outgoingPipelineBuilder.Build(), poolingSettings?.ObjectPoolSize ?? DefaultPoolingSettings.ObjectPoolSize, @@ -143,17 +135,6 @@ namespace NetSharp } /// <summary> - /// Sets the socket settings for the currently configured connection. - /// </summary> - /// <param name="settings">The socket settings to use.</param> - /// <returns>The builder instance for further configuration.</returns> - public ConnectionBuilder WithSocket(SocketSettings settings) - { - socketSettings = settings; - return this; - } - - /// <summary> /// Holds settings for configuring a connection's logging. /// </summary> public readonly struct LoggingSettings @@ -206,39 +187,5 @@ namespace NetSharp PreallocateBuffers = preallocateBuffers; } } - - /// <summary> - /// Holds settings fo configuring a connection's underlying socket. - /// </summary> - public readonly struct SocketSettings - { - /// <summary> - /// The address family for the socket underlying the connection. - /// </summary> - public readonly AddressFamily AddressFamily; - - /// <summary> - /// The socket type for the socket underlying the connection. - /// </summary> - public readonly ProtocolType ProtocolType; - - /// <summary> - /// The protocol type for the socket underlying the connection. - /// </summary> - public readonly SocketType SocketType; - - /// <summary> - /// Initialises a new instance of the <see cref="SocketSettings"/> struct. - /// </summary> - /// <param name="addressFamily">The address family for the underlying socket.</param> - /// <param name="socketType">The socket type for the underlying socket.</param> - /// <param name="protocolType">The protocol type for the underlying socket.</param> - public SocketSettings(AddressFamily addressFamily, SocketType socketType, ProtocolType protocolType) - { - AddressFamily = addressFamily; - SocketType = socketType; - ProtocolType = protocolType; - } - } } } \ No newline at end of file diff --git a/NetSharp/NetSharp/Extensions/ConnectionBuilderExtensions.cs b/NetSharp/NetSharp/Extensions/ConnectionBuilderExtensions.cs @@ -25,15 +25,5 @@ namespace NetSharp.Extensions public static ConnectionBuilder WithPooling(this ConnectionBuilder instance, int poolSize, bool preallocateBuffers) => instance.WithPooling(new ConnectionBuilder.PoolingSettings(poolSize, preallocateBuffers)); - - public static ConnectionBuilder WithSocket(this ConnectionBuilder instance, - AddressFamily addressFamily, SocketType socketType, ProtocolType protocolType) - => instance.WithSocket(new ConnectionBuilder.SocketSettings(addressFamily, socketType, protocolType)); - - public static ConnectionBuilder WithTcp(this ConnectionBuilder instance) - => instance.WithSocket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - - public static ConnectionBuilder WithUdp(this ConnectionBuilder instance) - => instance.WithSocket(AddressFamily.InterNetwork, SocketType.Dgram, ProtocolType.Udp); } } \ No newline at end of file diff --git a/NetSharp/NetSharp/Extensions/ConnectionExtensions.cs b/NetSharp/NetSharp/Extensions/ConnectionExtensions.cs @@ -20,11 +20,11 @@ namespace NetSharp.Extensions EndPoint remoteEndPoint, Memory<byte> inputBuffer, SocketFlags flags) => instance.ReceiveFromAsync(remoteEndPoint, inputBuffer, flags, Timeout.InfiniteTimeSpan); - public static Task<int> SendAsync(this Connection instance, + public static ValueTask<int> SendAsync(this Connection instance, Memory<byte> outputBuffer, SocketFlags flags) => instance.SendAsync(outputBuffer, flags, Timeout.InfiniteTimeSpan); - public static Task<int> SendToAsync(this Connection instance, + public static ValueTask<int> SendToAsync(this Connection instance, EndPoint remoteEndPoint, Memory<byte> outputBuffer, SocketFlags flags) => instance.SendToAsync(remoteEndPoint, outputBuffer, flags, Timeout.InfiniteTimeSpan); diff --git a/NetSharp/NetSharp/NetSharp.xml b/NetSharp/NetSharp/NetSharp.xml @@ -4,6 +4,26 @@ <name>NetSharp</name> </assembly> <members> + <member name="T:NetSharp.Connection"> + <summary> + Encapsulates a connection capable of receiving packets and responding to them with registered packet handlers. + </summary> + </member> + <member name="F:NetSharp.Connection.AnyRemoteEndPoint"> + <summary> + Represents any remote endpoint for datagram operations. + </summary> + </member> + <member name="F:NetSharp.Connection.incomingPacketPipeline"> + <summary> + Pipeline to convert incoming byte buffers to <see cref="T:NetSharp.Packets.NetworkPacket"/> instances. + </summary> + </member> + <member name="F:NetSharp.Connection.outgoingPacketPipeline"> + <summary> + Pipeline to convert outgoing <see cref="T:NetSharp.Packets.NetworkPacket"/> instances to a byte buffer for sending. + </summary> + </member> <member name="M:NetSharp.Connection.Finalize"> <summary> Destroys a <see cref="T:NetSharp.Connection"/> class instance, freeing all managed resources. @@ -14,9 +34,9 @@ Lock synchronisation object for the <see cref="F:NetSharp.Connection.logger"/> variable. </summary> </member> - <member name="F:NetSharp.Connection.ShutdownToken"> + <member name="F:NetSharp.Connection.ServerShutdownToken"> <summary> - Cancellation token which allows observing the shutdown of the server. It is set when <see cref="M:NetSharp.Connection.Shutdown"/> is called. + Cancellation token which allows observing the shutdown of the server. It is set when <see cref="M:NetSharp.Connection.ShutdownServer"/> is called. </summary> </member> <member name="F:NetSharp.Connection.logger"> @@ -30,13 +50,19 @@ </summary> <param name="disposing">Whether this method is called by <see cref="M:NetSharp.Connection.Dispose"/> or by the finaliser.</param> </member> + <member name="F:NetSharp.Connection.MaximumPacketBacklog"> + <summary> + The maximum number of packets that will be stored before older packets start to be dropped. + </summary> + TODO change this to a configurable builder option + </member> <member name="M:NetSharp.Connection.Dispose"> <inheritdoc /> </member> - <member name="M:NetSharp.Connection.RunAsync"> + <member name="M:NetSharp.Connection.RunServerAsync"> <summary> - Makes the connection listen for incoming client request packets, and handle them according to registered packet handler delegates. - This work can be cancelled by calling <see cref="M:NetSharp.Connection.Shutdown"/>. + Makes the connection listen for incoming request packets, and handle them according to registered packet handler delegates. + This work can be cancelled by calling <see cref="M:NetSharp.Connection.ShutdownServer"/>. </summary> <returns>The task representing the connection work.</returns> </member> @@ -48,7 +74,7 @@ <param name="loggingStream">The stream to which messages will be logged.</param> <param name="minimumLoggedSeverity">The minimum severity a message must be to be logged.</param> </member> - <member name="M:NetSharp.Connection.Shutdown"> + <member name="M:NetSharp.Connection.ShutdownServer"> <summary> Shuts down the connection, and releases managed and unmanaged resources. </summary> @@ -91,9 +117,6 @@ Returns a new <see cref="T:NetSharp.Connection"/> instance with the current configuration. </summary> <returns>The configured <see cref="T:NetSharp.Connection"/> instance.</returns> - <exception cref="T:System.ArgumentNullException"> - Thrown when <see cref="M:NetSharp.ConnectionBuilder.WithSocket(NetSharp.ConnectionBuilder.SocketSettings)"/> has not been called. - </exception> </member> <member name="M:NetSharp.ConnectionBuilder.WithIncomingPipelineStage(System.Func{System.Memory{System.Byte},System.Memory{System.Byte}}@,System.Int32)"> <summary> @@ -129,13 +152,6 @@ <param name="settings">The pooling settings to use.</param> <returns>The builder instance for further configuration.</returns> </member> - <member name="M:NetSharp.ConnectionBuilder.WithSocket(NetSharp.ConnectionBuilder.SocketSettings)"> - <summary> - Sets the socket settings for the currently configured connection. - </summary> - <param name="settings">The socket settings to use.</param> - <returns>The builder instance for further configuration.</returns> - </member> <member name="T:NetSharp.ConnectionBuilder.LoggingSettings"> <summary> Holds settings for configuring a connection's logging. @@ -180,34 +196,6 @@ <param name="poolSize">The number of objects that will be held in the object pools.</param> <param name="preallocateBuffers">Whether the buffers for receiving messages should be preallocated.</param> </member> - <member name="T:NetSharp.ConnectionBuilder.SocketSettings"> - <summary> - Holds settings fo configuring a connection's underlying socket. - </summary> - </member> - <member name="F:NetSharp.ConnectionBuilder.SocketSettings.AddressFamily"> - <summary> - The address family for the socket underlying the connection. - </summary> - </member> - <member name="F:NetSharp.ConnectionBuilder.SocketSettings.ProtocolType"> - <summary> - The socket type for the socket underlying the connection. - </summary> - </member> - <member name="F:NetSharp.ConnectionBuilder.SocketSettings.SocketType"> - <summary> - The protocol type for the socket underlying the connection. - </summary> - </member> - <member name="M:NetSharp.ConnectionBuilder.SocketSettings.#ctor(System.Net.Sockets.AddressFamily,System.Net.Sockets.SocketType,System.Net.Sockets.ProtocolType)"> - <summary> - Initialises a new instance of the <see cref="T:NetSharp.ConnectionBuilder.SocketSettings"/> struct. - </summary> - <param name="addressFamily">The address family for the underlying socket.</param> - <param name="socketType">The socket type for the underlying socket.</param> - <param name="protocolType">The protocol type for the underlying socket.</param> - </member> <member name="T:NetSharp.Deprecated.Client"> <summary> Provides methods for connecting to and talking with a <see cref="T:NetSharp.Deprecated.IServer"/> instance. @@ -1705,11 +1693,21 @@ </member> <member name="M:NetSharp.Packets.NetworkPacket.Serialise(NetSharp.Packets.NetworkPacket)"> <summary> - Serialises the given packet instance into a single byte buffer. + Serialises the given packet instance to a new byte buffer. </summary> <param name="instance">The packet instance to serialise.</param> <returns>The byte buffer that represents the packet instance.</returns> </member> + <member name="M:NetSharp.Packets.NetworkPacket.SerialiseToBuffer(System.Memory{System.Byte},NetSharp.Packets.NetworkPacket)"> + <summary> + Serialises the given packet instance into the given byte buffer. + </summary> + <param name="buffer"> + The buffer to which the instance should be serialised. Must be at least of size <see cref="F:NetSharp.Packets.NetworkPacket.PacketSize"/>. + </param> + <param name="instance">The packet instance to serialise.</param> + <exception cref="T:System.ArgumentException">Thrown if the given buffer is too small.</exception> + </member> <member name="F:NetSharp.Packets.NetworkPacketFooter.Size"> <summary> The number of bytes taken up by a packet footer. diff --git a/NetSharp/NetSharp/Packets/NetworkPacket.cs b/NetSharp/NetSharp/Packets/NetworkPacket.cs @@ -81,31 +81,46 @@ namespace NetSharp.Packets Span<byte> serialisedPacketFooter = buffer.Slice(HeaderSize + DataSegmentSize, FooterSize).Span; NetworkPacketFooter footer = NetworkPacketFooter.Deserialise(serialisedPacketFooter); - // data segment - Memory<byte> packetData = buffer.Slice(HeaderSize, DataSegmentSize); + Memory<byte> serialisedInstanceData = buffer.Slice(HeaderSize, DataSegmentSize); - return new NetworkPacket(packetData, header, footer); + return new NetworkPacket(serialisedInstanceData, header, footer); } /// <summary> - /// Serialises the given packet instance into a single byte buffer. + /// Serialises the given packet instance to a new byte buffer. /// </summary> /// <param name="instance">The packet instance to serialise.</param> /// <returns>The byte buffer that represents the packet instance.</returns> public static Memory<byte> Serialise(NetworkPacket instance) { byte[] buffer = new byte[PacketSize]; + SerialiseToBuffer(buffer, instance); + return buffer; + } - Span<byte> serialisedPacketHeader = new Span<byte>(buffer, 0, HeaderSize); + /// <summary> + /// Serialises the given packet instance into the given byte buffer. + /// </summary> + /// <param name="buffer"> + /// The buffer to which the instance should be serialised. Must be at least of size <see cref="PacketSize"/>. + /// </param> + /// <param name="instance">The packet instance to serialise.</param> + /// <exception cref="ArgumentException">Thrown if the given buffer is too small.</exception> + public static void SerialiseToBuffer(Memory<byte> buffer, NetworkPacket instance) + { + if (buffer.Length < PacketSize) + { + throw new ArgumentException("Given buffer is too small to serialise the packet instance into.", nameof(buffer)); + } + + Span<byte> serialisedPacketHeader = buffer.Slice(0, HeaderSize).Span; NetworkPacketHeader.Serialise(serialisedPacketHeader, instance.Header); - Span<byte> serialisedPacketFooter = new Span<byte>(buffer, HeaderSize + DataSegmentSize, FooterSize); + Span<byte> serialisedPacketFooter = buffer.Slice(HeaderSize + DataSegmentSize, FooterSize).Span; NetworkPacketFooter.Serialise(serialisedPacketFooter, instance.Footer); - Memory<byte> serialisedInstanceData = new Memory<byte>(buffer, HeaderSize, instance.Header.DataLength); + Memory<byte> serialisedInstanceData = buffer.Slice(HeaderSize, DataSegmentSize); instance.DataBuffer.CopyTo(serialisedInstanceData); - - return buffer; } } diff --git a/NetSharp/NetSharp/Sockets/SocketReader.cs b/NetSharp/NetSharp/Sockets/SocketReader.cs @@ -18,6 +18,8 @@ namespace NetSharp.Sockets private readonly int PacketBufferLength; private readonly ObjectPool<SocketAsyncEventArgs> receiveAsyncEventArgsPool; private readonly ArrayPool<byte> receiveBufferPool; + private readonly ObjectPool<SocketAsyncEventArgs> receiveFromAsyncEventArgsPool; + private readonly ArrayPool<byte> receiveFromBufferPool; private void HandleIOCompleted(object? sender, SocketAsyncEventArgs args) { @@ -76,8 +78,8 @@ namespace NetSharp.Sockets } } - receiveBufferPool.Return(asyncReceiveFromToken.RentedBuffer, true); - receiveAsyncEventArgsPool.Return(args); + receiveFromBufferPool.Return(asyncReceiveFromToken.RentedBuffer, true); + receiveFromAsyncEventArgsPool.Return(args); break; @@ -116,11 +118,21 @@ namespace NetSharp.Sockets new DefaultObjectPool<SocketAsyncEventArgs>(new DefaultPooledObjectPolicy<SocketAsyncEventArgs>(), maxPooledObjects); + receiveFromBufferPool = ArrayPool<byte>.Create(packetBufferLength, maxPooledObjects); + + receiveFromAsyncEventArgsPool = + new DefaultObjectPool<SocketAsyncEventArgs>(new DefaultPooledObjectPolicy<SocketAsyncEventArgs>(), + maxPooledObjects); + for (int i = 0; i < maxPooledObjects; i++) { - SocketAsyncEventArgs args = new SocketAsyncEventArgs(); - args.Completed += HandleIOCompleted; - receiveAsyncEventArgsPool.Return(args); + SocketAsyncEventArgs receiveArgs = new SocketAsyncEventArgs(); + receiveArgs.Completed += HandleIOCompleted; + receiveAsyncEventArgsPool.Return(receiveArgs); + + SocketAsyncEventArgs receiveFromArgs = new SocketAsyncEventArgs(); + receiveFromArgs.Completed += HandleIOCompleted; + receiveFromAsyncEventArgsPool.Return(receiveFromArgs); } } @@ -190,10 +202,10 @@ namespace NetSharp.Sockets { TaskCompletionSource<TransmissionResult> tcs = new TaskCompletionSource<TransmissionResult>(); - byte[] rentedReceiveFromBuffer = receiveBufferPool.Rent(PacketBufferLength); + byte[] rentedReceiveFromBuffer = receiveFromBufferPool.Rent(PacketBufferLength); Memory<byte> rentedReceiveFromBufferMemory = new Memory<byte>(rentedReceiveFromBuffer); - SocketAsyncEventArgs args = receiveAsyncEventArgsPool.Get(); + SocketAsyncEventArgs args = receiveFromAsyncEventArgsPool.Get(); args.SetBuffer(rentedReceiveFromBufferMemory); args.SocketFlags = socketFlags; args.RemoteEndPoint = remoteEndPoint; @@ -224,8 +236,8 @@ namespace NetSharp.Sockets TransmissionResult result = new TransmissionResult(args); - receiveBufferPool.Return(rentedReceiveFromBuffer, true); - receiveAsyncEventArgsPool.Return(args); + receiveFromBufferPool.Return(rentedReceiveFromBuffer, true); + receiveFromAsyncEventArgsPool.Return(args); return Task.FromResult(result); } diff --git a/NetSharp/NetSharp/Sockets/SocketWriter.cs b/NetSharp/NetSharp/Sockets/SocketWriter.cs @@ -17,6 +17,8 @@ namespace NetSharp.Sockets private readonly int PacketBufferLength; private readonly ObjectPool<SocketAsyncEventArgs> sendAsyncEventArgsPool; private readonly ArrayPool<byte> sendBufferPool; + private readonly ObjectPool<SocketAsyncEventArgs> sendToAsyncEventArgsPool; + private readonly ArrayPool<byte> sendToBufferPool; private void HandleIOCompleted(object? sender, SocketAsyncEventArgs args) { @@ -66,8 +68,8 @@ namespace NetSharp.Sockets } } - sendBufferPool.Return(asyncSendToToken.RentedBuffer, true); - sendAsyncEventArgsPool.Return(args); + sendToBufferPool.Return(asyncSendToToken.RentedBuffer, true); + sendToAsyncEventArgsPool.Return(args); break; default: @@ -103,11 +105,21 @@ namespace NetSharp.Sockets new DefaultObjectPool<SocketAsyncEventArgs>(new DefaultPooledObjectPolicy<SocketAsyncEventArgs>(), maxPooledObjects); + sendToBufferPool = ArrayPool<byte>.Create(packetBufferLength, maxPooledObjects); + + sendToAsyncEventArgsPool = + new DefaultObjectPool<SocketAsyncEventArgs>(new DefaultPooledObjectPolicy<SocketAsyncEventArgs>(), + maxPooledObjects); + for (int i = 0; i < maxPooledObjects; i++) { - SocketAsyncEventArgs args = new SocketAsyncEventArgs(); - args.Completed += HandleIOCompleted; - sendAsyncEventArgsPool.Return(args); + SocketAsyncEventArgs sendArgs = new SocketAsyncEventArgs(); + sendArgs.Completed += HandleIOCompleted; + sendAsyncEventArgsPool.Return(sendArgs); + + SocketAsyncEventArgs sendToArgs = new SocketAsyncEventArgs(); + sendToArgs.Completed += HandleIOCompleted; + sendToAsyncEventArgsPool.Return(sendArgs); } } @@ -119,7 +131,7 @@ namespace NetSharp.Sockets /// <param name="outputBuffer">The data buffer which should be sent.</param> /// <param name="cancellationToken">The cancellation token to observe for the operation.</param> /// <returns>The number of bytes of data which were written to the remote connection.</returns> - public Task<int> SendAsync(Socket socket, SocketFlags socketFlags, Memory<byte> outputBuffer, + public ValueTask<int> SendAsync(Socket socket, SocketFlags socketFlags, Memory<byte> outputBuffer, CancellationToken cancellationToken = default) { TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); @@ -153,14 +165,14 @@ namespace NetSharp.Sockets */ // if the send operation doesn't complete synchronously, return the awaitable task - if (socket.SendAsync(args)) return tcs.Task; + if (socket.SendAsync(args)) return new ValueTask<int>(tcs.Task); int result = args.BytesTransferred; sendBufferPool.Return(rentedSendBuffer, true); sendAsyncEventArgsPool.Return(args); - return Task.FromResult(result); + return new ValueTask<int>(result); } /// <summary> @@ -172,17 +184,17 @@ namespace NetSharp.Sockets /// <param name="outputBuffer">The data buffer which should be sent.</param> /// <param name="cancellationToken">The cancellation token to observe for the operation.</param> /// <returns>The number of bytes of data which were written to the remote endpoint.</returns> - public Task<int> SendToAsync(Socket socket, EndPoint remoteEndPoint, SocketFlags socketFlags, + public ValueTask<int> SendToAsync(Socket socket, EndPoint remoteEndPoint, SocketFlags socketFlags, Memory<byte> outputBuffer, CancellationToken cancellationToken = default) { TaskCompletionSource<int> tcs = new TaskCompletionSource<int>(); - byte[] rentedSendToBuffer = sendBufferPool.Rent(PacketBufferLength); + byte[] rentedSendToBuffer = sendToBufferPool.Rent(PacketBufferLength); Memory<byte> rentedSendToBufferMemory = new Memory<byte>(rentedSendToBuffer); outputBuffer.CopyTo(rentedSendToBufferMemory); - SocketAsyncEventArgs args = sendAsyncEventArgsPool.Get(); + SocketAsyncEventArgs args = sendToAsyncEventArgsPool.Get(); args.SetBuffer(rentedSendToBufferMemory); args.SocketFlags = socketFlags; args.RemoteEndPoint = remoteEndPoint; @@ -207,14 +219,14 @@ namespace NetSharp.Sockets */ // if the send operation doesn't complete synchronously, return the awaitable task - if (socket.SendToAsync(args)) return tcs.Task; + if (socket.SendToAsync(args)) return new ValueTask<int>(tcs.Task); int result = args.BytesTransferred; - sendBufferPool.Return(rentedSendToBuffer, true); - sendAsyncEventArgsPool.Return(args); + sendToBufferPool.Return(rentedSendToBuffer, true); + sendToAsyncEventArgsPool.Return(args); - return Task.FromResult(result); + return new ValueTask<int>(result); } } } \ No newline at end of file diff --git a/NetSharp/NetSharp/Utils/CryptographyHelpers.cs b/NetSharp/NetSharp/Utils/CryptographyHelpers.cs @@ -0,0 +1,96 @@ +using System; +using System.IO; +using System.Security.Cryptography; +using System.Text; + +namespace NetSharp.Utils +{ + internal static class CryptographyHelpers + { + #region Settings + + private static string _hash = "SHA1"; + private static int _iterations = 2; + private static int _keySize = 256; + private static string _salt = "aselrias38490a32"; // Random + private static string _vector = "8947az34awl34kjq"; // Random + + #endregion Settings + + public static string Decrypt(byte[] value, string password) + { + return Decrypt<AesManaged>(value, password); + } + + public static string Decrypt<T>(byte[] value, string password) where T : SymmetricAlgorithm, new() + { + byte[] vectorBytes = Encoding.ASCII.GetBytes(_vector); // GetBytes<ASCIIEncoding>(_vector); + byte[] saltBytes = Encoding.ASCII.GetBytes(_salt); // GetBytes<ASCIIEncoding>(_salt); + byte[] valueBytes = value; + + byte[] decrypted; + int decryptedByteCount = 0; + + using (T cipher = new T()) + { + PasswordDeriveBytes _passwordBytes = new PasswordDeriveBytes(password, saltBytes, _hash, _iterations); + byte[] keyBytes = _passwordBytes.GetBytes(_keySize / 8); + + cipher.Mode = CipherMode.CBC; + + try + { + using (ICryptoTransform decryptor = cipher.CreateDecryptor(keyBytes, vectorBytes)) + { + using MemoryStream from = new MemoryStream(valueBytes); + using CryptoStream reader = new CryptoStream(@from, decryptor, CryptoStreamMode.Read); + + decrypted = new byte[valueBytes.Length]; + decryptedByteCount = reader.Read(decrypted, 0, decrypted.Length); + } + } + catch (Exception ex) + { + return String.Empty; + } + + cipher.Clear(); + } + return Encoding.UTF8.GetString(decrypted, 0, decryptedByteCount); + } + + public static byte[] Encrypt(string value, string password) + { + return Encrypt<AesManaged>(value, password); + } + + public static byte[] Encrypt<T>(string value, string password) where T : SymmetricAlgorithm, new() + { + byte[] vectorBytes = Encoding.ASCII.GetBytes(_vector); // GetBytes<ASCIIEncoding>(_vector); + byte[] saltBytes = Encoding.ASCII.GetBytes(_salt); // GetBytes<ASCIIEncoding>(_salt); + byte[] valueBytes = Encoding.UTF8.GetBytes(value); // GetBytes<UTF8Encoding>(value); + + byte[] encrypted; + using (T cipher = new T()) + { + PasswordDeriveBytes _passwordBytes = + new PasswordDeriveBytes(password, saltBytes, _hash, _iterations); + byte[] keyBytes = _passwordBytes.GetBytes(_keySize / 8); + + cipher.Mode = CipherMode.CBC; + + using (ICryptoTransform encryptor = cipher.CreateEncryptor(keyBytes, vectorBytes)) + { + using MemoryStream to = new MemoryStream(); + using CryptoStream writer = new CryptoStream(to, encryptor, CryptoStreamMode.Write); + + writer.Write(valueBytes, 0, valueBytes.Length); + writer.FlushFinalBlock(); + encrypted = to.ToArray(); + } + cipher.Clear(); + } + return encrypted; + } + } +} +\ No newline at end of file diff --git a/NetSharp/NetSharpExamples/Program.cs b/NetSharp/NetSharpExamples/Program.cs @@ -8,6 +8,7 @@ using System.Threading; using System.Threading.Tasks; using NetSharp; using NetSharp.Extensions; +using NetSharp.Logging; using NetSharp.Packets; using NetSharp.Utils; @@ -35,16 +36,11 @@ namespace NetSharpExamples { TimeSpan socketTimeout = TimeSpan.FromSeconds(newtorkTimeout); - const int clientCount = 10; + const int clientCount = 1; const long sentPacketCount = 10_000; EndPoint serverEndPoint = new IPEndPoint(serverAddress, serverPort); - ConnectionBuilder clientBuilder = new ConnectionBuilder().WithUdp(); - - Connection ClientFactory() - { - return clientBuilder.Build(); - } + ConnectionBuilder clientBuilder = new ConnectionBuilder(); Console.WriteLine($"Testing client connections..."); @@ -54,14 +50,19 @@ namespace NetSharpExamples { Console.WriteLine($"Starting client {clientId}"); - using Connection client = ClientFactory(); + using Connection client = clientBuilder.Build(); client.TryBind(new IPEndPoint(IPAddress.Any, 0)); + //client.SetLoggingStream(Console.OpenStandardOutput()); //TimeSpan timeout = TimeSpan.FromMilliseconds(100); Stopwatch stopwatch = new Stopwatch(); byte[] message = Encoding.UTF8.GetBytes("Hello World!"); Memory<byte> messageBuffer = new Memory<byte>(message); + NetworkPacket requestPacket = new NetworkPacket(messageBuffer, message.Length, 1, NetworkErrorCode.Ok, false); + byte[] requestPacketBuffer = new byte[NetworkPacket.PacketSize]; + NetworkPacket.SerialiseToBuffer(requestPacketBuffer, requestPacket); + byte[] response = new byte[NetworkPacket.PacketSize]; Memory<byte> responseBuffer = new Memory<byte>(response); @@ -71,18 +72,20 @@ namespace NetSharpExamples try { stopwatch.Start(); - int sentBytes = await client.SendToAsync(serverEndPoint, messageBuffer, SocketFlags.None); + //int sentBytes = await client.SendToAsync(serverEndPoint, requestPacketBuffer, SocketFlags.None); + int sentBytes = await client.SendAsync(requestPacketBuffer, SocketFlags.None); stopwatch.Stop(); Interlocked.Increment(ref sentPackets); - //Console.WriteLine($"[Client {clientId}] Sent {sentBytes} bytes to {serverEndPoint}"); + Console.WriteLine($"[Client {clientId}] Sent {sentBytes} bytes to {serverEndPoint}"); stopwatch.Start(); - TransmissionResult result = await client.ReceiveFromAsync(serverEndPoint, responseBuffer, SocketFlags.None); + //TransmissionResult result = await client.ReceiveFromAsync(serverEndPoint, responseBuffer, SocketFlags.None); + TransmissionResult result = await client.ReceiveAsync(responseBuffer, SocketFlags.None); stopwatch.Stop(); Interlocked.Increment(ref receivedPackets); - //Console.WriteLine($"[Client {clientId}] Received {result.Count} bytes from {result.RemoteEndPoint}"); + Console.WriteLine($"[Client {clientId}] Received {result.Count} bytes from {result.RemoteEndPoint}"); //await Task.Delay(10); } @@ -110,31 +113,34 @@ namespace NetSharpExamples await using Stream serverOutputStream = File.OpenWrite(serverLogFile); EndPoint serverEndPoint = new IPEndPoint(serverAddress, serverPort); - ConnectionBuilder serverBuilder = new ConnectionBuilder().WithUdp(); + ConnectionBuilder serverBuilder = new ConnectionBuilder(); - using Connection server = serverBuilder.Build(); + using Connection server = serverBuilder.WithLogging(Console.OpenStandardOutput(), LogLevel.Info).Build(); server.TryBind(serverEndPoint); - //server.SetLoggingStream(Console.OpenStandardOutput(), LogLevel.Info); + //server.SetLoggingStream(Console.OpenStandardOutput()); //server.ChangeLoggingStream(serverOutputStream, LogLevel.Error); Console.WriteLine("Starting server..."); EndPoint nullEndPoint = new IPEndPoint(IPAddress.Any, 0); + await server.RunServerAsync(); + + /* byte[] request = new byte[NetworkPacket.PacketSize]; Memory<byte> requestBuffer = new Memory<byte>(request); while (true) { - TransmissionResult result = - await server.ReceiveFromAsync(nullEndPoint, requestBuffer, SocketFlags.None); + TransmissionResult result = await server.ReceiveFromAsync(nullEndPoint, requestBuffer, SocketFlags.None); - //Console.WriteLine($"[Server] Received {result.Count} bytes from {result.RemoteEndPoint}"); + Console.WriteLine($"[Server] Received {result.Count} bytes from {result.RemoteEndPoint}"); int sentBytes = await server.SendToAsync(result.RemoteEndPoint, requestBuffer, SocketFlags.None); - //Console.WriteLine($"[Server] Sent {sentBytes} bytes tp {result.RemoteEndPoint}"); + Console.WriteLine($"[Server] Sent {sentBytes} bytes tp {result.RemoteEndPoint}"); } + */ Console.WriteLine("Server stopped");