From 6812138eb18c8c07ada55778fcb324969847eff8 Mon Sep 17 00:00:00 2001 From: Mikhail Mitrofanov Date: Sat, 3 May 2025 14:32:21 +0300 Subject: [PATCH] Some strange IO Exception found --- Example.Frontend/Program.cs | 2 +- .../Backend/NetworkGatewayModule.cs | 18 +- .../ChannelInteractionModule.cs | 205 ++++++++++++------ .../Frontend/NetworkFrontendBridge.cs | 13 +- mROA/Implementation/StreamExtractor.cs | 95 -------- mROA/mROA.csproj | 2 +- 6 files changed, 154 insertions(+), 181 deletions(-) delete mode 100644 mROA/Implementation/StreamExtractor.cs diff --git a/Example.Frontend/Program.cs b/Example.Frontend/Program.cs index 8a5b163..cf79fd0 100644 --- a/Example.Frontend/Program.cs +++ b/Example.Frontend/Program.cs @@ -58,7 +58,7 @@ class Program Console.WriteLine("Printer created"); Thread.Sleep(100); - // frontendBridge.Obstacle(); + frontendBridge.Obstacle(); var name = disposingPrinter.GetName(); DemoCheck.BasicNonParamsCall = true; Console.WriteLine("Printer name : {0}", name); diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index b9f31af..1fe4169 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -75,17 +75,14 @@ namespace mROA.Implementation.Backend interaction!.Inject(_serialization); - - var streamExtractor = new StreamExtractor(client.GetStream(), _serialization); - interaction.IsConnected = () => streamExtractor.IsConnected; - streamExtractor.MessageReceived = message => - { - interaction.ReceiveChanel.Writer.WriteAsync(message); - }; + + var streamExtractor = new ChannelInteractionModule.StreamExtractor(client.GetStream(), _serialization); + interaction.IsConnected = () => streamExtractor.IsConnected; + streamExtractor.MessageReceived = message => { interaction.ReceiveChanel.Writer.WriteAsync(message); }; streamExtractor.SingleReceive(); var connectionRequest = interaction.GetNextMessageReceiving(false) .GetAwaiter().GetResult()!; - + switch (connectionRequest.MessageType) { case EMessageType.ClientConnect: @@ -101,10 +98,13 @@ namespace mROA.Implementation.Backend var recoveryRequest = _serialization!.Deserialize(connectionRequest.Data)!; var recoveryInteraction = _hub.GetInteraction(recoveryRequest.Id); + recoveryInteraction.IsConnected = () => streamExtractor.IsConnected; streamExtractor.MessageReceived = message => { recoveryInteraction.ReceiveChanel.Writer.WriteAsync(message); }; + Task.Run(async () => await streamExtractor.LoopedReceive()); + recoveryInteraction.Restart(false); Console.WriteLine("Connection recovery for client {0} finished", recoveryRequest.Id); @@ -116,7 +116,7 @@ namespace mROA.Implementation.Backend } } } - + private void ThrowIfNotInjected() { if (_hub is null) diff --git a/mROA/Implementation/ChannelInteractionModule.cs b/mROA/Implementation/ChannelInteractionModule.cs index 8671bf1..e706dfb 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -2,6 +2,7 @@ using System.Collections.Generic; using System.IO; using System.Linq; +using System.Threading; using System.Threading.Channels; using System.Threading.Tasks; using mROA.Abstract; @@ -28,24 +29,21 @@ namespace mROA.Implementation { SingleReader = false, SingleWriter = false, - // AllowSynchronousContinuations = true }); _receiveReader = ReceiveChanel.Reader; _outputTrustedChannel = Channel.CreateBounded(new BoundedChannelOptions(1) { SingleReader = true, SingleWriter = true, - // AllowSynchronousContinuations = true, - }); _trustedWriter = _outputTrustedChannel.Writer; _outputUntrustedChannel = Channel.CreateUnbounded(new UnboundedChannelOptions { SingleReader = true, SingleWriter = true, - // AllowSynchronousContinuations = true }); _untrustedWriter = _outputUntrustedChannel.Writer; + _reconnection = new TaskCompletionSource(); } public int ConnectionId { get; set; } @@ -84,7 +82,7 @@ namespace mROA.Implementation { return false; } - + await _trustedWriter.WriteAsync(messageHeader); return true; } @@ -125,13 +123,6 @@ namespace mROA.Implementation await _untrustedWriter.WriteAsync(messageHeader); } - public void HandleMessage(NetworkMessageHeader messageHeader) - { - _messageBuffer.Remove(messageHeader); - } - - // public NetworkMessageHeader[] UnhandledMessages => _messageBuffer.ToArray(); - public NetworkMessageHeader? FirstByFilter(Predicate predicate) { return _messageBuffer.FirstOrDefault(m => predicate(m)); @@ -139,44 +130,18 @@ namespace mROA.Implementation public event Action? OnDisconnected; - private async Task GetNextMessage() - { - if (_serialization == null) - throw new NullReferenceException("Serialization toolkit is null"); - - bool withError = false; - - while (true) - { - if (withError) - { - Console.WriteLine("Receive again"); - } - - try - { - var message = await _receiveReader.ReadAsync(); - return message; - } - catch (Exception) - { - if (!_isActive) - { - return NetworkMessageHeader.Null; - } - - withError = true; - await MakeRecovery("IN"); - } - } - } - public async Task Restart(bool sendRecovery) { if (sendRecovery) { await PostMessageAsync( new NetworkMessageHeader(_serialization!, new ClientRecovery(Math.Abs(ConnectionId)))); + var ping = await ReceiveChanel.Reader.ReadAsync(); + Console.WriteLine($"Ping received {ping.Id}"); + } + else + { + await _trustedWriter.WriteAsync(new NetworkMessageHeader()); } Console.WriteLine("Setting result for reconnection"); @@ -191,30 +156,25 @@ namespace mROA.Implementation private async Task MakeRecovery(string source) { //TODO переделать реконнект - - // Console.WriteLine("Staring recovery from {0}", source); - // - // lock (_reconnection) - // { - // Console.WriteLine("Got lock from {0}", source); - // - // Console.WriteLine("Call OnDisconnected from {0}", source); - // _isInReconnectionState = true; - // OnDisconnected?.Invoke(ConnectionId); - // } - // - // Console.WriteLine("Waiting for reconnect from {0}", source); - // if (!_reconnection.Task.IsCompleted && !_isConnected) - // { - // Console.WriteLine("Current connection state {0} from {1}", _isConnected, source); - // await _reconnection.Task; - // } - // - // Console.WriteLine("Reconnect finished from {0}", source); - // lock (_reconnection) - // { - // _isInReconnectionState = false; - // } + + Console.WriteLine("Staring recovery from {0}", source); + + lock (_reconnection) + { + Console.WriteLine("Got lock from {0}", source); + + Console.WriteLine("Call OnDisconnected from {0}", source); + OnDisconnected?.Invoke(ConnectionId); + } + + Console.WriteLine("Waiting for reconnect from {0}", source); + if (!_reconnection.Task.IsCompleted && !_isConnected) + { + Console.WriteLine("Current connection state {0} from {1}", _isConnected, source); + await _reconnection.Task; + } + + Console.WriteLine("Reconnect finished from {0}", source); } public void Dispose() @@ -226,5 +186,114 @@ namespace mROA.Implementation _currentReceiving?.Dispose(); } } + + public class StreamExtractor + { + private readonly Stream _ioStream; + private readonly ISerializationToolkit _serializationToolkit; + private const int BufferSize = ushort.MaxValue; + private readonly Memory _buffer = new byte[BufferSize]; + private bool _manualConnectionState = true; + + public readonly int Id = new Random().Next(); + + public StreamExtractor(Stream ioStream, ISerializationToolkit serializationToolkit) + { + _ioStream = ioStream; + _serializationToolkit = serializationToolkit; + } + + public Action MessageReceived = _ => { }; + + private ushort ReadMessageLength() + { + var firstBit = _ioStream.ReadByte(); + if (firstBit == -1) + { + _manualConnectionState = false; + throw new EndOfStreamException(); + } + + _manualConnectionState = true; + var secondBit = (byte)_ioStream.ReadByte(); + + var len = BitConverter.ToUInt16(new[] { (byte)firstBit, secondBit }); + + return len; + } + + public async Task SingleReceive(CancellationToken Token = default) + { +#if TRACE + Console.WriteLine($"[{Id}] Single receive started"); +#endif + var len = ReadMessageLength(); + var localSpan = _buffer[..len]; + + await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: Token); + + var message = _serializationToolkit.Deserialize(localSpan.Span); +#if TRACE + Console.WriteLine( + $"{DateTime.Now.TimeOfDay} [{Id}] Received Message {message.Id} - {message.MessageType}"); + TransmissionConfig.TotalTransmittedBytes += len; + Console.WriteLine($"Total received bytes are {TransmissionConfig.TotalTransmittedBytes}"); +#endif + MessageReceived(message); + } + + public async Task LoopedReceive(CancellationToken token = default) + { +#if TRACE + Console.WriteLine("LoopedReceive started"); +#endif + while (token.IsCancellationRequested == false && IsConnected) + { + await SingleReceive(token); + } + } + + public async Task Send(NetworkMessageHeader message, CancellationToken token = default) + { + try + { + var rawMessage = _serializationToolkit.Serialize(message); + var header = BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort)); + +#if TRACE + Console.WriteLine( + $"{DateTime.Now.TimeOfDay} [{Id}] Posting Message {message.Id} - {message.MessageType}"); + TransmissionConfig.TotalTransmittedBytes += rawMessage.Length; + Console.WriteLine($"Total received bytes are {TransmissionConfig.TotalTransmittedBytes}"); +#endif + + await _ioStream.WriteAsync(header, token); + await _ioStream.WriteAsync(rawMessage, token); +#if TRACE + Console.WriteLine( + $"{DateTime.Now.TimeOfDay} [{Id}] Posting finished {message.Id} - {message.MessageType}"); + +#endif + } + catch (Exception e) + { + Console.WriteLine(e); + throw; + } + } + + public async Task SendFromChannel(ChannelReader channel, + CancellationToken token = default) + { + while (token.IsCancellationRequested == false && IsConnected) + { + var message = await channel.ReadAsync(token); + await Send(message, token); + } + } + + + public bool IsConnected => _ioStream is { CanRead: true, CanWrite: true } && _manualConnectionState; + } } } \ No newline at end of file diff --git a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs index 4dab9de..fdbface 100644 --- a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs +++ b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs @@ -17,7 +17,7 @@ namespace mROA.Implementation.Frontend private TcpClient _tcpClient = new(); private IChannelInteractionModule? _interactionModule; private ISerializationToolkit? _serialization; - private StreamExtractor _currentExtractor; + private ChannelInteractionModule.StreamExtractor _currentExtractor; private CancellationTokenSource _rawExtractorCancellation; public NetworkFrontendBridge(IPEndPoint serverEndPoint) @@ -50,7 +50,7 @@ namespace mROA.Implementation.Frontend PrepareExtractor(); _interactionModule.IsConnected = () => _currentExtractor.IsConnected; - _interactionModule.OnDisconnected += id => { Reconnect(); }; + _interactionModule.OnDisconnected += _ => { Reconnect(); }; _interactionModule.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientConnect())).Wait(); @@ -64,8 +64,7 @@ namespace mROA.Implementation.Frontend } - var stopToken = _rawExtractorCancellation.Token; - Task.Run(async () => await _currentExtractor.LoopedReceive(stopToken)); + Task.Run(async () => await _currentExtractor.LoopedReceive(_rawExtractorCancellation.Token)); var assignment = _serialization.Deserialize(idMessage.Data)!; _interactionModule.ConnectionId = -assignment.Id; @@ -74,9 +73,9 @@ namespace mROA.Implementation.Frontend private void PrepareExtractor() { - _currentExtractor = new StreamExtractor(_tcpClient.GetStream(), _serialization); + _currentExtractor = new ChannelInteractionModule.StreamExtractor(_tcpClient.GetStream(), _serialization!); - _ = _currentExtractor.SendFromChannel(_interactionModule.TrustedPostChanel, + _ = _currentExtractor.SendFromChannel(_interactionModule!.TrustedPostChanel, _rawExtractorCancellation.Token); _currentExtractor.MessageReceived = message => { @@ -94,7 +93,7 @@ namespace mROA.Implementation.Frontend PrepareExtractor(); - _ = _currentExtractor.LoopedReceive(_rawExtractorCancellation.Token); + Task.Run(async () => await _currentExtractor.LoopedReceive(_rawExtractorCancellation.Token)); await _interactionModule.Restart(true); } diff --git a/mROA/Implementation/StreamExtractor.cs b/mROA/Implementation/StreamExtractor.cs deleted file mode 100644 index 5e748c0..0000000 --- a/mROA/Implementation/StreamExtractor.cs +++ /dev/null @@ -1,95 +0,0 @@ -using System; -using System.IO; -using System.Threading; -using System.Threading.Channels; -using System.Threading.Tasks; -using mROA.Abstract; - -namespace mROA.Implementation -{ - public class StreamExtractor - { - private readonly Stream _ioStream; - private readonly ISerializationToolkit _serializationToolkit; - private const int BufferSize = ushort.MaxValue; - private readonly Memory _buffer = new byte[BufferSize]; - private bool _manualConnectionState = true; - - public StreamExtractor(Stream ioStream, ISerializationToolkit serializationToolkit) - { - _ioStream = ioStream; - _serializationToolkit = serializationToolkit; - } - - public Action MessageReceived = _ => { }; - - private ushort ReadMessageLength() - { - var firstBit = _ioStream.ReadByte(); - if (firstBit == -1) - { - _manualConnectionState = false; - throw new EndOfStreamException(); - } - - _manualConnectionState = true; - var secondBit = (byte)_ioStream.ReadByte(); - - var len = BitConverter.ToUInt16(new[] { (byte)firstBit, secondBit }); - - return len; - } - - public async Task SingleReceive(CancellationToken Token = default) - { - var len = ReadMessageLength(); - var localSpan = _buffer[..len]; - - await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: Token); - - var message = _serializationToolkit.Deserialize(localSpan.Span); -#if TRACE - Console.WriteLine($"{DateTime.Now.TimeOfDay} Received Message {message.Id} - {message.MessageType}"); - TransmissionConfig.TotalTransmittedBytes += len; - Console.WriteLine($"Total received bytes are {TransmissionConfig.TotalTransmittedBytes}"); -#endif - MessageReceived(message); - } - - public async Task LoopedReceive(CancellationToken token = default) - { - while (token.IsCancellationRequested == false && IsConnected) - { - await SingleReceive(token); - } - } - - public async Task Send(NetworkMessageHeader message, CancellationToken token = default) - { - var rawMessage = _serializationToolkit.Serialize(message); - var header = BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort)); - -#if TRACE - Console.WriteLine($"{DateTime.Now.TimeOfDay} Posting Message {message.Id} - {message.MessageType}"); - TransmissionConfig.TotalTransmittedBytes += rawMessage.Length; - Console.WriteLine($"Total received bytes are {TransmissionConfig.TotalTransmittedBytes}"); -#endif - - await _ioStream.WriteAsync(header, token); - await _ioStream.WriteAsync(rawMessage, token); - } - - public async Task SendFromChannel(ChannelReader channel, - CancellationToken token = default) - { - while (token.IsCancellationRequested == false && IsConnected) - { - var message = await channel.ReadAsync(token); - await Send(message, token); - } - } - - - public bool IsConnected => _ioStream is { CanRead: true, CanWrite: true } && _manualConnectionState; - } -} \ No newline at end of file diff --git a/mROA/mROA.csproj b/mROA/mROA.csproj index 3c172de..e52150a 100644 --- a/mROA/mROA.csproj +++ b/mROA/mROA.csproj @@ -21,7 +21,7 @@ - + TRACE;