попытка буферизировать сообщения не увенчалась успехом, но может использоваться в будущем

This commit is contained in:
2025-02-18 14:56:31 +03:00
parent 6023f47212
commit 635fcbf914
5 changed files with 55 additions and 27 deletions
+4 -5
View File
@@ -18,8 +18,6 @@ builder.Modules.Add(new StaticRepresentationModuleProducer());
builder.Modules.Add(new RequestExtractor()); builder.Modules.Add(new RequestExtractor());
builder.Modules.Add(new BasicExecutionModule()); builder.Modules.Add(new BasicExecutionModule());
builder.Modules.Add(new CoCodegenMethodRepository()); builder.Modules.Add(new CoCodegenMethodRepository());
// builder.Modules.Add(new StreamBasedVirtualBackendInteractionModule());
// builder.Modules.Add(new JsonSerialisationModule());
builder.UseCollectableContextRepository(); builder.UseCollectableContextRepository();
builder.Build(); builder.Build();
@@ -28,7 +26,7 @@ TransmissionConfig.RealContextRepository = builder.GetModule<ContextRepository>(
TransmissionConfig.RemoteEndpointContextRepository = builder.GetModule<RemoteContextRepository>(); TransmissionConfig.RemoteEndpointContextRepository = builder.GetModule<RemoteContextRepository>();
builder.GetModule<NetworkFrontendBridge>()!.Connect(); builder.GetModule<NetworkFrontendBridge>()!.Connect();
// _ = builder.GetModule<RequestExtractor>()!.StartExtraction(); // _ = builder.GetModule<RequestExtractor>()!.StartExtraction();
Console.WriteLine(TransmissionConfig.OwnershipRepository!.GetOwnershipId()); Console.WriteLine(TransmissionConfig.OwnershipRepository!.GetOwnershipId());
var context = builder.GetModule<RemoteContextRepository>(); var context = builder.GetModule<RemoteContextRepository>();
@@ -36,6 +34,7 @@ var factory = context.GetSingleObject(typeof(IPrinterFactory)) as IPrinterFactor
//правильный порядок команд 8-5-10-7 //правильный порядок команд 8-5-10-7
var printer = factory.Create("Test"); var printer = factory.Create("Test");
Console.WriteLine("Printer created");
Thread.Sleep(100); Thread.Sleep(100);
var name = printer.Value.GetName(); var name = printer.Value.GetName();
@@ -43,13 +42,13 @@ Console.WriteLine("Printer name : {0}", name);
Thread.Sleep(100); Thread.Sleep(100);
factory.Register(new SharedObject<IPrinter>(new ClientBasedPrinter())); factory.Register(new SharedObject<IPrinter>(new ClientBasedPrinter()));
Console.WriteLine("Registered printer");
Thread.Sleep(100); Thread.Sleep(100);
Console.WriteLine("Registered printer");
var registred = factory.GetFirstPrinter(); var registred = factory.GetFirstPrinter();
Thread.Sleep(100);
Console.WriteLine("First printer"); Console.WriteLine("First printer");
Thread.Sleep(100);
Console.WriteLine(registred.Value); Console.WriteLine(registred.Value);
Console.WriteLine("Collecting all printers"); Console.WriteLine("Collecting all printers");
+2
View File
@@ -8,4 +8,6 @@ public interface INextGenerationInteractionModule : IInjectableModule
public Stream? BaseStream { get; set; } public Stream? BaseStream { get; set; }
Task<NetworkMessage> GetNextMessageReceiving(); Task<NetworkMessage> GetNextMessageReceiving();
Task PostMessage(NetworkMessage message); Task PostMessage(NetworkMessage message);
void HandleMessage(NetworkMessage message);
NetworkMessage[] UnhandledMessages { get; }
} }
@@ -61,7 +61,7 @@ public class RequestExtractor : IRequestExtractor
var request = var request =
_representationModule!.GetMessage<DefaultCallRequest>(messageType: MessageType.CallRequest); _representationModule!.GetMessage<DefaultCallRequest>(messageType: MessageType.CallRequest);
Console.WriteLine("Executing {0}", request.Id); // Console.WriteLine("Executing {0}", request.Id);
if (request.Parameter is not null) if (request.Parameter is not null)
{ {
@@ -1,6 +1,4 @@
using System.Security.Cryptography; using System.Text;
using System.Text;
using System.Text.Json;
using mROA.Abstract; using mROA.Abstract;
namespace mROA.Implementation; namespace mROA.Implementation;
@@ -15,7 +13,7 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
private readonly byte[] _buffer = new byte[BufferSize]; private readonly byte[] _buffer = new byte[BufferSize];
private List<NetworkMessage> _messageBuffer = []; private List<NetworkMessage> _messageBuffer = [];
public void Inject<T>(T dependency) public void Inject<T>(T dependency)
{ {
switch (dependency) switch (dependency)
@@ -32,27 +30,43 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
public Task<NetworkMessage> GetNextMessageReceiving() public Task<NetworkMessage> GetNextMessageReceiving()
{ {
Console.WriteLine(_currentReceiving?.Status); // Console.WriteLine(_currentReceiving?.Status);
if (_currentReceiving is { Status: TaskStatus.Running })
return _currentReceiving;
_currentReceiving = Task.Run(GetNextMessage); if (_currentReceiving == null)
return _currentReceiving; {
_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) public async Task PostMessage(NetworkMessage message)
{ {
if (BaseStream == null) if (BaseStream == null)
throw new NullReferenceException("BaseStream is 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); var rawMessage = _serialization.Serialize(message);
await BaseStream.WriteAsync(BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort))); await BaseStream.WriteAsync(BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort)));
await BaseStream.WriteAsync(rawMessage); await BaseStream.WriteAsync(rawMessage);
} }
public void HandleMessage(NetworkMessage message)
{
_messageBuffer.Remove(message);
}
public NetworkMessage[] UnhandledMessages => _messageBuffer.ToArray();
private NetworkMessage GetNextMessage() private NetworkMessage GetNextMessage()
{ {
if (BaseStream == null) if (BaseStream == null)
@@ -62,14 +76,14 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
throw new NullReferenceException("Serialization toolkit is null"); throw new NullReferenceException("Serialization toolkit is null");
Console.WriteLine("Receiving message"); // Console.WriteLine("Receiving message");
BaseStream.ReadExactly(_buffer, 0, 2); BaseStream.ReadExactly(_buffer, 0, 2);
var len = BitConverter.ToUInt16(_buffer, 0); var len = BitConverter.ToUInt16(_buffer, 0);
BaseStream.ReadExactly(_buffer, 0, len); BaseStream.ReadExactly(_buffer, 0, len);
Console.WriteLine("Receiving {0}", Encoding.Default.GetString(_buffer[..len])); // Console.WriteLine("Receiving {0}", Encoding.Default.GetString(_buffer[..len]));
var message = JsonSerializer.Deserialize<NetworkMessage>(Encoding.Default.GetString(_buffer[..len])); var message = _serialization.Deserialize<NetworkMessage>(_buffer[..len]);
_messageBuffer.Add(message!); _messageBuffer.Add(message!);
return message!; return message!;
+18 -5
View File
@@ -35,13 +35,26 @@ public class RepresentationModule : IRepresentationModule
public async Task<byte[]> GetRawMessage(Guid? requestId = null, MessageType? messageType = null) public async Task<byte[]> 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(); while (true)
if ((requestId is null || message.Id == requestId) && {
(messageType is null || message.SchemaId == messageType)) var message = await _interaction.GetNextMessageReceiving();
return message.Data; 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<T>(Guid id, MessageType messageType, T payload) public async Task PostCallMessageAsync<T>(Guid id, MessageType messageType, T payload)