Remove non-context serialization module

This commit is contained in:
2025-05-04 19:43:15 +03:00
parent eede4d28c4
commit 9575f745f9
19 changed files with 138 additions and 131 deletions
@@ -9,6 +9,7 @@ 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; }
@@ -1,9 +1,8 @@
using System; using System;
using mROA.Abstract;
namespace mROA.Cbor namespace mROA.Abstract
{ {
public interface IContextualSerializationToolKit : ISerializationToolkit public interface IContextualSerializationToolKit
{ {
byte[] Serialize(object objectToSerialize, IEndPointContext? context); byte[] Serialize(object objectToSerialize, IEndPointContext? context);
void Serialize(object objectToSerialize, Span<byte> destination, IEndPointContext? context); void Serialize(object objectToSerialize, Span<byte> destination, IEndPointContext? context);
+1 -1
View File
@@ -1,6 +1,6 @@
namespace mROA.Abstract namespace mROA.Abstract
{ {
public interface IEndPointContext public interface IEndPointContext : IInjectableModule
{ {
IContextRepository RealRepository { get; } IContextRepository RealRepository { get; }
IContextRepository RemoteRepository { get; } IContextRepository RemoteRepository { get; }
+11 -6
View File
@@ -11,15 +11,20 @@ namespace mROA.Abstract
int Id { get; } int Id { get; }
Task<(object? Deserialized, EMessageType MessageType)> GetSingle(Predicate<NetworkMessageHeader> rule, Task<(object? Deserialized, EMessageType MessageType)> GetSingle(Predicate<NetworkMessageHeader> rule,
CancellationToken token, IEndPointContext? context, CancellationToken token = default,
params Func<NetworkMessageHeader, Type?>[] converter); params Func<NetworkMessageHeader, Type?>[] converter);
IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(Predicate<NetworkMessageHeader> rule, CancellationToken token, IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(Predicate<NetworkMessageHeader> rule,
IEndPointContext? context, CancellationToken token = default,
params Func<NetworkMessageHeader, Type?>[] converter); params Func<NetworkMessageHeader, Type?>[] converter);
Task PostCallMessageAsync<T>(Guid id, EMessageType eMessageType, T payload) where T : notnull; Task PostCallMessageAsync<T>(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context)
void PostCallMessage<T>(Guid id, EMessageType eMessageType, T payload) where T : notnull; where T : notnull;
Task PostCallMessageUntrustedAsync<T>(Guid id, EMessageType eMessageType, T payload) where T : notnull;
void PostCallMessage(Guid id, EMessageType eMessageType, object payload, Type payloadType); void PostCallMessage<T>(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context)
where T : notnull;
Task PostCallMessageUntrustedAsync<T>(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context)
where T : notnull;
} }
} }
+11 -11
View File
@@ -2,15 +2,15 @@
namespace mROA.Abstract namespace mROA.Abstract
{ {
public interface ISerializationToolkit : IInjectableModule // public interface IContextualSerializationToolKit : IInjectableModule
{ // {
byte[] Serialize<T>(T objectToSerialize); // byte[] Serialize<T>(T objectToSerialize);
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); // T? Deserialize<T>(Span<byte> rawData);
object? Deserialize(Span<byte> rawData, Type type); // 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);
} // }
} }
@@ -9,7 +9,7 @@ namespace mROA.Implementation.Backend
{ {
private ICancellationRepository? _cancellationRepo; private ICancellationRepository? _cancellationRepo;
private IMethodRepository? _methodRepo; private IMethodRepository? _methodRepo;
private ISerializationToolkit? _serialization; private IContextualSerializationToolKit? _serialization;
public void Inject<T>(T dependency) public void Inject<T>(T dependency)
{ {
@@ -21,7 +21,7 @@ namespace mROA.Implementation.Backend
case ICancellationRepository cancellationRepo: case ICancellationRepository cancellationRepo:
_cancellationRepo = cancellationRepo; _cancellationRepo = cancellationRepo;
break; break;
case ISerializationToolkit serializationToolkit: case IContextualSerializationToolKit serializationToolkit:
_serialization = serializationToolkit; _serialization = serializationToolkit;
break; break;
} }
+2 -2
View File
@@ -7,7 +7,7 @@ namespace mROA.Implementation.Backend
public class ConnectionHub : IConnectionHub public class ConnectionHub : IConnectionHub
{ {
private readonly Dictionary<int, IChannelInteractionModule> _connections = new(); private readonly Dictionary<int, IChannelInteractionModule> _connections = new();
private ISerializationToolkit? _serializationToolkit; private IContextualSerializationToolKit? _serializationToolkit;
public void RegisterInteraction(IChannelInteractionModule interaction) public void RegisterInteraction(IChannelInteractionModule interaction)
{ {
@@ -31,7 +31,7 @@ namespace mROA.Implementation.Backend
public void Inject<T>(T dependency) public void Inject<T>(T dependency)
{ {
if (dependency is ISerializationToolkit serializationToolkit) if (dependency is IContextualSerializationToolKit serializationToolkit)
_serializationToolkit = serializationToolkit; _serializationToolkit = serializationToolkit;
} }
} }
@@ -10,7 +10,7 @@ namespace mROA.Implementation.Backend
private IContextRepository? _contextRepository; private IContextRepository? _contextRepository;
private IContextRepository? _remoteContextRepository; private IContextRepository? _remoteContextRepository;
private IMethodRepository? _methodRepository; private IMethodRepository? _methodRepository;
private ISerializationToolkit? _serializationToolkit; private IContextualSerializationToolKit? _serializationToolkit;
private IExecuteModule? _executeModule; private IExecuteModule? _executeModule;
private readonly Type _extractorType; private readonly Type _extractorType;
@@ -37,7 +37,7 @@ namespace mROA.Implementation.Backend
case IMethodRepository methodRepository: case IMethodRepository methodRepository:
_methodRepository = methodRepository; _methodRepository = methodRepository;
break; break;
case ISerializationToolkit serializationToolkit: case IContextualSerializationToolKit serializationToolkit:
_serializationToolkit = serializationToolkit; _serializationToolkit = serializationToolkit;
break; break;
case IExecuteModule executeModule: case IExecuteModule executeModule:
@@ -15,7 +15,7 @@ namespace mROA.Implementation.Backend
private readonly Type? _interactionModuleType; private readonly Type? _interactionModuleType;
private readonly TcpListener _tcpListener; private readonly TcpListener _tcpListener;
private IConnectionHub? _hub; private IConnectionHub? _hub;
private ISerializationToolkit? _serialization; private IContextualSerializationToolKit? _serialization;
private Dictionary<int, CancellationTokenSource> _extractorsCTS = new(); private Dictionary<int, CancellationTokenSource> _extractorsCTS = new();
@@ -58,7 +58,7 @@ namespace mROA.Implementation.Backend
case IConnectionHub interactionModule: case IConnectionHub interactionModule:
_hub = interactionModule; _hub = interactionModule;
break; break;
case ISerializationToolkit serializationToolkit: case IContextualSerializationToolKit serializationToolkit:
_serialization = serializationToolkit; _serialization = serializationToolkit;
break; break;
} }
+2 -2
View File
@@ -16,7 +16,7 @@ namespace mROA.Implementation.Backend
private UdpClient _client; private UdpClient _client;
private Dictionary<IPEndPoint, int> _reservedPorts = new(); private Dictionary<IPEndPoint, int> _reservedPorts = new();
private CancellationTokenSource _tokenSource = new(); private CancellationTokenSource _tokenSource = new();
private ISerializationToolkit _serializationToolkit; private IContextualSerializationToolKit _serializationToolkit;
public UdpGateway(IPEndPoint listeningEndpoint) public UdpGateway(IPEndPoint listeningEndpoint)
{ {
@@ -31,7 +31,7 @@ namespace mROA.Implementation.Backend
case IConnectionHub hub: case IConnectionHub hub:
_hub = hub; _hub = hub;
break; break;
case ISerializationToolkit serializationToolkit: case IContextualSerializationToolKit serializationToolkit:
_serializationToolkit = serializationToolkit; _serializationToolkit = serializationToolkit;
break; break;
} }
+13 -16
View File
@@ -1,7 +1,5 @@
using System; using System;
using System.Collections.Generic;
using System.IO; using System.IO;
using System.Linq;
using System.Threading; using System.Threading;
using System.Threading.Channels; using System.Threading.Channels;
using System.Threading.Tasks; using System.Threading.Tasks;
@@ -16,13 +14,12 @@ 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 readonly List<NetworkMessageHeader> _messageBuffer = new(128);
private Task<NetworkMessageHeader>? _currentReceiving; private Task<NetworkMessageHeader>? _currentReceiving;
private ISerializationToolkit? _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; }
public ChannelInteractionModule() public ChannelInteractionModule()
{ {
ReceiveChanel = Channel.CreateUnbounded<NetworkMessageHeader>(new UnboundedChannelOptions ReceiveChanel = Channel.CreateUnbounded<NetworkMessageHeader>(new UnboundedChannelOptions
@@ -58,12 +55,15 @@ namespace mROA.Implementation
{ {
switch (dependency) switch (dependency)
{ {
case ISerializationToolkit toolkit: case IContextualSerializationToolKit toolkit:
_serialization = toolkit; _serialization = toolkit;
break; break;
case IIdentityGenerator identityGenerator: case IIdentityGenerator identityGenerator:
ConnectionId = identityGenerator.GetNextIdentity(); ConnectionId = identityGenerator.GetNextIdentity();
break; break;
case IEndPointContext endpointContext:
Context = endpointContext;
break;
} }
} }
@@ -122,11 +122,6 @@ namespace mROA.Implementation
await _untrustedWriter.WriteAsync(messageHeader); await _untrustedWriter.WriteAsync(messageHeader);
} }
public NetworkMessageHeader? FirstByFilter(Predicate<NetworkMessageHeader> predicate)
{
return _messageBuffer.FirstOrDefault(m => predicate(m));
}
public event Action<int>? OnDisconnected; public event Action<int>? OnDisconnected;
public async Task Restart(bool sendRecovery) public async Task Restart(bool sendRecovery)
@@ -189,17 +184,18 @@ namespace mROA.Implementation
public class StreamExtractor public class StreamExtractor
{ {
private readonly Stream _ioStream; private readonly Stream _ioStream;
private readonly ISerializationToolkit _serializationToolkit; private readonly IContextualSerializationToolKit _serializationToolkit;
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;
public readonly int Id = new Random().Next(); public readonly int Id = new Random().Next();
public StreamExtractor(Stream ioStream, ISerializationToolkit serializationToolkit) public StreamExtractor(Stream ioStream, IContextualSerializationToolKit serializationToolkit, IEndPointContext context)
{ {
_ioStream = ioStream; _ioStream = ioStream;
_serializationToolkit = serializationToolkit; _serializationToolkit = serializationToolkit;
Context = context;
} }
public Action<NetworkMessageHeader> MessageReceived = _ => { }; public Action<NetworkMessageHeader> MessageReceived = _ => { };
@@ -231,7 +227,7 @@ namespace mROA.Implementation
await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: Token); await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: Token);
var message = _serializationToolkit.Deserialize<NetworkMessageHeader>(localSpan.Span); var message = _serializationToolkit.Deserialize<NetworkMessageHeader>(localSpan, Context);
#if TRACE #if TRACE
Console.WriteLine( Console.WriteLine(
$"{DateTime.Now.TimeOfDay} [{Id}] Received Message {message.Id} - {message.MessageType}"); $"{DateTime.Now.TimeOfDay} [{Id}] Received Message {message.Id} - {message.MessageType}");
@@ -254,9 +250,10 @@ namespace mROA.Implementation
public async Task Send(NetworkMessageHeader message, CancellationToken token = default) public async Task Send(NetworkMessageHeader message, CancellationToken token = default)
{ {
try try
{ {
var rawMessage = _serializationToolkit.Serialize(message); 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));
#if TRACE #if TRACE
+4
View File
@@ -16,5 +16,9 @@ namespace mROA.Implementation
// ReSharper disable once UnusedMember.Global // ReSharper disable once UnusedMember.Global
set { OwnerFunc = () => value; } set { OwnerFunc = () => value; }
} }
public void Inject<T>(T dependency)
{
}
} }
} }
@@ -16,7 +16,7 @@ namespace mROA.Implementation.Frontend
private readonly IPEndPoint _serverEndPoint; private readonly IPEndPoint _serverEndPoint;
private TcpClient _tcpClient = new(); private TcpClient _tcpClient = new();
private IChannelInteractionModule? _interactionModule; private IChannelInteractionModule? _interactionModule;
private ISerializationToolkit? _serialization; private IContextualSerializationToolKit? _serialization;
private ChannelInteractionModule.StreamExtractor _currentExtractor; private ChannelInteractionModule.StreamExtractor _currentExtractor;
private CancellationTokenSource _rawExtractorCancellation; private CancellationTokenSource _rawExtractorCancellation;
@@ -33,7 +33,7 @@ namespace mROA.Implementation.Frontend
case ChannelInteractionModule interactionModule: case ChannelInteractionModule interactionModule:
_interactionModule = interactionModule; _interactionModule = interactionModule;
break; break;
case ISerializationToolkit toolkit: case IContextualSerializationToolKit toolkit:
_serialization = toolkit; _serialization = toolkit;
break; break;
} }
@@ -16,7 +16,7 @@ namespace mROA.Implementation.Frontend
private IContextRepository? _realContextRepository; private IContextRepository? _realContextRepository;
private IContextRepository? _remoteContextRepository; private IContextRepository? _remoteContextRepository;
private IRepresentationModule? _representationModule; private IRepresentationModule? _representationModule;
private ISerializationToolkit? _serializationToolkit; private IContextualSerializationToolKit? _serializationToolkit;
public void Inject<T>(T dependency) public void Inject<T>(T dependency)
{ {
@@ -38,7 +38,7 @@ namespace mROA.Implementation.Frontend
case IRepresentationModule representationModule: case IRepresentationModule representationModule:
_representationModule = representationModule; _representationModule = representationModule;
break; break;
case ISerializationToolkit serializationToolkit: case IContextualSerializationToolKit serializationToolkit:
_serializationToolkit = serializationToolkit; _serializationToolkit = serializationToolkit;
break; break;
} }
@@ -9,7 +9,7 @@ namespace mROA.Implementation.Frontend
{ {
public class UdpUntrustedInteraction : IUntrustedInteractionModule public class UdpUntrustedInteraction : IUntrustedInteractionModule
{ {
private ISerializationToolkit _serializationToolkit; private IContextualSerializationToolKit _serializationToolkit;
private IChannelInteractionModule _channelInteractionModule; private IChannelInteractionModule _channelInteractionModule;
private CancellationTokenSource _tokenSource = new CancellationTokenSource(); private CancellationTokenSource _tokenSource = new CancellationTokenSource();
@@ -77,7 +77,7 @@ namespace mROA.Implementation.Frontend
case IChannelInteractionModule channelModule: case IChannelInteractionModule channelModule:
_channelInteractionModule = channelModule; _channelInteractionModule = channelModule;
break; break;
case ISerializationToolkit serializationToolkit: case IContextualSerializationToolKit serializationToolkit:
_serializationToolkit = serializationToolkit; _serializationToolkit = serializationToolkit;
break; break;
} }
+54 -54
View File
@@ -4,58 +4,58 @@ using mROA.Abstract;
namespace mROA.Implementation namespace mROA.Implementation
{ {
public class JsonSerializationToolkit : ISerializationToolkit // public class JsonSerializationToolkit : IContextualSerializationToolKit
{ // {
public byte[] Serialize<T>(T objectToSerialize) // public byte[] Serialize<T>(T objectToSerialize)
{ // {
return JsonSerializer.SerializeToUtf8Bytes(objectToSerialize); // return JsonSerializer.SerializeToUtf8Bytes(objectToSerialize);
} // }
//
public byte[] Serialize(object objectToSerialize, Type type) // public byte[] Serialize(object objectToSerialize, Type type)
{ // {
return JsonSerializer.SerializeToUtf8Bytes(objectToSerialize, type); // return JsonSerializer.SerializeToUtf8Bytes(objectToSerialize, type);
} // }
//
public T? Deserialize<T>(byte[] rawData) // public T? Deserialize<T>(byte[] rawData)
{ // {
return JsonSerializer.Deserialize<T>(rawData); // return JsonSerializer.Deserialize<T>(rawData);
} // }
//
public object? Deserialize(byte[] rawData, Type type) // public object? Deserialize(byte[] rawData, Type type)
{ // {
return JsonSerializer.Deserialize(rawData, type); // return JsonSerializer.Deserialize(rawData, type);
} // }
//
public T? Deserialize<T>(Span<byte> rawData) // public T? Deserialize<T>(Span<byte> rawData)
{ // {
return JsonSerializer.Deserialize<T>(rawData); // return JsonSerializer.Deserialize<T>(rawData);
} // }
//
public object? Deserialize(Span<byte> rawData, Type type) // public object? Deserialize(Span<byte> rawData, Type type)
{ // {
return JsonSerializer.Deserialize(rawData, type); // return JsonSerializer.Deserialize(rawData, type);
} // }
//
public T Cast<T>(object nonCasted) // public T Cast<T>(object nonCasted)
{ // {
return nonCasted switch // return nonCasted switch
{ // {
JsonElement jsonElement => jsonElement.Deserialize<T>()!, // JsonElement jsonElement => jsonElement.Deserialize<T>()!,
T casted => casted, // T casted => casted,
_ => throw new JsonException("Cannot cast object to type " + typeof(T).FullName) // _ => throw new JsonException("Cannot cast object to type " + typeof(T).FullName)
}; // };
} // }
//
public object Cast(object nonCasted, Type type) // public object Cast(object nonCasted, Type type)
{ // {
if (nonCasted is JsonElement jsonElement) // if (nonCasted is JsonElement jsonElement)
return jsonElement.Deserialize(type)!; // return jsonElement.Deserialize(type)!;
//
throw new JsonException("Cannot cast object to type " + type.FullName); // throw new JsonException("Cannot cast object to type " + type.FullName);
} // }
//
public void Inject<T>(T dependency) // public void Inject<T>(T dependency)
{ // {
} // }
} // }
} }
+1 -1
View File
@@ -34,7 +34,7 @@ namespace mROA.Implementation
MessageType = EMessageType.Unknown; MessageType = EMessageType.Unknown;
Data = Array.Empty<byte>(); Data = Array.Empty<byte>();
} }
public NetworkMessageHeader(ISerializationToolkit serializationToolkit, INetworkMessage networkMessage) public NetworkMessageHeader(IContextualSerializationToolKit serializationToolkit, INetworkMessage networkMessage)
{ {
MessageType = networkMessage.MessageType; MessageType = networkMessage.MessageType;
Data = serializationToolkit.Serialize(networkMessage); Data = serializationToolkit.Serialize(networkMessage);
+4 -1
View File
@@ -1,4 +1,5 @@
using System; using System;
using System.Net;
using System.Threading; using System.Threading;
using System.Threading.Tasks; using System.Threading.Tasks;
using mROA.Abstract; using mROA.Abstract;
@@ -10,6 +11,7 @@ namespace mROA.Implementation
{ {
public abstract class RemoteObjectBase : IDisposable public abstract class RemoteObjectBase : IDisposable
{ {
private readonly IEndPointContext _context;
public bool Equals(RemoteObjectBase other) public bool Equals(RemoteObjectBase other)
{ {
return _identifier.Equals(other._identifier); return _identifier.Equals(other._identifier);
@@ -31,10 +33,11 @@ namespace mROA.Implementation
private readonly ComplexObjectIdentifier _identifier; private readonly ComplexObjectIdentifier _identifier;
private readonly IRepresentationModule _representationModule; private readonly IRepresentationModule _representationModule;
protected RemoteObjectBase(int id, IRepresentationModule representationModule) protected RemoteObjectBase(int id, IRepresentationModule representationModule, IEndPointContext context)
{ {
_identifier = new ComplexObjectIdentifier { ContextId = id, OwnerId = representationModule.Id }; _identifier = new ComplexObjectIdentifier { ContextId = id, OwnerId = representationModule.Id };
_representationModule = representationModule; _representationModule = representationModule;
_context = context;
} }
public int Id => _identifier.ContextId; public int Id => _identifier.ContextId;
+17 -19
View File
@@ -11,13 +11,13 @@ namespace mROA.Implementation
public class RepresentationModule : IRepresentationModule public class RepresentationModule : IRepresentationModule
{ {
private IChannelInteractionModule? _interaction; private IChannelInteractionModule? _interaction;
private ISerializationToolkit? _serialization; private IContextualSerializationToolKit? _serialization;
public void Inject<T>(T dependency) public void Inject<T>(T dependency)
{ {
switch (dependency) switch (dependency)
{ {
case ISerializationToolkit toolkit: case IContextualSerializationToolKit toolkit:
_serialization = toolkit; _serialization = toolkit;
break; break;
case IChannelInteractionModule interactionModule: case IChannelInteractionModule interactionModule:
@@ -31,7 +31,7 @@ namespace mROA.Implementation
.ConnectionId; .ConnectionId;
public async Task<(object? Deserialized, EMessageType MessageType)> GetSingle( public async Task<(object? Deserialized, EMessageType MessageType)> GetSingle(
Predicate<NetworkMessageHeader> rule, Predicate<NetworkMessageHeader> rule, IEndPointContext? context,
CancellationToken token = default, params Func<NetworkMessageHeader, Type?>[] converter) CancellationToken token = default, params Func<NetworkMessageHeader, Type?>[] converter)
{ {
var writer = _interaction.ReceiveChanel.Writer; var writer = _interaction.ReceiveChanel.Writer;
@@ -47,7 +47,7 @@ namespace mROA.Implementation
} }
var type = converter.Select(i => i(message)).First(i => i != null)!; var type = converter.Select(i => i(message)).First(i => i != null)!;
var deserialized = _serialization.Deserialize(message.Data, type); var deserialized = _serialization.Deserialize(message.Data, type, context);
return (deserialized, message.MessageType); return (deserialized, message.MessageType);
} }
@@ -55,7 +55,7 @@ namespace mROA.Implementation
} }
public async IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream( public async IAsyncEnumerable<(object parced, EMessageType originalType)> GetStream(
Predicate<NetworkMessageHeader> rule, [EnumeratorCancellation] CancellationToken token = default, Predicate<NetworkMessageHeader> rule, IEndPointContext? context, [EnumeratorCancellation] CancellationToken token = default,
params Func<NetworkMessageHeader, Type?>[] converter) params Func<NetworkMessageHeader, Type?>[] converter)
{ {
var writer = _interaction?.ReceiveChanel.Writer; var writer = _interaction?.ReceiveChanel.Writer;
@@ -68,43 +68,41 @@ namespace mROA.Implementation
} }
var type = converter.Select(i => i(message)).First(i => i != null)!; var type = converter.Select(i => i(message)).First(i => i != null)!;
var deserialized = _serialization.Deserialize(message.Data, type); var deserialized = _serialization.Deserialize(message.Data, type, context);
yield return (deserialized, message.MessageType)!; yield return (deserialized, message.MessageType)!;
} }
} }
public async Task PostCallMessageAsync<T>(Guid id, EMessageType eMessageType, T payload) where T : notnull public async Task PostCallMessageAsync<T>(Guid id, EMessageType eMessageType, T payload,
IEndPointContext? context) where T : notnull
{ {
await PostCallMessageAsync(id, eMessageType, payload, typeof(T)); await PostCallMessageAsync(id, eMessageType, payload, context);
} }
public async Task PostCallMessageAsync(Guid id, EMessageType eMessageType, object payload, Type payloadType) public async Task PostCallMessageAsync(Guid id, EMessageType eMessageType, object payload,
IEndPointContext? context)
{ {
if (_interaction == null) if (_interaction == null)
throw new NullReferenceException("Interaction toolkit is not initialized"); throw new NullReferenceException("Interaction toolkit is not initialized");
if (_serialization == null) if (_serialization == null)
throw new NullReferenceException("Serialization toolkit is not initialized"); throw new NullReferenceException("Serialization toolkit is not initialized");
var serialized = _serialization.Serialize(payload, payloadType); var serialized = _serialization.Serialize(payload, context);
await _interaction.PostMessageAsync(new NetworkMessageHeader await _interaction.PostMessageAsync(new NetworkMessageHeader
{ Id = id, MessageType = eMessageType, Data = serialized }); { Id = id, MessageType = eMessageType, Data = serialized });
} }
public void PostCallMessage<T>(Guid id, EMessageType eMessageType, T payload) where T : notnull public void PostCallMessage<T>(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull
{ {
PostCallMessageAsync(id, eMessageType, payload).GetAwaiter().GetResult(); PostCallMessageAsync(id, eMessageType, payload, context).GetAwaiter().GetResult();
} }
public async Task PostCallMessageUntrustedAsync<T>(Guid id, EMessageType eMessageType, T payload) where T : notnull public async Task PostCallMessageUntrustedAsync<T>(Guid id, EMessageType eMessageType, T payload,
IEndPointContext? context) where T : notnull
{ {
var serialized = _serialization.Serialize(payload, typeof(T)); var serialized = _serialization.Serialize(payload, context);
await _interaction.PostMessageUntrustedAsync(new NetworkMessageHeader await _interaction.PostMessageUntrustedAsync(new NetworkMessageHeader
{ Id = id, MessageType = eMessageType, Data = serialized }); { Id = id, MessageType = eMessageType, Data = serialized });
} }
public void PostCallMessage(Guid id, EMessageType eMessageType, object payload, Type payloadType)
{
PostCallMessageAsync(id, eMessageType, payload, payloadType).GetAwaiter().GetResult();
}
} }
} }