From 71cc47056961eaaa3e41d7a64db1f133898bc0ca Mon Sep 17 00:00:00 2001 From: Alexandru Macocian Date: Fri, 10 Apr 2020 21:44:40 +0200 Subject: [PATCH] - server now reads continiously instead of polling - server increases reading delay when client is inactive - new procedure to check if the client is connected or not - ftphandler sets affinity to connected client --- MTSC-TestServer/Program.cs | 11 ++-- MTSC.UnitTests/E2ETests.cs | 69 ++++++++++++-------- MTSC.UnitTests/MultipartModule.cs | 20 ++++++ MTSC/MTSC.csproj | 6 +- MTSC/ServerSide/ClientData.cs | 4 +- MTSC/ServerSide/Handlers/FtpHandler.cs | 1 + MTSC/ServerSide/Server.cs | 88 ++++++++++++++++++-------- 7 files changed, 137 insertions(+), 62 deletions(-) create mode 100644 MTSC.UnitTests/MultipartModule.cs diff --git a/MTSC-TestServer/Program.cs b/MTSC-TestServer/Program.cs index 359a2c5..7e29b21 100644 --- a/MTSC-TestServer/Program.cs +++ b/MTSC-TestServer/Program.cs @@ -1,4 +1,6 @@ using MTSC.Common.Ftp.FtpModules; +using MTSC.Common.Http.RoutingModules; +using MTSC.Common.Http.ServerModules; using MTSC.Exceptions; using MTSC.Logging; using MTSC.ServerSide; @@ -13,19 +15,18 @@ namespace MTSC_TestServer { static void Main(string[] args) { - Server server = new Server(443); + Server server = new Server(800); RSACryptoServiceProvider rsa = new RSACryptoServiceProvider(1024); server - .WithCertificate(new X509Certificate2("powershellcert.pfx", "123")) + //.WithCertificate(new X509Certificate2("powershellcert.pfx", "123")) .AddLogger(new ConsoleLogger()) .AddLogger(new DebugConsoleLogger()) .AddServerUsageMonitor(new TickrateEnforcer() .SetTicksPerSecond(60) .SetSilent(true)) .AddExceptionHandler(new ExceptionConsoleLogger()) - //.AddHandler(new WebsocketHandler()) - //.AddHandler(new HttpHandler().AddHttpModule(new FileServerModule()) - // .AddHttpModule(new PostModule())) + .AddHandler(new HttpRoutingHandler() + .AddRoute(MTSC.Common.Http.HttpMessage.HttpMethods.Get, "hello", new Http200Module())) .AddHandler(new FtpHandler() .AddModule(new AuthenticationModule()) .AddModule(new SystModule()) diff --git a/MTSC.UnitTests/E2ETests.cs b/MTSC.UnitTests/E2ETests.cs index 4773e43..33598c1 100644 --- a/MTSC.UnitTests/E2ETests.cs +++ b/MTSC.UnitTests/E2ETests.cs @@ -54,6 +54,7 @@ namespace MTSC.UnitTests .AddRoute(HttpMessage.HttpMethods.Get, "echo", new EchoModule()) .AddRoute(HttpMessage.HttpMethods.Post, "echo", new EchoModule()) .AddRoute(HttpMessage.HttpMethods.Get, "long-running", new LongRunningModule()) + .AddRoute(HttpMessage.HttpMethods.Post, "multipart", new MultipartModule()) .WithFragmentsExpirationTime(TimeSpan.FromMilliseconds(3000)) .WithMaximumSize(250000)) .AddLogger(new ConsoleLogger()) @@ -79,7 +80,7 @@ namespace MTSC.UnitTests } var result = longRunningTask.Result; Assert.AreEqual(result.StatusCode, System.Net.HttpStatusCode.OK); - Assert.IsTrue(responses > 50); + Assert.IsTrue(responses > 5); } [TestMethod] @@ -91,6 +92,23 @@ namespace MTSC.UnitTests Assert.AreEqual(result.StatusCode, System.Net.HttpStatusCode.OK); } + [TestMethod] + public void MultipleRequestsShouldRespond() + { + HttpClient httpClient = new HttpClient(); + httpClient.BaseAddress = new Uri("https://localhost:800"); + for (int i = 0; i < 100; i++) + { + var result = httpClient.GetAsync("").GetAwaiter().GetResult(); + Assert.AreEqual(result.StatusCode, System.Net.HttpStatusCode.OK); + } + for (int i = 0; i < 100; i++) + { + var result = httpClient.GetAsync("echo").GetAwaiter().GetResult(); + Assert.AreEqual(result.StatusCode, System.Net.HttpStatusCode.OK); + } + } + [TestMethod] public void SendFragmentedHttpMessage() { @@ -221,10 +239,34 @@ namespace MTSC.UnitTests using(var sc = new ByteArrayContent(Encoding.UTF8.GetBytes(s))) { var response = httpClient.PostAsync("echo", sc).Result; - Assert.AreEqual(response.StatusCode, HttpStatusCode.OK); + Assert.AreEqual(HttpStatusCode.OK, response.StatusCode); Assert.AreEqual(response.Content.ReadAsStringAsync().Result, s); } } + + [TestMethod] + public void UploadFileShouldSucceed() + { + byte[] bytes = new byte[120000]; + + for(int i = 0; i < bytes.Length; i++) + { + bytes[i] = 43; + } + + using (var client = new HttpClient()) + { + using (var content = new MultipartFormDataContent("Upload----" + DateTime.Now.ToString())) + { + content.Add(new StreamContent(new MemoryStream(bytes)), "file", "upload.zip"); + + using (var message = client.PostAsync("https://localhost:800/multipart", content).GetAwaiter().GetResult()) + { + + } + } + } + } [TestMethod] public void GetWithQueryHttp() @@ -255,29 +297,6 @@ namespace MTSC.UnitTests Assert.AreEqual(resultString, "Hello world!"); } - [TestMethod] - public void HTTPStressTest() - { - var httpClient = new HttpClient(); - httpClient.BaseAddress = new Uri("https://localhost:800"); - for(int i = 0; i < stressIterations; i++) - { - var startTime = DateTime.Now; - var tasks = new Task[Environment.ProcessorCount]; - for(int j = 0; j < Environment.ProcessorCount; j++) - { - tasks[j] = httpClient.GetAsync(""); - } - Task.WaitAll(tasks); - var duration = DateTime.Now - startTime; - foreach(Task t in tasks) - { - Assert.AreEqual(t.Result.StatusCode, System.Net.HttpStatusCode.OK); - } - TestContext.WriteLine($"{i}: Processed {tasks.Length} requests in {duration.TotalMilliseconds} ms."); - } - } - [ClassCleanup] public static void CleanupServer() { diff --git a/MTSC.UnitTests/MultipartModule.cs b/MTSC.UnitTests/MultipartModule.cs new file mode 100644 index 0000000..6d92d3b --- /dev/null +++ b/MTSC.UnitTests/MultipartModule.cs @@ -0,0 +1,20 @@ +using MTSC.Common.Http; +using MTSC.Common.Http.RoutingModules; +using MTSC.ServerSide; +using System; +using System.Threading.Tasks; + +namespace MTSC.UnitTests +{ + public class MultipartModule : HttpRouteBase + { + public override Task HandleRequest(HttpRequest request, ClientData client, Server server) + { + if (request.Form.Count > 0) + { + return Task.FromResult(new HttpResponse { StatusCode = HttpMessage.StatusCodes.OK }); + } + return Task.FromResult(new HttpResponse { StatusCode = HttpMessage.StatusCodes.BadRequest }); + } + } +} diff --git a/MTSC/MTSC.csproj b/MTSC/MTSC.csproj index 123c936..d516d34 100644 --- a/MTSC/MTSC.csproj +++ b/MTSC/MTSC.csproj @@ -5,12 +5,12 @@ netcoreapp2.1;net48;netstandard2.0;netcoreapp3.0;netcoreapp3.1 - 2.6 + 2.7 Alexandru-Victor Macocian MTSC Modular TCP Server and Client - 0.2.6.0 - 0.2.6.0 + 0.2.7.0 + 0.2.7.0 true AnyCPU;x64 https://github.com/AlexMacocian/MTSC diff --git a/MTSC/ServerSide/ClientData.cs b/MTSC/ServerSide/ClientData.cs index 93859fb..c3d6d92 100644 --- a/MTSC/ServerSide/ClientData.cs +++ b/MTSC/ServerSide/ClientData.cs @@ -27,13 +27,13 @@ namespace MTSC.ServerSide /// public IHandler Affinity { get; private set; } - IConsumerQueue IQueueHolder.ConsumerQueue { get => messageQueue; } - public bool ToBeRemoved { get; set; } = false; public SslStream SslStream { get; set; } = null; public SafeNetworkStream SafeNetworkStream { get; } public ResourceDictionary Resources { get; set; } = new ResourceDictionary(); + IConsumerQueue IQueueHolder.ConsumerQueue => messageQueue; + public ClientData(TcpClient client) { this.TcpClient = client; diff --git a/MTSC/ServerSide/Handlers/FtpHandler.cs b/MTSC/ServerSide/Handlers/FtpHandler.cs index 9388daa..3352b28 100644 --- a/MTSC/ServerSide/Handlers/FtpHandler.cs +++ b/MTSC/ServerSide/Handlers/FtpHandler.cs @@ -58,6 +58,7 @@ namespace MTSC.ServerSide.Handlers { client.Resources.SetResource(State.Initialized); server.QueueMessage(client, Encoding.ASCII.GetBytes(welcomeMessage)); + client.SetAffinity(this); } }); return false; diff --git a/MTSC/ServerSide/Server.cs b/MTSC/ServerSide/Server.cs index a26c992..0660bb9 100644 --- a/MTSC/ServerSide/Server.cs +++ b/MTSC/ServerSide/Server.cs @@ -8,6 +8,7 @@ using MTSC.ServerSide.UsageMonitors; using System; using System.Collections.Concurrent; using System.Collections.Generic; +using System.IO; using System.Linq; using System.Net; using System.Net.Security; @@ -433,11 +434,6 @@ namespace MTSC.ServerSide } } - /* - * Gather all messages from clients and put them in a queue - */ - GatherReceivedMessages(); - /* * Call the scheduler to handle all received messages and distribute them to the handlers */ @@ -565,6 +561,18 @@ namespace MTSC.ServerSide { toRemove.Add(client); } + else + { + if (client.TcpClient.Client.Poll(0, SelectMode.SelectRead)) + { + byte[] buff = new byte[1]; + if (client.TcpClient.Client.Receive(buff, SocketFlags.Peek) == 0) + { + // Client disconnected + toRemove.Add(client); + } + } + } } foreach (ClientData client in toRemove) { @@ -610,36 +618,47 @@ namespace MTSC.ServerSide } } - private void GatherReceivedMessages() + private void WaitAndCollectMessages(ClientData client) { - foreach(var client in Clients) + Task.Run(async () => { - if (!client.TcpClient.Connected) + int messageCount = 0; + var buffer = new byte[8192]; + int increasingDelay = 100; + while (true) { - client.ToBeRemoved = true; - } - try - { - if (!client.ToBeRemoved && client.TcpClient.Available > 0) + if(client.TcpClient.Available == 0) { - Message message = CommunicationPrimitives.GetMessage(client); - (client as IActiveClient).UpdateLastReceivedMessage(); - LogDebug("Received message from " + client.TcpClient.Client.RemoteEndPoint.ToString() + - "\nMessage length: " + message.MessageLength); - (client as IQueueHolder).Enqueue(message); + await Task.Delay(Math.Min(increasingDelay++, 1000)); + continue; } - } - catch(Exception e) - { - foreach (IExceptionHandler exceptionHandler in exceptionHandlers) + increasingDelay = 100; + Stream stream; + if (client.SslStream != null) { - if (exceptionHandler.HandleException(e)) + stream = client.SslStream; + } + else + { + stream = client.TcpClient.GetStream(); + } + try + { + var byteCount = await stream.ReadAsync(buffer, 0, buffer.Length); + if (byteCount > 0) { - break; + Message message = new Message((uint)byteCount, buffer.Take(byteCount).ToArray()); + (client as IQueueHolder).Enqueue(message); + messageCount++; } } + catch (Exception) + { + client.ToBeRemoved = true; + return; + } } - } + }); } private void HandleClientMessages(ClientData client, IConsumerQueue messages) @@ -747,9 +766,24 @@ namespace MTSC.ServerSide this.EncryptionPolicy); client.SslStream = sslStream; - sslStream.AuthenticateAsServerAsync(this.certificate, this.RequestClientCertificate, this.SslProtocols, false).Wait(SslAuthenticationTimeout); + if(sslStream.AuthenticateAsServerAsync(this.certificate, this.RequestClientCertificate, this.SslProtocols, false).Wait(SslAuthenticationTimeout)) + { + /* + * Client authenticated in the alloted time + */ + this._ProducerClientQueue.Enqueue(client); + WaitAndCollectMessages(client); + } + else + { + client.Dispose(); + } + } + else + { + this._ProducerClientQueue.Enqueue(client); + WaitAndCollectMessages(client); } - this._ProducerClientQueue.Enqueue(client); } catch (Exception e) {