From 635fcbf914c1a9c7b3f2a5377793a9b670522bce Mon Sep 17 00:00:00 2001 From: Mikhail Mitrofanov Date: Tue, 18 Feb 2025 14:56:31 +0300 Subject: [PATCH] =?UTF-8?q?=D0=BF=D0=BE=D0=BF=D1=8B=D1=82=D0=BA=D0=B0=20?= =?UTF-8?q?=D0=B1=D1=83=D1=84=D0=B5=D1=80=D0=B8=D0=B7=D0=B8=D1=80=D0=BE?= =?UTF-8?q?=D0=B2=D0=B0=D1=82=D1=8C=20=D1=81=D0=BE=D0=BE=D0=B1=D1=89=D0=B5?= =?UTF-8?q?=D0=BD=D0=B8=D1=8F=20=D0=BD=D0=B5=20=D1=83=D0=B2=D0=B5=D0=BD?= =?UTF-8?q?=D1=87=D0=B0=D0=BB=D0=B0=D1=81=D1=8C=20=D1=83=D1=81=D0=BF=D0=B5?= =?UTF-8?q?=D1=85=D0=BE=D0=BC,=20=D0=BD=D0=BE=20=D0=BC=D0=BE=D0=B6=D0=B5?= =?UTF-8?q?=D1=82=20=D0=B8=D1=81=D0=BF=D0=BE=D0=BB=D1=8C=D0=B7=D0=BE=D0=B2?= =?UTF-8?q?=D0=B0=D1=82=D1=8C=D1=81=D1=8F=20=D0=B2=20=D0=B1=D1=83=D0=B4?= =?UTF-8?q?=D1=83=D1=89=D0=B5=D0=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Example.Frontend/Program.cs | 9 ++-- mROA/Abstract/IInteractionModule.cs | 2 + .../Frontend/RequestExtractor.cs | 2 +- .../NextGenerationInteractionModule.cs | 46 ++++++++++++------- mROA/Implementation/RepresentationModule.cs | 23 ++++++++-- 5 files changed, 55 insertions(+), 27 deletions(-) diff --git a/Example.Frontend/Program.cs b/Example.Frontend/Program.cs index 10ae566..234cdaa 100644 --- a/Example.Frontend/Program.cs +++ b/Example.Frontend/Program.cs @@ -18,8 +18,6 @@ builder.Modules.Add(new StaticRepresentationModuleProducer()); builder.Modules.Add(new RequestExtractor()); builder.Modules.Add(new BasicExecutionModule()); builder.Modules.Add(new CoCodegenMethodRepository()); -// builder.Modules.Add(new StreamBasedVirtualBackendInteractionModule()); -// builder.Modules.Add(new JsonSerialisationModule()); builder.UseCollectableContextRepository(); builder.Build(); @@ -28,7 +26,7 @@ TransmissionConfig.RealContextRepository = builder.GetModule( TransmissionConfig.RemoteEndpointContextRepository = builder.GetModule(); builder.GetModule()!.Connect(); - // _ = builder.GetModule()!.StartExtraction(); +// _ = builder.GetModule()!.StartExtraction(); Console.WriteLine(TransmissionConfig.OwnershipRepository!.GetOwnershipId()); var context = builder.GetModule(); @@ -36,6 +34,7 @@ var factory = context.GetSingleObject(typeof(IPrinterFactory)) as IPrinterFactor //правильный порядок команд 8-5-10-7 var printer = factory.Create("Test"); +Console.WriteLine("Printer created"); Thread.Sleep(100); var name = printer.Value.GetName(); @@ -43,13 +42,13 @@ Console.WriteLine("Printer name : {0}", name); Thread.Sleep(100); factory.Register(new SharedObject(new ClientBasedPrinter())); +Console.WriteLine("Registered printer"); Thread.Sleep(100); -Console.WriteLine("Registered printer"); var registred = factory.GetFirstPrinter(); -Thread.Sleep(100); Console.WriteLine("First printer"); +Thread.Sleep(100); Console.WriteLine(registred.Value); Console.WriteLine("Collecting all printers"); diff --git a/mROA/Abstract/IInteractionModule.cs b/mROA/Abstract/IInteractionModule.cs index 96a91b3..e2697c9 100644 --- a/mROA/Abstract/IInteractionModule.cs +++ b/mROA/Abstract/IInteractionModule.cs @@ -8,4 +8,6 @@ public interface INextGenerationInteractionModule : IInjectableModule public Stream? BaseStream { get; set; } Task GetNextMessageReceiving(); Task PostMessage(NetworkMessage message); + void HandleMessage(NetworkMessage message); + NetworkMessage[] UnhandledMessages { get; } } \ No newline at end of file diff --git a/mROA/Implementation/Frontend/RequestExtractor.cs b/mROA/Implementation/Frontend/RequestExtractor.cs index e50d997..57afbcd 100644 --- a/mROA/Implementation/Frontend/RequestExtractor.cs +++ b/mROA/Implementation/Frontend/RequestExtractor.cs @@ -61,7 +61,7 @@ public class RequestExtractor : IRequestExtractor var request = _representationModule!.GetMessage(messageType: MessageType.CallRequest); - Console.WriteLine("Executing {0}", request.Id); + // Console.WriteLine("Executing {0}", request.Id); if (request.Parameter is not null) { diff --git a/mROA/Implementation/NextGenerationInteractionModule.cs b/mROA/Implementation/NextGenerationInteractionModule.cs index 918315b..6210c38 100644 --- a/mROA/Implementation/NextGenerationInteractionModule.cs +++ b/mROA/Implementation/NextGenerationInteractionModule.cs @@ -1,6 +1,4 @@ -using System.Security.Cryptography; -using System.Text; -using System.Text.Json; +using System.Text; using mROA.Abstract; namespace mROA.Implementation; @@ -15,7 +13,7 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule private readonly byte[] _buffer = new byte[BufferSize]; private List _messageBuffer = []; - + public void Inject(T dependency) { switch (dependency) @@ -32,27 +30,43 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule public Task GetNextMessageReceiving() { - Console.WriteLine(_currentReceiving?.Status); - if (_currentReceiving is { Status: TaskStatus.Running }) - return _currentReceiving; + // Console.WriteLine(_currentReceiving?.Status); - _currentReceiving = Task.Run(GetNextMessage); - return _currentReceiving; + if (_currentReceiving == null) + { + _currentReceiving = Task.Run(GetNextMessage); + return _currentReceiving; + } + lock (_currentReceiving) + { + if (_currentReceiving is { Status: TaskStatus.Running }) + return _currentReceiving; + + _currentReceiving = Task.Run(GetNextMessage); + return _currentReceiving; + } } public async Task PostMessage(NetworkMessage message) { if (BaseStream == null) throw new NullReferenceException("BaseStream is null"); - - Console.WriteLine("Sending {0}", JsonSerializer.Serialize(message)); - + // Console.WriteLine("Sending {0}", JsonSerializer.Serialize(message)); + + var rawMessage = _serialization.Serialize(message); await BaseStream.WriteAsync(BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort))); await BaseStream.WriteAsync(rawMessage); } + public void HandleMessage(NetworkMessage message) + { + _messageBuffer.Remove(message); + } + + public NetworkMessage[] UnhandledMessages => _messageBuffer.ToArray(); + private NetworkMessage GetNextMessage() { if (BaseStream == null) @@ -62,14 +76,14 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule throw new NullReferenceException("Serialization toolkit is null"); - Console.WriteLine("Receiving message"); + // Console.WriteLine("Receiving message"); BaseStream.ReadExactly(_buffer, 0, 2); var len = BitConverter.ToUInt16(_buffer, 0); BaseStream.ReadExactly(_buffer, 0, len); - Console.WriteLine("Receiving {0}", Encoding.Default.GetString(_buffer[..len])); - - var message = JsonSerializer.Deserialize(Encoding.Default.GetString(_buffer[..len])); + // Console.WriteLine("Receiving {0}", Encoding.Default.GetString(_buffer[..len])); + + var message = _serialization.Deserialize(_buffer[..len]); _messageBuffer.Add(message!); return message!; diff --git a/mROA/Implementation/RepresentationModule.cs b/mROA/Implementation/RepresentationModule.cs index 6b6c266..094866c 100644 --- a/mROA/Implementation/RepresentationModule.cs +++ b/mROA/Implementation/RepresentationModule.cs @@ -35,13 +35,26 @@ public class RepresentationModule : IRepresentationModule public async Task GetRawMessage(Guid? requestId = null, MessageType? messageType = null) { - while (true) + var fromBuffer = + _interaction.UnhandledMessages.FirstOrDefault(message => + (requestId is null || message.Id == requestId) && + (messageType is null || message.SchemaId == messageType)); + if (fromBuffer == null) { - var message = await _interaction.GetNextMessageReceiving(); - if ((requestId is null || message.Id == requestId) && - (messageType is null || message.SchemaId == messageType)) - return message.Data; + while (true) + { + var message = await _interaction.GetNextMessageReceiving(); + if ((requestId is null || message.Id == requestId) && + (messageType is null || message.SchemaId == messageType)) + { + _interaction.HandleMessage(message); + return message.Data; + } + } } + + _interaction.HandleMessage(fromBuffer); + return fromBuffer.Data; } public async Task PostCallMessageAsync(Guid id, MessageType messageType, T payload)