From abaac48246dc0f460fc7dfc037efa0c3344cd587 Mon Sep 17 00:00:00 2001 From: Mikhail Mitrofanov Date: Sun, 6 Apr 2025 20:51:12 +0300 Subject: [PATCH] Reconnection base --- Example.Frontend/Program.cs | 3 +- mROA.Test/NextGenTest.cs | 2 +- mROA/Abstract/IFrontendBridge.cs | 7 +- mROA/Abstract/IInteractionModule.cs | 4 +- .../Backend/NetworkGatewayModule.cs | 4 +- mROA/Implementation/ClientDisconnect.cs | 7 ++ mROA/Implementation/EMessageType.cs | 3 +- .../Frontend/NetworkFrontendBridge.cs | 38 +++++++++-- mROA/Implementation/IdAssignment.cs | 5 ++ .../NextGenerationInteractionModule.cs | 64 +++++++++++++------ mROA/Implementation/RepresentationModule.cs | 2 +- 11 files changed, 106 insertions(+), 33 deletions(-) create mode 100644 mROA/Implementation/ClientDisconnect.cs diff --git a/Example.Frontend/Program.cs b/Example.Frontend/Program.cs index 156026d..75291dc 100644 --- a/Example.Frontend/Program.cs +++ b/Example.Frontend/Program.cs @@ -5,6 +5,7 @@ using System.Threading; using System.Threading.Tasks; using Example.Frontend; using Example.Shared; +using mROA.Abstract; using mROA.Cbor; using mROA.Codegen; using mROA.Implementation; @@ -38,7 +39,7 @@ class Program TransmissionConfig.RealContextRepository = builder.GetModule(); TransmissionConfig.RemoteEndpointContextRepository = builder.GetModule(); - builder.GetModule()!.Connect(); + builder.GetModule()!.Connect(); _ = builder.GetModule()!.StartExtraction(); Console.WriteLine(TransmissionConfig.OwnershipRepository.GetOwnershipId()); var context = builder.GetModule(); diff --git a/mROA.Test/NextGenTest.cs b/mROA.Test/NextGenTest.cs index 02c2fe1..31e2347 100644 --- a/mROA.Test/NextGenTest.cs +++ b/mROA.Test/NextGenTest.cs @@ -37,7 +37,7 @@ namespace mROA.Test foreach (var guid in guids) { - _interactionModuleB.PostMessage(new NetworkMessageHeader { Id = guid, Data = "Hello user"u8.ToArray() }); + _interactionModuleB.PostMessageAsync(new NetworkMessageHeader { Id = guid, Data = "Hello user"u8.ToArray() }); } }); diff --git a/mROA/Abstract/IFrontendBridge.cs b/mROA/Abstract/IFrontendBridge.cs index 909f486..53217d7 100644 --- a/mROA/Abstract/IFrontendBridge.cs +++ b/mROA/Abstract/IFrontendBridge.cs @@ -1,6 +1,11 @@ +using System; + namespace mROA.Abstract { - public interface IFrontendBridge : IInjectableModule + public interface IFrontendBridge : IInjectableModule, IDisposable { + void Connect(); + void Obstacle(); + void Disconnect(); } } \ No newline at end of file diff --git a/mROA/Abstract/IInteractionModule.cs b/mROA/Abstract/IInteractionModule.cs index 45ee3c7..c622df6 100644 --- a/mROA/Abstract/IInteractionModule.cs +++ b/mROA/Abstract/IInteractionModule.cs @@ -10,9 +10,11 @@ namespace mROA.Abstract int ConnectionId { get; set; } public Stream? BaseStream { get; set; } Task GetNextMessageReceiving(); - Task PostMessage(NetworkMessageHeader messageHeader); + Task PostMessageAsync(NetworkMessageHeader messageHeader); void HandleMessage(NetworkMessageHeader messageHeader); NetworkMessageHeader[] UnhandledMessages { get; } NetworkMessageHeader? FirstByFilter(Predicate predicate); + event Action OnDisconected; + } } \ No newline at end of file diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index d33b4d0..9573178 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -79,7 +79,7 @@ namespace mROA.Implementation.Backend if (connectionRequest.MessageType == EMessageType.ClientConnect) { - interaction.PostMessage(new NetworkMessageHeader(_serialization!, + interaction.PostMessageAsync(new NetworkMessageHeader(_serialization!, new IdAssignment { Id = -interaction.ConnectionId })); _hub!.RegisterInteraction(interaction); Console.WriteLine("Client registered"); @@ -98,7 +98,7 @@ namespace mROA.Implementation.Backend interaction.BaseStream = client.GetStream(); - interaction.PostMessage(new NetworkMessageHeader(_serialization!, + interaction.PostMessageAsync(new NetworkMessageHeader(_serialization!, new IdAssignment { Id = -interaction.ConnectionId })); _hub!.RegisterInteraction(interaction); Console.WriteLine("Client registered"); diff --git a/mROA/Implementation/ClientDisconnect.cs b/mROA/Implementation/ClientDisconnect.cs new file mode 100644 index 0000000..805e17d --- /dev/null +++ b/mROA/Implementation/ClientDisconnect.cs @@ -0,0 +1,7 @@ +namespace mROA.Implementation +{ + public class ClientDisconnect : INetworkMessage + { + public EMessageType MessageType => EMessageType.ClientDisconnect; + } +} \ No newline at end of file diff --git a/mROA/Implementation/EMessageType.cs b/mROA/Implementation/EMessageType.cs index 8bb556b..bfb6304 100644 --- a/mROA/Implementation/EMessageType.cs +++ b/mROA/Implementation/EMessageType.cs @@ -10,6 +10,7 @@ namespace mROA.Implementation CancelRequest, EventRequest, ClientRecovery, - ClientConnect + ClientConnect, + ClientDisconnect, } } \ No newline at end of file diff --git a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs index 2fb5cc6..526fdaa 100644 --- a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs +++ b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs @@ -1,6 +1,7 @@ using System; using System.Net; using System.Net.Sockets; +using System.Threading.Tasks; using mROA.Abstract; using Exception = System.Exception; @@ -39,21 +40,50 @@ namespace mROA.Implementation.Frontend throw new NullReferenceException("Serialization toolkit is not initialized"); _tcpClient.Connect(_ipEndPoint); + _interactionModule.BaseStream = _tcpClient.GetStream(); - _ = _interactionModule.PostMessage(new NetworkMessageHeader(_serialization, new ClientConnect())); + _interactionModule.OnDisconected += async id => + { + await Reconect(); + }; + _ = _interactionModule.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientConnect())); var welcomeMessage = _interactionModule.GetNextMessageReceiving().GetAwaiter().GetResult(); if (welcomeMessage.MessageType != EMessageType.IdAssigning) { throw new Exception( $"Incorrect message type. Must be IdAssigning, current : {welcomeMessage.MessageType.ToString()}"); } - - - + + var assignment = _serialization.Deserialize(welcomeMessage.Data)!; _interactionModule.ConnectionId = -assignment.Id; TransmissionConfig.OwnershipRepository = new StaticOwnershipRepository(assignment.Id); } + + private async Task Reconect() + { + _tcpClient.Connect(_ipEndPoint); + _interactionModule.BaseStream = _tcpClient.GetStream(); + await _interactionModule.Restart(); + } + + public void Obstacle() + { + _interactionModule!.BaseStream!.Close(); + _tcpClient.Close(); + } + + public void Disconnect() + { + _ = _interactionModule!.PostMessageAsync(new NetworkMessageHeader(_serialization!, new ClientDisconnect())); + _interactionModule.BaseStream!.Close(); + _tcpClient.Dispose(); + } + + public void Dispose() + { + Disconnect(); + } } } \ No newline at end of file diff --git a/mROA/Implementation/IdAssignment.cs b/mROA/Implementation/IdAssignment.cs index 924eb94..f3b3f09 100644 --- a/mROA/Implementation/IdAssignment.cs +++ b/mROA/Implementation/IdAssignment.cs @@ -9,6 +9,11 @@ namespace mROA.Implementation public class ClientRecovery : INetworkMessage { + public ClientRecovery(int id) + { + Id = id; + } + public int Id { get; set; } public EMessageType MessageType => EMessageType.ClientRecovery; } diff --git a/mROA/Implementation/NextGenerationInteractionModule.cs b/mROA/Implementation/NextGenerationInteractionModule.cs index 3bb89df..c27fa52 100644 --- a/mROA/Implementation/NextGenerationInteractionModule.cs +++ b/mROA/Implementation/NextGenerationInteractionModule.cs @@ -2,6 +2,7 @@ using System.Collections.Generic; using System.IO; using System.Linq; +using System.Net.Sockets; using System.Threading.Tasks; using mROA.Abstract; @@ -15,6 +16,9 @@ namespace mROA.Implementation private Task? _currentReceiving; private ISerializationToolkit? _serialization; private Stream? _baseStream; + + private TaskCompletionSource _reconection = new(); + public int ConnectionId { get; set; } public Stream? BaseStream @@ -24,8 +28,8 @@ namespace mROA.Implementation { if (_baseStream is null) { - } + _baseStream = value; } } @@ -51,7 +55,7 @@ namespace mROA.Implementation return _currentReceiving; } - public async Task PostMessage(NetworkMessageHeader messageHeader) + public async Task PostMessageAsync(NetworkMessageHeader messageHeader) { if (BaseStream == null) throw new NullReferenceException("BaseStream is null"); @@ -81,6 +85,8 @@ namespace mROA.Implementation return _messageBuffer.FirstOrDefault(m => predicate(m)); } + public event Action? OnDisconected; + private async Task GetNextMessage() { if (BaseStream == null) @@ -89,37 +95,53 @@ namespace mROA.Implementation if (_serialization == null) throw new NullReferenceException("Serialization toolkit is null"); - - try + while (true) { - // Console.WriteLine("Receiving message"); - var firstBit = (byte)BaseStream.ReadByte(); - var secondBit = (byte)BaseStream.ReadByte(); + try + { + return await Receive(); + } + catch (Exception) + { + OnDisconected!.Invoke(ConnectionId); + _ = await _reconection.Task; + } + } + } - var len = BitConverter.ToUInt16(new[] { firstBit, secondBit }); - var localSpan = _buffer[..len]; + private ushort ReadMessageLength() + { + var firstBit = (byte)BaseStream.ReadByte(); + var secondBit = (byte)BaseStream.ReadByte(); - await BaseStream.ReadExactlyAsync(localSpan); + var len = BitConverter.ToUInt16(new[] { firstBit, secondBit }); - // Console.WriteLine("Receiving {0}", Encoding.Default.GetString(_buffer[..len])); + return len; + } - var message = _serialization.Deserialize(localSpan.Span); + private async Task Receive() + { + var len = ReadMessageLength(); + var localSpan = _buffer[..len]; + + await BaseStream.ReadExactlyAsync(localSpan); + + var message = _serialization.Deserialize(localSpan.Span); #if TRACE Console.WriteLine($"{DateTime.Now.TimeOfDay} Received Message {message.Id} - {message.SchemaId}"); TransmissionConfig.TotalTransmittedBytes += len; Console.WriteLine($"Total recieced bytes are {TransmissionConfig.TotalTransmittedBytes}"); #endif - _messageBuffer.Add(message); - _currentReceiving = Task.Run(async () => await GetNextMessage()); + _messageBuffer.Add(message); + _currentReceiving = Task.Run(async () => await GetNextMessage()); - return message; - } - catch (Exception e) - { - Console.WriteLine(e); - throw; - } + return message; + } + public async Task Restart() + { + await PostMessageAsync(new NetworkMessageHeader(_serialization!, new ClientRecovery(Math.Abs(ConnectionId)))); + _reconection.SetResult(BaseStream!); } } } \ No newline at end of file diff --git a/mROA/Implementation/RepresentationModule.cs b/mROA/Implementation/RepresentationModule.cs index e5c1d67..8b72ccf 100644 --- a/mROA/Implementation/RepresentationModule.cs +++ b/mROA/Implementation/RepresentationModule.cs @@ -97,7 +97,7 @@ namespace mROA.Implementation #endif var serialized = _serialization.Serialize(payload, payloadType); - await _interaction.PostMessage(new NetworkMessageHeader + await _interaction.PostMessageAsync(new NetworkMessageHeader { Id = id, MessageType = eMessageType, Data = serialized }); }