From 5e9c2cbb7b49cb2246793c8ebb4395986a1bce3c Mon Sep 17 00:00:00 2001 From: Mikhail Mitrofanov Date: Sun, 16 Feb 2025 11:45:08 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9D=D0=BE=D0=B2=D0=B0=D1=8F=20=D1=81=D0=B8?= =?UTF-8?q?=D1=81=D1=82=D0=B5=D0=BC=D0=B0=20=D1=81=D0=B5=D1=82=D0=B5=D0=B2?= =?UTF-8?q?=D0=BE=D0=B3=D0=BE=20=D0=B2=D0=B7=D0=B0=D0=B8=D0=BC=D0=BE=D0=B4?= =?UTF-8?q?=D0=B5=D0=B9=D1=81=D1=82=D0=B2=D0=B8=D1=8F=20=D0=BE=D0=B1=D0=BB?= =?UTF-8?q?=D0=BE=D0=B6=D0=B5=D0=BD=D0=B0=20=D1=82=D0=B5=D1=81=D1=82=D0=B0?= =?UTF-8?q?=D0=BC=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- mROA.Test/NextGenTest.cs | 71 +++++++++++++++++++++++++ mROA/Abstract/IInteractionModule.cs | 12 +++++ mROA/NextGenerationInteractionModule.cs | 54 +++++++++++++++++++ 3 files changed, 137 insertions(+) create mode 100644 mROA.Test/NextGenTest.cs create mode 100644 mROA/NextGenerationInteractionModule.cs diff --git a/mROA.Test/NextGenTest.cs b/mROA.Test/NextGenTest.cs new file mode 100644 index 0000000..f57dc5c --- /dev/null +++ b/mROA.Test/NextGenTest.cs @@ -0,0 +1,71 @@ +using System.Net; +using System.Net.Sockets; +using System.Text; +using mROA.Implementation; + +namespace mROA.Test; + +public class NextGenTest +{ + private TcpListener _listener; + private NextGenerationInteractionModule _interactionModuleA; + private NextGenerationInteractionModule _interactionModuleB; + private Guid[] guids = [Guid.NewGuid(), Guid.NewGuid(), Guid.NewGuid()]; + + [SetUp] + public void Setup() + { + _listener = new TcpListener(IPAddress.Loopback, 4567); + _interactionModuleA = new NextGenerationInteractionModule(); + _interactionModuleA.Inject(new JsonSerializationToolkit()); + _interactionModuleB = new NextGenerationInteractionModule(); + _interactionModuleB.Inject(new JsonSerializationToolkit()); + + } + + [Test] + public void MultithreadedTest() + { + + Task.Run(() => + { + _listener.Start(); + _interactionModuleB.BaseStream = _listener.AcceptTcpClient().GetStream(); + + foreach (var guid in guids) + { + _interactionModuleB.PostMessage(new NetworkMessage { Id = guid, Data = "Hello user"u8.ToArray() }); + } + }); + + var client = new TcpClient(); + client.Connect(IPAddress.Loopback, 4567); + _interactionModuleA.BaseStream = client.GetStream(); + + var tasks = guids.Select(ReadStream); + + Task.WaitAll(tasks.ToArray()); + Assert.Pass(); + } + + private async Task ReadStream(Guid current) + { + var msg = await _interactionModuleA.GetNextMessageReceiving(); + Console.WriteLine( + $"{Environment.CurrentManagedThreadId} Received message: {Encoding.Default.GetString(new JsonSerializationToolkit().Serialize(msg))}"); + while (msg.Id != current) + { + msg = await _interactionModuleA.GetNextMessageReceiving(); + Console.WriteLine( + $"{Environment.CurrentManagedThreadId} Received message: {Encoding.Default.GetString(new JsonSerializationToolkit().Serialize(msg))}"); + } + + Console.WriteLine($"{Environment.CurrentManagedThreadId} Good message received"); + } + + [TearDown] + public void TearDown() + { + _listener.Dispose(); + } +} \ No newline at end of file diff --git a/mROA/Abstract/IInteractionModule.cs b/mROA/Abstract/IInteractionModule.cs index b0400a6..5d49bf1 100644 --- a/mROA/Abstract/IInteractionModule.cs +++ b/mROA/Abstract/IInteractionModule.cs @@ -1,3 +1,5 @@ +using mROA.Implementation; + namespace mROA.Abstract; public interface IInteractionModule : IInjectableModule @@ -11,4 +13,14 @@ public interface IInteractionModule : IInjectableModule public Task ReceiveMessage(); public void PostMessage(byte[] message); } +} + +public interface INextGenerationInteractionModule : IInjectableModule +{ + int ConntectionId { get; } + public Stream? BaseStream { get; set; } + Task GetNextMessageReceiving(); + void PostMessage(NetworkMessage message); + NetworkMessage[] UnhandledMessages { get; } + void HandleMessage(NetworkMessage msg); } \ No newline at end of file diff --git a/mROA/NextGenerationInteractionModule.cs b/mROA/NextGenerationInteractionModule.cs new file mode 100644 index 0000000..07404da --- /dev/null +++ b/mROA/NextGenerationInteractionModule.cs @@ -0,0 +1,54 @@ +using mROA.Abstract; +using mROA.Implementation; + +namespace mROA; + +public class NextGenerationInteractionModule : INextGenerationInteractionModule +{ + private ISerializationToolkit _serialization; + public int ConntectionId { get; set; } + public Stream BaseStream { get; set; } + private Task? _currentReceiving; + private const int BufferSize = ushort.MaxValue; + private byte[] _buffer = new byte[BufferSize]; + private List _unhandledMessages = new(8); + + public void Inject(T dependency) + { + if (dependency is ISerializationToolkit toolkit) + _serialization = toolkit; + } + + + public Task GetNextMessageReceiving() + { + if (_currentReceiving is { IsCompleted: false }) + return _currentReceiving; + + _currentReceiving = GetNextMessage(); + return _currentReceiving; + } + + public void PostMessage(NetworkMessage message) + { + var rawMessage = _serialization.Serialize(message); + BaseStream.Write(BitConverter.GetBytes((ushort)rawMessage.Length), 0, sizeof(ushort)); + BaseStream.Write(rawMessage, 0, rawMessage.Length); + } + + public NetworkMessage[] UnhandledMessages => _unhandledMessages.ToArray(); + + public void HandleMessage(NetworkMessage msg) + { + _unhandledMessages.Remove(msg); + } + + private async Task GetNextMessage() + { + await BaseStream.ReadExactlyAsync(_buffer, 0, 2); + var len = BitConverter.ToUInt16(_buffer, 0); + await BaseStream.ReadExactlyAsync(_buffer, 0, len); + + return _serialization.Deserialize(_buffer[..len])!; + } +} \ No newline at end of file