Abstract distribution model

This commit is contained in:
2025-08-09 13:03:20 +03:00
parent c29db73a0e
commit b409e5ef08
6 changed files with 120 additions and 16 deletions
+1
View File
@@ -27,6 +27,7 @@ class Program
var listening = new IPEndPoint(IPAddress.Any, 4567); var listening = new IPEndPoint(IPAddress.Any, 4567);
builder.Services.Configure<GatewayOptions>(options => options.Endpoint = listening); builder.Services.Configure<GatewayOptions>(options => options.Endpoint = listening);
builder.Services.Configure<DistributionOptions>(o => o.DistributionType = EDistributionType.ExtractorFirst); builder.Services.Configure<DistributionOptions>(o => o.DistributionType = EDistributionType.ExtractorFirst);
builder.Services.AddSingleton<IDistributionModule, ExtractorFirstDistributionModule>();
builder.Services.AddSingleton<HubRequestExtractor>(); builder.Services.AddSingleton<HubRequestExtractor>();
builder.Services.AddSingleton<IExecuteModule, BasicExecutionModule>(); builder.Services.AddSingleton<IExecuteModule, BasicExecutionModule>();
+30
View File
@@ -0,0 +1,30 @@
using BenchmarkDotNet.Attributes;
namespace mROA.Benchmark;
[MemoryDiagnoser]
[DisassemblyDiagnoser]
public class InvokePerformance
{
public Func<int, int> 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);
+1 -1
View File
@@ -5,4 +5,4 @@ using mROA.Benchmark;
Console.WriteLine("Hello, World!"); Console.WriteLine("Hello, World!");
BenchmarkRunner.Run<TaskWaiting>(); BenchmarkRunner.Run<InvokePerformance>();
+10
View File
@@ -0,0 +1,10 @@
using System;
using mROA.Implementation;
namespace mROA.Abstract
{
public interface IDistributionModule
{
Action<NetworkMessage> GetDistributionAction(int clientId);
}
}
@@ -14,15 +14,14 @@ namespace mROA.Implementation.Backend
private readonly TcpListener _tcpListener; private readonly TcpListener _tcpListener;
private readonly IConnectionHub _hub; private readonly IConnectionHub _hub;
private readonly HubRequestExtractor _hre; private readonly HubRequestExtractor _hre;
private readonly DistributionOptions _distribution; private readonly IDistributionModule _distribution;
private readonly IContextualSerializationToolKit _serialization; private readonly IContextualSerializationToolKit _serialization;
private readonly Dictionary<int, CancellationTokenSource> _extractorsTokenSources = new(); private readonly Dictionary<int, CancellationTokenSource> _extractorsTokenSources = new();
private readonly ICallIndexProvider _callIndexProvider; private readonly ICallIndexProvider _callIndexProvider;
private readonly IIdentityGenerator _identityGenerator; private readonly IIdentityGenerator _identityGenerator;
public NetworkGatewayModule(IOptions<GatewayOptions> options, IIdentityGenerator identityGenerator, public NetworkGatewayModule(IOptions<GatewayOptions> options, IIdentityGenerator identityGenerator,
IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, HubRequestExtractor hre, IDistributionModule distribution)
IOptions<DistributionOptions> distribution, HubRequestExtractor hre)
{ {
_tcpListener = new(options.Value.Endpoint); _tcpListener = new(options.Value.Endpoint);
_identityGenerator = identityGenerator; _identityGenerator = identityGenerator;
@@ -30,7 +29,7 @@ namespace mROA.Implementation.Backend
_callIndexProvider = callIndexProvider; _callIndexProvider = callIndexProvider;
_hub = hub; _hub = hub;
_hre = hre; _hre = hre;
_distribution = distribution.Value; _distribution = distribution;
} }
public void Run() public void Run()
@@ -107,12 +106,10 @@ namespace mROA.Implementation.Backend
_extractorsTokenSources[interaction.ConnectionId] = cts; _extractorsTokenSources[interaction.ConnectionId] = cts;
_hub.RegisterInteraction(interaction); _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, private void BindRequestFirstDistribution(IEndPointContext context, IChannelInteractionModule interaction,
@@ -156,12 +153,8 @@ namespace mROA.Implementation.Backend
}; };
_ = streamExtractor.SendFromChannel(recoveryInteraction.TrustedPostChanel, cts.Token); _ = streamExtractor.SendFromChannel(recoveryInteraction.TrustedPostChanel, cts.Token);
if (_distribution.DistributionType == EDistributionType.ExtractorFirst) streamExtractor.MessageReceived = _distribution.GetDistributionAction(recoveryInteraction.ConnectionId);
{
BindRequestFirstDistribution(recoveryInteraction.Context, recoveryInteraction, streamExtractor,
_hre[recoveryInteraction.ConnectionId]);
}
Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token).ConfigureAwait(false)); Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token).ConfigureAwait(false));
recoveryInteraction.Restart(false); recoveryInteraction.Restart(false);
@@ -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<NetworkMessage> 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<NetworkMessage> 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);
};
}
}
}