From be12f1da6508cb90d684adc3386f2c10c35e53e8 Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Sat, 26 Jul 2025 20:44:19 +0300 Subject: [PATCH 1/4] Preparing for new request id --- mROA.Cbor/CborSerializationToolkit.cs | 2 +- mROA.Cbor/IOrdinaryStructureParser.cs | 4 +- mROA/Abstract/IChannelInteractionModule.cs | 12 +-- mROA/Abstract/IRepresentationModule.cs | 8 +- mROA/Abstract/IRequestExtractor.cs | 4 +- .../Backend/NetworkGatewayModule.cs | 6 +- mROA/Implementation/Backend/UdpGateway.cs | 2 +- .../ChannelInteractionModule.cs | 75 +++++++++---------- mROA/Implementation/EMessageType.cs | 2 +- .../Frontend/NetworkFrontendBridge.cs | 4 +- .../Frontend/RequestExtractor.cs | 4 +- .../Frontend/UdpUntrustedInteraction.cs | 4 +- ...workMessageHeader.cs => NetworkMessage.cs} | 35 +++++++-- mROA/Implementation/RepresentationModule.cs | 12 +-- mROA/Implementation/RequestId.cs | 18 +++++ 15 files changed, 116 insertions(+), 76 deletions(-) rename mROA/Implementation/{NetworkMessageHeader.cs => NetworkMessage.cs} (50%) create mode 100644 mROA/Implementation/RequestId.cs diff --git a/mROA.Cbor/CborSerializationToolkit.cs b/mROA.Cbor/CborSerializationToolkit.cs index c5cd523..f50fe7c 100644 --- a/mROA.Cbor/CborSerializationToolkit.cs +++ b/mROA.Cbor/CborSerializationToolkit.cs @@ -27,7 +27,7 @@ namespace mROA.Cbor private bool FindParser(Type t, out IOrdinaryStructureParser parser) { - if (t == typeof(NetworkMessageHeader)) + if (t == typeof(NetworkMessage)) { parser = _parsers[0]; return true; diff --git a/mROA.Cbor/IOrdinaryStructureParser.cs b/mROA.Cbor/IOrdinaryStructureParser.cs index 8faa205..2d3711a 100644 --- a/mROA.Cbor/IOrdinaryStructureParser.cs +++ b/mROA.Cbor/IOrdinaryStructureParser.cs @@ -16,7 +16,7 @@ namespace mROA.Cbor { public void Write(CborWriter writer, object value, IEndPointContext context, CborSerializationToolkit serialization) { - var v = value as NetworkMessageHeader; + var v = value as NetworkMessage; writer.WriteStartArray(3); writer.WriteByteString(v.Id.ToByteArray()); writer.WriteInt32((int)v.MessageType); @@ -27,7 +27,7 @@ namespace mROA.Cbor public object Read(CborReader reader, IEndPointContext context, CborSerializationToolkit serialization) { reader.ReadStartArray(); - var value = new NetworkMessageHeader + var value = new NetworkMessage { Id = new Guid(reader.ReadByteString()), MessageType = (EMessageType)reader.ReadInt32(), diff --git a/mROA/Abstract/IChannelInteractionModule.cs b/mROA/Abstract/IChannelInteractionModule.cs index e7054c0..0542417 100644 --- a/mROA/Abstract/IChannelInteractionModule.cs +++ b/mROA/Abstract/IChannelInteractionModule.cs @@ -9,13 +9,13 @@ namespace mROA.Abstract { int ConnectionId { get; set; } IEndPointContext Context { get; set; } - Channel ReceiveChanel { get; } - ChannelReader TrustedPostChanel { get; } - ChannelReader UntrustedPostChanel { get; } + Channel ReceiveChanel { get; } + ChannelReader TrustedPostChanel { get; } + ChannelReader UntrustedPostChanel { get; } Func IsConnected { get; set; } - ValueTask GetNextMessageReceiving(); - Task PostMessageAsync(NetworkMessageHeader messageHeader); - Task PostMessageUntrustedAsync(NetworkMessageHeader messageHeader); + ValueTask GetNextMessageReceiving(); + Task PostMessageAsync(NetworkMessage message); + Task PostMessageUntrustedAsync(NetworkMessage message); event Action OnDisconnected; Task Restart(bool sendRecovery); void PassReconnection(); diff --git a/mROA/Abstract/IRepresentationModule.cs b/mROA/Abstract/IRepresentationModule.cs index 7fd4bf5..45d44d3 100644 --- a/mROA/Abstract/IRepresentationModule.cs +++ b/mROA/Abstract/IRepresentationModule.cs @@ -11,13 +11,13 @@ namespace mROA.Abstract int Id { get; } IEndPointContext Context { get; } - Task<(object? Deserialized, EMessageType MessageType)> GetSingle(Predicate rule, + Task<(object? Deserialized, EMessageType MessageType)> GetSingle(Predicate rule, IEndPointContext? context, CancellationToken token = default, - params Func[] converter); + params Func[] converter); - IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(Predicate rule, + IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(Predicate rule, IEndPointContext? context, CancellationToken token = default, - params Func[] converter); + params Func[] converter); Task PostCallMessageAsync(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull; diff --git a/mROA/Abstract/IRequestExtractor.cs b/mROA/Abstract/IRequestExtractor.cs index d4a3c44..6fc3816 100644 --- a/mROA/Abstract/IRequestExtractor.cs +++ b/mROA/Abstract/IRequestExtractor.cs @@ -8,7 +8,7 @@ namespace mROA.Abstract { Task StartExtraction(); void PushMessage(object parced, EMessageType originalType); - Predicate Rule { get; } - Func[] Converters { get; } + Predicate Rule { get; } + Func[] Converters { get; } } } \ No newline at end of file diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index 1308bb7..6d4bdfb 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -98,14 +98,14 @@ namespace mROA.Implementation.Backend } private void HandleNewClient(EndPointContext context, ChannelInteractionModule interaction, - ChannelInteractionModule.StreamExtractor streamExtractor, CancellationTokenSource cts, NetworkMessageHeader connectionHeader) + ChannelInteractionModule.StreamExtractor streamExtractor, CancellationTokenSource cts, NetworkMessage connection) { context.HostId = 0; context.OwnerId = -interaction.ConnectionId; interaction.Context = context; Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token)); _ = streamExtractor.SendFromChannel(interaction.TrustedPostChanel, cts.Token); - interaction.PostMessageAsync(new NetworkMessageHeader(_serialization, + interaction.PostMessageAsync(new NetworkMessage(_serialization, new IdAssignment { Id = interaction.ConnectionId }, null)); _extractorsTokenSources[interaction.ConnectionId] = cts; @@ -143,7 +143,7 @@ namespace mROA.Implementation.Backend }; } - private void RecoverDisconnectedClient(NetworkMessageHeader connectionRequest, + private void RecoverDisconnectedClient(NetworkMessage connectionRequest, ChannelInteractionModule.StreamExtractor streamExtractor, CancellationTokenSource cts) { diff --git a/mROA/Implementation/Backend/UdpGateway.cs b/mROA/Implementation/Backend/UdpGateway.cs index 12a5297..623858b 100644 --- a/mROA/Implementation/Backend/UdpGateway.cs +++ b/mROA/Implementation/Backend/UdpGateway.cs @@ -41,7 +41,7 @@ namespace mROA.Implementation.Backend while (token.IsCancellationRequested == false) { var incoming = await _client.ReceiveAsync(); - var parsed = _serializationToolkit.Deserialize(incoming.Buffer, null); + var parsed = _serializationToolkit.Deserialize(incoming.Buffer, null); try { int channelId; diff --git a/mROA/Implementation/ChannelInteractionModule.cs b/mROA/Implementation/ChannelInteractionModule.cs index 154c5fa..d084ad4 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -9,11 +9,11 @@ namespace mROA.Implementation { public class ChannelInteractionModule : IChannelInteractionModule { - private readonly ChannelReader _receiveReader; - private readonly ChannelWriter _trustedWriter; - private readonly ChannelWriter _untrustedWriter; - private readonly Channel _outputTrustedChannel; - private readonly Channel _outputUntrustedChannel; + private readonly ChannelReader _receiveReader; + private readonly ChannelWriter _trustedWriter; + private readonly ChannelWriter _untrustedWriter; + private readonly Channel _outputTrustedChannel; + private readonly Channel _outputUntrustedChannel; private readonly IContextualSerializationToolKit _serialization; private bool _isConnected = true; private bool _isActive = true; @@ -28,19 +28,19 @@ namespace mROA.Implementation public ChannelInteractionModule(IContextualSerializationToolKit serialization) { _serialization = serialization; - ReceiveChanel = Channel.CreateUnbounded(new UnboundedChannelOptions + ReceiveChanel = Channel.CreateUnbounded(new UnboundedChannelOptions { SingleReader = false, SingleWriter = false, }); _receiveReader = ReceiveChanel.Reader; - _outputTrustedChannel = Channel.CreateUnbounded(new UnboundedChannelOptions + _outputTrustedChannel = Channel.CreateUnbounded(new UnboundedChannelOptions { SingleReader = true, SingleWriter = true }); _trustedWriter = _outputTrustedChannel.Writer; - _outputUntrustedChannel = Channel.CreateUnbounded(new UnboundedChannelOptions + _outputUntrustedChannel = Channel.CreateUnbounded(new UnboundedChannelOptions { SingleReader = true, SingleWriter = true, @@ -52,33 +52,33 @@ namespace mROA.Implementation public int ConnectionId { get; set; } public IEndPointContext Context { get; set; } - public Channel ReceiveChanel { get; } + public Channel ReceiveChanel { get; } - public ChannelReader TrustedPostChanel => _outputTrustedChannel.Reader; - public ChannelReader UntrustedPostChanel => _outputUntrustedChannel.Reader; + public ChannelReader TrustedPostChanel => _outputTrustedChannel.Reader; + public ChannelReader UntrustedPostChanel => _outputUntrustedChannel.Reader; public Func IsConnected { get; set; } = () => false; - public ValueTask GetNextMessageReceiving() + public ValueTask GetNextMessageReceiving() { return _receiveReader.ReadAsync(); } - private async ValueTask PostMessageInternal(NetworkMessageHeader messageHeader) + private async ValueTask PostMessageInternal(NetworkMessage message) { if (!IsConnected()) { return false; } - await _trustedWriter.WriteAsync(messageHeader); + await _trustedWriter.WriteAsync(message); return true; } - public async Task PostMessageAsync(NetworkMessageHeader messageHeader) + public async Task PostMessageAsync(NetworkMessage message) { while (true) { - if (await PostMessageInternal(messageHeader)) + if (await PostMessageInternal(message)) break; if (!_isActive) @@ -91,9 +91,9 @@ namespace mROA.Implementation } } - public async Task PostMessageUntrustedAsync(NetworkMessageHeader messageHeader) + public async Task PostMessageUntrustedAsync(NetworkMessage message) { - await _untrustedWriter.WriteAsync(messageHeader); + await _untrustedWriter.WriteAsync(message); } public event Action? OnDisconnected; @@ -103,12 +103,12 @@ namespace mROA.Implementation if (sendRecovery) { await PostMessageAsync( - new NetworkMessageHeader(_serialization, new ClientRecovery(Math.Abs(ConnectionId)), Context)); + new NetworkMessage(_serialization, new ClientRecovery(Math.Abs(ConnectionId)), Context)); await ReceiveChanel.Reader.ReadAsync(); } else { - await _trustedWriter.WriteAsync(new NetworkMessageHeader()); + await _trustedWriter.WriteAsync(new NetworkMessage()); } PassReconnection(); @@ -141,7 +141,7 @@ namespace mROA.Implementation public class StreamExtractor { - private const int BufferSize = ushort.MaxValue; + private const int BufferSize = ushort.MaxValue + 2; private readonly Stream _ioStream; private readonly IContextualSerializationToolKit _serializationToolkit; @@ -158,25 +158,22 @@ namespace mROA.Implementation _lenBuffer = new byte[2]; } - public Action MessageReceived = _ => { }; - - private async Task ReadMessageLength() - { - await _ioStream.ReadAsync(_lenBuffer); - - var len = BitConverter.ToUInt16(_lenBuffer); - - return len; - } - + public Action MessageReceived = _ => { }; + public async Task SingleReceive(CancellationToken token = default) { - var len = await ReadMessageLength(); - var localSpan = _buffer[..len]; + int firstRead = _ioStream.Read(_buffer.Span); - await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: token); - var message = _serializationToolkit.Deserialize(localSpan, _context); - // _logger.LogTrace("RECV {0}", message.ToString()); + var metadata = new NetworkMessage.NetworkMessageMeta(_buffer.Span); + + var localSpan = _buffer[2..len]; + if (firstRead - 2 != len) + { + localSpan = _buffer[(len + 2)..]; + await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: token); + } + + var message = _serializationToolkit.Deserialize(localSpan, _context); MessageReceived(message); } @@ -188,7 +185,7 @@ namespace mROA.Implementation } } - private async Task Send(NetworkMessageHeader message, CancellationToken token = default) + private async Task Send(NetworkMessage message, CancellationToken token = default) { var bodySpan = _buffer[2..]; var len = _serializationToolkit.Serialize(message, bodySpan.Span, _context); @@ -199,7 +196,7 @@ namespace mROA.Implementation // _logger.LogTrace("SEND {0}", message.ToString()); } - public async Task SendFromChannel(ChannelReader channel, + public async Task SendFromChannel(ChannelReader channel, CancellationToken token = default) { while (token.IsCancellationRequested == false && IsConnected) diff --git a/mROA/Implementation/EMessageType.cs b/mROA/Implementation/EMessageType.cs index 205e85d..c0b5980 100644 --- a/mROA/Implementation/EMessageType.cs +++ b/mROA/Implementation/EMessageType.cs @@ -1,6 +1,6 @@ namespace mROA.Implementation { - public enum EMessageType + public enum EMessageType : byte { Unknown, FinishedCommandExecution, diff --git a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs index 7203285..cfd7f4a 100644 --- a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs +++ b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs @@ -39,7 +39,7 @@ namespace mROA.Implementation.Frontend _interactionModule.IsConnected = () => _currentExtractor.IsConnected; _interactionModule.OnDisconnected += _ => { Reconnect().ConfigureAwait(false); }; - _interactionModule.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientConnect(), _context)) + _interactionModule.PostMessageAsync(new NetworkMessage(_serialization, new ClientConnect(), _context)) .Wait(); _ = _currentExtractor.SingleReceive().ConfigureAwait(false); @@ -95,7 +95,7 @@ namespace mROA.Implementation.Frontend public void Disconnect() { - _ = _interactionModule.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientDisconnect(), + _ = _interactionModule.PostMessageAsync(new NetworkMessage(_serialization, new ClientDisconnect(), _context)); _interactionModule.Dispose(); _tcpClient.Dispose(); diff --git a/mROA/Implementation/Frontend/RequestExtractor.cs b/mROA/Implementation/Frontend/RequestExtractor.cs index 1db877d..f13afe1 100644 --- a/mROA/Implementation/Frontend/RequestExtractor.cs +++ b/mROA/Implementation/Frontend/RequestExtractor.cs @@ -56,11 +56,11 @@ namespace mROA.Implementation.Frontend } } - public Predicate Rule { get; } = m => + public Predicate Rule { get; } = m => m.MessageType is EMessageType.CallRequest or EMessageType.CancelRequest or EMessageType.EventRequest or EMessageType.ClientDisconnect; - public Func[] Converters { get; } = + public Func[] Converters { get; } = { m => m.MessageType == EMessageType.CallRequest ? typeof(DefaultCallRequest) : null, m => m.MessageType == EMessageType.CancelRequest ? typeof(CancelRequest) : null, diff --git a/mROA/Implementation/Frontend/UdpUntrustedInteraction.cs b/mROA/Implementation/Frontend/UdpUntrustedInteraction.cs index b663248..13a6bce 100644 --- a/mROA/Implementation/Frontend/UdpUntrustedInteraction.cs +++ b/mROA/Implementation/Frontend/UdpUntrustedInteraction.cs @@ -44,7 +44,7 @@ namespace mROA.Implementation.Frontend while (token.IsCancellationRequested == false) { var message = new Memory((await udpClient.ReceiveAsync()).Buffer); - var parsed = _serializationToolkit.Deserialize(message, _context); + var parsed = _serializationToolkit.Deserialize(message, _context); await writer.WriteAsync(parsed, token); } @@ -52,7 +52,7 @@ namespace mROA.Implementation.Frontend private async Task Posting(UdpClient udpClient, CancellationToken token) { - var initMessage = new NetworkMessageHeader + var initMessage = new NetworkMessage { MessageType = EMessageType.UntrustedConnect, Id = Guid.NewGuid(), Data = BitConverter.GetBytes(_channelInteractionModule.ConnectionId) diff --git a/mROA/Implementation/NetworkMessageHeader.cs b/mROA/Implementation/NetworkMessage.cs similarity index 50% rename from mROA/Implementation/NetworkMessageHeader.cs rename to mROA/Implementation/NetworkMessage.cs index 39715d5..6053b3e 100644 --- a/mROA/Implementation/NetworkMessageHeader.cs +++ b/mROA/Implementation/NetworkMessage.cs @@ -3,9 +3,9 @@ using mROA.Abstract; namespace mROA.Implementation { - public class NetworkMessageHeader + public class NetworkMessage { - private bool Equals(NetworkMessageHeader other) + private bool Equals(NetworkMessage other) { return Id.Equals(other.Id) && MessageType == other.MessageType; } @@ -14,15 +14,15 @@ namespace mROA.Implementation { if (obj is null) return false; if (ReferenceEquals(this, obj)) return true; - return obj.GetType() == GetType() && Equals((NetworkMessageHeader)obj); + return obj.GetType() == GetType() && Equals((NetworkMessage)obj); } - public NetworkMessageHeader() + public NetworkMessage() { Data = Array.Empty(); } - public NetworkMessageHeader(IContextualSerializationToolKit serializationToolkit, + public NetworkMessage(IContextualSerializationToolKit serializationToolkit, INetworkMessage networkMessage, IEndPointContext? context) { MessageType = networkMessage.MessageType; @@ -40,5 +40,30 @@ namespace mROA.Implementation { return $" {Id}:{MessageType} [{Data.Length}]"; } + + public struct NetworkMessageMeta + { + public byte Type; + public Guid Id; + public ushort BodyLength; + + public NetworkMessageMeta(ReadOnlySpan metadata) + { + Type = metadata[0]; + Id = new Guid(metadata[1..17]); + BodyLength = BitConverter.ToUInt16(metadata[17..]); + } + + public NetworkMessage ToMessage(ReadOnlySpan memory) + { + var data = memory[19..][..BodyLength]; + return new NetworkMessage + { + Data = data.ToArray(), + Id = Id, + MessageType = (EMessageType)Type + }; + } + } } } \ No newline at end of file diff --git a/mROA/Implementation/RepresentationModule.cs b/mROA/Implementation/RepresentationModule.cs index ccb20b8..586cb7e 100644 --- a/mROA/Implementation/RepresentationModule.cs +++ b/mROA/Implementation/RepresentationModule.cs @@ -28,8 +28,8 @@ namespace mROA.Implementation public IEndPointContext Context => _interaction.Context; public async Task<(object? Deserialized, EMessageType MessageType)> GetSingle( - Predicate rule, IEndPointContext? context, - CancellationToken token = default, params Func[] converter) + Predicate rule, IEndPointContext? context, + CancellationToken token = default, params Func[] converter) { var writer = _interaction.ReceiveChanel.Writer; var reader = _interaction.ReceiveChanel.Reader; @@ -52,9 +52,9 @@ namespace mROA.Implementation } public async IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream( - Predicate rule, IEndPointContext? context, + Predicate rule, IEndPointContext? context, [EnumeratorCancellation] CancellationToken token = default, - params Func[] converter) + params Func[] converter) { var writer = _interaction.ReceiveChanel.Writer; await foreach (var message in _interaction.ReceiveChanel.Reader.ReadAllAsync(token)) @@ -83,7 +83,7 @@ namespace mROA.Implementation IEndPointContext? context) where T : notnull { var serialized = _serialization.Serialize(payload, context); - await _interaction.PostMessageAsync(new NetworkMessageHeader + await _interaction.PostMessageAsync(new NetworkMessage { Id = id, MessageType = eMessageType, Data = serialized }); } @@ -97,7 +97,7 @@ namespace mROA.Implementation IEndPointContext? context) where T : notnull { var serialized = _serialization.Serialize(payload, context); - await _interaction.PostMessageUntrustedAsync(new NetworkMessageHeader + await _interaction.PostMessageUntrustedAsync(new NetworkMessage { Id = id, MessageType = eMessageType, Data = serialized }); } } diff --git a/mROA/Implementation/RequestId.cs b/mROA/Implementation/RequestId.cs new file mode 100644 index 0000000..e4b49f7 --- /dev/null +++ b/mROA/Implementation/RequestId.cs @@ -0,0 +1,18 @@ +using System; + +namespace mROA.Implementation +{ + public struct RequestId + { + public ulong P0; + public ulong P1; + public RequestId Generate() + { + var guid = Guid.NewGuid(); + var bytes = guid.ToByteArray(); + var high = BitConverter.ToUInt64(bytes, 0); + var low = BitConverter.ToUInt64(bytes, 8); + return new RequestId{ P0 = high, P1 = low}; + } + } +} \ No newline at end of file From bb3c4887df827501d50653c9b54700cc8b4c4666 Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Sun, 27 Jul 2025 00:14:15 +0300 Subject: [PATCH 2/4] Works, but slow for many connections --- Example.Load/Program.cs | 98 ++++++++++--------- mROA.Cbor/CborSerializationToolkit.cs | 10 +- mROA.Cbor/IOrdinaryStructureParser.cs | 8 +- mROA.Codegen/RemoteTypeBinder.cstmpl | 2 +- mROA.Test/Identifier.cs | 24 +++++ mROA/Abstract/ICancellationRepository.cs | 7 +- mROA/Abstract/ICommandExecution.cs | 2 +- mROA/Abstract/IRepresentationModule.cs | 6 +- .../Backend/BasicExecutionModule.cs | 6 +- .../Backend/NetworkGatewayModule.cs | 4 +- mROA/Implementation/CallRequest.cs | 6 +- mROA/Implementation/CancellationRepository.cs | 8 +- .../ChannelInteractionModule.cs | 41 ++++---- .../CommandExecution/AsyncCommandExecution.cs | 2 +- .../ExceptionCommandExecution.cs | 2 +- .../CommandExecution/FinalCommandExecution.cs | 4 +- .../Frontend/NetworkFrontendBridge.cs | 4 +- .../Frontend/RemoteException.cs | 2 +- .../Frontend/UdpUntrustedInteraction.cs | 2 +- mROA/Implementation/NetworkMessage.cs | 25 ++--- mROA/Implementation/RemoteObjectBase.cs | 6 +- mROA/Implementation/RepresentationModule.cs | 6 +- mROA/Implementation/RequestContext.cs | 4 +- mROA/Implementation/RequestId.cs | 60 ++++++++++-- 24 files changed, 209 insertions(+), 130 deletions(-) diff --git a/Example.Load/Program.cs b/Example.Load/Program.cs index 295fb0f..160e454 100644 --- a/Example.Load/Program.cs +++ b/Example.Load/Program.cs @@ -35,65 +35,75 @@ File.AppendAllText("results.txt", $"[SINGLE CBOR WRITER ALLOC] {totalRequests}\r async Task> GetLoadEndpoints(int count) { - var loads = new List(); - for (int i = 0; i < count; i++) + try { - var builder = Host.CreateApplicationBuilder(new HostApplicationBuilderSettings { DisableDefaults = true }); - builder.Services.AddSingleton(); - builder.Services.AddSingleton(); - builder.Services.AddSingleton(provider => + var loads = new List(); + for (int i = 0; i < count; i++) { - var repo = new InstanceRepository(provider.GetService()); - repo.FillSingletons(typeof(Program).Assembly); - return repo; - }); + Console.WriteLine($"Initializing {i}"); + var builder = Host.CreateApplicationBuilder(new HostApplicationBuilderSettings { DisableDefaults = true }); + builder.Services.AddSingleton(); + builder.Services.AddSingleton(); + builder.Services.AddSingleton(provider => + { + var repo = new InstanceRepository(provider.GetService()); + repo.FillSingletons(typeof(Program).Assembly); + return repo; + }); - builder.Services.AddSingleton(); - builder.Services.AddSingleton(); - // builder.Services.AddSingleton(); - builder.Services.AddSingleton(); - var serverEndPoint = new IPEndPoint(IPAddress.Loopback, 4567); - builder.Services.AddSingleton(); - builder.Services.AddOptions(); - builder.Services.Configure(options => options.Endpoint = serverEndPoint); - builder.Services.AddSingleton(); - // builder.Services.AddSingleton(); - builder.Services.AddSingleton(); + builder.Services.AddSingleton(); + builder.Services.AddSingleton(); + // builder.Services.AddSingleton(); + builder.Services.AddSingleton(); + var serverEndPoint = new IPEndPoint(IPAddress.Loopback, 4567); + builder.Services.AddSingleton(); + builder.Services.AddOptions(); + builder.Services.Configure(options => options.Endpoint = serverEndPoint); + builder.Services.AddSingleton(); + // builder.Services.AddSingleton(); + builder.Services.AddSingleton(); - builder.Services.AddSingleton(p => - { - var methodRepo = new CollectableMethodRepository(); - methodRepo.AppendInvokers(new GeneratedInvokersCollection()); - return methodRepo; - }); - builder.Services.AddSingleton(); - builder.Services.AddSingleton(); + builder.Services.AddSingleton(p => + { + var methodRepo = new CollectableMethodRepository(); + methodRepo.AppendInvokers(new GeneratedInvokersCollection()); + return methodRepo; + }); + builder.Services.AddSingleton(); + builder.Services.AddSingleton(); - var app = builder.Build(); + var app = builder.Build(); - var frontendBridge = app.Services.GetService()!; - await frontendBridge.Connect(); - // _ = app.Services.GetService()!.StartExtraction(); - // _ = app.Services.GetService().Start(serverEndPoint); - var context = app.Services.GetService(); + Console.WriteLine($"Connecting {i}"); + var frontendBridge = app.Services.GetService()!; + await frontendBridge.Connect(); + // _ = app.Services.GetService()!.StartExtraction(); + // _ = app.Services.GetService().Start(serverEndPoint); + var context = app.Services.GetService(); + Console.WriteLine($"Connected {i}"); + var singletonObject = + context.GetSingletonObject( + app.Services.GetService()); + loads.Add(singletonObject); + } - var singletonObject = - context.GetSingletonObject( - app.Services.GetService()); - loads.Add(singletonObject); + return loads; + } + catch (Exception e) + { + Console.WriteLine(e); + throw; } - return loads; } async Task Requests(CancellationToken token, int id, ILoadTest load) { try { - - int count = 0; - while (true){ + while (true) + { if (token.IsCancellationRequested) { break; @@ -101,7 +111,7 @@ async Task Requests(CancellationToken token, int id, ILoadTest load) await load.Next(2); count++; - } + } Console.WriteLine(id); return count; diff --git a/mROA.Cbor/CborSerializationToolkit.cs b/mROA.Cbor/CborSerializationToolkit.cs index f50fe7c..0506e65 100644 --- a/mROA.Cbor/CborSerializationToolkit.cs +++ b/mROA.Cbor/CborSerializationToolkit.cs @@ -116,9 +116,9 @@ namespace mROA.Cbor return preParsed.ToObject(type, context); - if (type == typeof(Guid)) + if (type == typeof(RequestId)) { - return new Guid((byte[])nonCasted); + return new RequestId((byte[])nonCasted); } return Convert.ChangeType(nonCasted, type); @@ -164,7 +164,7 @@ namespace mROA.Cbor case DateTimeOffset dto: writer.WriteDateTimeOffset(dto); break; - case Guid g: + case RequestId g: writer.WriteByteString(g.ToByteArray()); break; case byte[] bytes: @@ -270,8 +270,8 @@ namespace mROA.Cbor return reader.ReadUInt64(); case CborReaderState.ByteString: - if (type == typeof(Guid)) - return new Guid(reader.ReadByteString()); + if (type == typeof(RequestId)) + return new RequestId(reader.ReadByteString()); return reader.ReadByteString(); case CborReaderState.TextString: return reader.ReadTextString(); diff --git a/mROA.Cbor/IOrdinaryStructureParser.cs b/mROA.Cbor/IOrdinaryStructureParser.cs index 2d3711a..10ff38d 100644 --- a/mROA.Cbor/IOrdinaryStructureParser.cs +++ b/mROA.Cbor/IOrdinaryStructureParser.cs @@ -29,7 +29,7 @@ namespace mROA.Cbor reader.ReadStartArray(); var value = new NetworkMessage { - Id = new Guid(reader.ReadByteString()), + Id = new RequestId(reader.ReadByteString()), MessageType = (EMessageType)reader.ReadInt32(), Data = reader.ReadByteString() }; @@ -57,7 +57,7 @@ namespace mROA.Cbor reader.ReadStartArray(); var value = new DefaultCallRequest { - Id = new Guid(reader.ReadByteString()), + Id = new RequestId(reader.ReadByteString()), CommandId = reader.ReadInt32(), ObjectId = (ComplexObjectIdentifier)ComplexObjectIdentifierParser.Instance.Read(reader, context, serialization), Parameters = serialization.ReadData(reader, typeof(object[]), context) as object[] @@ -102,7 +102,7 @@ namespace mROA.Cbor reader.ReadStartArray(); var result = new FinalCommandExecution { - Id = new Guid(reader.ReadByteString()), + Id = new RequestId(reader.ReadByteString()), Result = serialization.ReadData(reader, typeof(object), context), }; reader.ReadEndArray(); @@ -125,7 +125,7 @@ namespace mROA.Cbor reader.ReadStartArray(); var result = new FinalCommandExecution { - Id = new Guid(reader.ReadByteString()) + Id = new RequestId(reader.ReadByteString()) }; reader.ReadEndArray(); return result; diff --git a/mROA.Codegen/RemoteTypeBinder.cstmpl b/mROA.Codegen/RemoteTypeBinder.cstmpl index 8735464..78df6c3 100644 --- a/mROA.Codegen/RemoteTypeBinder.cstmpl +++ b/mROA.Codegen/RemoteTypeBinder.cstmpl @@ -31,7 +31,7 @@ namespace mROA.Codegen Console.WriteLine("Sending event..."); var request = new DefaultCallRequest { - Id = Guid.NewGuid(), + Id = RequestId.Generate(), CommandId = , ObjectId = new ComplexObjectIdentifier(index, ownerId), Parameters = new object[] { } diff --git a/mROA.Test/Identifier.cs b/mROA.Test/Identifier.cs index 6b929bc..333ca33 100644 --- a/mROA.Test/Identifier.cs +++ b/mROA.Test/Identifier.cs @@ -20,4 +20,28 @@ public class Identifier Assert.Fail(); } } + + [Test] + public void RequestIdTest() + { + var id = RequestId.Generate(); + Assert.Pass(id.ToString()); + } + + [Test] + public void EqualsTest() + { + var id = RequestId.Generate(); + var id2 = new RequestId { P0 = id.P0, P1 = id.P1 }; + Assert.That(id2, Is.EqualTo(id)); + } + + [Test] + public void ByteString() + { + var id = RequestId.Generate(); + var binary = id.ToByteArray(); + var reverced = new RequestId(binary); + Assert.That(reverced, Is.EqualTo(id)); + } } \ No newline at end of file diff --git a/mROA/Abstract/ICancellationRepository.cs b/mROA/Abstract/ICancellationRepository.cs index 22189cb..502d1ea 100644 --- a/mROA/Abstract/ICancellationRepository.cs +++ b/mROA/Abstract/ICancellationRepository.cs @@ -1,12 +1,13 @@ using System; using System.Threading; +using mROA.Implementation; namespace mROA.Abstract { public interface ICancellationRepository { - void RegisterCancellation(Guid id, CancellationTokenSource cts); - CancellationTokenSource? GetCancellation(Guid id); - void FreeCancelation(Guid id); + void RegisterCancellation(RequestId id, CancellationTokenSource cts); + CancellationTokenSource? GetCancellation(RequestId id); + void FreeCancellation(RequestId id); } } \ No newline at end of file diff --git a/mROA/Abstract/ICommandExecution.cs b/mROA/Abstract/ICommandExecution.cs index e546c68..cf4ec0b 100644 --- a/mROA/Abstract/ICommandExecution.cs +++ b/mROA/Abstract/ICommandExecution.cs @@ -5,6 +5,6 @@ namespace mROA.Abstract { public interface ICommandExecution : INetworkMessage { - Guid Id { get; set; } + RequestId Id { get; set; } } } \ No newline at end of file diff --git a/mROA/Abstract/IRepresentationModule.cs b/mROA/Abstract/IRepresentationModule.cs index 45d44d3..3446380 100644 --- a/mROA/Abstract/IRepresentationModule.cs +++ b/mROA/Abstract/IRepresentationModule.cs @@ -19,13 +19,13 @@ namespace mROA.Abstract IEndPointContext? context, CancellationToken token = default, params Func[] converter); - Task PostCallMessageAsync(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context) + Task PostCallMessageAsync(RequestId id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull; - void PostCallMessage(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context) + void PostCallMessage(RequestId id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull; - Task PostCallMessageUntrustedAsync(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context) + Task PostCallMessageUntrustedAsync(RequestId id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull; } } \ No newline at end of file diff --git a/mROA/Implementation/Backend/BasicExecutionModule.cs b/mROA/Implementation/Backend/BasicExecutionModule.cs index 91c5a1d..d72ec7c 100644 --- a/mROA/Implementation/Backend/BasicExecutionModule.cs +++ b/mROA/Implementation/Backend/BasicExecutionModule.cs @@ -114,7 +114,7 @@ namespace mROA.Implementation.Backend if (cts == null) throw new NullReferenceException("Can't find cancellation for this request"); cts.Cancel(); - _cancellationRepo.FreeCancelation(command.Id); + _cancellationRepo.FreeCancellation(command.Id); return new FinalCommandExecution { @@ -164,7 +164,7 @@ namespace mROA.Implementation.Backend { Id = command.Id }; - _cancellationRepo.FreeCancelation(command.Id); + _cancellationRepo.FreeCancellation(command.Id); if (invoker.IsTrusted) @@ -195,7 +195,7 @@ namespace mROA.Implementation.Backend Id = command.Id, Result = finalResult }; - _cancellationRepo.FreeCancelation(command.Id); + _cancellationRepo.FreeCancellation(command.Id); representationModule.PostCallMessage(command.Id, EMessageType.FinishedCommandExecution, payload, context); diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index 6d4bdfb..e569185 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -56,7 +56,7 @@ namespace mROA.Implementation.Backend while (true) { var client = await _tcpListener.AcceptTcpClientAsync(); - _ = HandleConnection(client).ConfigureAwait(false); + _ = HandleConnection(client); } } @@ -70,7 +70,7 @@ namespace mROA.Implementation.Backend CallIndexProvider = _callIndexProvider }; var streamExtractor = - new ChannelInteractionModule.StreamExtractor(client.GetStream(), _serialization, context); + new ChannelInteractionModule.StreamExtractor(client.GetStream()); interaction.IsConnected = () => streamExtractor.IsConnected; streamExtractor.MessageReceived = async message => { diff --git a/mROA/Implementation/CallRequest.cs b/mROA/Implementation/CallRequest.cs index cf4834b..fcf9131 100644 --- a/mROA/Implementation/CallRequest.cs +++ b/mROA/Implementation/CallRequest.cs @@ -4,7 +4,7 @@ namespace mROA.Implementation { public interface ICallRequest { - Guid Id { get; } + RequestId Id { get; } int CommandId { get; } ComplexObjectIdentifier ObjectId { get; } object?[]? Parameters { get; } @@ -12,7 +12,7 @@ namespace mROA.Implementation public struct DefaultCallRequest : ICallRequest { - public Guid Id { get; set; } + public RequestId Id { get; set; } public int CommandId { get; set; } public ComplexObjectIdentifier ObjectId { get; set; } @@ -26,7 +26,7 @@ namespace mROA.Implementation public class CancelRequest : ICallRequest { - public Guid Id { get; set; } + public RequestId Id { get; set; } public int CommandId { get; set; } = -2; public ComplexObjectIdentifier ObjectId { get; set; } = ComplexObjectIdentifier.Null; public object?[]? Parameters { get; set; } = null; diff --git a/mROA/Implementation/CancellationRepository.cs b/mROA/Implementation/CancellationRepository.cs index 92db694..9b6be1b 100644 --- a/mROA/Implementation/CancellationRepository.cs +++ b/mROA/Implementation/CancellationRepository.cs @@ -8,19 +8,19 @@ namespace mROA.Implementation { public class CancellationRepository : ICancellationRepository { - private readonly ConcurrentDictionary _cancellations = new(); + private readonly ConcurrentDictionary _cancellations = new(); - public void RegisterCancellation(Guid id, CancellationTokenSource cts) + public void RegisterCancellation(RequestId id, CancellationTokenSource cts) { _cancellations.TryAdd(id, cts); } - public CancellationTokenSource? GetCancellation(Guid id) + public CancellationTokenSource? GetCancellation(RequestId id) { return _cancellations.GetValueOrDefault(id); } - public void FreeCancelation(Guid id) + public void FreeCancellation(RequestId id) { _cancellations.Remove(id, out _); } diff --git a/mROA/Implementation/ChannelInteractionModule.cs b/mROA/Implementation/ChannelInteractionModule.cs index d084ad4..63b4778 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -1,5 +1,6 @@ using System; using System.IO; +using System.Runtime.InteropServices; using System.Threading; using System.Threading.Channels; using System.Threading.Tasks; @@ -141,39 +142,34 @@ namespace mROA.Implementation public class StreamExtractor { - private const int BufferSize = ushort.MaxValue + 2; + private const int BufferSize = ushort.MaxValue + 19; private readonly Stream _ioStream; - private readonly IContextualSerializationToolKit _serializationToolkit; private readonly Memory _buffer = new byte[BufferSize]; - private readonly IEndPointContext _context; - private readonly byte[] _lenBuffer; - public StreamExtractor(Stream ioStream, IContextualSerializationToolKit serializationToolkit, - IEndPointContext context) + public StreamExtractor(Stream ioStream) { _ioStream = ioStream; - _serializationToolkit = serializationToolkit; - _context = context; - _lenBuffer = new byte[2]; } public Action MessageReceived = _ => { }; public async Task SingleReceive(CancellationToken token = default) { - int firstRead = _ioStream.Read(_buffer.Span); - - var metadata = new NetworkMessage.NetworkMessageMeta(_buffer.Span); + var firstRead = _ioStream.Read(_buffer.Span); - var localSpan = _buffer[2..len]; - if (firstRead - 2 != len) + var meta = MemoryMarshal.Read(_buffer.Span); + + var len = meta.BodyLength; + var readLen = firstRead - 19; + + if (readLen != len) { - localSpan = _buffer[(len + 2)..]; - await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: token); + var lastPart = _buffer[firstRead..(len + 19)]; + await _ioStream.ReadExactlyAsync(lastPart, cancellationToken: token); } - - var message = _serializationToolkit.Deserialize(localSpan, _context); + + var message = meta.ToMessage(_buffer.Span); MessageReceived(message); } @@ -187,11 +183,10 @@ namespace mROA.Implementation private async Task Send(NetworkMessage message, CancellationToken token = default) { - var bodySpan = _buffer[2..]; - var len = _serializationToolkit.Serialize(message, bodySpan.Span, _context); - var header = BitConverter.GetBytes((ushort)len); - header.CopyTo(_buffer); - var sendingSpan = _buffer[..(len + 2)]; + var meta = message.ToMeta(); + MemoryMarshal.Write(_buffer.Span, ref meta); + message.Data.CopyTo(_buffer.Span[19..]); + var sendingSpan = _buffer[..(19 + meta.BodyLength)]; await _ioStream.WriteAsync(sendingSpan, token); // _logger.LogTrace("SEND {0}", message.ToString()); } diff --git a/mROA/Implementation/CommandExecution/AsyncCommandExecution.cs b/mROA/Implementation/CommandExecution/AsyncCommandExecution.cs index deb3299..768668a 100644 --- a/mROA/Implementation/CommandExecution/AsyncCommandExecution.cs +++ b/mROA/Implementation/CommandExecution/AsyncCommandExecution.cs @@ -5,7 +5,7 @@ namespace mROA.Implementation.CommandExecution { public class AsyncCommandExecution : ICommandExecution { - public Guid Id { get; set; } + public RequestId Id { get; set; } public EMessageType MessageType => EMessageType.Unknown; } } \ No newline at end of file diff --git a/mROA/Implementation/CommandExecution/ExceptionCommandExecution.cs b/mROA/Implementation/CommandExecution/ExceptionCommandExecution.cs index 3199141..29d9144 100644 --- a/mROA/Implementation/CommandExecution/ExceptionCommandExecution.cs +++ b/mROA/Implementation/CommandExecution/ExceptionCommandExecution.cs @@ -6,7 +6,7 @@ namespace mROA.Implementation.CommandExecution { public class ExceptionCommandExecution : ICommandExecution { - public Guid Id { get; set; } + public RequestId Id { get; set; } public EMessageType MessageType => EMessageType.ExceptionCommandExecution; public string Exception { get; set; } diff --git a/mROA/Implementation/CommandExecution/FinalCommandExecution.cs b/mROA/Implementation/CommandExecution/FinalCommandExecution.cs index 5911109..cef5580 100644 --- a/mROA/Implementation/CommandExecution/FinalCommandExecution.cs +++ b/mROA/Implementation/CommandExecution/FinalCommandExecution.cs @@ -5,13 +5,13 @@ namespace mROA.Implementation.CommandExecution { public struct FinalCommandExecution : ICommandExecution { - public Guid Id { get; set; } + public RequestId Id { get; set; } public EMessageType MessageType => EMessageType.FinishedCommandExecution; } public struct FinalCommandExecution : ICommandExecution { - public Guid Id { get; set; } + public RequestId Id { get; set; } public EMessageType MessageType => EMessageType.FinishedCommandExecution; public T? Result { get; set; } } diff --git a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs index cfd7f4a..5a8572b 100644 --- a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs +++ b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs @@ -28,7 +28,7 @@ namespace mROA.Implementation.Frontend _serialization = serialization; _interactionModule = interactionModule; _rawExtractorCancellation = new CancellationTokenSource(); - _currentExtractor = new ChannelInteractionModule.StreamExtractor(Stream.Null, _serialization, context); + _currentExtractor = new ChannelInteractionModule.StreamExtractor(Stream.Null); } public async Task Connect() @@ -63,7 +63,7 @@ namespace mROA.Implementation.Frontend private void PrepareExtractor() { _currentExtractor = - new ChannelInteractionModule.StreamExtractor(_tcpClient.GetStream(), _serialization, _context); + new ChannelInteractionModule.StreamExtractor(_tcpClient.GetStream()); _ = _currentExtractor.SendFromChannel(_interactionModule.TrustedPostChanel, _rawExtractorCancellation.Token); diff --git a/mROA/Implementation/Frontend/RemoteException.cs b/mROA/Implementation/Frontend/RemoteException.cs index e63e8b3..873c836 100644 --- a/mROA/Implementation/Frontend/RemoteException.cs +++ b/mROA/Implementation/Frontend/RemoteException.cs @@ -4,7 +4,7 @@ namespace mROA.Implementation.Frontend { public class RemoteException : Exception { - public Guid CallRequestId; + public RequestId CallRequestId; private readonly string _error; public RemoteException(string error) diff --git a/mROA/Implementation/Frontend/UdpUntrustedInteraction.cs b/mROA/Implementation/Frontend/UdpUntrustedInteraction.cs index 13a6bce..c3a63ff 100644 --- a/mROA/Implementation/Frontend/UdpUntrustedInteraction.cs +++ b/mROA/Implementation/Frontend/UdpUntrustedInteraction.cs @@ -54,7 +54,7 @@ namespace mROA.Implementation.Frontend { var initMessage = new NetworkMessage { - MessageType = EMessageType.UntrustedConnect, Id = Guid.NewGuid(), + MessageType = EMessageType.UntrustedConnect, Id = RequestId.Generate(), Data = BitConverter.GetBytes(_channelInteractionModule.ConnectionId) }; diff --git a/mROA/Implementation/NetworkMessage.cs b/mROA/Implementation/NetworkMessage.cs index 6053b3e..3a7e5ed 100644 --- a/mROA/Implementation/NetworkMessage.cs +++ b/mROA/Implementation/NetworkMessage.cs @@ -27,10 +27,10 @@ namespace mROA.Implementation { MessageType = networkMessage.MessageType; Data = serializationToolkit.Serialize(networkMessage, context); - Id = Guid.NewGuid(); + Id = RequestId.Generate(); } - public Guid Id { get; set; } + public RequestId Id { get; set; } public EMessageType MessageType { get; set; } @@ -40,19 +40,22 @@ namespace mROA.Implementation { return $" {Id}:{MessageType} [{Data.Length}]"; } + + public NetworkMessageMeta ToMeta() + { + return new NetworkMessageMeta + { + BodyLength = (ushort)(Data == null ? 0 : Data.Length), + Type = (byte)MessageType, + Id = Id + }; + } public struct NetworkMessageMeta { - public byte Type; - public Guid Id; + public RequestId Id; public ushort BodyLength; - - public NetworkMessageMeta(ReadOnlySpan metadata) - { - Type = metadata[0]; - Id = new Guid(metadata[1..17]); - BodyLength = BitConverter.ToUInt16(metadata[17..]); - } + public byte Type; public NetworkMessage ToMessage(ReadOnlySpan memory) { diff --git a/mROA/Implementation/RemoteObjectBase.cs b/mROA/Implementation/RemoteObjectBase.cs index bc70cf1..9968aea 100644 --- a/mROA/Implementation/RemoteObjectBase.cs +++ b/mROA/Implementation/RemoteObjectBase.cs @@ -57,7 +57,7 @@ namespace mROA.Implementation { var request = new DefaultCallRequest { - Id = Guid.NewGuid(), CommandId = methodId, ObjectId = _identifier, Parameters = parameters + Id = RequestId.Generate(), CommandId = methodId, ObjectId = _identifier, Parameters = parameters }; await _representationModule.PostCallMessageAsync(request.Id, EMessageType.CallRequest, request, _context); @@ -102,7 +102,7 @@ namespace mROA.Implementation { var request = new DefaultCallRequest { - Id = Guid.NewGuid(), CommandId = methodId, ObjectId = _identifier, Parameters = parameters + Id = RequestId.Generate(), CommandId = methodId, ObjectId = _identifier, Parameters = parameters }; await _representationModule.PostCallMessageAsync(request.Id, EMessageType.CallRequest, request, _context); @@ -149,7 +149,7 @@ namespace mROA.Implementation { var request = new DefaultCallRequest { - Id = Guid.NewGuid(), CommandId = methodId, ObjectId = _identifier, Parameters = parameters + Id = RequestId.Generate(), CommandId = methodId, ObjectId = _identifier, Parameters = parameters }; await _representationModule.PostCallMessageUntrustedAsync(request.Id, EMessageType.CallRequest, request, _context); diff --git a/mROA/Implementation/RepresentationModule.cs b/mROA/Implementation/RepresentationModule.cs index 586cb7e..3452277 100644 --- a/mROA/Implementation/RepresentationModule.cs +++ b/mROA/Implementation/RepresentationModule.cs @@ -79,7 +79,7 @@ namespace mROA.Implementation } } - public async Task PostCallMessageAsync(Guid id, EMessageType eMessageType, T payload, + public async Task PostCallMessageAsync(RequestId id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull { var serialized = _serialization.Serialize(payload, context); @@ -87,13 +87,13 @@ namespace mROA.Implementation { Id = id, MessageType = eMessageType, Data = serialized }); } - public void PostCallMessage(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context) + public void PostCallMessage(RequestId id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull { PostCallMessageAsync(id, eMessageType, payload, context).GetAwaiter().GetResult(); } - public async Task PostCallMessageUntrustedAsync(Guid id, EMessageType eMessageType, T payload, + public async Task PostCallMessageUntrustedAsync(RequestId id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull { var serialized = _serialization.Serialize(payload, context); diff --git a/mROA/Implementation/RequestContext.cs b/mROA/Implementation/RequestContext.cs index 81db0d6..8ac4cae 100644 --- a/mROA/Implementation/RequestContext.cs +++ b/mROA/Implementation/RequestContext.cs @@ -5,9 +5,9 @@ namespace mROA.Implementation public struct RequestContext { public int OwnerId { get; } - public Guid RequestId { get; } + public RequestId RequestId { get; } - public RequestContext(Guid requestId, int ownerId) + public RequestContext(RequestId requestId, int ownerId) { RequestId = requestId; OwnerId = ownerId; diff --git a/mROA/Implementation/RequestId.cs b/mROA/Implementation/RequestId.cs index e4b49f7..3f19867 100644 --- a/mROA/Implementation/RequestId.cs +++ b/mROA/Implementation/RequestId.cs @@ -1,18 +1,64 @@ using System; +using System.Collections.Generic; +using System.Runtime.InteropServices; namespace mROA.Implementation { - public struct RequestId + public struct RequestId : IEquatable { + private static Random _random = new(); + public ulong P0; public ulong P1; - public RequestId Generate() + + public static RequestId Generate() { - var guid = Guid.NewGuid(); - var bytes = guid.ToByteArray(); - var high = BitConverter.ToUInt64(bytes, 0); - var low = BitConverter.ToUInt64(bytes, 8); - return new RequestId{ P0 = high, P1 = low}; + var high = (ulong)(ushort)_random.Next() << 32 | (ulong)_random.Next(); + var low = (ulong)(ushort)_random.Next() << 32 | (ulong)_random.Next(); + return new RequestId { P0 = high, P1 = low }; + } + + public RequestId(byte[] bytes) + { + P0 = BitConverter.ToUInt64(bytes); + P1 = BitConverter.ToUInt64(bytes, 8); + } + + public bool Equals(RequestId other) + { + return P0 == other.P0 && P1 == other.P1; + } + + public override bool Equals(object? obj) + { + return obj is RequestId other && Equals(other); + } + + public static bool operator ==(RequestId r1, RequestId r2) + { + return r1.Equals(r2); + } + + public static bool operator !=(RequestId r1, RequestId r2) + { + return !(r1 == r2); + } + + public override string ToString() + { + return $"{P0:X}{P1:X}"; + } + + public override int GetHashCode() + { + return HashCode.Combine(P0, P1); + } + + public byte[] ToByteArray() + { + var array = new byte[16]; + MemoryMarshal.Write(array, ref this); + return array; } } } \ No newline at end of file From 6c55567d3f345d3d34e3f69e7783def7741b0f42 Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Sun, 27 Jul 2025 11:52:16 +0300 Subject: [PATCH 3/4] fast new message header --- Example.Load/Program.cs | 2 +- mROA/Implementation/ChannelInteractionModule.cs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/Example.Load/Program.cs b/Example.Load/Program.cs index 160e454..a0d0b94 100644 --- a/Example.Load/Program.cs +++ b/Example.Load/Program.cs @@ -31,7 +31,7 @@ Console.WriteLine("End waiting"); var totalRequests = tasks.Sum(i => i.Result); Console.WriteLine($"Total requests: {totalRequests:N0}"); Console.WriteLine($"Results: {totalRequests / time.TotalSeconds:N} RPS"); -File.AppendAllText("results.txt", $"[SINGLE CBOR WRITER ALLOC] {totalRequests}\r\n"); +File.AppendAllText("results.txt", $"[HEADER REMAKE] {totalRequests}\r\n"); async Task> GetLoadEndpoints(int count) { diff --git a/mROA/Implementation/ChannelInteractionModule.cs b/mROA/Implementation/ChannelInteractionModule.cs index 63b4778..6b41ec0 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -156,7 +156,7 @@ namespace mROA.Implementation public async Task SingleReceive(CancellationToken token = default) { - var firstRead = _ioStream.Read(_buffer.Span); + var firstRead = await _ioStream.ReadAsync(_buffer, token); var meta = MemoryMarshal.Read(_buffer.Span); From 965c8eb8238d6c4172205141db7b7133a5d63e00 Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Sun, 27 Jul 2025 23:46:54 +0300 Subject: [PATCH 4/4] serialization improves and benchmarks --- Example.Load/Program.cs | 2 +- mROA.Benchmark/CborTest.cs | 93 +++++++++++++++++++ mROA.Benchmark/Program.cs | 75 +++++++-------- mROA.Benchmark/RequestWriter.cs | 57 ++++++++++++ mROA.Benchmark/mROA.Benchmark.csproj | 17 +++- mROA.Cbor/CborExtensions.cs | 19 ++++ mROA.Cbor/CborSerializationToolkit.cs | 17 ++-- mROA.Cbor/IOrdinaryStructureParser.cs | 34 ++----- mROA.Cbor/mROA.Cbor.csproj | 1 + mROA.sln | 12 +-- .../Backend/BasicExecutionModule.cs | 4 +- .../Backend/MultiClientInstanceRepository.cs | 1 + mROA/Implementation/NetworkMessage.cs | 5 +- 13 files changed, 243 insertions(+), 94 deletions(-) create mode 100644 mROA.Benchmark/CborTest.cs create mode 100644 mROA.Benchmark/RequestWriter.cs create mode 100644 mROA.Cbor/CborExtensions.cs diff --git a/Example.Load/Program.cs b/Example.Load/Program.cs index a0d0b94..280f5a4 100644 --- a/Example.Load/Program.cs +++ b/Example.Load/Program.cs @@ -31,7 +31,7 @@ Console.WriteLine("End waiting"); var totalRequests = tasks.Sum(i => i.Result); Console.WriteLine($"Total requests: {totalRequests:N0}"); Console.WriteLine($"Results: {totalRequests / time.TotalSeconds:N} RPS"); -File.AppendAllText("results.txt", $"[HEADER REMAKE] {totalRequests}\r\n"); +File.AppendAllText("results.txt", $"[FAST ID] {totalRequests}\r\n"); async Task> GetLoadEndpoints(int count) { diff --git a/mROA.Benchmark/CborTest.cs b/mROA.Benchmark/CborTest.cs new file mode 100644 index 0000000..d66ba1e --- /dev/null +++ b/mROA.Benchmark/CborTest.cs @@ -0,0 +1,93 @@ +using System.Formats.Cbor; +using BenchmarkDotNet.Attributes; + +namespace mROA.Benchmark; + + +[MemoryDiagnoser] +public class CborTest +{ + private CborWriter _writer; + private CborReader _reader; + + private Memory FlatEncoded; + private Memory ArrayEncoded; + public CborTest() + { + _writer = new CborWriter(initialCapacity:512); + + _writer.WriteStartArray(2); + _writer.WriteByteString([1, 2, 3, 4, 5, 6, 7, 8]); + _writer.WriteStartArray(2); + _writer.WriteInt32(12); + _writer.WriteTextString("tralala"); + _writer.WriteEndArray(); + _writer.WriteEndArray(); + + ArrayEncoded = _writer.Encode(); + + _writer.Reset(); + + _writer.WriteStartArray(3); + _writer.WriteByteString([1, 2, 3, 4, 5, 6, 7, 8]); + _writer.WriteInt32(12); + _writer.WriteTextString("tralala"); + _writer.WriteEndArray(); + FlatEncoded = _writer.Encode(); + } + + [Benchmark] + public int FlatWrite() + { + _writer.Reset(); + _writer.WriteStartArray(3); + _writer.WriteByteString([1, 2, 3, 4, 5, 6, 7, 8]); + _writer.WriteInt32(12); + _writer.WriteTextString("tralala"); + _writer.WriteEndArray(); + var len = _writer.Encode(FlatEncoded.Span); + return len; + } + + [Benchmark] + public int ArrayWrite() + { + _writer.Reset(); + _writer.WriteStartArray(2); + _writer.WriteByteString([1, 2, 3, 4, 5, 6, 7, 8]); + _writer.WriteStartArray(2); + _writer.WriteInt32(12); + _writer.WriteTextString("tralala"); + _writer.WriteEndArray(); + _writer.WriteEndArray(); + var len = _writer.Encode(ArrayEncoded.Span); + return len; + } + [Benchmark] + public int FlatRead() + { + var reader = new CborReader(FlatEncoded); + reader.ReadStartArray(); + var arr = reader.ReadByteString(); + var i = reader.ReadInt32(); + var text = reader.ReadTextString(); + reader.ReadEndArray(); + + return arr.Length; + } + + [Benchmark] + public int ArrayRead() + { + var reader = new CborReader(ArrayEncoded); + reader.ReadStartArray(); + var arr = reader.ReadByteString(); + reader.ReadStartArray(); + var i = reader.ReadInt32(); + var text = reader.ReadTextString(); + reader.ReadEndArray(); + reader.ReadEndArray(); + + return arr.Length; + } +} \ No newline at end of file diff --git a/mROA.Benchmark/Program.cs b/mROA.Benchmark/Program.cs index 217a6d7..1f030d2 100644 --- a/mROA.Benchmark/Program.cs +++ b/mROA.Benchmark/Program.cs @@ -1,51 +1,44 @@ -using System.Collections.Generic; -using System.Linq; -using BenchmarkDotNet.Attributes; +// See https://aka.ms/new-console-template for more information -namespace mROA.Benchmark +using System.Formats.Cbor; +using System.Runtime.CompilerServices; +using System.Runtime.InteropServices; +using BenchmarkDotNet.Running; +using mROA.Benchmark; +using mROA.Implementation; + +Console.WriteLine("Hello, World!"); +var test = new CborTest(); +test.ArrayWrite(); +BenchmarkRunner.Run(); + +public static class CborExtensions { - class Program + public static unsafe void WriteToCbor(this RequestId id, CborWriter writer) { - static void Main(string[] args) - { - // Console.WriteLine("Hello, World!"); - // var summary = BenchmarkRunner.Run(); - } + Span span = stackalloc byte[16]; + MemoryMarshal.Write(span, ref id); + writer.WriteByteString(span); } - public class CollectionsSpeed + [MethodImpl(MethodImplOptions.AggressiveInlining)] + public static unsafe void WriteToCborInline(this RequestId id, CborWriter writer) { - private const int N = 1000; + Span span = stackalloc byte[16]; + MemoryMarshal.Write(span, ref id); + writer.WriteByteString(span); + } - private readonly List _immutable; - private readonly int[] _array; - - public CollectionsSpeed() - { - _array = Enumerable.Range(0, N).ToArray(); - // _immutable = [.._array]; - } - - [Benchmark] - public int DefaultArray() - { - var sum = 0; - for (int i = 0; i < N; i++) - { - sum += _array[i]; - } - return sum; - } + public static void WriteToDest(this RequestId id, Span destination) + { + MemoryMarshal.Write(destination, ref id); + } - [Benchmark] - public int ImmutableArray() - { - var sum = 0; - for (int i = 0; i < N; i++) - { - sum += _immutable[i]; - } - return sum; - } + [MethodImpl(MethodImplOptions.AggressiveOptimization)] + public static unsafe void WriteToCborOpt(this RequestId id, CborWriter writer) + { + Span span = stackalloc byte[16]; + MemoryMarshal.Write(span, ref id); + writer.WriteByteString(span); } } \ No newline at end of file diff --git a/mROA.Benchmark/RequestWriter.cs b/mROA.Benchmark/RequestWriter.cs new file mode 100644 index 0000000..4ccd89f --- /dev/null +++ b/mROA.Benchmark/RequestWriter.cs @@ -0,0 +1,57 @@ +using System.Formats.Cbor; +using BenchmarkDotNet.Attributes; +using mROA.Implementation; + +[MemoryDiagnoser] +public class RequestWriter +{ + private const int N = 1000; + public RequestId Id = RequestId.Generate(); + private CborWriter _writer; + + public RequestWriter() + { + _writer = new CborWriter(initialCapacity: 512); + } + + [Benchmark(Baseline = true)] + public int DefaultCbor() + { + _writer.Reset(); + _writer.WriteByteString(Id.ToByteArray()); + return _writer.BytesWritten; + } + [Benchmark] + public int DirectCbor() + { + _writer.Reset(); + Id.WriteToCbor(_writer); + return _writer.BytesWritten; + } + + [Benchmark] + public int Stackalloc() + { + _writer.Reset(); + Span span = stackalloc byte[16]; + Id.WriteToDest(span); + _writer.WriteByteString(span); + return _writer.BytesWritten; + } + + [Benchmark] + public int DirectCborInline() + { + _writer.Reset(); + Id.WriteToCborInline(_writer); + return _writer.BytesWritten; + } + + [Benchmark] + public int DirectCborOpt() + { + _writer.Reset(); + Id.WriteToCborOpt(_writer); + return _writer.BytesWritten; + } +} \ No newline at end of file diff --git a/mROA.Benchmark/mROA.Benchmark.csproj b/mROA.Benchmark/mROA.Benchmark.csproj index 5697c66..a7f69a0 100644 --- a/mROA.Benchmark/mROA.Benchmark.csproj +++ b/mROA.Benchmark/mROA.Benchmark.csproj @@ -2,13 +2,24 @@ Exe - netstandard2.1 - + net9.0 + enable enable + true - + + + + + + ..\..\..\..\.nuget\packages\system.formats.cbor\9.0.7\lib\net9.0\System.Formats.Cbor.dll + + + + + diff --git a/mROA.Cbor/CborExtensions.cs b/mROA.Cbor/CborExtensions.cs new file mode 100644 index 0000000..1685ea3 --- /dev/null +++ b/mROA.Cbor/CborExtensions.cs @@ -0,0 +1,19 @@ +using System; +using System.Formats.Cbor; +using System.Runtime.CompilerServices; +using System.Runtime.InteropServices; +using mROA.Implementation; + +namespace mROA.Cbor +{ + public static class CborExtensions + { + [MethodImpl(MethodImplOptions.AggressiveInlining)] + public static unsafe void WriteToCborInline(this RequestId id, CborWriter writer) + { + Span span = stackalloc byte[16]; + MemoryMarshal.Write(span, ref id); + writer.WriteByteString(span); + } + } +} \ No newline at end of file diff --git a/mROA.Cbor/CborSerializationToolkit.cs b/mROA.Cbor/CborSerializationToolkit.cs index 0506e65..4036166 100644 --- a/mROA.Cbor/CborSerializationToolkit.cs +++ b/mROA.Cbor/CborSerializationToolkit.cs @@ -18,7 +18,7 @@ namespace mROA.Cbor private readonly IOrdinaryStructureParser[] _parsers = { - new NetworkMessageHeaderParser(), new DefaultCallRequestParser(), new FinalCommandExecutionParser(), + new DefaultCallRequestParser(), new FinalCommandExecutionParser(), new FinalCommandExecutionResultlessParser() }; @@ -27,27 +27,21 @@ namespace mROA.Cbor private bool FindParser(Type t, out IOrdinaryStructureParser parser) { - if (t == typeof(NetworkMessage)) + if (t == typeof(DefaultCallRequest)) { parser = _parsers[0]; return true; } - if (t == typeof(DefaultCallRequest)) + if (t == typeof(FinalCommandExecution)) { parser = _parsers[1]; return true; } - if (t == typeof(FinalCommandExecution)) - { - parser = _parsers[2]; - return true; - } - if (t == typeof(FinalCommandExecution)) { - parser = _parsers[3]; + parser = _parsers[2]; return true; } @@ -165,7 +159,8 @@ namespace mROA.Cbor writer.WriteDateTimeOffset(dto); break; case RequestId g: - writer.WriteByteString(g.ToByteArray()); + // writer.WriteByteString(g.ToByteArray()); + g.WriteToCborInline(writer); break; case byte[] bytes: writer.WriteByteString(bytes); diff --git a/mROA.Cbor/IOrdinaryStructureParser.cs b/mROA.Cbor/IOrdinaryStructureParser.cs index 10ff38d..e39e8e0 100644 --- a/mROA.Cbor/IOrdinaryStructureParser.cs +++ b/mROA.Cbor/IOrdinaryStructureParser.cs @@ -12,38 +12,14 @@ namespace mROA.Cbor object Read(CborReader reader, IEndPointContext context, CborSerializationToolkit serialization); } - public class NetworkMessageHeaderParser : IOrdinaryStructureParser - { - public void Write(CborWriter writer, object value, IEndPointContext context, CborSerializationToolkit serialization) - { - var v = value as NetworkMessage; - writer.WriteStartArray(3); - writer.WriteByteString(v.Id.ToByteArray()); - writer.WriteInt32((int)v.MessageType); - writer.WriteByteString(v.Data); - writer.WriteEndArray(); - } - - public object Read(CborReader reader, IEndPointContext context, CborSerializationToolkit serialization) - { - reader.ReadStartArray(); - var value = new NetworkMessage - { - Id = new RequestId(reader.ReadByteString()), - MessageType = (EMessageType)reader.ReadInt32(), - Data = reader.ReadByteString() - }; - return value; - } - } - public class DefaultCallRequestParser : IOrdinaryStructureParser { public void Write(CborWriter writer, object value, IEndPointContext context, CborSerializationToolkit serialization) { var v = (DefaultCallRequest)value; writer.WriteStartArray(4); - writer.WriteByteString(v.Id.ToByteArray()); + v.Id.WriteToCborInline(writer); + // writer.WriteByteString(v.Id.ToByteArray()); writer.WriteInt32(v.CommandId); writer.WriteStartArray(1); writer.WriteUInt64(v.ObjectId.Flat); @@ -92,7 +68,8 @@ namespace mROA.Cbor { var v = (FinalCommandExecution)value; writer.WriteStartArray(2); - writer.WriteByteString(v.Id.ToByteArray()); + // writer.WriteByteString(v.Id.ToByteArray()); + v.Id.WriteToCborInline(writer); serialization.WriteData(v.Result, writer, context); writer.WriteEndArray(); } @@ -116,7 +93,8 @@ namespace mROA.Cbor { var v = (FinalCommandExecution)value; writer.WriteStartArray(1); - writer.WriteByteString(v.Id.ToByteArray()); + // writer.WriteByteString(v.Id.ToByteArray()); + v.Id.WriteToCborInline(writer); writer.WriteEndArray(); } diff --git a/mROA.Cbor/mROA.Cbor.csproj b/mROA.Cbor/mROA.Cbor.csproj index 53687e0..6fe99f0 100644 --- a/mROA.Cbor/mROA.Cbor.csproj +++ b/mROA.Cbor/mROA.Cbor.csproj @@ -7,6 +7,7 @@ 2.0.7 9 mroaLogo.png + true diff --git a/mROA.sln b/mROA.sln index b6d9335..97ede35 100644 --- a/mROA.sln +++ b/mROA.sln @@ -17,8 +17,6 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Example.Shared", "Example.S EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Example.Frontend", "Example.Frontend\Example.Frontend.csproj", "{9BD25A13-3165-47C0-9EAA-5C59EC490E32}" EndProject -Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "mROA.Benchmark", "mROA.Benchmark\mROA.Benchmark.csproj", "{6868F42B-E30D-4040-AD4A-BC2A2E76D03A}" -EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "mROA.Cbor", "mROA.Cbor\mROA.Cbor.csproj", "{6211E4EA-13FD-4EB1-8E6A-C0173DD0784A}" EndProject Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "TotalDemo", "TotalDemo", "{FC4FA752-10A7-4D78-A7C2-9BCD8A81FB5E}" @@ -31,6 +29,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Functionality.Shared", "Fun EndProject Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Test", "Test", "{8C20901F-B416-4ABC-8AA4-9059646B081B}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "mROA.Benchmark", "mROA.Benchmark\mROA.Benchmark.csproj", "{E021E8B3-56C2-400E-A05E-523CF7831189}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -61,10 +61,6 @@ Global {9BD25A13-3165-47C0-9EAA-5C59EC490E32}.Debug|Any CPU.Build.0 = Debug|Any CPU {9BD25A13-3165-47C0-9EAA-5C59EC490E32}.Release|Any CPU.ActiveCfg = Release|Any CPU {9BD25A13-3165-47C0-9EAA-5C59EC490E32}.Release|Any CPU.Build.0 = Release|Any CPU - {6868F42B-E30D-4040-AD4A-BC2A2E76D03A}.Debug|Any CPU.ActiveCfg = Debug|Any CPU - {6868F42B-E30D-4040-AD4A-BC2A2E76D03A}.Debug|Any CPU.Build.0 = Debug|Any CPU - {6868F42B-E30D-4040-AD4A-BC2A2E76D03A}.Release|Any CPU.ActiveCfg = Release|Any CPU - {6868F42B-E30D-4040-AD4A-BC2A2E76D03A}.Release|Any CPU.Build.0 = Release|Any CPU {6211E4EA-13FD-4EB1-8E6A-C0173DD0784A}.Debug|Any CPU.ActiveCfg = Debug|Any CPU {6211E4EA-13FD-4EB1-8E6A-C0173DD0784A}.Debug|Any CPU.Build.0 = Debug|Any CPU {6211E4EA-13FD-4EB1-8E6A-C0173DD0784A}.Release|Any CPU.ActiveCfg = Release|Any CPU @@ -81,6 +77,10 @@ Global {D9D28596-E10C-4A98-A2AA-573219467506}.Debug|Any CPU.Build.0 = Debug|Any CPU {D9D28596-E10C-4A98-A2AA-573219467506}.Release|Any CPU.ActiveCfg = Release|Any CPU {D9D28596-E10C-4A98-A2AA-573219467506}.Release|Any CPU.Build.0 = Release|Any CPU + {E021E8B3-56C2-400E-A05E-523CF7831189}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {E021E8B3-56C2-400E-A05E-523CF7831189}.Debug|Any CPU.Build.0 = Debug|Any CPU + {E021E8B3-56C2-400E-A05E-523CF7831189}.Release|Any CPU.ActiveCfg = Release|Any CPU + {E021E8B3-56C2-400E-A05E-523CF7831189}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE diff --git a/mROA/Implementation/Backend/BasicExecutionModule.cs b/mROA/Implementation/Backend/BasicExecutionModule.cs index d72ec7c..53b8937 100644 --- a/mROA/Implementation/Backend/BasicExecutionModule.cs +++ b/mROA/Implementation/Backend/BasicExecutionModule.cs @@ -168,7 +168,7 @@ namespace mROA.Implementation.Backend if (invoker.IsTrusted) - representationModule.PostCallMessage(command.Id, EMessageType.FinishedCommandExecution, + representationModule.PostCallMessageAsync(command.Id, EMessageType.FinishedCommandExecution, payload, context); }); @@ -197,7 +197,7 @@ namespace mROA.Implementation.Backend }; _cancellationRepo.FreeCancellation(command.Id); - representationModule.PostCallMessage(command.Id, EMessageType.FinishedCommandExecution, + representationModule.PostCallMessageAsync(command.Id, EMessageType.FinishedCommandExecution, payload, context); }); diff --git a/mROA/Implementation/Backend/MultiClientInstanceRepository.cs b/mROA/Implementation/Backend/MultiClientInstanceRepository.cs index 0488e21..f5d6cd2 100644 --- a/mROA/Implementation/Backend/MultiClientInstanceRepository.cs +++ b/mROA/Implementation/Backend/MultiClientInstanceRepository.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Collections.Specialized; using mROA.Abstract; namespace mROA.Implementation.Backend diff --git a/mROA/Implementation/NetworkMessage.cs b/mROA/Implementation/NetworkMessage.cs index 3a7e5ed..719bd42 100644 --- a/mROA/Implementation/NetworkMessage.cs +++ b/mROA/Implementation/NetworkMessage.cs @@ -35,7 +35,8 @@ namespace mROA.Implementation public EMessageType MessageType { get; set; } public byte[] Data { get; set; } - + public object Serialized { get; set; } + public IEndPointContext Context { get; set; } public override string ToString() { return $" {Id}:{MessageType} [{Data.Length}]"; @@ -57,7 +58,7 @@ namespace mROA.Implementation public ushort BodyLength; public byte Type; - public NetworkMessage ToMessage(ReadOnlySpan memory) + public NetworkMessage ToMessage(Span memory) { var data = memory[19..][..BodyLength]; return new NetworkMessage