Preparing for new request id
This commit is contained in:
@@ -27,7 +27,7 @@ namespace mROA.Cbor
|
|||||||
|
|
||||||
private bool FindParser(Type t, out IOrdinaryStructureParser parser)
|
private bool FindParser(Type t, out IOrdinaryStructureParser parser)
|
||||||
{
|
{
|
||||||
if (t == typeof(NetworkMessageHeader))
|
if (t == typeof(NetworkMessage))
|
||||||
{
|
{
|
||||||
parser = _parsers[0];
|
parser = _parsers[0];
|
||||||
return true;
|
return true;
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ namespace mROA.Cbor
|
|||||||
{
|
{
|
||||||
public void Write(CborWriter writer, object value, IEndPointContext context, CborSerializationToolkit serialization)
|
public void Write(CborWriter writer, object value, IEndPointContext context, CborSerializationToolkit serialization)
|
||||||
{
|
{
|
||||||
var v = value as NetworkMessageHeader;
|
var v = value as NetworkMessage;
|
||||||
writer.WriteStartArray(3);
|
writer.WriteStartArray(3);
|
||||||
writer.WriteByteString(v.Id.ToByteArray());
|
writer.WriteByteString(v.Id.ToByteArray());
|
||||||
writer.WriteInt32((int)v.MessageType);
|
writer.WriteInt32((int)v.MessageType);
|
||||||
@@ -27,7 +27,7 @@ namespace mROA.Cbor
|
|||||||
public object Read(CborReader reader, IEndPointContext context, CborSerializationToolkit serialization)
|
public object Read(CborReader reader, IEndPointContext context, CborSerializationToolkit serialization)
|
||||||
{
|
{
|
||||||
reader.ReadStartArray();
|
reader.ReadStartArray();
|
||||||
var value = new NetworkMessageHeader
|
var value = new NetworkMessage
|
||||||
{
|
{
|
||||||
Id = new Guid(reader.ReadByteString()),
|
Id = new Guid(reader.ReadByteString()),
|
||||||
MessageType = (EMessageType)reader.ReadInt32(),
|
MessageType = (EMessageType)reader.ReadInt32(),
|
||||||
|
|||||||
@@ -9,13 +9,13 @@ namespace mROA.Abstract
|
|||||||
{
|
{
|
||||||
int ConnectionId { get; set; }
|
int ConnectionId { get; set; }
|
||||||
IEndPointContext Context { get; set; }
|
IEndPointContext Context { get; set; }
|
||||||
Channel<NetworkMessageHeader> ReceiveChanel { get; }
|
Channel<NetworkMessage> ReceiveChanel { get; }
|
||||||
ChannelReader<NetworkMessageHeader> TrustedPostChanel { get; }
|
ChannelReader<NetworkMessage> TrustedPostChanel { get; }
|
||||||
ChannelReader<NetworkMessageHeader> UntrustedPostChanel { get; }
|
ChannelReader<NetworkMessage> UntrustedPostChanel { get; }
|
||||||
Func<bool> IsConnected { get; set; }
|
Func<bool> IsConnected { get; set; }
|
||||||
ValueTask<NetworkMessageHeader> GetNextMessageReceiving();
|
ValueTask<NetworkMessage> GetNextMessageReceiving();
|
||||||
Task PostMessageAsync(NetworkMessageHeader messageHeader);
|
Task PostMessageAsync(NetworkMessage message);
|
||||||
Task PostMessageUntrustedAsync(NetworkMessageHeader messageHeader);
|
Task PostMessageUntrustedAsync(NetworkMessage message);
|
||||||
event Action<int> OnDisconnected;
|
event Action<int> OnDisconnected;
|
||||||
Task Restart(bool sendRecovery);
|
Task Restart(bool sendRecovery);
|
||||||
void PassReconnection();
|
void PassReconnection();
|
||||||
|
|||||||
@@ -11,13 +11,13 @@ namespace mROA.Abstract
|
|||||||
int Id { get; }
|
int Id { get; }
|
||||||
IEndPointContext Context { get; }
|
IEndPointContext Context { get; }
|
||||||
|
|
||||||
Task<(object? Deserialized, EMessageType MessageType)> GetSingle(Predicate<NetworkMessageHeader> rule,
|
Task<(object? Deserialized, EMessageType MessageType)> GetSingle(Predicate<NetworkMessage> rule,
|
||||||
IEndPointContext? context, CancellationToken token = default,
|
IEndPointContext? context, CancellationToken token = default,
|
||||||
params Func<NetworkMessageHeader, Type?>[] converter);
|
params Func<NetworkMessage, Type?>[] converter);
|
||||||
|
|
||||||
IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(Predicate<NetworkMessageHeader> rule,
|
IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(Predicate<NetworkMessage> rule,
|
||||||
IEndPointContext? context, CancellationToken token = default,
|
IEndPointContext? context, CancellationToken token = default,
|
||||||
params Func<NetworkMessageHeader, Type?>[] converter);
|
params Func<NetworkMessage, Type?>[] converter);
|
||||||
|
|
||||||
Task PostCallMessageAsync<T>(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context)
|
Task PostCallMessageAsync<T>(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context)
|
||||||
where T : notnull;
|
where T : notnull;
|
||||||
|
|||||||
@@ -8,7 +8,7 @@ namespace mROA.Abstract
|
|||||||
{
|
{
|
||||||
Task StartExtraction();
|
Task StartExtraction();
|
||||||
void PushMessage(object parced, EMessageType originalType);
|
void PushMessage(object parced, EMessageType originalType);
|
||||||
Predicate<NetworkMessageHeader> Rule { get; }
|
Predicate<NetworkMessage> Rule { get; }
|
||||||
Func<NetworkMessageHeader, Type?>[] Converters { get; }
|
Func<NetworkMessage, Type?>[] Converters { get; }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -98,14 +98,14 @@ namespace mROA.Implementation.Backend
|
|||||||
}
|
}
|
||||||
|
|
||||||
private void HandleNewClient(EndPointContext context, ChannelInteractionModule interaction,
|
private void HandleNewClient(EndPointContext context, ChannelInteractionModule interaction,
|
||||||
ChannelInteractionModule.StreamExtractor streamExtractor, CancellationTokenSource cts, NetworkMessageHeader connectionHeader)
|
ChannelInteractionModule.StreamExtractor streamExtractor, CancellationTokenSource cts, NetworkMessage connection)
|
||||||
{
|
{
|
||||||
context.HostId = 0;
|
context.HostId = 0;
|
||||||
context.OwnerId = -interaction.ConnectionId;
|
context.OwnerId = -interaction.ConnectionId;
|
||||||
interaction.Context = context;
|
interaction.Context = context;
|
||||||
Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token));
|
Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token));
|
||||||
_ = streamExtractor.SendFromChannel(interaction.TrustedPostChanel, cts.Token);
|
_ = streamExtractor.SendFromChannel(interaction.TrustedPostChanel, cts.Token);
|
||||||
interaction.PostMessageAsync(new NetworkMessageHeader(_serialization,
|
interaction.PostMessageAsync(new NetworkMessage(_serialization,
|
||||||
new IdAssignment { Id = interaction.ConnectionId }, null));
|
new IdAssignment { Id = interaction.ConnectionId }, null));
|
||||||
_extractorsTokenSources[interaction.ConnectionId] = cts;
|
_extractorsTokenSources[interaction.ConnectionId] = cts;
|
||||||
|
|
||||||
@@ -143,7 +143,7 @@ namespace mROA.Implementation.Backend
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
private void RecoverDisconnectedClient(NetworkMessageHeader connectionRequest,
|
private void RecoverDisconnectedClient(NetworkMessage connectionRequest,
|
||||||
ChannelInteractionModule.StreamExtractor streamExtractor,
|
ChannelInteractionModule.StreamExtractor streamExtractor,
|
||||||
CancellationTokenSource cts)
|
CancellationTokenSource cts)
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -41,7 +41,7 @@ namespace mROA.Implementation.Backend
|
|||||||
while (token.IsCancellationRequested == false)
|
while (token.IsCancellationRequested == false)
|
||||||
{
|
{
|
||||||
var incoming = await _client.ReceiveAsync();
|
var incoming = await _client.ReceiveAsync();
|
||||||
var parsed = _serializationToolkit.Deserialize<NetworkMessageHeader>(incoming.Buffer, null);
|
var parsed = _serializationToolkit.Deserialize<NetworkMessage>(incoming.Buffer, null);
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
int channelId;
|
int channelId;
|
||||||
|
|||||||
@@ -9,11 +9,11 @@ namespace mROA.Implementation
|
|||||||
{
|
{
|
||||||
public class ChannelInteractionModule : IChannelInteractionModule
|
public class ChannelInteractionModule : IChannelInteractionModule
|
||||||
{
|
{
|
||||||
private readonly ChannelReader<NetworkMessageHeader> _receiveReader;
|
private readonly ChannelReader<NetworkMessage> _receiveReader;
|
||||||
private readonly ChannelWriter<NetworkMessageHeader> _trustedWriter;
|
private readonly ChannelWriter<NetworkMessage> _trustedWriter;
|
||||||
private readonly ChannelWriter<NetworkMessageHeader> _untrustedWriter;
|
private readonly ChannelWriter<NetworkMessage> _untrustedWriter;
|
||||||
private readonly Channel<NetworkMessageHeader> _outputTrustedChannel;
|
private readonly Channel<NetworkMessage> _outputTrustedChannel;
|
||||||
private readonly Channel<NetworkMessageHeader> _outputUntrustedChannel;
|
private readonly Channel<NetworkMessage> _outputUntrustedChannel;
|
||||||
private readonly IContextualSerializationToolKit _serialization;
|
private readonly IContextualSerializationToolKit _serialization;
|
||||||
private bool _isConnected = true;
|
private bool _isConnected = true;
|
||||||
private bool _isActive = true;
|
private bool _isActive = true;
|
||||||
@@ -28,19 +28,19 @@ namespace mROA.Implementation
|
|||||||
public ChannelInteractionModule(IContextualSerializationToolKit serialization)
|
public ChannelInteractionModule(IContextualSerializationToolKit serialization)
|
||||||
{
|
{
|
||||||
_serialization = serialization;
|
_serialization = serialization;
|
||||||
ReceiveChanel = Channel.CreateUnbounded<NetworkMessageHeader>(new UnboundedChannelOptions
|
ReceiveChanel = Channel.CreateUnbounded<NetworkMessage>(new UnboundedChannelOptions
|
||||||
{
|
{
|
||||||
SingleReader = false,
|
SingleReader = false,
|
||||||
SingleWriter = false,
|
SingleWriter = false,
|
||||||
});
|
});
|
||||||
_receiveReader = ReceiveChanel.Reader;
|
_receiveReader = ReceiveChanel.Reader;
|
||||||
_outputTrustedChannel = Channel.CreateUnbounded<NetworkMessageHeader>(new UnboundedChannelOptions
|
_outputTrustedChannel = Channel.CreateUnbounded<NetworkMessage>(new UnboundedChannelOptions
|
||||||
{
|
{
|
||||||
SingleReader = true,
|
SingleReader = true,
|
||||||
SingleWriter = true
|
SingleWriter = true
|
||||||
});
|
});
|
||||||
_trustedWriter = _outputTrustedChannel.Writer;
|
_trustedWriter = _outputTrustedChannel.Writer;
|
||||||
_outputUntrustedChannel = Channel.CreateUnbounded<NetworkMessageHeader>(new UnboundedChannelOptions
|
_outputUntrustedChannel = Channel.CreateUnbounded<NetworkMessage>(new UnboundedChannelOptions
|
||||||
{
|
{
|
||||||
SingleReader = true,
|
SingleReader = true,
|
||||||
SingleWriter = true,
|
SingleWriter = true,
|
||||||
@@ -52,33 +52,33 @@ namespace mROA.Implementation
|
|||||||
public int ConnectionId { get; set; }
|
public int ConnectionId { get; set; }
|
||||||
|
|
||||||
public IEndPointContext Context { get; set; }
|
public IEndPointContext Context { get; set; }
|
||||||
public Channel<NetworkMessageHeader> ReceiveChanel { get; }
|
public Channel<NetworkMessage> ReceiveChanel { get; }
|
||||||
|
|
||||||
public ChannelReader<NetworkMessageHeader> TrustedPostChanel => _outputTrustedChannel.Reader;
|
public ChannelReader<NetworkMessage> TrustedPostChanel => _outputTrustedChannel.Reader;
|
||||||
public ChannelReader<NetworkMessageHeader> UntrustedPostChanel => _outputUntrustedChannel.Reader;
|
public ChannelReader<NetworkMessage> UntrustedPostChanel => _outputUntrustedChannel.Reader;
|
||||||
public Func<bool> IsConnected { get; set; } = () => false;
|
public Func<bool> IsConnected { get; set; } = () => false;
|
||||||
|
|
||||||
public ValueTask<NetworkMessageHeader> GetNextMessageReceiving()
|
public ValueTask<NetworkMessage> GetNextMessageReceiving()
|
||||||
{
|
{
|
||||||
return _receiveReader.ReadAsync();
|
return _receiveReader.ReadAsync();
|
||||||
}
|
}
|
||||||
|
|
||||||
private async ValueTask<bool> PostMessageInternal(NetworkMessageHeader messageHeader)
|
private async ValueTask<bool> PostMessageInternal(NetworkMessage message)
|
||||||
{
|
{
|
||||||
if (!IsConnected())
|
if (!IsConnected())
|
||||||
{
|
{
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
await _trustedWriter.WriteAsync(messageHeader);
|
await _trustedWriter.WriteAsync(message);
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
public async Task PostMessageAsync(NetworkMessageHeader messageHeader)
|
public async Task PostMessageAsync(NetworkMessage message)
|
||||||
{
|
{
|
||||||
while (true)
|
while (true)
|
||||||
{
|
{
|
||||||
if (await PostMessageInternal(messageHeader))
|
if (await PostMessageInternal(message))
|
||||||
break;
|
break;
|
||||||
|
|
||||||
if (!_isActive)
|
if (!_isActive)
|
||||||
@@ -91,9 +91,9 @@ namespace mROA.Implementation
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public async Task PostMessageUntrustedAsync(NetworkMessageHeader messageHeader)
|
public async Task PostMessageUntrustedAsync(NetworkMessage message)
|
||||||
{
|
{
|
||||||
await _untrustedWriter.WriteAsync(messageHeader);
|
await _untrustedWriter.WriteAsync(message);
|
||||||
}
|
}
|
||||||
|
|
||||||
public event Action<int>? OnDisconnected;
|
public event Action<int>? OnDisconnected;
|
||||||
@@ -103,12 +103,12 @@ namespace mROA.Implementation
|
|||||||
if (sendRecovery)
|
if (sendRecovery)
|
||||||
{
|
{
|
||||||
await PostMessageAsync(
|
await PostMessageAsync(
|
||||||
new NetworkMessageHeader(_serialization, new ClientRecovery(Math.Abs(ConnectionId)), Context));
|
new NetworkMessage(_serialization, new ClientRecovery(Math.Abs(ConnectionId)), Context));
|
||||||
await ReceiveChanel.Reader.ReadAsync();
|
await ReceiveChanel.Reader.ReadAsync();
|
||||||
}
|
}
|
||||||
else
|
else
|
||||||
{
|
{
|
||||||
await _trustedWriter.WriteAsync(new NetworkMessageHeader());
|
await _trustedWriter.WriteAsync(new NetworkMessage());
|
||||||
}
|
}
|
||||||
|
|
||||||
PassReconnection();
|
PassReconnection();
|
||||||
@@ -141,7 +141,7 @@ namespace mROA.Implementation
|
|||||||
|
|
||||||
public class StreamExtractor
|
public class StreamExtractor
|
||||||
{
|
{
|
||||||
private const int BufferSize = ushort.MaxValue;
|
private const int BufferSize = ushort.MaxValue + 2;
|
||||||
|
|
||||||
private readonly Stream _ioStream;
|
private readonly Stream _ioStream;
|
||||||
private readonly IContextualSerializationToolKit _serializationToolkit;
|
private readonly IContextualSerializationToolKit _serializationToolkit;
|
||||||
@@ -158,25 +158,22 @@ namespace mROA.Implementation
|
|||||||
_lenBuffer = new byte[2];
|
_lenBuffer = new byte[2];
|
||||||
}
|
}
|
||||||
|
|
||||||
public Action<NetworkMessageHeader> MessageReceived = _ => { };
|
public Action<NetworkMessage> MessageReceived = _ => { };
|
||||||
|
|
||||||
private async Task<ushort> ReadMessageLength()
|
|
||||||
{
|
|
||||||
await _ioStream.ReadAsync(_lenBuffer);
|
|
||||||
|
|
||||||
var len = BitConverter.ToUInt16(_lenBuffer);
|
|
||||||
|
|
||||||
return len;
|
|
||||||
}
|
|
||||||
|
|
||||||
public async Task SingleReceive(CancellationToken token = default)
|
public async Task SingleReceive(CancellationToken token = default)
|
||||||
{
|
{
|
||||||
var len = await ReadMessageLength();
|
int firstRead = _ioStream.Read(_buffer.Span);
|
||||||
var localSpan = _buffer[..len];
|
|
||||||
|
|
||||||
|
var metadata = new NetworkMessage.NetworkMessageMeta(_buffer.Span);
|
||||||
|
|
||||||
|
var localSpan = _buffer[2..len];
|
||||||
|
if (firstRead - 2 != len)
|
||||||
|
{
|
||||||
|
localSpan = _buffer[(len + 2)..];
|
||||||
await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: token);
|
await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: token);
|
||||||
var message = _serializationToolkit.Deserialize<NetworkMessageHeader>(localSpan, _context);
|
}
|
||||||
// _logger.LogTrace("RECV {0}", message.ToString());
|
|
||||||
|
var message = _serializationToolkit.Deserialize<NetworkMessage>(localSpan, _context);
|
||||||
MessageReceived(message);
|
MessageReceived(message);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -188,7 +185,7 @@ namespace mROA.Implementation
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private async Task Send(NetworkMessageHeader message, CancellationToken token = default)
|
private async Task Send(NetworkMessage message, CancellationToken token = default)
|
||||||
{
|
{
|
||||||
var bodySpan = _buffer[2..];
|
var bodySpan = _buffer[2..];
|
||||||
var len = _serializationToolkit.Serialize(message, bodySpan.Span, _context);
|
var len = _serializationToolkit.Serialize(message, bodySpan.Span, _context);
|
||||||
@@ -199,7 +196,7 @@ namespace mROA.Implementation
|
|||||||
// _logger.LogTrace("SEND {0}", message.ToString());
|
// _logger.LogTrace("SEND {0}", message.ToString());
|
||||||
}
|
}
|
||||||
|
|
||||||
public async Task SendFromChannel(ChannelReader<NetworkMessageHeader> channel,
|
public async Task SendFromChannel(ChannelReader<NetworkMessage> channel,
|
||||||
CancellationToken token = default)
|
CancellationToken token = default)
|
||||||
{
|
{
|
||||||
while (token.IsCancellationRequested == false && IsConnected)
|
while (token.IsCancellationRequested == false && IsConnected)
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
namespace mROA.Implementation
|
namespace mROA.Implementation
|
||||||
{
|
{
|
||||||
public enum EMessageType
|
public enum EMessageType : byte
|
||||||
{
|
{
|
||||||
Unknown,
|
Unknown,
|
||||||
FinishedCommandExecution,
|
FinishedCommandExecution,
|
||||||
|
|||||||
@@ -39,7 +39,7 @@ namespace mROA.Implementation.Frontend
|
|||||||
_interactionModule.IsConnected = () => _currentExtractor.IsConnected;
|
_interactionModule.IsConnected = () => _currentExtractor.IsConnected;
|
||||||
_interactionModule.OnDisconnected += _ => { Reconnect().ConfigureAwait(false); };
|
_interactionModule.OnDisconnected += _ => { Reconnect().ConfigureAwait(false); };
|
||||||
|
|
||||||
_interactionModule.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientConnect(), _context))
|
_interactionModule.PostMessageAsync(new NetworkMessage(_serialization, new ClientConnect(), _context))
|
||||||
.Wait();
|
.Wait();
|
||||||
|
|
||||||
_ = _currentExtractor.SingleReceive().ConfigureAwait(false);
|
_ = _currentExtractor.SingleReceive().ConfigureAwait(false);
|
||||||
@@ -95,7 +95,7 @@ namespace mROA.Implementation.Frontend
|
|||||||
|
|
||||||
public void Disconnect()
|
public void Disconnect()
|
||||||
{
|
{
|
||||||
_ = _interactionModule.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientDisconnect(),
|
_ = _interactionModule.PostMessageAsync(new NetworkMessage(_serialization, new ClientDisconnect(),
|
||||||
_context));
|
_context));
|
||||||
_interactionModule.Dispose();
|
_interactionModule.Dispose();
|
||||||
_tcpClient.Dispose();
|
_tcpClient.Dispose();
|
||||||
|
|||||||
@@ -56,11 +56,11 @@ namespace mROA.Implementation.Frontend
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public Predicate<NetworkMessageHeader> Rule { get; } = m =>
|
public Predicate<NetworkMessage> Rule { get; } = m =>
|
||||||
m.MessageType is EMessageType.CallRequest or EMessageType.CancelRequest
|
m.MessageType is EMessageType.CallRequest or EMessageType.CancelRequest
|
||||||
or EMessageType.EventRequest or EMessageType.ClientDisconnect;
|
or EMessageType.EventRequest or EMessageType.ClientDisconnect;
|
||||||
|
|
||||||
public Func<NetworkMessageHeader, Type?>[] Converters { get; } =
|
public Func<NetworkMessage, Type?>[] Converters { get; } =
|
||||||
{
|
{
|
||||||
m => m.MessageType == EMessageType.CallRequest ? typeof(DefaultCallRequest) : null,
|
m => m.MessageType == EMessageType.CallRequest ? typeof(DefaultCallRequest) : null,
|
||||||
m => m.MessageType == EMessageType.CancelRequest ? typeof(CancelRequest) : null,
|
m => m.MessageType == EMessageType.CancelRequest ? typeof(CancelRequest) : null,
|
||||||
|
|||||||
@@ -44,7 +44,7 @@ namespace mROA.Implementation.Frontend
|
|||||||
while (token.IsCancellationRequested == false)
|
while (token.IsCancellationRequested == false)
|
||||||
{
|
{
|
||||||
var message = new Memory<byte>((await udpClient.ReceiveAsync()).Buffer);
|
var message = new Memory<byte>((await udpClient.ReceiveAsync()).Buffer);
|
||||||
var parsed = _serializationToolkit.Deserialize<NetworkMessageHeader>(message, _context);
|
var parsed = _serializationToolkit.Deserialize<NetworkMessage>(message, _context);
|
||||||
|
|
||||||
await writer.WriteAsync(parsed, token);
|
await writer.WriteAsync(parsed, token);
|
||||||
}
|
}
|
||||||
@@ -52,7 +52,7 @@ namespace mROA.Implementation.Frontend
|
|||||||
|
|
||||||
private async Task Posting(UdpClient udpClient, CancellationToken token)
|
private async Task Posting(UdpClient udpClient, CancellationToken token)
|
||||||
{
|
{
|
||||||
var initMessage = new NetworkMessageHeader
|
var initMessage = new NetworkMessage
|
||||||
{
|
{
|
||||||
MessageType = EMessageType.UntrustedConnect, Id = Guid.NewGuid(),
|
MessageType = EMessageType.UntrustedConnect, Id = Guid.NewGuid(),
|
||||||
Data = BitConverter.GetBytes(_channelInteractionModule.ConnectionId)
|
Data = BitConverter.GetBytes(_channelInteractionModule.ConnectionId)
|
||||||
|
|||||||
+30
-5
@@ -3,9 +3,9 @@ using mROA.Abstract;
|
|||||||
|
|
||||||
namespace mROA.Implementation
|
namespace mROA.Implementation
|
||||||
{
|
{
|
||||||
public class NetworkMessageHeader
|
public class NetworkMessage
|
||||||
{
|
{
|
||||||
private bool Equals(NetworkMessageHeader other)
|
private bool Equals(NetworkMessage other)
|
||||||
{
|
{
|
||||||
return Id.Equals(other.Id) && MessageType == other.MessageType;
|
return Id.Equals(other.Id) && MessageType == other.MessageType;
|
||||||
}
|
}
|
||||||
@@ -14,15 +14,15 @@ namespace mROA.Implementation
|
|||||||
{
|
{
|
||||||
if (obj is null) return false;
|
if (obj is null) return false;
|
||||||
if (ReferenceEquals(this, obj)) return true;
|
if (ReferenceEquals(this, obj)) return true;
|
||||||
return obj.GetType() == GetType() && Equals((NetworkMessageHeader)obj);
|
return obj.GetType() == GetType() && Equals((NetworkMessage)obj);
|
||||||
}
|
}
|
||||||
|
|
||||||
public NetworkMessageHeader()
|
public NetworkMessage()
|
||||||
{
|
{
|
||||||
Data = Array.Empty<byte>();
|
Data = Array.Empty<byte>();
|
||||||
}
|
}
|
||||||
|
|
||||||
public NetworkMessageHeader(IContextualSerializationToolKit serializationToolkit,
|
public NetworkMessage(IContextualSerializationToolKit serializationToolkit,
|
||||||
INetworkMessage networkMessage, IEndPointContext? context)
|
INetworkMessage networkMessage, IEndPointContext? context)
|
||||||
{
|
{
|
||||||
MessageType = networkMessage.MessageType;
|
MessageType = networkMessage.MessageType;
|
||||||
@@ -40,5 +40,30 @@ namespace mROA.Implementation
|
|||||||
{
|
{
|
||||||
return $" {Id}:{MessageType} [{Data.Length}]";
|
return $" {Id}:{MessageType} [{Data.Length}]";
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public struct NetworkMessageMeta
|
||||||
|
{
|
||||||
|
public byte Type;
|
||||||
|
public Guid Id;
|
||||||
|
public ushort BodyLength;
|
||||||
|
|
||||||
|
public NetworkMessageMeta(ReadOnlySpan<byte> metadata)
|
||||||
|
{
|
||||||
|
Type = metadata[0];
|
||||||
|
Id = new Guid(metadata[1..17]);
|
||||||
|
BodyLength = BitConverter.ToUInt16(metadata[17..]);
|
||||||
|
}
|
||||||
|
|
||||||
|
public NetworkMessage ToMessage(ReadOnlySpan<byte> memory)
|
||||||
|
{
|
||||||
|
var data = memory[19..][..BodyLength];
|
||||||
|
return new NetworkMessage
|
||||||
|
{
|
||||||
|
Data = data.ToArray(),
|
||||||
|
Id = Id,
|
||||||
|
MessageType = (EMessageType)Type
|
||||||
|
};
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -28,8 +28,8 @@ namespace mROA.Implementation
|
|||||||
public IEndPointContext Context => _interaction.Context;
|
public IEndPointContext Context => _interaction.Context;
|
||||||
|
|
||||||
public async Task<(object? Deserialized, EMessageType MessageType)> GetSingle(
|
public async Task<(object? Deserialized, EMessageType MessageType)> GetSingle(
|
||||||
Predicate<NetworkMessageHeader> rule, IEndPointContext? context,
|
Predicate<NetworkMessage> rule, IEndPointContext? context,
|
||||||
CancellationToken token = default, params Func<NetworkMessageHeader, Type?>[] converter)
|
CancellationToken token = default, params Func<NetworkMessage, Type?>[] converter)
|
||||||
{
|
{
|
||||||
var writer = _interaction.ReceiveChanel.Writer;
|
var writer = _interaction.ReceiveChanel.Writer;
|
||||||
var reader = _interaction.ReceiveChanel.Reader;
|
var reader = _interaction.ReceiveChanel.Reader;
|
||||||
@@ -52,9 +52,9 @@ namespace mROA.Implementation
|
|||||||
}
|
}
|
||||||
|
|
||||||
public async IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(
|
public async IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(
|
||||||
Predicate<NetworkMessageHeader> rule, IEndPointContext? context,
|
Predicate<NetworkMessage> rule, IEndPointContext? context,
|
||||||
[EnumeratorCancellation] CancellationToken token = default,
|
[EnumeratorCancellation] CancellationToken token = default,
|
||||||
params Func<NetworkMessageHeader, Type?>[] converter)
|
params Func<NetworkMessage, Type?>[] converter)
|
||||||
{
|
{
|
||||||
var writer = _interaction.ReceiveChanel.Writer;
|
var writer = _interaction.ReceiveChanel.Writer;
|
||||||
await foreach (var message in _interaction.ReceiveChanel.Reader.ReadAllAsync(token))
|
await foreach (var message in _interaction.ReceiveChanel.Reader.ReadAllAsync(token))
|
||||||
@@ -83,7 +83,7 @@ namespace mROA.Implementation
|
|||||||
IEndPointContext? context) where T : notnull
|
IEndPointContext? context) where T : notnull
|
||||||
{
|
{
|
||||||
var serialized = _serialization.Serialize(payload, context);
|
var serialized = _serialization.Serialize(payload, context);
|
||||||
await _interaction.PostMessageAsync(new NetworkMessageHeader
|
await _interaction.PostMessageAsync(new NetworkMessage
|
||||||
{ Id = id, MessageType = eMessageType, Data = serialized });
|
{ Id = id, MessageType = eMessageType, Data = serialized });
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -97,7 +97,7 @@ namespace mROA.Implementation
|
|||||||
IEndPointContext? context) where T : notnull
|
IEndPointContext? context) where T : notnull
|
||||||
{
|
{
|
||||||
var serialized = _serialization.Serialize(payload, context);
|
var serialized = _serialization.Serialize(payload, context);
|
||||||
await _interaction.PostMessageUntrustedAsync(new NetworkMessageHeader
|
await _interaction.PostMessageUntrustedAsync(new NetworkMessage
|
||||||
{ Id = id, MessageType = eMessageType, Data = serialized });
|
{ Id = id, MessageType = eMessageType, Data = serialized });
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,18 @@
|
|||||||
|
using System;
|
||||||
|
|
||||||
|
namespace mROA.Implementation
|
||||||
|
{
|
||||||
|
public struct RequestId
|
||||||
|
{
|
||||||
|
public ulong P0;
|
||||||
|
public ulong P1;
|
||||||
|
public RequestId Generate()
|
||||||
|
{
|
||||||
|
var guid = Guid.NewGuid();
|
||||||
|
var bytes = guid.ToByteArray();
|
||||||
|
var high = BitConverter.ToUInt64(bytes, 0);
|
||||||
|
var low = BitConverter.ToUInt64(bytes, 8);
|
||||||
|
return new RequestId{ P0 = high, P1 = low};
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user