From b409e5ef08ec5172f53bd64151fdc0bcab49290c Mon Sep 17 00:00:00 2001 From: Mikhail Mitrofanov Date: Sat, 9 Aug 2025 13:03:20 +0300 Subject: [PATCH] Abstract distribution model --- Example.Backend/Program.cs | 1 + mROA.Benchmark/InvokePerformance.cs | 30 ++++++++ mROA.Benchmark/Program.cs | 2 +- mROA/Abstract/IDistributionModule.cs | 10 +++ .../Backend/NetworkGatewayModule.cs | 23 +++--- mROA/Implementation/DistributionModules.cs | 70 +++++++++++++++++++ 6 files changed, 120 insertions(+), 16 deletions(-) create mode 100644 mROA.Benchmark/InvokePerformance.cs create mode 100644 mROA/Abstract/IDistributionModule.cs create mode 100644 mROA/Implementation/DistributionModules.cs diff --git a/Example.Backend/Program.cs b/Example.Backend/Program.cs index e80c1d7..5ef1be4 100644 --- a/Example.Backend/Program.cs +++ b/Example.Backend/Program.cs @@ -27,6 +27,7 @@ class Program var listening = new IPEndPoint(IPAddress.Any, 4567); builder.Services.Configure(options => options.Endpoint = listening); builder.Services.Configure(o => o.DistributionType = EDistributionType.ExtractorFirst); + builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); diff --git a/mROA.Benchmark/InvokePerformance.cs b/mROA.Benchmark/InvokePerformance.cs new file mode 100644 index 0000000..0c07908 --- /dev/null +++ b/mROA.Benchmark/InvokePerformance.cs @@ -0,0 +1,30 @@ +using BenchmarkDotNet.Attributes; + +namespace mROA.Benchmark; + +[MemoryDiagnoser] +[DisassemblyDiagnoser] +public class InvokePerformance +{ + public Func A = x => + { + var result = x * 5 + x * x / 5; + return result; + }; + + public SomeMath B = x => x * 5 + x * x / 5; + + [Benchmark(Baseline = true)] + public int ActionInvoke() + { + return A(132); + } + + [Benchmark] + public int DelegateInvoke() + { + return A(132); + } +} + +public delegate int SomeMath(int x); \ No newline at end of file diff --git a/mROA.Benchmark/Program.cs b/mROA.Benchmark/Program.cs index 5150ab8..fdde9af 100644 --- a/mROA.Benchmark/Program.cs +++ b/mROA.Benchmark/Program.cs @@ -5,4 +5,4 @@ using mROA.Benchmark; Console.WriteLine("Hello, World!"); -BenchmarkRunner.Run(); \ No newline at end of file +BenchmarkRunner.Run(); \ No newline at end of file diff --git a/mROA/Abstract/IDistributionModule.cs b/mROA/Abstract/IDistributionModule.cs new file mode 100644 index 0000000..e06beb4 --- /dev/null +++ b/mROA/Abstract/IDistributionModule.cs @@ -0,0 +1,10 @@ +using System; +using mROA.Implementation; + +namespace mROA.Abstract +{ + public interface IDistributionModule + { + Action GetDistributionAction(int clientId); + } +} \ No newline at end of file diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index 2905d04..ecbf97b 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -14,15 +14,14 @@ namespace mROA.Implementation.Backend private readonly TcpListener _tcpListener; private readonly IConnectionHub _hub; private readonly HubRequestExtractor _hre; - private readonly DistributionOptions _distribution; + private readonly IDistributionModule _distribution; private readonly IContextualSerializationToolKit _serialization; private readonly Dictionary _extractorsTokenSources = new(); private readonly ICallIndexProvider _callIndexProvider; private readonly IIdentityGenerator _identityGenerator; public NetworkGatewayModule(IOptions options, IIdentityGenerator identityGenerator, - IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, - IOptions distribution, HubRequestExtractor hre) + IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, HubRequestExtractor hre, IDistributionModule distribution) { _tcpListener = new(options.Value.Endpoint); _identityGenerator = identityGenerator; @@ -30,7 +29,7 @@ namespace mROA.Implementation.Backend _callIndexProvider = callIndexProvider; _hub = hub; _hre = hre; - _distribution = distribution.Value; + _distribution = distribution; } public void Run() @@ -107,12 +106,10 @@ namespace mROA.Implementation.Backend _extractorsTokenSources[interaction.ConnectionId] = cts; _hub.RegisterInteraction(interaction); - var requestExtractor = _hre.HubOnOnConnected(new RepresentationModule(interaction, _serialization.Clone())); + _hre.HubOnOnConnected(new RepresentationModule(interaction, _serialization.Clone())); + + streamExtractor.MessageReceived = _distribution.GetDistributionAction(interaction.ConnectionId); - if (_distribution.DistributionType != EDistributionType.Channeled) - { - BindRequestFirstDistribution(context, interaction, streamExtractor, requestExtractor); - } } private void BindRequestFirstDistribution(IEndPointContext context, IChannelInteractionModule interaction, @@ -156,12 +153,8 @@ namespace mROA.Implementation.Backend }; _ = streamExtractor.SendFromChannel(recoveryInteraction.TrustedPostChanel, cts.Token); - if (_distribution.DistributionType == EDistributionType.ExtractorFirst) - { - BindRequestFirstDistribution(recoveryInteraction.Context, recoveryInteraction, streamExtractor, - _hre[recoveryInteraction.ConnectionId]); - } - + streamExtractor.MessageReceived = _distribution.GetDistributionAction(recoveryInteraction.ConnectionId); + Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token).ConfigureAwait(false)); recoveryInteraction.Restart(false); diff --git a/mROA/Implementation/DistributionModules.cs b/mROA/Implementation/DistributionModules.cs new file mode 100644 index 0000000..66d5b1e --- /dev/null +++ b/mROA/Implementation/DistributionModules.cs @@ -0,0 +1,70 @@ +using System; +using System.Threading.Tasks; +using mROA.Abstract; +using mROA.Implementation.Backend; + +namespace mROA.Implementation +{ + public class ChannelDistributionModule : IDistributionModule + { + private readonly IConnectionHub _hub; + + public ChannelDistributionModule(IConnectionHub hub) + { + _hub = hub; + } + + public Action GetDistributionAction(int clientId) + { + var writer = _hub.GetInteraction(clientId).ReceiveChanel.Writer; + return message => + { + writer.TryWrite(message); + }; + } + + } + + public class ExtractorFirstDistributionModule : IDistributionModule + { + private readonly IConnectionHub _hub; + private readonly IContextualSerializationToolKit _serialization; + private readonly HubRequestExtractor _extractorHub; + + public ExtractorFirstDistributionModule(HubRequestExtractor extractorHub, IConnectionHub hub, IContextualSerializationToolKit serialization) + { + _extractorHub = extractorHub; + _hub = hub; + _serialization = serialization; + } + + public Action GetDistributionAction(int clientId) + { + var interaction = _hub.GetInteraction(clientId); + var writer = interaction.ReceiveChanel.Writer; + var context = interaction.Context; + var requestExtractor = _extractorHub[clientId]; + var converters = requestExtractor.Converters; + var serialization = _serialization.Clone(); + return message => + { + if (requestExtractor.Rule(message)) + { + for (var i = 0; i < converters.Length; i++) + { + var func = converters[i]; + if (func(message) is not { } t) continue; + + var deserialized = serialization.Deserialize(message.Data, t, context)!; + Task.Run(() => requestExtractor.PushMessage(deserialized, message.MessageType)); + break; + } + + return; + } + + writer.WriteAsync(message).ConfigureAwait(false); + }; + } + } +} \ No newline at end of file