From 17a6d3aad6193be2940be47c8f5301e206b34f74 Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Mon, 14 Jul 2025 21:03:12 +0300 Subject: [PATCH] New distribution model base --- Example.Backend/Program.cs | 1 + mROA/Abstract/IPrimaryMessageDistributior.cs | 15 ++++ mROA/Abstract/IRequestExtractor.cs | 5 ++ .../Backend/BasicExecutionModule.cs | 13 ---- .../Backend/HubRequestExtractor.cs | 2 +- .../Backend/NetworkGatewayModule.cs | 8 +- .../ChannelDistributorFactory.cs | 35 +++++++++ .../ChannelInteractionModule.cs | 3 - .../CreativeRepresentationModuleProducer.cs | 3 - .../Frontend/NetworkFrontendBridge.cs | 5 -- .../Frontend/RequestExtractor.cs | 73 +++++++++---------- .../RemoteInstanceRepository.cs | 5 -- mROA/Implementation/RepresentationModule.cs | 5 -- 13 files changed, 94 insertions(+), 79 deletions(-) create mode 100644 mROA/Abstract/IPrimaryMessageDistributior.cs create mode 100644 mROA/Implementation/ChannelDistributorFactory.cs diff --git a/Example.Backend/Program.cs b/Example.Backend/Program.cs index b31d198..5711803 100644 --- a/Example.Backend/Program.cs +++ b/Example.Backend/Program.cs @@ -40,6 +40,7 @@ class Program return repo; })); + builder.Services.AddSingleton(); builder.Services.AddSingleton(p => { var methodRepo = new CollectableMethodRepository(); diff --git a/mROA/Abstract/IPrimaryMessageDistributior.cs b/mROA/Abstract/IPrimaryMessageDistributior.cs new file mode 100644 index 0000000..aa0295c --- /dev/null +++ b/mROA/Abstract/IPrimaryMessageDistributior.cs @@ -0,0 +1,15 @@ +using System.Threading.Tasks; +using mROA.Implementation; + +namespace mROA.Abstract +{ + public interface IPrimaryMessageDistributior + { + Task Distribute(NetworkMessageHeader message); + } + + public interface IMessageDistributorFactory + { + IPrimaryMessageDistributior Produce(int clientId); + } +} \ No newline at end of file diff --git a/mROA/Abstract/IRequestExtractor.cs b/mROA/Abstract/IRequestExtractor.cs index f64e8a7..2cc176f 100644 --- a/mROA/Abstract/IRequestExtractor.cs +++ b/mROA/Abstract/IRequestExtractor.cs @@ -1,9 +1,14 @@ +using System; using System.Threading.Tasks; +using mROA.Implementation; namespace mROA.Abstract { public interface IRequestExtractor { Task StartExtraction(); + void PushMessage(object parced, EMessageType originalType); + Predicate Rule { get; } + Func[] Converters { get; } } } \ No newline at end of file diff --git a/mROA/Implementation/Backend/BasicExecutionModule.cs b/mROA/Implementation/Backend/BasicExecutionModule.cs index 31f2735..b32cdfd 100644 --- a/mROA/Implementation/Backend/BasicExecutionModule.cs +++ b/mROA/Implementation/Backend/BasicExecutionModule.cs @@ -23,7 +23,6 @@ namespace mROA.Implementation.Backend { try { - ThrowIfNotInjected(instanceRepository); if (command is CancelRequest) { return CancelExecution(command); @@ -96,18 +95,6 @@ namespace mROA.Implementation.Backend return castedParams; } - private void ThrowIfNotInjected(IInstanceRepository instanceRepository) - { - if (_cancellationRepo is null) - throw new NullReferenceException("Method repository was not defined"); - - if (_methodRepo is null) - throw new NullReferenceException("Method repository was not defined"); - - if (instanceRepository is null) - throw new NullReferenceException("Context repository was not defined"); - } - private FinalCommandExecution CancelExecution(ICallRequest command) { var cts = _cancellationRepo.GetCancellation(command.Id); diff --git a/mROA/Implementation/Backend/HubRequestExtractor.cs b/mROA/Implementation/Backend/HubRequestExtractor.cs index 2b09ae2..7d5a2cb 100644 --- a/mROA/Implementation/Backend/HubRequestExtractor.cs +++ b/mROA/Implementation/Backend/HubRequestExtractor.cs @@ -37,7 +37,7 @@ namespace mROA.Implementation.Backend private IRequestExtractor CreateExtractor(IRepresentationModule interaction) { - var extractor = new RequestExtractor(_executeModule, _methodRepository, interaction, _serializationToolkit, interaction.Context); + var extractor = new RequestExtractor(_executeModule, interaction, interaction.Context); var context = interaction.Context; if (_contextRepository is IContextRepositoryHub contextHub) context.RealRepository = contextHub.GetRepository(interaction.Id); diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index 8661876..a8f485e 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -11,21 +11,21 @@ namespace mROA.Implementation.Backend { public class NetworkGatewayModule : IGatewayModule { - private readonly IServiceProvider _serviceProvider; private readonly TcpListener _tcpListener; private readonly IConnectionHub _hub; private readonly IContextualSerializationToolKit _serialization; private readonly Dictionary _extractorsCTS = new(); private ICallIndexProvider _callIndexProvider; private readonly IIdentityGenerator _identityGenerator; - public NetworkGatewayModule(IOptions options, IServiceProvider service, IIdentityGenerator identityGenerator, IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub) + private readonly IMessageDistributorFactory _distributorFactory; + public NetworkGatewayModule(IOptions options, IIdentityGenerator identityGenerator, IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, IMessageDistributorFactory distributorFactory) { _tcpListener = new(options.Value.Endpoint); - _serviceProvider = service; _identityGenerator = identityGenerator; _serialization = serialization; _callIndexProvider = callIndexProvider; _hub = hub; + _distributorFactory = distributorFactory; } public void Run() @@ -74,7 +74,7 @@ namespace mROA.Implementation.Backend interaction.PostMessageAsync(new NetworkMessageHeader(_serialization!, new IdAssignment { Id = interaction.ConnectionId }, null)); _extractorsCTS[interaction.ConnectionId] = cts; - _hub!.RegisterInteraction(interaction); + _hub.RegisterInteraction(interaction); Console.WriteLine("Client registered"); break; case EMessageType.ClientRecovery: diff --git a/mROA/Implementation/ChannelDistributorFactory.cs b/mROA/Implementation/ChannelDistributorFactory.cs new file mode 100644 index 0000000..5998694 --- /dev/null +++ b/mROA/Implementation/ChannelDistributorFactory.cs @@ -0,0 +1,35 @@ +using System.Threading.Channels; +using System.Threading.Tasks; +using mROA.Abstract; + +namespace mROA.Implementation +{ + public class ChannelDistributorFactory : IMessageDistributorFactory + { + private readonly IConnectionHub _hub; + + public ChannelDistributorFactory(IConnectionHub hub) + { + _hub = hub; + } + public IPrimaryMessageDistributior Produce(int clientId) + { + return new ChannelMessageDistributor(_hub.GetInteraction(clientId).ReceiveChanel.Writer); + } + } + + public class ChannelMessageDistributor : IPrimaryMessageDistributior + { + private readonly ChannelWriter _writer; + + public ChannelMessageDistributor(ChannelWriter writer) + { + _writer = writer; + } + + public async Task Distribute(NetworkMessageHeader message) + { + await _writer.WriteAsync(message); + } + } +} \ No newline at end of file diff --git a/mROA/Implementation/ChannelInteractionModule.cs b/mROA/Implementation/ChannelInteractionModule.cs index 300f2c9..0c6b827 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -82,9 +82,6 @@ namespace mROA.Implementation public async Task PostMessageAsync(NetworkMessageHeader messageHeader) { - if (_serialization == null) - throw new NullReferenceException("Serialization toolkit is not initialized"); - while (true) { if (await PostMessageInternal(messageHeader)) diff --git a/mROA/Implementation/CreativeRepresentationModuleProducer.cs b/mROA/Implementation/CreativeRepresentationModuleProducer.cs index fb9bb2c..66356a1 100644 --- a/mROA/Implementation/CreativeRepresentationModuleProducer.cs +++ b/mROA/Implementation/CreativeRepresentationModuleProducer.cs @@ -18,9 +18,6 @@ namespace mROA.Implementation public IRepresentationModule Produce(int id) { - if (_hub == null) - throw new NullReferenceException("Interaction module is null"); - var interaction = _hub.GetInteraction(id); var produced = new RepresentationModule(interaction, _serialization); diff --git a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs index 7b096ea..e3f3d8e 100644 --- a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs +++ b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs @@ -31,11 +31,6 @@ namespace mROA.Implementation.Frontend public async Task Connect() { - if (_interactionModule is null) - throw new Exception("Interaction module was not injected"); - if (_serialization == null) - throw new NullReferenceException("Serialization toolkit is not initialized"); - _tcpClient.Connect(_serverEndPoint); _tcpClient.NoDelay = true; PrepareExtractor(); diff --git a/mROA/Implementation/Frontend/RequestExtractor.cs b/mROA/Implementation/Frontend/RequestExtractor.cs index eb240e0..838f4e2 100644 --- a/mROA/Implementation/Frontend/RequestExtractor.cs +++ b/mROA/Implementation/Frontend/RequestExtractor.cs @@ -9,72 +9,65 @@ namespace mROA.Implementation.Frontend { public class RequestExtractor : IRequestExtractor { + private IExecuteModule _executeModule; - private IMethodRepository _methodRepository; - private IRepresentationModule _representationModule; - private IContextualSerializationToolKit _serializationToolkit; private IEndPointContext _context; - public RequestExtractor(IExecuteModule executeModule, IMethodRepository methodRepository, IRepresentationModule representationModule, IContextualSerializationToolKit serializationToolkit, IEndPointContext context) + public RequestExtractor(IExecuteModule executeModule, IRepresentationModule representationModule, IEndPointContext context) { _executeModule = executeModule; - _methodRepository = methodRepository; _representationModule = representationModule; - _serializationToolkit = serializationToolkit; _context = context; } public async Task StartExtraction() { - ThrowIfNotInjected(); - var streamTokenSource = new CancellationTokenSource(); - var query = _representationModule.GetStream(m => - m.MessageType is EMessageType.CallRequest or EMessageType.CancelRequest - or EMessageType.EventRequest or EMessageType.ClientDisconnect, _context, + var query = _representationModule.GetStream(Rule, _context, streamTokenSource.Token, - m => m.MessageType == EMessageType.CallRequest ? typeof(DefaultCallRequest) : null, - m => m.MessageType == EMessageType.CancelRequest ? typeof(CancelRequest) : null, - m => m.MessageType == EMessageType.EventRequest ? typeof(DefaultCallRequest) : null, - m => m.MessageType == EMessageType.ClientDisconnect ? typeof(ClientDisconnect) : null); + Converters); await foreach (var command in query) { - switch (command.originalType) - { - case EMessageType.CallRequest: - HandleCallRequest((DefaultCallRequest)command.parced); - break; - case EMessageType.ClientDisconnect: - return; - case EMessageType.EventRequest: - HandleEventRequest((DefaultCallRequest)command.parced); - break; - case EMessageType.CancelRequest: - HandleCancelRequest((command.parced as CancelRequest)!); - break; - default: - continue; - } + PushMessage(command.parced, command.originalType); } } - private void ThrowIfNotInjected() + public void PushMessage(object parced, EMessageType originalType) { - if (_serializationToolkit == null) - throw new NullReferenceException("Serializing toolkit is null."); - if (_executeModule == null) - throw new NullReferenceException("Execute module is null."); - if (_representationModule == null) - throw new NullReferenceException("Representation module is null."); - if (_methodRepository == null) - throw new NullReferenceException("Method repository is null."); + switch (originalType) + { + case EMessageType.CallRequest: + HandleCallRequest((DefaultCallRequest)parced); + break; + case EMessageType.ClientDisconnect: + return; + case EMessageType.EventRequest: + HandleEventRequest((DefaultCallRequest)parced); + break; + case EMessageType.CancelRequest: + HandleCancelRequest((parced as CancelRequest)!); + break; + default: + return; + } } + public Predicate Rule { get; } = m => + m.MessageType is EMessageType.CallRequest or EMessageType.CancelRequest + or EMessageType.EventRequest or EMessageType.ClientDisconnect; + + public Func[] Converters { get; } = { + m => m.MessageType == EMessageType.CallRequest ? typeof(DefaultCallRequest) : null, + m => m.MessageType == EMessageType.CancelRequest ? typeof(CancelRequest) : null, + m => m.MessageType == EMessageType.EventRequest ? typeof(DefaultCallRequest) : null, + m => m.MessageType == EMessageType.ClientDisconnect ? typeof(ClientDisconnect) : null + }; + private void HandleCancelRequest(CancelRequest req) { _executeModule.Execute(req, _context.RealRepository, _representationModule, _context); diff --git a/mROA/Implementation/RemoteInstanceRepository.cs b/mROA/Implementation/RemoteInstanceRepository.cs index 72e21fe..451058c 100644 --- a/mROA/Implementation/RemoteInstanceRepository.cs +++ b/mROA/Implementation/RemoteInstanceRepository.cs @@ -34,8 +34,6 @@ namespace mROA.Implementation var index = _producedProxys.Find(i => i.Identifier.Equals(id)); if (index is not null) return (T)(index as object); - if (_representationProducer == null) - throw new NullReferenceException("representation producer is not initialized"); if (!_callIndexProvider.Activators.TryGetValue(typeof(T), out var remoteType)) throw new NotSupportedException(); @@ -56,9 +54,6 @@ namespace mROA.Implementation public object GetSingletonObject(Type type, IEndPointContext context) { - if (_representationProducer == null) - throw new NullReferenceException("representation producer is not initialized"); - var representationModule = _representationProducer.Produce(context.OwnerId); diff --git a/mROA/Implementation/RepresentationModule.cs b/mROA/Implementation/RepresentationModule.cs index 2b25dd9..9ecb7a5 100644 --- a/mROA/Implementation/RepresentationModule.cs +++ b/mROA/Implementation/RepresentationModule.cs @@ -82,11 +82,6 @@ namespace mROA.Implementation public async Task PostCallMessageAsync(Guid id, EMessageType eMessageType, T payload, IEndPointContext? context) where T : notnull { - if (_interaction == null) - throw new NullReferenceException("Interaction toolkit is not initialized"); - if (_serialization == null) - throw new NullReferenceException("Serialization toolkit is not initialized"); - var serialized = _serialization.Serialize(payload, context); await _interaction.PostMessageAsync(new NetworkMessageHeader { Id = id, MessageType = eMessageType, Data = serialized });