Производительность выше крыши
This commit is contained in:
@@ -71,6 +71,7 @@ var x = 0;
|
||||
for (int i = 0; i < iterations; i++)
|
||||
{
|
||||
x = loadSingleton.Next(x);
|
||||
// Console.WriteLine(x);
|
||||
}
|
||||
|
||||
timer.Stop();
|
||||
|
||||
@@ -31,7 +31,6 @@ public class NextGenTest
|
||||
{
|
||||
_listener.Start();
|
||||
_interactionModuleB.BaseStream = _listener.AcceptTcpClient().GetStream();
|
||||
|
||||
foreach (var guid in guids)
|
||||
{
|
||||
_interactionModuleB.PostMessage(new NetworkMessage { Id = guid, Data = "Hello user"u8.ToArray() });
|
||||
|
||||
@@ -5,7 +5,7 @@ namespace mROA.Abstract;
|
||||
public interface INextGenerationInteractionModule : IInjectableModule
|
||||
{
|
||||
int ConnectionId { get; }
|
||||
public Stream? BaseStream { get; set; }
|
||||
public Stream BaseStream { get; set; }
|
||||
NetworkMessage[] UnhandledMessages { get; }
|
||||
NetworkMessage LastMessage { get; }
|
||||
EventWaitHandle CurrentReceivingHandle { get; }
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
using System.Windows.Input;
|
||||
using mROA.Implementation;
|
||||
using mROA.Implementation.CommandExecution;
|
||||
|
||||
@@ -20,9 +21,10 @@ public interface ISerialisationModule : IInjectableModule
|
||||
public interface IRepresentationModule : IInjectableModule
|
||||
{
|
||||
int Id { get; }
|
||||
Task<T> GetMessageAsync<T>(Guid? requestId = null, MessageType? messageType = null);
|
||||
Task<T> GetMessageAsync<T>(Guid? requestId = null, MessageType? messageType = null, CancellationToken token = default);
|
||||
T GetMessage<T>(Guid? requestId = null, MessageType? messageType = null);
|
||||
Task<byte[]> GetRawMessage(Guid? requestId = null, MessageType? messageType = null);
|
||||
T GetMessage<T>(Predicate<NetworkMessage> filter);
|
||||
Task<byte[]> GetRawMessage(Predicate<NetworkMessage> filter, CancellationToken token = default);
|
||||
|
||||
Task PostCallMessageAsync<T>(Guid id, MessageType messageType, T payload) where T : notnull;
|
||||
Task PostCallMessageAsync(Guid id, MessageType messageType, object payload, Type payloadType);
|
||||
|
||||
@@ -68,7 +68,7 @@ public class NetworkGatewayModule : IGatewayModule
|
||||
interaction!.Inject(_serialization);
|
||||
|
||||
interaction.BaseStream = client.GetStream();
|
||||
|
||||
interaction.StartInfiniteReceiving();
|
||||
interaction.PostMessage(new NetworkMessage
|
||||
{
|
||||
Id = Guid.NewGuid(), SchemaId = MessageType.IdAssigning,
|
||||
|
||||
@@ -32,12 +32,16 @@ public class NetworkFrontendBridge(IPEndPoint ipEndPoint) : IFrontendBridge
|
||||
|
||||
_tcpClient.Connect(ipEndPoint);
|
||||
_interactionModule.BaseStream = _tcpClient.GetStream();
|
||||
var welcomeMessage = _interactionModule.GetNextMessageReceiving().GetAwaiter().GetResult();
|
||||
_interactionModule.StartInfiniteReceiving();
|
||||
|
||||
var handle = _interactionModule.CurrentReceivingHandle;
|
||||
handle.WaitOne();
|
||||
var welcomeMessage = _interactionModule.LastMessage;
|
||||
if (welcomeMessage.SchemaId != MessageType.IdAssigning)
|
||||
{
|
||||
throw new Exception($"Incorrect message type. Must be IdAssigning, current : {welcomeMessage.SchemaId.ToString()}");
|
||||
}
|
||||
|
||||
_interactionModule.HandleMessage(welcomeMessage);
|
||||
TransmissionConfig.OwnershipRepository = new StaticOwnershipRepository(_serialization.Deserialize<IdAssingnment>(welcomeMessage.Data)!.Id);
|
||||
}
|
||||
}
|
||||
@@ -56,13 +56,13 @@ public class RequestExtractor : IRequestExtractor
|
||||
|
||||
try
|
||||
{
|
||||
var lastCommandId = Guid.Empty;
|
||||
while (true)
|
||||
{
|
||||
|
||||
var request =
|
||||
_representationModule!.GetMessage<DefaultCallRequest>(messageType: MessageType.CallRequest);
|
||||
|
||||
// Console.WriteLine("Executing {0}", request.Id);
|
||||
|
||||
_representationModule!.GetMessage<DefaultCallRequest>(m => m.Id != lastCommandId && m.SchemaId == MessageType.CallRequest);
|
||||
lastCommandId = request.Id;
|
||||
if (request.Parameter is not null)
|
||||
{
|
||||
var parameterType = _methodRepository!.GetMethod(request.CommandId).GetParameters().First()
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
using System.Text.Json.Serialization;
|
||||
|
||||
// ReSharper disable UnusedMember.Global
|
||||
|
||||
namespace mROA.Implementation;
|
||||
@@ -7,13 +8,25 @@ public class NetworkMessage
|
||||
{
|
||||
public static readonly NetworkMessage Null = new() { SchemaId = MessageType.Unknown, Id = Guid.Empty, Data = [] };
|
||||
public Guid Id { get; init; }
|
||||
|
||||
[JsonConverter(typeof(JsonStringEnumConverter))]
|
||||
public MessageType SchemaId { get; init; }
|
||||
|
||||
public required byte[] Data { get; init; }
|
||||
|
||||
public bool IsValidMessage(Guid? requestId = null, MessageType? messageType = null)
|
||||
{
|
||||
return (requestId is null || Id == requestId) &&
|
||||
(messageType is null || SchemaId == messageType);
|
||||
}
|
||||
}
|
||||
|
||||
public enum MessageType
|
||||
{
|
||||
Unknown, FinishedCommandExecution, ExceptionCommandExecution, AsyncCancelCommandExecution, CallRequest, IdAssigning
|
||||
Unknown,
|
||||
FinishedCommandExecution,
|
||||
ExceptionCommandExecution,
|
||||
AsyncCancelCommandExecution,
|
||||
CallRequest,
|
||||
IdAssigning
|
||||
}
|
||||
@@ -6,14 +6,14 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
|
||||
{
|
||||
private ISerializationToolkit? _serialization;
|
||||
public int ConnectionId { get; private set; }
|
||||
public Stream? BaseStream { get; set; }
|
||||
public Stream BaseStream { get; set; } = Stream.Null;
|
||||
private Task<NetworkMessage>? _currentReceiving;
|
||||
private const int BufferSize = ushort.MaxValue;
|
||||
private readonly Memory<byte> _buffer = new byte[BufferSize];
|
||||
private readonly List<NetworkMessage> _messageBuffer = new (128);
|
||||
public NetworkMessage[] UnhandledMessages => _messageBuffer.ToArray();
|
||||
public NetworkMessage LastMessage { get; private set; } = NetworkMessage.Null;
|
||||
public EventWaitHandle CurrentReceivingHandle { get; private set; } = new(false, EventResetMode.ManualReset);
|
||||
public EventWaitHandle CurrentReceivingHandle { get; private set; } = new(true, EventResetMode.ManualReset);
|
||||
|
||||
public void Inject<T>(T dependency)
|
||||
{
|
||||
@@ -31,14 +31,25 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
|
||||
|
||||
public async void StartInfiniteReceiving()
|
||||
{
|
||||
try
|
||||
{
|
||||
CurrentReceivingHandle = new EventWaitHandle(false, EventResetMode.ManualReset);
|
||||
|
||||
await Task.Yield();
|
||||
while (true)
|
||||
{
|
||||
var message = ReceiveMessage();
|
||||
LastMessage = message;
|
||||
_messageBuffer.Add(message);
|
||||
CurrentReceivingHandle.Set();
|
||||
CurrentReceivingHandle = new EventWaitHandle(false, EventResetMode.ManualReset);
|
||||
}
|
||||
}
|
||||
catch (Exception)
|
||||
{
|
||||
// ignored
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -71,9 +82,9 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
|
||||
_messageBuffer.Remove(message);
|
||||
}
|
||||
|
||||
public NetworkMessage? FirstByFilter(Predicate<NetworkMessage> predicate)
|
||||
public NetworkMessage FirstByFilter(Predicate<NetworkMessage> predicate)
|
||||
{
|
||||
return _messageBuffer.FirstOrDefault(m => predicate(m));
|
||||
return _messageBuffer.FirstOrDefault(m => predicate(m)) ?? NetworkMessage.Null;
|
||||
}
|
||||
|
||||
private NetworkMessage GetNextMessage()
|
||||
@@ -84,18 +95,17 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
|
||||
if (_serialization == null)
|
||||
throw new NullReferenceException("Serialization toolkit is null");
|
||||
|
||||
|
||||
var message = ReceiveMessage();
|
||||
|
||||
_messageBuffer.Add(message);
|
||||
_currentReceiving = Task.Run(GetNextMessage);
|
||||
|
||||
return message!;
|
||||
return message;
|
||||
}
|
||||
|
||||
private NetworkMessage ReceiveMessage()
|
||||
{
|
||||
var len = BitConverter.ToUInt16([(byte)BaseStream!.ReadByte(), (byte)BaseStream.ReadByte()]);
|
||||
var len = BitConverter.ToUInt16([(byte)BaseStream.ReadByte(), (byte)BaseStream.ReadByte()]);
|
||||
var localSpan = _buffer.Span.Slice(0, len);
|
||||
BaseStream.ReadExactly(localSpan);
|
||||
|
||||
|
||||
@@ -16,17 +16,26 @@ public abstract class RemoteObjectBase(int id, IRepresentationModule representat
|
||||
{ CommandId = methodId, ObjectId = id, Parameter = parameter, ParameterType = parameter?.GetType() };
|
||||
await representationModule.PostCallMessageAsync(request.Id, MessageType.CallRequest, request);
|
||||
|
||||
var localTokenSource = new CancellationTokenSource();
|
||||
|
||||
var successResponse =
|
||||
representationModule.GetMessageAsync<FinalCommandExecution<T>>(
|
||||
messageType: MessageType.FinishedCommandExecution, requestId: request.Id);
|
||||
messageType: MessageType.FinishedCommandExecution, requestId: request.Id, token: localTokenSource.Token);
|
||||
var errorResponse =
|
||||
representationModule.GetMessageAsync<ExceptionCommandExecution>(
|
||||
messageType: MessageType.ExceptionCommandExecution, requestId: request.Id);
|
||||
messageType: MessageType.ExceptionCommandExecution, requestId: request.Id, token: localTokenSource.Token);
|
||||
|
||||
Task.WaitAny(successResponse, errorResponse);
|
||||
|
||||
if (successResponse.IsCompletedSuccessfully)
|
||||
return successResponse.Result.Result!;
|
||||
|
||||
|
||||
if (successResponse.IsCompletedSuccessfully)
|
||||
{
|
||||
await localTokenSource.CancelAsync();
|
||||
return successResponse.Result.Result!;
|
||||
}
|
||||
|
||||
await localTokenSource.CancelAsync();
|
||||
throw errorResponse.Result.GetException();
|
||||
}
|
||||
|
||||
|
||||
@@ -22,47 +22,74 @@ public class RepresentationModule : IRepresentationModule
|
||||
|
||||
public int Id => (_interaction ?? throw new NullReferenceException("Interaction is not initialized")).ConnectionId;
|
||||
|
||||
public async Task<T> GetMessageAsync<T>(Guid? requestId, MessageType? messageType)
|
||||
public async Task<T> GetMessageAsync<T>(Guid? requestId, MessageType? messageType,
|
||||
CancellationToken token = default)
|
||||
{
|
||||
if (_serialization == null)
|
||||
throw new NullReferenceException("Serialization toolkit is not initialized");
|
||||
|
||||
return _serialization.Deserialize<T>(await GetRawMessage(requestId, messageType))!;
|
||||
return _serialization.Deserialize<T>(await GetRawMessage(m => m.IsValidMessage(requestId, messageType), token))!;
|
||||
}
|
||||
|
||||
public T GetMessage<T>(Guid? requestId = null, MessageType? messageType = null)
|
||||
{
|
||||
return GetMessage<T>(m => m.IsValidMessage(requestId, messageType));
|
||||
}
|
||||
|
||||
public T GetMessage<T>(Predicate<NetworkMessage> filter)
|
||||
{
|
||||
if (_serialization == null)
|
||||
throw new NullReferenceException("Serialization toolkit is not initialized");
|
||||
|
||||
return _serialization.Deserialize<T>(GetRawMessage(requestId, messageType).GetAwaiter().GetResult())!;
|
||||
return _serialization.Deserialize<T>(GetRawMessage(filter).GetAwaiter().GetResult())!;
|
||||
}
|
||||
|
||||
public async Task<byte[]> GetRawMessage(Guid? requestId = null, MessageType? messageType = null)
|
||||
public async Task<byte[]> GetRawMessage(Predicate<NetworkMessage> filter, CancellationToken token = default)
|
||||
{
|
||||
await Task.Yield();
|
||||
|
||||
if (_interaction == null)
|
||||
throw new NullReferenceException("Interaction toolkit is not initialized");
|
||||
|
||||
var fromBuffer =
|
||||
_interaction.FirstByFilter(message =>
|
||||
(requestId is null || message.Id == requestId) &&
|
||||
(messageType is null || message.SchemaId == messageType));
|
||||
// Console.WriteLine(
|
||||
// $"{DateTime.Now.TimeOfDay} {Environment.CurrentManagedThreadId} Representation : Reading");
|
||||
|
||||
if (fromBuffer == null)
|
||||
if (filter(_interaction.LastMessage))
|
||||
{
|
||||
while (true)
|
||||
{
|
||||
var message = await _interaction.GetNextMessageReceiving();
|
||||
if ((requestId is not null && message.Id != requestId) ||
|
||||
(messageType is not null && message.SchemaId != messageType)) continue;
|
||||
_interaction.HandleMessage(_interaction.LastMessage);
|
||||
return _interaction.LastMessage.Data;
|
||||
}
|
||||
|
||||
var message = _interaction.FirstByFilter(filter);
|
||||
|
||||
if (message != NetworkMessage.Null)
|
||||
{
|
||||
_interaction.HandleMessage(message);
|
||||
return message.Data;
|
||||
}
|
||||
|
||||
|
||||
while (!token.IsCancellationRequested)
|
||||
{
|
||||
// Console.WriteLine(
|
||||
// $"{DateTime.Now.TimeOfDay} {Environment.CurrentManagedThreadId} Representation : Receiving message...");
|
||||
var handle = _interaction.CurrentReceivingHandle;
|
||||
handle.WaitOne();
|
||||
|
||||
message = _interaction.LastMessage;
|
||||
|
||||
// Console.WriteLine(
|
||||
// $"{DateTime.Now.TimeOfDay} {Environment.CurrentManagedThreadId} Representation : Message received {message.SchemaId} - {message.Id}");
|
||||
if (!filter(message)) continue;
|
||||
|
||||
message = _interaction.LastMessage;
|
||||
// Console.WriteLine(
|
||||
// $"{DateTime.Now.TimeOfDay} {Environment.CurrentManagedThreadId} Representation : Message received Successfully {message.SchemaId} - {message.Id}");
|
||||
_interaction.HandleMessage(message);
|
||||
return message.Data;
|
||||
}
|
||||
|
||||
_interaction.HandleMessage(fromBuffer);
|
||||
return fromBuffer.Data;
|
||||
return [];
|
||||
}
|
||||
|
||||
public async Task PostCallMessageAsync<T>(Guid id, MessageType messageType, T payload) where T : notnull
|
||||
|
||||
Reference in New Issue
Block a user