попытка ускорить выполнение
This commit is contained in:
@@ -10,4 +10,5 @@ public interface INextGenerationInteractionModule : IInjectableModule
|
|||||||
Task PostMessage(NetworkMessage message);
|
Task PostMessage(NetworkMessage message);
|
||||||
void HandleMessage(NetworkMessage message);
|
void HandleMessage(NetworkMessage message);
|
||||||
NetworkMessage[] UnhandledMessages { get; }
|
NetworkMessage[] UnhandledMessages { get; }
|
||||||
|
NetworkMessage? FirstByFilter(Predicate<NetworkMessage> predicate);
|
||||||
}
|
}
|
||||||
@@ -6,6 +6,8 @@ public interface ISerializationToolkit : IInjectableModule
|
|||||||
byte[] Serialize(object objectToSerialize, Type type);
|
byte[] Serialize(object objectToSerialize, Type type);
|
||||||
T? Deserialize<T>(byte[] rawData);
|
T? Deserialize<T>(byte[] rawData);
|
||||||
object? Deserialize(byte[] rawData, Type type);
|
object? Deserialize(byte[] rawData, Type type);
|
||||||
|
T? Deserialize<T>(Span<byte> rawData);
|
||||||
|
object? Deserialize(Span<byte> rawData, Type type);
|
||||||
T Cast<T>(object nonCasted);
|
T Cast<T>(object nonCasted);
|
||||||
object Cast(object nonCasted, Type type);
|
object Cast(object nonCasted, Type type);
|
||||||
|
|
||||||
|
|||||||
@@ -25,6 +25,16 @@ public class JsonSerializationToolkit : ISerializationToolkit
|
|||||||
return JsonSerializer.Deserialize(rawData, type);
|
return JsonSerializer.Deserialize(rawData, type);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public T? Deserialize<T>(Span<byte> rawData)
|
||||||
|
{
|
||||||
|
return JsonSerializer.Deserialize<T>(rawData);
|
||||||
|
}
|
||||||
|
|
||||||
|
public object? Deserialize(Span<byte> rawData, Type type)
|
||||||
|
{
|
||||||
|
return JsonSerializer.Deserialize(rawData, type);
|
||||||
|
}
|
||||||
|
|
||||||
public T Cast<T>(object nonCasted)
|
public T Cast<T>(object nonCasted)
|
||||||
{
|
{
|
||||||
return nonCasted switch
|
return nonCasted switch
|
||||||
|
|||||||
@@ -9,8 +9,8 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
|
|||||||
public Stream? BaseStream { get; set; }
|
public Stream? BaseStream { get; set; }
|
||||||
private Task<NetworkMessage>? _currentReceiving;
|
private Task<NetworkMessage>? _currentReceiving;
|
||||||
private const int BufferSize = ushort.MaxValue;
|
private const int BufferSize = ushort.MaxValue;
|
||||||
private readonly byte[] _buffer = new byte[BufferSize];
|
private readonly Memory<byte> _buffer = new byte[BufferSize];
|
||||||
private readonly List<NetworkMessage> _messageBuffer = [];
|
private readonly List<NetworkMessage> _messageBuffer = new (128);
|
||||||
|
|
||||||
|
|
||||||
public void Inject<T>(T dependency)
|
public void Inject<T>(T dependency)
|
||||||
@@ -29,8 +29,6 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
|
|||||||
|
|
||||||
public Task<NetworkMessage> GetNextMessageReceiving()
|
public Task<NetworkMessage> GetNextMessageReceiving()
|
||||||
{
|
{
|
||||||
// Console.WriteLine(_currentReceiving?.Status);
|
|
||||||
|
|
||||||
if (_currentReceiving != null) return _currentReceiving;
|
if (_currentReceiving != null) return _currentReceiving;
|
||||||
_currentReceiving = Task.Run(GetNextMessage);
|
_currentReceiving = Task.Run(GetNextMessage);
|
||||||
return _currentReceiving;
|
return _currentReceiving;
|
||||||
@@ -59,6 +57,10 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
|
|||||||
}
|
}
|
||||||
|
|
||||||
public NetworkMessage[] UnhandledMessages => _messageBuffer.ToArray();
|
public NetworkMessage[] UnhandledMessages => _messageBuffer.ToArray();
|
||||||
|
public NetworkMessage? FirstByFilter(Predicate<NetworkMessage> predicate)
|
||||||
|
{
|
||||||
|
return _messageBuffer.FirstOrDefault(m => predicate(m));
|
||||||
|
}
|
||||||
|
|
||||||
private NetworkMessage GetNextMessage()
|
private NetworkMessage GetNextMessage()
|
||||||
{
|
{
|
||||||
@@ -70,15 +72,14 @@ public class NextGenerationInteractionModule : INextGenerationInteractionModule
|
|||||||
|
|
||||||
|
|
||||||
// Console.WriteLine("Receiving message");
|
// Console.WriteLine("Receiving message");
|
||||||
BaseStream.ReadExactly(_buffer, 0, 2);
|
var len = BitConverter.ToUInt16([(byte)BaseStream.ReadByte(), (byte)BaseStream.ReadByte()]);
|
||||||
var len = BitConverter.ToUInt16(_buffer, 0);
|
var localSpan = _buffer.Span.Slice(0, len);
|
||||||
BaseStream.ReadExactly(_buffer, 0, len);
|
BaseStream.ReadExactly(localSpan);
|
||||||
|
|
||||||
// Console.WriteLine("Receiving {0}", Encoding.Default.GetString(_buffer[..len]));
|
// Console.WriteLine("Receiving {0}", Encoding.Default.GetString(_buffer[..len]));
|
||||||
|
|
||||||
var message = _serialization.Deserialize<NetworkMessage>(_buffer[..len]);
|
var message = _serialization.Deserialize<NetworkMessage>(localSpan);
|
||||||
_messageBuffer.Add(message!);
|
_messageBuffer.Add(message!);
|
||||||
|
|
||||||
_currentReceiving = Task.Run(GetNextMessage);
|
_currentReceiving = Task.Run(GetNextMessage);
|
||||||
|
|
||||||
return message!;
|
return message!;
|
||||||
|
|||||||
@@ -44,9 +44,10 @@ public class RepresentationModule : IRepresentationModule
|
|||||||
throw new NullReferenceException("Interaction toolkit is not initialized");
|
throw new NullReferenceException("Interaction toolkit is not initialized");
|
||||||
|
|
||||||
var fromBuffer =
|
var fromBuffer =
|
||||||
_interaction.UnhandledMessages.FirstOrDefault(message =>
|
_interaction.FirstByFilter(message =>
|
||||||
(requestId is null || message.Id == requestId) &&
|
(requestId is null || message.Id == requestId) &&
|
||||||
(messageType is null || message.SchemaId == messageType));
|
(messageType is null || message.SchemaId == messageType));
|
||||||
|
|
||||||
if (fromBuffer == null)
|
if (fromBuffer == null)
|
||||||
{
|
{
|
||||||
while (true)
|
while (true)
|
||||||
|
|||||||
Reference in New Issue
Block a user