From 1ecca099d86cd057aa8ae5d7255f70bff667faf3 Mon Sep 17 00:00:00 2001 From: Mikhail Mitrofanov Date: Fri, 5 Jun 2026 09:24:44 +0300 Subject: [PATCH] fixing reading message crash --- .../Backend/NetworkGatewayModule.cs | 1 + .../ChannelInteractionModule.cs | 30 ++++++++++++------- 2 files changed, 20 insertions(+), 11 deletions(-) diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index ecbf97b..9b307d5 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -57,6 +57,7 @@ namespace mROA.Implementation.Backend private async Task HandleConnection(TcpClient client) { + client.NoDelay = true; Console.WriteLine($"Client connected from {client.Client.RemoteEndPoint}"); var interaction = new ChannelInteractionModule(_serialization, _identityGenerator); diff --git a/mROA/Implementation/ChannelInteractionModule.cs b/mROA/Implementation/ChannelInteractionModule.cs index 25f6d1f..4a742b7 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -145,30 +145,34 @@ namespace mROA.Implementation private const int BufferSize = ushort.MaxValue + 19; private readonly Stream _ioStream; private readonly Memory _buffer = new byte[BufferSize]; + public StreamExtractor(Stream ioStream) { _ioStream = ioStream; } + public Action MessageReceived = _ => { }; + public async Task SingleReceive(CancellationToken token = default) { - var firstRead = await _ioStream.ReadAsync(_buffer, token); - var meta = MemoryMarshal.Read(_buffer.Span); + _ = await _ioStream.ReadExactlyAsync(_buffer[..19], token); - var len = meta.BodyLength; - var readLen = firstRead - 19; + var meta = MemoryMarshal.Read(_buffer.Span); - if (readLen != len) - { - var lastPart = _buffer[firstRead..(len + 19)]; + var len = meta.BodyLength; + + var range = 19..(len + 19); + // Console.WriteLine($"{range} {_buffer.Length}"); + var lastPart = _buffer[range]; await _ioStream.ReadExactlyAsync(lastPart, cancellationToken: token); - } + + var message = meta.ToMessage(_buffer.Span); + // Console.WriteLine("RECV " + message); - var message = meta.ToMessage(_buffer.Span); - - MessageReceived(message); + MessageReceived(message); } + public async Task LoopedReceive(CancellationToken token = default) { while (token.IsCancellationRequested == false && IsConnected) @@ -176,14 +180,17 @@ namespace mROA.Implementation await SingleReceive(token); } } + private async Task Send(NetworkMessage message, CancellationToken token = default) { var meta = message.ToMeta(); MemoryMarshal.Write(_buffer.Span, ref meta); message.Data.CopyTo(_buffer.Span[19..]); var sendingSpan = _buffer[..(19 + meta.BodyLength)]; + // Console.WriteLine("SEND " + message); await _ioStream.WriteAsync(sendingSpan, token); } + public async Task SendFromChannel(ChannelReader channel, CancellationToken token = default) { @@ -193,6 +200,7 @@ namespace mROA.Implementation await Send(message, token); } } + public bool IsConnected => _ioStream is { CanRead: true, CanWrite: true }; } }