Reconnection base
This commit is contained in:
@@ -5,6 +5,7 @@ using System.Threading;
|
|||||||
using System.Threading.Tasks;
|
using System.Threading.Tasks;
|
||||||
using Example.Frontend;
|
using Example.Frontend;
|
||||||
using Example.Shared;
|
using Example.Shared;
|
||||||
|
using mROA.Abstract;
|
||||||
using mROA.Cbor;
|
using mROA.Cbor;
|
||||||
using mROA.Codegen;
|
using mROA.Codegen;
|
||||||
using mROA.Implementation;
|
using mROA.Implementation;
|
||||||
@@ -38,7 +39,7 @@ class Program
|
|||||||
TransmissionConfig.RealContextRepository = builder.GetModule<ContextRepository>();
|
TransmissionConfig.RealContextRepository = builder.GetModule<ContextRepository>();
|
||||||
TransmissionConfig.RemoteEndpointContextRepository = builder.GetModule<RemoteContextRepository>();
|
TransmissionConfig.RemoteEndpointContextRepository = builder.GetModule<RemoteContextRepository>();
|
||||||
|
|
||||||
builder.GetModule<NetworkFrontendBridge>()!.Connect();
|
builder.GetModule<IFrontendBridge>()!.Connect();
|
||||||
_ = builder.GetModule<RequestExtractor>()!.StartExtraction();
|
_ = builder.GetModule<RequestExtractor>()!.StartExtraction();
|
||||||
Console.WriteLine(TransmissionConfig.OwnershipRepository.GetOwnershipId());
|
Console.WriteLine(TransmissionConfig.OwnershipRepository.GetOwnershipId());
|
||||||
var context = builder.GetModule<RemoteContextRepository>();
|
var context = builder.GetModule<RemoteContextRepository>();
|
||||||
|
|||||||
@@ -37,7 +37,7 @@ namespace mROA.Test
|
|||||||
|
|
||||||
foreach (var guid in guids)
|
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() });
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,11 @@
|
|||||||
|
using System;
|
||||||
|
|
||||||
namespace mROA.Abstract
|
namespace mROA.Abstract
|
||||||
{
|
{
|
||||||
public interface IFrontendBridge : IInjectableModule
|
public interface IFrontendBridge : IInjectableModule, IDisposable
|
||||||
{
|
{
|
||||||
|
void Connect();
|
||||||
|
void Obstacle();
|
||||||
|
void Disconnect();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -10,9 +10,11 @@ namespace mROA.Abstract
|
|||||||
int ConnectionId { get; set; }
|
int ConnectionId { get; set; }
|
||||||
public Stream? BaseStream { get; set; }
|
public Stream? BaseStream { get; set; }
|
||||||
Task<NetworkMessageHeader> GetNextMessageReceiving();
|
Task<NetworkMessageHeader> GetNextMessageReceiving();
|
||||||
Task PostMessage(NetworkMessageHeader messageHeader);
|
Task PostMessageAsync(NetworkMessageHeader messageHeader);
|
||||||
void HandleMessage(NetworkMessageHeader messageHeader);
|
void HandleMessage(NetworkMessageHeader messageHeader);
|
||||||
NetworkMessageHeader[] UnhandledMessages { get; }
|
NetworkMessageHeader[] UnhandledMessages { get; }
|
||||||
NetworkMessageHeader? FirstByFilter(Predicate<NetworkMessageHeader> predicate);
|
NetworkMessageHeader? FirstByFilter(Predicate<NetworkMessageHeader> predicate);
|
||||||
|
event Action<int> OnDisconected;
|
||||||
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -79,7 +79,7 @@ namespace mROA.Implementation.Backend
|
|||||||
|
|
||||||
if (connectionRequest.MessageType == EMessageType.ClientConnect)
|
if (connectionRequest.MessageType == EMessageType.ClientConnect)
|
||||||
{
|
{
|
||||||
interaction.PostMessage(new NetworkMessageHeader(_serialization!,
|
interaction.PostMessageAsync(new NetworkMessageHeader(_serialization!,
|
||||||
new IdAssignment { Id = -interaction.ConnectionId }));
|
new IdAssignment { Id = -interaction.ConnectionId }));
|
||||||
_hub!.RegisterInteraction(interaction);
|
_hub!.RegisterInteraction(interaction);
|
||||||
Console.WriteLine("Client registered");
|
Console.WriteLine("Client registered");
|
||||||
@@ -98,7 +98,7 @@ namespace mROA.Implementation.Backend
|
|||||||
|
|
||||||
interaction.BaseStream = client.GetStream();
|
interaction.BaseStream = client.GetStream();
|
||||||
|
|
||||||
interaction.PostMessage(new NetworkMessageHeader(_serialization!,
|
interaction.PostMessageAsync(new NetworkMessageHeader(_serialization!,
|
||||||
new IdAssignment { Id = -interaction.ConnectionId }));
|
new IdAssignment { Id = -interaction.ConnectionId }));
|
||||||
_hub!.RegisterInteraction(interaction);
|
_hub!.RegisterInteraction(interaction);
|
||||||
Console.WriteLine("Client registered");
|
Console.WriteLine("Client registered");
|
||||||
|
|||||||
@@ -0,0 +1,7 @@
|
|||||||
|
namespace mROA.Implementation
|
||||||
|
{
|
||||||
|
public class ClientDisconnect : INetworkMessage
|
||||||
|
{
|
||||||
|
public EMessageType MessageType => EMessageType.ClientDisconnect;
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -10,6 +10,7 @@ namespace mROA.Implementation
|
|||||||
CancelRequest,
|
CancelRequest,
|
||||||
EventRequest,
|
EventRequest,
|
||||||
ClientRecovery,
|
ClientRecovery,
|
||||||
ClientConnect
|
ClientConnect,
|
||||||
|
ClientDisconnect,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1,6 +1,7 @@
|
|||||||
using System;
|
using System;
|
||||||
using System.Net;
|
using System.Net;
|
||||||
using System.Net.Sockets;
|
using System.Net.Sockets;
|
||||||
|
using System.Threading.Tasks;
|
||||||
using mROA.Abstract;
|
using mROA.Abstract;
|
||||||
using Exception = System.Exception;
|
using Exception = System.Exception;
|
||||||
|
|
||||||
@@ -39,21 +40,50 @@ namespace mROA.Implementation.Frontend
|
|||||||
throw new NullReferenceException("Serialization toolkit is not initialized");
|
throw new NullReferenceException("Serialization toolkit is not initialized");
|
||||||
|
|
||||||
_tcpClient.Connect(_ipEndPoint);
|
_tcpClient.Connect(_ipEndPoint);
|
||||||
|
|
||||||
_interactionModule.BaseStream = _tcpClient.GetStream();
|
_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();
|
var welcomeMessage = _interactionModule.GetNextMessageReceiving().GetAwaiter().GetResult();
|
||||||
if (welcomeMessage.MessageType != EMessageType.IdAssigning)
|
if (welcomeMessage.MessageType != EMessageType.IdAssigning)
|
||||||
{
|
{
|
||||||
throw new Exception(
|
throw new Exception(
|
||||||
$"Incorrect message type. Must be IdAssigning, current : {welcomeMessage.MessageType.ToString()}");
|
$"Incorrect message type. Must be IdAssigning, current : {welcomeMessage.MessageType.ToString()}");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
var assignment = _serialization.Deserialize<IdAssignment>(welcomeMessage.Data)!;
|
var assignment = _serialization.Deserialize<IdAssignment>(welcomeMessage.Data)!;
|
||||||
_interactionModule.ConnectionId = -assignment.Id;
|
_interactionModule.ConnectionId = -assignment.Id;
|
||||||
TransmissionConfig.OwnershipRepository = new StaticOwnershipRepository(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();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -9,6 +9,11 @@ namespace mROA.Implementation
|
|||||||
|
|
||||||
public class ClientRecovery : INetworkMessage
|
public class ClientRecovery : INetworkMessage
|
||||||
{
|
{
|
||||||
|
public ClientRecovery(int id)
|
||||||
|
{
|
||||||
|
Id = id;
|
||||||
|
}
|
||||||
|
|
||||||
public int Id { get; set; }
|
public int Id { get; set; }
|
||||||
public EMessageType MessageType => EMessageType.ClientRecovery;
|
public EMessageType MessageType => EMessageType.ClientRecovery;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2,6 +2,7 @@
|
|||||||
using System.Collections.Generic;
|
using System.Collections.Generic;
|
||||||
using System.IO;
|
using System.IO;
|
||||||
using System.Linq;
|
using System.Linq;
|
||||||
|
using System.Net.Sockets;
|
||||||
using System.Threading.Tasks;
|
using System.Threading.Tasks;
|
||||||
using mROA.Abstract;
|
using mROA.Abstract;
|
||||||
|
|
||||||
@@ -15,6 +16,9 @@ namespace mROA.Implementation
|
|||||||
private Task<NetworkMessageHeader>? _currentReceiving;
|
private Task<NetworkMessageHeader>? _currentReceiving;
|
||||||
private ISerializationToolkit? _serialization;
|
private ISerializationToolkit? _serialization;
|
||||||
private Stream? _baseStream;
|
private Stream? _baseStream;
|
||||||
|
|
||||||
|
private TaskCompletionSource<Stream> _reconection = new();
|
||||||
|
|
||||||
public int ConnectionId { get; set; }
|
public int ConnectionId { get; set; }
|
||||||
|
|
||||||
public Stream? BaseStream
|
public Stream? BaseStream
|
||||||
@@ -24,8 +28,8 @@ namespace mROA.Implementation
|
|||||||
{
|
{
|
||||||
if (_baseStream is null)
|
if (_baseStream is null)
|
||||||
{
|
{
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
_baseStream = value;
|
_baseStream = value;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -51,7 +55,7 @@ namespace mROA.Implementation
|
|||||||
return _currentReceiving;
|
return _currentReceiving;
|
||||||
}
|
}
|
||||||
|
|
||||||
public async Task PostMessage(NetworkMessageHeader messageHeader)
|
public async Task PostMessageAsync(NetworkMessageHeader messageHeader)
|
||||||
{
|
{
|
||||||
if (BaseStream == null)
|
if (BaseStream == null)
|
||||||
throw new NullReferenceException("BaseStream is null");
|
throw new NullReferenceException("BaseStream is null");
|
||||||
@@ -81,6 +85,8 @@ namespace mROA.Implementation
|
|||||||
return _messageBuffer.FirstOrDefault(m => predicate(m));
|
return _messageBuffer.FirstOrDefault(m => predicate(m));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public event Action<int>? OnDisconected;
|
||||||
|
|
||||||
private async Task<NetworkMessageHeader> GetNextMessage()
|
private async Task<NetworkMessageHeader> GetNextMessage()
|
||||||
{
|
{
|
||||||
if (BaseStream == null)
|
if (BaseStream == null)
|
||||||
@@ -89,37 +95,53 @@ namespace mROA.Implementation
|
|||||||
if (_serialization == null)
|
if (_serialization == null)
|
||||||
throw new NullReferenceException("Serialization toolkit is null");
|
throw new NullReferenceException("Serialization toolkit is null");
|
||||||
|
|
||||||
|
while (true)
|
||||||
try
|
|
||||||
{
|
{
|
||||||
// Console.WriteLine("Receiving message");
|
try
|
||||||
var firstBit = (byte)BaseStream.ReadByte();
|
{
|
||||||
var secondBit = (byte)BaseStream.ReadByte();
|
return await Receive();
|
||||||
|
}
|
||||||
|
catch (Exception)
|
||||||
|
{
|
||||||
|
OnDisconected!.Invoke(ConnectionId);
|
||||||
|
_ = await _reconection.Task;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
var len = BitConverter.ToUInt16(new[] { firstBit, secondBit });
|
private ushort ReadMessageLength()
|
||||||
var localSpan = _buffer[..len];
|
{
|
||||||
|
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<NetworkMessageHeader>(localSpan.Span);
|
private async Task<NetworkMessageHeader> Receive()
|
||||||
|
{
|
||||||
|
var len = ReadMessageLength();
|
||||||
|
var localSpan = _buffer[..len];
|
||||||
|
|
||||||
|
await BaseStream.ReadExactlyAsync(localSpan);
|
||||||
|
|
||||||
|
var message = _serialization.Deserialize<NetworkMessageHeader>(localSpan.Span);
|
||||||
#if TRACE
|
#if TRACE
|
||||||
Console.WriteLine($"{DateTime.Now.TimeOfDay} Received Message {message.Id} - {message.SchemaId}");
|
Console.WriteLine($"{DateTime.Now.TimeOfDay} Received Message {message.Id} - {message.SchemaId}");
|
||||||
TransmissionConfig.TotalTransmittedBytes += len;
|
TransmissionConfig.TotalTransmittedBytes += len;
|
||||||
Console.WriteLine($"Total recieced bytes are {TransmissionConfig.TotalTransmittedBytes}");
|
Console.WriteLine($"Total recieced bytes are {TransmissionConfig.TotalTransmittedBytes}");
|
||||||
#endif
|
#endif
|
||||||
_messageBuffer.Add(message);
|
_messageBuffer.Add(message);
|
||||||
_currentReceiving = Task.Run(async () => await GetNextMessage());
|
_currentReceiving = Task.Run(async () => await GetNextMessage());
|
||||||
|
|
||||||
return message;
|
return message;
|
||||||
}
|
}
|
||||||
catch (Exception e)
|
|
||||||
{
|
|
||||||
Console.WriteLine(e);
|
|
||||||
throw;
|
|
||||||
}
|
|
||||||
|
|
||||||
|
public async Task Restart()
|
||||||
|
{
|
||||||
|
await PostMessageAsync(new NetworkMessageHeader(_serialization!, new ClientRecovery(Math.Abs(ConnectionId))));
|
||||||
|
_reconection.SetResult(BaseStream!);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -97,7 +97,7 @@ namespace mROA.Implementation
|
|||||||
#endif
|
#endif
|
||||||
|
|
||||||
var serialized = _serialization.Serialize(payload, payloadType);
|
var serialized = _serialization.Serialize(payload, payloadType);
|
||||||
await _interaction.PostMessage(new NetworkMessageHeader
|
await _interaction.PostMessageAsync(new NetworkMessageHeader
|
||||||
{ Id = id, MessageType = eMessageType, Data = serialized });
|
{ Id = id, MessageType = eMessageType, Data = serialized });
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user