From ab40cad3ee0f242e112341c1f61ca696ba61fc21 Mon Sep 17 00:00:00 2001 From: Mikhail Mitrofanov Date: Tue, 8 Apr 2025 22:42:00 +0300 Subject: [PATCH] Reconnection DONE!!! --- mROA/Abstract/IInteractionModule.cs | 4 +- .../Backend/NetworkGatewayModule.cs | 13 +-- .../Frontend/NetworkFrontendBridge.cs | 4 +- mROA/Implementation/NetworkMessageHeader.cs | 1 + .../NextGenerationInteractionModule.cs | 107 +++++++----------- 5 files changed, 54 insertions(+), 75 deletions(-) diff --git a/mROA/Abstract/IInteractionModule.cs b/mROA/Abstract/IInteractionModule.cs index 0b5e657..5459266 100644 --- a/mROA/Abstract/IInteractionModule.cs +++ b/mROA/Abstract/IInteractionModule.cs @@ -9,13 +9,13 @@ namespace mROA.Abstract { int ConnectionId { get; set; } public Stream? BaseStream { get; set; } - Task GetNextMessageReceiving(); + public IntPtr StreamHandle { get; set; } + Task GetNextMessageReceiving(bool infinite = true); Task PostMessageAsync(NetworkMessageHeader messageHeader); void HandleMessage(NetworkMessageHeader messageHeader); NetworkMessageHeader[] UnhandledMessages { get; } NetworkMessageHeader? FirstByFilter(Predicate predicate); event Action OnDisconected; Task Restart(bool sendRecovery); - } } \ No newline at end of file diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index 8e7106c..303cffc 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -1,10 +1,8 @@ using System; -using System.Linq; using System.Net; using System.Net.Sockets; using System.Threading.Tasks; using mROA.Abstract; -using static System.Byte; namespace mROA.Implementation.Backend { @@ -67,17 +65,18 @@ namespace mROA.Implementation.Backend while (true) { var client = _tcpListener.AcceptTcpClient(); - Console.WriteLine($"Client connected from {client.Client.RemoteEndPoint}"); var interaction = Activator.CreateInstance(_interactionModuleType!) as INextGenerationInteractionModule; + foreach (var injectableModule in _injectableModules!) interaction!.Inject(injectableModule); interaction!.Inject(_serialization); - interaction.BaseStream = client.GetStream(); + interaction.StreamHandle = client.Client.Handle; - var connectionRequest = interaction.GetNextMessageReceiving().GetAwaiter().GetResult()!; + var connectionRequest = interaction.GetNextMessageReceiving(false) + .GetAwaiter().GetResult()!; switch (connectionRequest.MessageType) { @@ -92,10 +91,10 @@ namespace mROA.Implementation.Backend var recoveryRequest = _serialization!.Deserialize(connectionRequest.Data)!; var recoveryInteraction = _hub.GetInteraction(recoveryRequest.Id); recoveryInteraction.BaseStream = client.GetStream(); + recoveryInteraction.StreamHandle = client.Client.Handle; - recoveryInteraction.Restart(false); - Console.WriteLine($"Client {recoveryRequest.Id} reconnected"); + Console.WriteLine("Connection recovery for client {0} finished", recoveryRequest.Id); break; } default: diff --git a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs index 6aca2a1..97bba05 100644 --- a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs +++ b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs @@ -71,8 +71,8 @@ namespace mROA.Implementation.Frontend public void Obstacle() { - _interactionModule!.BaseStream!.Close(); - _tcpClient.Close(); + _interactionModule!.BaseStream!.Dispose(); + _tcpClient.Dispose(); } public void Disconnect() diff --git a/mROA/Implementation/NetworkMessageHeader.cs b/mROA/Implementation/NetworkMessageHeader.cs index 8efc141..8eb5149 100644 --- a/mROA/Implementation/NetworkMessageHeader.cs +++ b/mROA/Implementation/NetworkMessageHeader.cs @@ -18,6 +18,7 @@ namespace mROA.Implementation { MessageType = networkMessage.MessageType; Data = serializationToolkit.Serialize(networkMessage); + Id = Guid.NewGuid(); } public Guid Id { get; set; } diff --git a/mROA/Implementation/NextGenerationInteractionModule.cs b/mROA/Implementation/NextGenerationInteractionModule.cs index f3816c1..254457c 100644 --- a/mROA/Implementation/NextGenerationInteractionModule.cs +++ b/mROA/Implementation/NextGenerationInteractionModule.cs @@ -31,16 +31,11 @@ namespace mROA.Implementation public Stream? BaseStream { get => _baseStream; - set - { - if (_baseStream is null) - { - } - - _baseStream = value; - } + set => _baseStream = value; } + public IntPtr StreamHandle { get; set; } + public void Inject(T dependency) { @@ -55,25 +50,28 @@ namespace mROA.Implementation } } - public Task GetNextMessageReceiving() + public Task GetNextMessageReceiving(bool infinite = true) { + if (!infinite) return Receive().AsTask(); if (_currentReceiving != null) return _currentReceiving; _currentReceiving = Task.Run(async () => await GetNextMessage()); return _currentReceiving; + } #pragma warning disable CS8602 // Dereference of a possibly null reference. private async ValueTask PostMessageInternal(NetworkMessageHeader messageHeader) { - #if TRACE - Console.WriteLine($"{DateTime.Now.TimeOfDay} Posting message: {messageHeader.Id} - {messageHeader.MessageType} to {ConnectionId}"); + Console.WriteLine( + $"{DateTime.Now.TimeOfDay} Posting message to {StreamHandle}: {messageHeader.Id} - {messageHeader.MessageType} to {ConnectionId}"); #endif - + var rawMessage = _serialization.Serialize(messageHeader); var header = BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort)); if (!_baseStream.CanWrite) return false; + await BaseStream.WriteAsync(header); await BaseStream.WriteAsync(rawMessage); return true; @@ -103,24 +101,8 @@ namespace mROA.Implementation break; _isConnected = false; + withError = true; await MakeRecovery("OUT"); - // Console.WriteLine("Try to get lock from post"); - // lock (_reconection) - // { - // Console.WriteLine("Got lock from post"); - // if (!_isRecovering) - // { - // _isRecovering = true; - // Console.WriteLine("Disconnect invoke for post"); - // OnDisconected?.Invoke(ConnectionId); - // Console.WriteLine("Disconnect invoked for post"); - // } - // } - // - // withError = true; - // Console.WriteLine("Start waiting for recovery from post"); - // _ = await _reconection.Task; - // Console.WriteLine("Connection recovered from post"); } } @@ -152,42 +134,26 @@ namespace mROA.Implementation { if (withError) { - Console.WriteLine("Recieve again"); + Console.WriteLine("Receive again"); } try { - return await Receive(); + var message = await Receive(); + _currentReceiving = Task.Run(async () => await GetNextMessage()); + + return message; } catch (Exception ex) { + withError = true; await MakeRecovery("IN"); - // Console.WriteLine("Try to get lock from receive"); - // lock (_reconection) - // { - // Console.WriteLine("Got lock from receive"); - // - // if (!_isRecovering) - // { - // _isRecovering = true; - // Console.WriteLine("Disconnect invoke"); - // OnDisconected?.Invoke(ConnectionId); - // Console.WriteLine("Disconnect invoked for receive"); - // - // } - // } - // - // withError = true; - // Console.WriteLine("Start waiting for recovery from receive"); - // _ = await _reconection.Task; - // Console.WriteLine("Connection recovered"); } } } private ushort ReadMessageLength() { - var firstBit = BaseStream.ReadByte(); if (firstBit == -1) { @@ -212,13 +178,11 @@ namespace mROA.Implementation var message = _serialization.Deserialize(localSpan.Span); #if TRACE - Console.WriteLine($"{DateTime.Now.TimeOfDay} Received Message {message.Id} - {message.MessageType}"); + Console.WriteLine($"{DateTime.Now.TimeOfDay} Received Message from {StreamHandle} {message.Id} - {message.MessageType}"); TransmissionConfig.TotalTransmittedBytes += len; Console.WriteLine($"Total received bytes are {TransmissionConfig.TotalTransmittedBytes}"); #endif _messageBuffer.Add(message); - _currentReceiving = Task.Run(async () => await GetNextMessage()); - return message; } @@ -228,32 +192,42 @@ namespace mROA.Implementation { await PostMessageAsync( new NetworkMessageHeader(_serialization!, new ClientRecovery(Math.Abs(ConnectionId)))); - - var confirmByte = BaseStream.ReadByte(); - Console.WriteLine("Reconnection byte {0}", confirmByte); + var iTest = _baseStream.ReadByte(); + var bTest = (byte)iTest; + _baseStream.WriteByte(bTest); } else { - BaseStream.WriteByte(byte.MaxValue); + const byte confirmByte = 128; + _baseStream.WriteByte(confirmByte); + var iPong = _baseStream.ReadByte(); + var bPong = (byte)iPong; + if (confirmByte != bPong) + { + Console.WriteLine("Incorrect byte"); + } } Console.WriteLine("Setting result for reconnection"); - _reconnection.SetResult(BaseStream!); - Console.WriteLine("Set result for reconnection successfull"); + var setting = _reconnection.TrySetResult(BaseStream); + _isInReconnectionState = false; + _isConnected = true; + Console.WriteLine($"Set result for reconnection {setting}"); _reconnection = new TaskCompletionSource(); } private async Task MakeRecovery(string source) { - Console.WriteLine("Staring recovery from {0}", source); + lock (_reconnection) { Console.WriteLine("Got lock from {0}", source); if (_isConnected || _isInReconnectionState) { - Console.WriteLine($"{_isConnected} {_isInReconnectionState} {!_baseStream.CanRead} {!_baseStream.CanWrite}"); + Console.WriteLine( + $"{source} {_isConnected} {_isInReconnectionState} {!_baseStream.CanRead} {!_baseStream.CanWrite}"); return; } @@ -261,9 +235,14 @@ namespace mROA.Implementation _isInReconnectionState = true; OnDisconected?.Invoke(ConnectionId); } - + Console.WriteLine("Waiting for reconnect from {0}", source); - await _reconnection.Task; + 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) {