Some strange IO Exception found
This commit is contained in:
@@ -58,7 +58,7 @@ class Program
|
|||||||
Console.WriteLine("Printer created");
|
Console.WriteLine("Printer created");
|
||||||
Thread.Sleep(100);
|
Thread.Sleep(100);
|
||||||
|
|
||||||
// frontendBridge.Obstacle();
|
frontendBridge.Obstacle();
|
||||||
var name = disposingPrinter.GetName();
|
var name = disposingPrinter.GetName();
|
||||||
DemoCheck.BasicNonParamsCall = true;
|
DemoCheck.BasicNonParamsCall = true;
|
||||||
Console.WriteLine("Printer name : {0}", name);
|
Console.WriteLine("Printer name : {0}", name);
|
||||||
|
|||||||
@@ -76,12 +76,9 @@ namespace mROA.Implementation.Backend
|
|||||||
interaction!.Inject(_serialization);
|
interaction!.Inject(_serialization);
|
||||||
|
|
||||||
|
|
||||||
var streamExtractor = new StreamExtractor(client.GetStream(), _serialization);
|
var streamExtractor = new ChannelInteractionModule.StreamExtractor(client.GetStream(), _serialization);
|
||||||
interaction.IsConnected = () => streamExtractor.IsConnected;
|
interaction.IsConnected = () => streamExtractor.IsConnected;
|
||||||
streamExtractor.MessageReceived = message =>
|
streamExtractor.MessageReceived = message => { interaction.ReceiveChanel.Writer.WriteAsync(message); };
|
||||||
{
|
|
||||||
interaction.ReceiveChanel.Writer.WriteAsync(message);
|
|
||||||
};
|
|
||||||
streamExtractor.SingleReceive();
|
streamExtractor.SingleReceive();
|
||||||
var connectionRequest = interaction.GetNextMessageReceiving(false)
|
var connectionRequest = interaction.GetNextMessageReceiving(false)
|
||||||
.GetAwaiter().GetResult()!;
|
.GetAwaiter().GetResult()!;
|
||||||
@@ -101,10 +98,13 @@ namespace mROA.Implementation.Backend
|
|||||||
var recoveryRequest = _serialization!.Deserialize<ClientRecovery>(connectionRequest.Data)!;
|
var recoveryRequest = _serialization!.Deserialize<ClientRecovery>(connectionRequest.Data)!;
|
||||||
var recoveryInteraction = _hub.GetInteraction(recoveryRequest.Id);
|
var recoveryInteraction = _hub.GetInteraction(recoveryRequest.Id);
|
||||||
|
|
||||||
|
recoveryInteraction.IsConnected = () => streamExtractor.IsConnected;
|
||||||
streamExtractor.MessageReceived = message =>
|
streamExtractor.MessageReceived = message =>
|
||||||
{
|
{
|
||||||
recoveryInteraction.ReceiveChanel.Writer.WriteAsync(message);
|
recoveryInteraction.ReceiveChanel.Writer.WriteAsync(message);
|
||||||
};
|
};
|
||||||
|
Task.Run(async () => await streamExtractor.LoopedReceive());
|
||||||
|
|
||||||
|
|
||||||
recoveryInteraction.Restart(false);
|
recoveryInteraction.Restart(false);
|
||||||
Console.WriteLine("Connection recovery for client {0} finished", recoveryRequest.Id);
|
Console.WriteLine("Connection recovery for client {0} finished", recoveryRequest.Id);
|
||||||
|
|||||||
@@ -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.Threading;
|
||||||
using System.Threading.Channels;
|
using System.Threading.Channels;
|
||||||
using System.Threading.Tasks;
|
using System.Threading.Tasks;
|
||||||
using mROA.Abstract;
|
using mROA.Abstract;
|
||||||
@@ -28,24 +29,21 @@ namespace mROA.Implementation
|
|||||||
{
|
{
|
||||||
SingleReader = false,
|
SingleReader = false,
|
||||||
SingleWriter = false,
|
SingleWriter = false,
|
||||||
// AllowSynchronousContinuations = true
|
|
||||||
});
|
});
|
||||||
_receiveReader = ReceiveChanel.Reader;
|
_receiveReader = ReceiveChanel.Reader;
|
||||||
_outputTrustedChannel = Channel.CreateBounded<NetworkMessageHeader>(new BoundedChannelOptions(1)
|
_outputTrustedChannel = Channel.CreateBounded<NetworkMessageHeader>(new BoundedChannelOptions(1)
|
||||||
{
|
{
|
||||||
SingleReader = true,
|
SingleReader = true,
|
||||||
SingleWriter = true,
|
SingleWriter = true,
|
||||||
// AllowSynchronousContinuations = true,
|
|
||||||
|
|
||||||
});
|
});
|
||||||
_trustedWriter = _outputTrustedChannel.Writer;
|
_trustedWriter = _outputTrustedChannel.Writer;
|
||||||
_outputUntrustedChannel = Channel.CreateUnbounded<NetworkMessageHeader>(new UnboundedChannelOptions
|
_outputUntrustedChannel = Channel.CreateUnbounded<NetworkMessageHeader>(new UnboundedChannelOptions
|
||||||
{
|
{
|
||||||
SingleReader = true,
|
SingleReader = true,
|
||||||
SingleWriter = true,
|
SingleWriter = true,
|
||||||
// AllowSynchronousContinuations = true
|
|
||||||
});
|
});
|
||||||
_untrustedWriter = _outputUntrustedChannel.Writer;
|
_untrustedWriter = _outputUntrustedChannel.Writer;
|
||||||
|
_reconnection = new TaskCompletionSource<Stream>();
|
||||||
}
|
}
|
||||||
|
|
||||||
public int ConnectionId { get; set; }
|
public int ConnectionId { get; set; }
|
||||||
@@ -125,13 +123,6 @@ namespace mROA.Implementation
|
|||||||
await _untrustedWriter.WriteAsync(messageHeader);
|
await _untrustedWriter.WriteAsync(messageHeader);
|
||||||
}
|
}
|
||||||
|
|
||||||
public void HandleMessage(NetworkMessageHeader messageHeader)
|
|
||||||
{
|
|
||||||
_messageBuffer.Remove(messageHeader);
|
|
||||||
}
|
|
||||||
|
|
||||||
// public NetworkMessageHeader[] UnhandledMessages => _messageBuffer.ToArray();
|
|
||||||
|
|
||||||
public NetworkMessageHeader? FirstByFilter(Predicate<NetworkMessageHeader> predicate)
|
public NetworkMessageHeader? FirstByFilter(Predicate<NetworkMessageHeader> predicate)
|
||||||
{
|
{
|
||||||
return _messageBuffer.FirstOrDefault(m => predicate(m));
|
return _messageBuffer.FirstOrDefault(m => predicate(m));
|
||||||
@@ -139,44 +130,18 @@ namespace mROA.Implementation
|
|||||||
|
|
||||||
public event Action<int>? OnDisconnected;
|
public event Action<int>? OnDisconnected;
|
||||||
|
|
||||||
private async Task<NetworkMessageHeader> 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)
|
public async Task Restart(bool sendRecovery)
|
||||||
{
|
{
|
||||||
if (sendRecovery)
|
if (sendRecovery)
|
||||||
{
|
{
|
||||||
await PostMessageAsync(
|
await PostMessageAsync(
|
||||||
new NetworkMessageHeader(_serialization!, new ClientRecovery(Math.Abs(ConnectionId))));
|
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");
|
Console.WriteLine("Setting result for reconnection");
|
||||||
@@ -192,29 +157,24 @@ namespace mROA.Implementation
|
|||||||
{
|
{
|
||||||
//TODO переделать реконнект
|
//TODO переделать реконнект
|
||||||
|
|
||||||
// Console.WriteLine("Staring recovery from {0}", source);
|
Console.WriteLine("Staring recovery from {0}", source);
|
||||||
//
|
|
||||||
// lock (_reconnection)
|
lock (_reconnection)
|
||||||
// {
|
{
|
||||||
// Console.WriteLine("Got lock from {0}", source);
|
Console.WriteLine("Got lock from {0}", source);
|
||||||
//
|
|
||||||
// Console.WriteLine("Call OnDisconnected from {0}", source);
|
Console.WriteLine("Call OnDisconnected from {0}", source);
|
||||||
// _isInReconnectionState = true;
|
OnDisconnected?.Invoke(ConnectionId);
|
||||||
// OnDisconnected?.Invoke(ConnectionId);
|
}
|
||||||
// }
|
|
||||||
//
|
Console.WriteLine("Waiting for reconnect from {0}", source);
|
||||||
// Console.WriteLine("Waiting for reconnect from {0}", source);
|
if (!_reconnection.Task.IsCompleted && !_isConnected)
|
||||||
// if (!_reconnection.Task.IsCompleted && !_isConnected)
|
{
|
||||||
// {
|
Console.WriteLine("Current connection state {0} from {1}", _isConnected, source);
|
||||||
// Console.WriteLine("Current connection state {0} from {1}", _isConnected, source);
|
await _reconnection.Task;
|
||||||
// await _reconnection.Task;
|
}
|
||||||
// }
|
|
||||||
//
|
Console.WriteLine("Reconnect finished from {0}", source);
|
||||||
// Console.WriteLine("Reconnect finished from {0}", source);
|
|
||||||
// lock (_reconnection)
|
|
||||||
// {
|
|
||||||
// _isInReconnectionState = false;
|
|
||||||
// }
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public void Dispose()
|
public void Dispose()
|
||||||
@@ -226,5 +186,114 @@ namespace mROA.Implementation
|
|||||||
_currentReceiving?.Dispose();
|
_currentReceiving?.Dispose();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public class StreamExtractor
|
||||||
|
{
|
||||||
|
private readonly Stream _ioStream;
|
||||||
|
private readonly ISerializationToolkit _serializationToolkit;
|
||||||
|
private const int BufferSize = ushort.MaxValue;
|
||||||
|
private readonly Memory<byte> _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<NetworkMessageHeader> 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<NetworkMessageHeader>(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<NetworkMessageHeader> 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;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -17,7 +17,7 @@ namespace mROA.Implementation.Frontend
|
|||||||
private TcpClient _tcpClient = new();
|
private TcpClient _tcpClient = new();
|
||||||
private IChannelInteractionModule? _interactionModule;
|
private IChannelInteractionModule? _interactionModule;
|
||||||
private ISerializationToolkit? _serialization;
|
private ISerializationToolkit? _serialization;
|
||||||
private StreamExtractor _currentExtractor;
|
private ChannelInteractionModule.StreamExtractor _currentExtractor;
|
||||||
private CancellationTokenSource _rawExtractorCancellation;
|
private CancellationTokenSource _rawExtractorCancellation;
|
||||||
|
|
||||||
public NetworkFrontendBridge(IPEndPoint serverEndPoint)
|
public NetworkFrontendBridge(IPEndPoint serverEndPoint)
|
||||||
@@ -50,7 +50,7 @@ namespace mROA.Implementation.Frontend
|
|||||||
|
|
||||||
PrepareExtractor();
|
PrepareExtractor();
|
||||||
_interactionModule.IsConnected = () => _currentExtractor.IsConnected;
|
_interactionModule.IsConnected = () => _currentExtractor.IsConnected;
|
||||||
_interactionModule.OnDisconnected += id => { Reconnect(); };
|
_interactionModule.OnDisconnected += _ => { Reconnect(); };
|
||||||
|
|
||||||
_interactionModule.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientConnect())).Wait();
|
_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(_rawExtractorCancellation.Token));
|
||||||
Task.Run(async () => await _currentExtractor.LoopedReceive(stopToken));
|
|
||||||
|
|
||||||
var assignment = _serialization.Deserialize<IdAssignment>(idMessage.Data)!;
|
var assignment = _serialization.Deserialize<IdAssignment>(idMessage.Data)!;
|
||||||
_interactionModule.ConnectionId = -assignment.Id;
|
_interactionModule.ConnectionId = -assignment.Id;
|
||||||
@@ -74,9 +73,9 @@ namespace mROA.Implementation.Frontend
|
|||||||
|
|
||||||
private void PrepareExtractor()
|
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);
|
_rawExtractorCancellation.Token);
|
||||||
_currentExtractor.MessageReceived = message =>
|
_currentExtractor.MessageReceived = message =>
|
||||||
{
|
{
|
||||||
@@ -94,7 +93,7 @@ namespace mROA.Implementation.Frontend
|
|||||||
|
|
||||||
PrepareExtractor();
|
PrepareExtractor();
|
||||||
|
|
||||||
_ = _currentExtractor.LoopedReceive(_rawExtractorCancellation.Token);
|
Task.Run(async () => await _currentExtractor.LoopedReceive(_rawExtractorCancellation.Token));
|
||||||
|
|
||||||
await _interactionModule.Restart(true);
|
await _interactionModule.Restart(true);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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<byte> _buffer = new byte[BufferSize];
|
|
||||||
private bool _manualConnectionState = true;
|
|
||||||
|
|
||||||
public StreamExtractor(Stream ioStream, ISerializationToolkit serializationToolkit)
|
|
||||||
{
|
|
||||||
_ioStream = ioStream;
|
|
||||||
_serializationToolkit = serializationToolkit;
|
|
||||||
}
|
|
||||||
|
|
||||||
public Action<NetworkMessageHeader> 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<NetworkMessageHeader>(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<NetworkMessageHeader> 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;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+1
-1
@@ -21,7 +21,7 @@
|
|||||||
</PropertyGroup>
|
</PropertyGroup>
|
||||||
|
|
||||||
<PropertyGroup Condition=" '$(Configuration)' == 'Debug' ">
|
<PropertyGroup Condition=" '$(Configuration)' == 'Debug' ">
|
||||||
<DefineConstants></DefineConstants>
|
<DefineConstants>TRACE;</DefineConstants>
|
||||||
</PropertyGroup>
|
</PropertyGroup>
|
||||||
|
|
||||||
<ItemGroup>
|
<ItemGroup>
|
||||||
|
|||||||
Reference in New Issue
Block a user