Base of synced distribution model written

This commit is contained in:
2025-07-19 16:01:41 +03:00
parent 0c648f0af0
commit 79d0e4b5b1
4 changed files with 49 additions and 9 deletions
+1 -1
View File
@@ -24,7 +24,7 @@ class Program
builder.Services.AddOptions();
var listening = new IPEndPoint(IPAddress.Any, 4567);
builder.Services.Configure<GatewayOptions>(options => options.Endpoint = listening);
builder.Services.Configure<DistributionOptions>(o => o.DistributionType = EDistributionType.Channeled);
builder.Services.Configure<DistributionOptions>(o => o.DistributionType = EDistributionType.ExtractorFirst);
builder.Services.AddSingleton<HubRequestExtractor>();
builder.Services.AddSingleton<IExecuteModule, BasicExecutionModule>();
@@ -1,3 +1,4 @@
using System.Collections.Generic;
using Microsoft.Extensions.Options;
using mROA.Abstract;
using mROA.Implementation.Frontend;
@@ -10,6 +11,7 @@ namespace mROA.Implementation.Backend
private readonly IInstanceRepository _remoteContextRepository;
private readonly IExecuteModule _executeModule;
private readonly DistributionOptions _mode;
private Dictionary<int, IRequestExtractor> _producedExtractors = new();
public HubRequestExtractor(IRealStoreInstanceRepository contextRepository,
IInstanceRepository remoteContextRepository, IExecuteModule executeModule,
@@ -21,6 +23,8 @@ namespace mROA.Implementation.Backend
_mode = mode.Value;
}
public IRequestExtractor this[int id] => _producedExtractors[id];
public IRequestExtractor HubOnOnConnected(IRepresentationModule interaction)
{
var extractor = CreateExtractor(interaction);
@@ -29,6 +33,7 @@ namespace mROA.Implementation.Backend
extractor.StartExtraction().ContinueWith(_ => OnDisconnected(interaction));
}
_producedExtractors[interaction.Id] = extractor;
return extractor;
}
@@ -19,7 +19,10 @@ namespace mROA.Implementation.Backend
private readonly Dictionary<int, CancellationTokenSource> _extractorsTokenSources = new();
private readonly ICallIndexProvider _callIndexProvider;
private readonly IIdentityGenerator _identityGenerator;
public NetworkGatewayModule(IOptions<GatewayOptions> options, IIdentityGenerator identityGenerator, IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, IOptions<DistributionOptions> distribution, HubRequestExtractor hre)
public NetworkGatewayModule(IOptions<GatewayOptions> options, IIdentityGenerator identityGenerator,
IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub,
IOptions<DistributionOptions> distribution, HubRequestExtractor hre)
{
_tcpListener = new(options.Value.Endpoint);
_identityGenerator = identityGenerator;
@@ -102,15 +105,42 @@ namespace mROA.Implementation.Backend
_extractorsTokenSources[interaction.ConnectionId] = cts;
_hub.RegisterInteraction(interaction);
_hre.HubOnOnConnected(new RepresentationModule(interaction, _serialization));
var requestExtractor = _hre.HubOnOnConnected(new RepresentationModule(interaction, _serialization));
if (_distribution.DistributionType != EDistributionType.Channeled)
{
BindRequestFirstDistribution(context, interaction, streamExtractor, requestExtractor);
}
}
private void RecoverDisconnectedClient(NetworkMessageHeader connectionRequest, ChannelInteractionModule.StreamExtractor streamExtractor,
private void BindRequestFirstDistribution(IEndPointContext context, IChannelInteractionModule interaction,
ChannelInteractionModule.StreamExtractor streamExtractor, IRequestExtractor requestExtractor)
{
var converters = requestExtractor.Converters;
streamExtractor.MessageReceived = message =>
{
if (requestExtractor.Rule(message))
{
for (int i = 0; i < converters.Length; i++)
{
var func = converters[i];
if (func(message) is { } t)
{
var deserialized = _serialization.Deserialize(message.Data, t, context);
requestExtractor.PushMessage(deserialized, message.MessageType);
break;
}
}
return;
}
interaction.ReceiveChanel.Writer.WriteAsync(message).ConfigureAwait(false);
};
}
private void RecoverDisconnectedClient(NetworkMessageHeader connectionRequest,
ChannelInteractionModule.StreamExtractor streamExtractor,
CancellationTokenSource cts)
{
var recoveryRequest = _serialization.Deserialize<ClientRecovery>(connectionRequest.Data, null);
@@ -125,6 +155,11 @@ namespace mROA.Implementation.Backend
};
_ = streamExtractor.SendFromChannel(recoveryInteraction.TrustedPostChanel, cts.Token);
if (_distribution.DistributionType == EDistributionType.ExtractorFirst)
{
BindRequestFirstDistribution(recoveryInteraction.Context, recoveryInteraction, streamExtractor, _hre[recoveryInteraction.ConnectionId]);
}
Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token).ConfigureAwait(false));
recoveryInteraction.Restart(false);