ChannelInteractionModule.cs refactor

This commit is contained in:
2025-05-07 23:34:12 +03:00
parent 2f65fb24c0
commit 323e9c9003
2 changed files with 16 additions and 27 deletions
@@ -9,7 +9,6 @@ namespace mROA.Abstract
public interface IChannelInteractionModule : IInjectableModule, IDisposable public interface IChannelInteractionModule : IInjectableModule, IDisposable
{ {
int ConnectionId { get; set; } int ConnectionId { get; set; }
IEndPointContext Context { get; set; }
Channel<NetworkMessageHeader> ReceiveChanel { get; } Channel<NetworkMessageHeader> ReceiveChanel { get; }
ChannelReader<NetworkMessageHeader> TrustedPostChanel { get; } ChannelReader<NetworkMessageHeader> TrustedPostChanel { get; }
ChannelReader<NetworkMessageHeader> UntrustedPostChanel { get; } ChannelReader<NetworkMessageHeader> UntrustedPostChanel { get; }
+16 -26
View File
@@ -14,12 +14,11 @@ namespace mROA.Implementation
private readonly ChannelWriter<NetworkMessageHeader> _untrustedWriter; private readonly ChannelWriter<NetworkMessageHeader> _untrustedWriter;
private readonly Channel<NetworkMessageHeader> _outputTrustedChannel; private readonly Channel<NetworkMessageHeader> _outputTrustedChannel;
private readonly Channel<NetworkMessageHeader> _outputUntrustedChannel; private readonly Channel<NetworkMessageHeader> _outputUntrustedChannel;
private Task<NetworkMessageHeader>? _currentReceiving;
private IContextualSerializationToolKit? _serialization; private IContextualSerializationToolKit? _serialization;
private bool _isConnected = true; private bool _isConnected = true;
private bool _isActive = true; private bool _isActive = true;
private TaskCompletionSource<Stream> _reconnection; private TaskCompletionSource<Stream> _reconnection;
public IEndPointContext Context { get; set; } private IEndPointContext? _context;
public ChannelInteractionModule() public ChannelInteractionModule()
{ {
@@ -50,7 +49,7 @@ namespace mROA.Implementation
public ChannelReader<NetworkMessageHeader> TrustedPostChanel => _outputTrustedChannel.Reader; public ChannelReader<NetworkMessageHeader> TrustedPostChanel => _outputTrustedChannel.Reader;
public ChannelReader<NetworkMessageHeader> UntrustedPostChanel => _outputUntrustedChannel.Reader; public ChannelReader<NetworkMessageHeader> UntrustedPostChanel => _outputUntrustedChannel.Reader;
public Func<bool> IsConnected { get; set; } public Func<bool> IsConnected { get; set; } = () => false;
public void Inject<T>(T dependency) public void Inject<T>(T dependency)
{ {
@@ -63,7 +62,7 @@ namespace mROA.Implementation
ConnectionId = identityGenerator.GetNextIdentity(); ConnectionId = identityGenerator.GetNextIdentity();
break; break;
case IEndPointContext endpointContext: case IEndPointContext endpointContext:
Context = endpointContext; _context = endpointContext;
break; break;
} }
} }
@@ -93,8 +92,7 @@ namespace mROA.Implementation
{ {
if (_serialization == null) if (_serialization == null)
throw new NullReferenceException("Serialization toolkit is not initialized"); throw new NullReferenceException("Serialization toolkit is not initialized");
bool withError = false;
while (true) while (true)
{ {
if (await PostMessageInternal(messageHeader)) if (await PostMessageInternal(messageHeader))
@@ -106,8 +104,7 @@ namespace mROA.Implementation
} }
_isConnected = false; _isConnected = false;
withError = true; await MakeRecovery();
await MakeRecovery("OUT");
} }
} }
@@ -123,23 +120,21 @@ namespace mROA.Implementation
if (sendRecovery) if (sendRecovery)
{ {
await PostMessageAsync( await PostMessageAsync(
new NetworkMessageHeader(_serialization!, new ClientRecovery(Math.Abs(ConnectionId)), Context)); new NetworkMessageHeader(_serialization!, new ClientRecovery(Math.Abs(ConnectionId)), _context));
var ping = await ReceiveChanel.Reader.ReadAsync(); await ReceiveChanel.Reader.ReadAsync();
} }
else else
{ {
await _trustedWriter.WriteAsync(new NetworkMessageHeader()); await _trustedWriter.WriteAsync(new NetworkMessageHeader());
} }
var setting = _reconnection.TrySetResult(null); _reconnection.TrySetResult(Stream.Null);
// _isInReconnectionState = false;
_isConnected = true; _isConnected = true;
_reconnection = new TaskCompletionSource<Stream>(); _reconnection = new TaskCompletionSource<Stream>();
} }
private async Task MakeRecovery(string source) private async Task MakeRecovery()
{ {
//TODO переделать реконнект
lock (_reconnection) lock (_reconnection)
{ {
OnDisconnected?.Invoke(ConnectionId); OnDisconnected?.Invoke(ConnectionId);
@@ -154,10 +149,6 @@ namespace mROA.Implementation
public void Dispose() public void Dispose()
{ {
_isActive = false; _isActive = false;
if (_currentReceiving is { IsCompleted: true })
{
_currentReceiving?.Dispose();
}
} }
public class StreamExtractor public class StreamExtractor
@@ -167,15 +158,14 @@ namespace mROA.Implementation
private const int BufferSize = ushort.MaxValue; private const int BufferSize = ushort.MaxValue;
private readonly Memory<byte> _buffer = new byte[BufferSize]; private readonly Memory<byte> _buffer = new byte[BufferSize];
private bool _manualConnectionState = true; private bool _manualConnectionState = true;
private IEndPointContext Context; private readonly IEndPointContext _context;
public readonly int Id = new Random().Next();
public StreamExtractor(Stream ioStream, IContextualSerializationToolKit serializationToolkit, public StreamExtractor(Stream ioStream, IContextualSerializationToolKit serializationToolkit,
IEndPointContext context) IEndPointContext context)
{ {
_ioStream = ioStream; _ioStream = ioStream;
_serializationToolkit = serializationToolkit; _serializationToolkit = serializationToolkit;
Context = context; _context = context;
} }
public Action<NetworkMessageHeader> MessageReceived = _ => { }; public Action<NetworkMessageHeader> MessageReceived = _ => { };
@@ -197,14 +187,14 @@ namespace mROA.Implementation
return len; return len;
} }
public async Task SingleReceive(CancellationToken Token = default) public async Task SingleReceive(CancellationToken token = default)
{ {
var len = ReadMessageLength(); var len = ReadMessageLength();
var localSpan = _buffer[..len]; var localSpan = _buffer[..len];
await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: Token); await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: token);
var message = _serializationToolkit.Deserialize<NetworkMessageHeader>(localSpan, Context); var message = _serializationToolkit.Deserialize<NetworkMessageHeader>(localSpan, _context);
MessageReceived(message); MessageReceived(message);
} }
@@ -216,9 +206,9 @@ namespace mROA.Implementation
} }
} }
public async Task Send(NetworkMessageHeader message, CancellationToken token = default) private async Task Send(NetworkMessageHeader message, CancellationToken token = default)
{ {
var rawMessage = _serializationToolkit.Serialize(message, Context); var rawMessage = _serializationToolkit.Serialize(message, _context);
var header = BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort)); var header = BitConverter.GetBytes((ushort)rawMessage.Length).AsMemory(0, sizeof(ushort));
await _ioStream.WriteAsync(header, token); await _ioStream.WriteAsync(header, token);