улучшен новый модуль взаимодействия и написан новый модуль репрезентации

This commit is contained in:
2025-02-16 17:05:16 +03:00
parent 5e9c2cbb7b
commit 5c643dd2c9
8 changed files with 73 additions and 68 deletions
+1
View File
@@ -46,6 +46,7 @@ public class NextGenTest
Task.WaitAll(tasks.ToArray()); Task.WaitAll(tasks.ToArray());
Assert.Pass(); Assert.Pass();
} }
private async Task ReadStream(Guid current) private async Task ReadStream(Guid current)
+1 -16
View File
@@ -2,25 +2,10 @@ using mROA.Implementation;
namespace mROA.Abstract; namespace mROA.Abstract;
public interface IInteractionModule : IInjectableModule
{
void SendTo(int clientId, byte[] message);
void RegisterSource(Stream stream);
Stream GetSource(int clientId);
public interface IFrontendInteractionModule : IInjectableModule
{
int ClientId { get; }
public Task<byte[]> ReceiveMessage();
public void PostMessage(byte[] message);
}
}
public interface INextGenerationInteractionModule : IInjectableModule public interface INextGenerationInteractionModule : IInjectableModule
{ {
int ConntectionId { get; } int ConntectionId { get; }
public Stream? BaseStream { get; set; } public Stream? BaseStream { get; set; }
Task<NetworkMessage> GetNextMessageReceiving(); Task<NetworkMessage> GetNextMessageReceiving();
void PostMessage(NetworkMessage message); Task PostMessage(NetworkMessage message);
NetworkMessage[] UnhandledMessages { get; }
void HandleMessage(NetworkMessage msg);
} }
+9 -1
View File
@@ -1,3 +1,4 @@
using System.Windows.Input;
using mROA.Implementation; using mROA.Implementation;
namespace mROA.Abstract; namespace mROA.Abstract;
@@ -12,7 +13,14 @@ public interface ISerialisationModule : IInjectableModule
int ClientId { get; } int ClientId { get; }
Task<T> GetNextCommandExecution<T>(Guid requestId) where T : ICommandExecution; Task<T> GetNextCommandExecution<T>(Guid requestId) where T : ICommandExecution;
Task<FinalCommandExecution<T>> GetFinalCommandExecution<T>(Guid requestId); Task<FinalCommandExecution<T>> GetFinalCommandExecution<T>(Guid requestId);
void PostCallRequest(ICallRequest callRequest); void PostCallRequest(ICallRequest callRequest);
} }
}
public interface IRepresentationModule : IInjectableModule
{
int Id { get; }
Task<T> GetMessage<T>(Guid? requestId, MessageType? messageType);
Task PostCallMessage<T>(Guid id, MessageType messageType, T payload);
Task PostCallMessage(Guid id, MessageType messageType, object payload, Type payloadType);
} }
@@ -1,33 +0,0 @@
namespace mROA.Implementation;
public class AsyncNetworkInteractionModule
{
private ISerializationToolkit _iSerializationToolkit;
private const int BufferSize = ushort.MaxValue;
private Stream? _stream;
private List<NetworkMessage> _unhandledMessages = new(8);
private byte[] _buffer = new byte[BufferSize];
private Task? _currentReadTask;
public Task ReceiveNext()
{
_currentReadTask ??= StartReceiving();
return _currentReadTask;
}
private async Task StartReceiving()
{
await _stream.ReadExactlyAsync(_buffer, 0, 2);
var len = BitConverter.ToUInt16(_buffer, 0);
await _stream.ReadExactlyAsync(_buffer, 0, len);
}
public NetworkMessage GetLastMessage()
{
return _unhandledMessages.Last();
}
}
+2 -1
View File
@@ -1,4 +1,5 @@
using mROA.Abstract; using System.Text.Json;
using mROA.Abstract;
namespace mROA.Implementation; namespace mROA.Implementation;
@@ -1,17 +1,15 @@
using mROA.Abstract; using mROA.Abstract;
using mROA.Implementation;
namespace mROA; namespace mROA.Implementation;
public class NextGenerationInteractionModule : INextGenerationInteractionModule public class NextGenerationInteractionModule : INextGenerationInteractionModule
{ {
private ISerializationToolkit _serialization; private ISerializationToolkit? _serialization;
public int ConntectionId { get; set; } public int ConntectionId { get; set; }
public Stream BaseStream { get; set; } public Stream? BaseStream { get; set; }
private Task<NetworkMessage>? _currentReceiving; private Task<NetworkMessage>? _currentReceiving;
private const int BufferSize = ushort.MaxValue; private const int BufferSize = ushort.MaxValue;
private byte[] _buffer = new byte[BufferSize]; private readonly byte[] _buffer = new byte[BufferSize];
private List<NetworkMessage> _unhandledMessages = new(8);
public void Inject<T>(T dependency) public void Inject<T>(T dependency)
{ {
@@ -29,22 +27,24 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
return _currentReceiving; return _currentReceiving;
} }
public void PostMessage(NetworkMessage message) public async Task PostMessage(NetworkMessage message)
{ {
if (BaseStream == null)
throw new NullReferenceException("BaseStream is null");
var rawMessage = _serialization.Serialize(message); var rawMessage = _serialization.Serialize(message);
BaseStream.Write(BitConverter.GetBytes((ushort)rawMessage.Length), 0, sizeof(ushort)); await BaseStream.WriteAsync(BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort)));
BaseStream.Write(rawMessage, 0, rawMessage.Length); await BaseStream.WriteAsync(rawMessage);
}
public NetworkMessage[] UnhandledMessages => _unhandledMessages.ToArray();
public void HandleMessage(NetworkMessage msg)
{
_unhandledMessages.Remove(msg);
} }
private async Task<NetworkMessage> GetNextMessage() private async Task<NetworkMessage> GetNextMessage()
{ {
if (BaseStream == null)
throw new NullReferenceException("BaseStream is null");
if (_serialization == null)
throw new NullReferenceException("Serialization toolkit is null");
await BaseStream.ReadExactlyAsync(_buffer, 0, 2); await BaseStream.ReadExactlyAsync(_buffer, 0, 2);
var len = BitConverter.ToUInt16(_buffer, 0); var len = BitConverter.ToUInt16(_buffer, 0);
await BaseStream.ReadExactlyAsync(_buffer, 0, len); await BaseStream.ReadExactlyAsync(_buffer, 0, len);
@@ -1,5 +1,4 @@
using System.Collections.Frozen; using System.Collections.Frozen;
using System.Collections.Immutable;
using mROA.Abstract; using mROA.Abstract;
namespace mROA.Implementation; namespace mROA.Implementation;
@@ -0,0 +1,44 @@
using mROA.Abstract;
namespace mROA.Implementation;
public class RepresentationModule : IRepresentationModule
{
private ISerializationToolkit _serialization;
private INextGenerationInteractionModule _interaction;
public void Inject<T>(T dependency)
{
switch (dependency)
{
case ISerializationToolkit toolkit:
_serialization = toolkit;
break;
case INextGenerationInteractionModule interactionModule:
_interaction = interactionModule;
break;
}
}
public int Id { get; set; }
public async Task<T> GetMessage<T>(Guid? requestId, MessageType? messageType)
{
while (true)
{
var message = await _interaction.GetNextMessageReceiving();
if ((requestId is null || message.Id == requestId) && (messageType is null || message.SchemaId == messageType))
return _serialization.Deserialize<T>(message.Data)!;
}
}
public async Task PostCallMessage<T>(Guid id, MessageType messageType, T payload)
{
await PostCallMessage(id, messageType, payload, typeof(T));
}
public async Task PostCallMessage(Guid id, MessageType messageType, object payload, Type payloadType)
{
await _interaction.PostMessage(new NetworkMessage {Id = id, SchemaId = messageType, Data = _serialization.Serialize(payload, payloadType)});
}
}