commit 18ba6e50c408176e86b5dea2e0278ecc86d6dc11
parent 049f3cb441d7921fbef7d8c46858bbb8c1df9a45
Author: Mikolaj Lenczewski <mikolaj.lenczewski308@gmail.com>
Date: Thu, 9 Apr 2020 16:07:41 +0100
Made datagram server handle requests, not clients, in a fire-and-forget manner
Diffstat:
4 files changed, 126 insertions(+), 8 deletions(-)
diff --git a/NetSharp/NetSharp/Sockets/Datagram/DatagramSocketServer.cs b/NetSharp/NetSharp/Sockets/Datagram/DatagramSocketServer.cs
@@ -75,7 +75,8 @@ namespace NetSharp.Sockets.Datagram
TransmissionResult sendResult =
await SocketAsyncOperations
- .SendToAsync(clientArgs, connection, clientEndPoint, SocketFlags.None, responseBufferMemory, cancellationToken);
+ .SendToAsync(clientArgs, connection, clientEndPoint, SocketFlags.None, responseBufferMemory,
+ cancellationToken);
#if DEBUG
lock (typeof(Console))
@@ -86,21 +87,109 @@ namespace NetSharp.Sockets.Datagram
#endif
}
}
- catch (OperationCanceledException) { }
+ catch (OperationCanceledException)
+ {
+ Console.WriteLine($"Client task for {clientArgs.RemoteEndPoint} cancelled!");
+ }
finally
{
DestroyConnectionArgs(clientArgs);
}
}
+ private readonly struct ClientRequest
+ {
+ public readonly NetworkPacket RequestPacket;
+
+ public readonly EndPoint ClientEndPoint;
+
+ public readonly CancellationToken CancellationToken;
+
+ public ClientRequest(in NetworkPacket requestPacket, in EndPoint clientEndPoint, in CancellationToken cancellationToken)
+ {
+ RequestPacket = requestPacket;
+
+ ClientEndPoint = clientEndPoint;
+
+ CancellationToken = cancellationToken;
+ }
+ }
+
+ private async Task HandleClientRequest(object clientRequestObj)
+ {
+ SocketAsyncEventArgs clientArgs = SocketArgsPool.Get();
+
+ byte[] responseBuffer = BufferPool.Rent(NetworkPacket.TotalSize);
+ Memory<byte> responseBufferMemory = new Memory<byte>(responseBuffer);
+
+ ClientRequest clientRequest = (ClientRequest) clientRequestObj;
+
+ NetworkPacket request = clientRequest.RequestPacket;
+ EndPoint remoteEndPoint = clientRequest.ClientEndPoint;
+ CancellationToken cancellationToken = clientRequest.CancellationToken;
+
+ // TODO implement actual request handling, besides just an echo
+ NetworkPacket response = request;
+
+ NetworkPacket.Serialise(response, responseBufferMemory);
+
+ TransmissionResult sendResult =
+ await SocketAsyncOperations
+ .SendToAsync(clientArgs, connection, remoteEndPoint, SocketFlags.None, responseBufferMemory, cancellationToken)
+ .ConfigureAwait(false);
+
+#if DEBUG
+ lock (typeof(Console))
+ {
+ Console.WriteLine($"[Server] Sent {sendResult.Count} bytes to {sendResult.RemoteEndPoint}");
+ Console.WriteLine($"[Server] >>>> {Encoding.UTF8.GetString(sendResult.Buffer.Span)}");
+ }
+#endif
+
+ BufferPool.Return(responseBuffer, true);
+
+ SocketArgsPool.Return(clientArgs);
+ }
+
public override async Task RunAsync(CancellationToken cancellationToken = default)
{
+ EndPoint remoteEndPoint = new IPEndPoint(IPAddress.Any, 0);
+ using SocketAsyncEventArgs remoteArgs = GenerateConnectionArgs(remoteEndPoint);
+
byte[] requestBuffer = new byte[NetworkPacket.TotalSize];
Memory<byte> requestBufferMemory = new Memory<byte>(requestBuffer);
+ while (!cancellationToken.IsCancellationRequested)
+ {
+ TransmissionResult receiveResult =
+ await SocketAsyncOperations
+ .ReceiveFromAsync(remoteArgs, connection, remoteEndPoint, SocketFlags.None, requestBufferMemory, cancellationToken)
+ .ConfigureAwait(false);
+
+ EndPoint clientEndPoint = receiveResult.RemoteEndPoint;
+
+#if DEBUG
+ lock (typeof(Console))
+ {
+ Console.WriteLine($"[Server] Received {receiveResult.Count} bytes from {receiveResult.RemoteEndPoint}");
+ Console.WriteLine($"[Server] <<<< {Encoding.UTF8.GetString(receiveResult.Buffer.Span)}");
+ }
+#endif
+
+ NetworkPacket requestPacket = NetworkPacket.Deserialise(requestBufferMemory);
+
+ ClientRequest request = new ClientRequest(in requestPacket, in clientEndPoint, in cancellationToken);
+
+ Task _ = HandleClientRequest(request);
+ }
+
+ /*
EndPoint remoteEndPoint = new IPEndPoint(IPAddress.Any, 0);
using SocketAsyncEventArgs remoteArgs = GenerateConnectionArgs(remoteEndPoint);
+ byte[] requestBuffer = new byte[NetworkPacket.TotalSize];
+ Memory<byte> requestBufferMemory = new Memory<byte>(requestBuffer);
+
while (!cancellationToken.IsCancellationRequested)
{
remoteArgs.RemoteEndPoint = remoteEndPoint;
@@ -133,10 +222,11 @@ namespace NetSharp.Sockets.Datagram
ConnectedClientHandlerTasks[clientEndPoint] = HandleClient(clientArgs, cancellationToken);
}
- NetworkPacket requestPacket = NetworkPacket.Deserialise(requestBufferMemory);
+ NetworkPacket request = NetworkPacket.Deserialise(requestBufferMemory);
- await connectedClientTokens[clientEndPoint].PacketWriter.WriteAsync(requestPacket, cancellationToken);
+ await connectedClientTokens[clientEndPoint].PacketWriter.WriteAsync(request, cancellationToken);
}
+ */
/*
while (true)
diff --git a/NetSharp/NetSharp/Sockets/SocketServer.cs b/NetSharp/NetSharp/Sockets/SocketServer.cs
@@ -6,6 +6,7 @@ using System.Net;
using System.Net.Sockets;
using System.Threading;
using System.Threading.Tasks;
+using Microsoft.Extensions.ObjectPool;
namespace NetSharp.Sockets
{
@@ -15,11 +16,15 @@ namespace NetSharp.Sockets
protected readonly ArrayPool<byte> BufferPool;
+ protected readonly ObjectPool<SocketAsyncEventArgs> SocketArgsPool;
+
protected SocketServer(in AddressFamily connectionAddressFamily, in SocketType connectionSocketType, in ProtocolType connectionProtocolType)
: base(in connectionAddressFamily, in connectionSocketType, in connectionProtocolType)
{
BufferPool = ArrayPool<byte>.Create(NetworkPacket.TotalSize, 10);
+ SocketArgsPool = new DefaultObjectPool<SocketAsyncEventArgs>(new PooledSocketAsyncEventArgsPolicy());
+
ConnectedClientHandlerTasks = new ConcurrentDictionary<EndPoint, Task>();
}
@@ -32,4 +37,21 @@ namespace NetSharp.Sockets
public abstract Task RunAsync(CancellationToken cancellationToken = default);
}
+
+ public class PooledSocketAsyncEventArgsPolicy : IPooledObjectPolicy<SocketAsyncEventArgs>
+ {
+ public SocketAsyncEventArgs Create()
+ {
+ SocketAsyncEventArgs args = new SocketAsyncEventArgs();
+
+ args.Completed += SocketAsyncOperations.HandleIoCompleted;
+
+ return args;
+ }
+
+ public bool Return(SocketAsyncEventArgs obj)
+ {
+ return true;
+ }
+ }
}
\ No newline at end of file
diff --git a/NetSharp/NetSharp/Sockets/Stream/StreamSocketServer.cs b/NetSharp/NetSharp/Sockets/Stream/StreamSocketServer.cs
@@ -75,7 +75,8 @@ namespace NetSharp.Sockets.Stream
{
TransmissionResult receiveResult =
await SocketAsyncOperations
- .ReceiveAsync(clientArgs, clientSocket, clientEndPoint, SocketFlags.None, requestBufferMemory, cancellationToken)
+ .ReceiveAsync(clientArgs, clientSocket, clientEndPoint, SocketFlags.None,
+ requestBufferMemory, cancellationToken)
.ConfigureAwait(false);
if (receiveResult.Count == 0)
@@ -99,7 +100,8 @@ namespace NetSharp.Sockets.Stream
TransmissionResult sendResult =
await SocketAsyncOperations
- .SendAsync(clientArgs, clientSocket, clientEndPoint, SocketFlags.None, responseBufferMemory, cancellationToken)
+ .SendAsync(clientArgs, clientSocket, clientEndPoint, SocketFlags.None, responseBufferMemory,
+ cancellationToken)
.ConfigureAwait(false);
#if DEBUG
@@ -111,7 +113,10 @@ namespace NetSharp.Sockets.Stream
#endif
}
}
- catch (OperationCanceledException) { }
+ catch (OperationCanceledException)
+ {
+ Console.WriteLine($"Client task for {clientArgs.RemoteEndPoint} cancelled!");
+ }
finally
{
DestroyConnectionArgs(clientArgs);
diff --git a/NetSharp/NetSharpExamples/Program.cs b/NetSharp/NetSharpExamples/Program.cs
@@ -113,9 +113,10 @@ namespace NetSharpExamples
}
#endif
+ EndPoint serverEndPoint = ServerEndPoint;
+
rttStopwatch.Start();
bandwidthStopwatch.Start();
- EndPoint serverEndPoint = ServerEndPoint;
#if TCP
int receiveResult = client.ReceiveBytes(responseBufferMemory);