From 85cf1bc48ab57a74fe1966e212afc656dd269d56 Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Fri, 18 Jul 2025 13:49:36 +0300 Subject: [PATCH 1/6] Refactoring and configure await false --- Example.Backend/Program.cs | 22 +++++------- Example.Frontend/Program.cs | 3 ++ mROA/Abstract/IPrimaryMessageDistributior.cs | 15 -------- mROA/Implementation/Backend/ConnectionHub.cs | 8 ----- .../Backend/HubRequestExtractor.cs | 31 ++++++++-------- .../Backend/InstanceRepository.cs | 2 +- .../Backend/NetworkGatewayModule.cs | 26 +++++++++----- mROA/Implementation/Backend/UdpGateway.cs | 10 +++--- .../ChannelDistributorFactory.cs | 35 ------------------- .../ChannelInteractionModule.cs | 6 ---- .../CollectableMethodRepository.cs | 2 +- .../CreativeRepresentationModuleProducer.cs | 11 +++--- mROA/Implementation/DistributionOptions.cs | 13 +++++++ .../Frontend/NetworkFrontendBridge.cs | 10 +++--- .../Frontend/RequestExtractor.cs | 8 ++--- .../RemoteInstanceRepository.cs | 14 ++++---- mROA/Implementation/RepresentationModule.cs | 4 +-- 17 files changed, 87 insertions(+), 133 deletions(-) delete mode 100644 mROA/Abstract/IPrimaryMessageDistributior.cs delete mode 100644 mROA/Implementation/ChannelDistributorFactory.cs create mode 100644 mROA/Implementation/DistributionOptions.cs diff --git a/Example.Backend/Program.cs b/Example.Backend/Program.cs index 5711803..9763db9 100644 --- a/Example.Backend/Program.cs +++ b/Example.Backend/Program.cs @@ -1,5 +1,4 @@ using System; -using System.Linq; using System.Net; using Example.Backend; using Example.Shared; @@ -21,13 +20,13 @@ class Program builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); + builder.Services.AddOptions(); var listening = new IPEndPoint(IPAddress.Any, 4567); - builder.Services.Configure(options => options.Endpoint = listening); - + builder.Services.Configure(o => o.DistributionType = EDistributionType.Channeled); + builder.Services.AddSingleton(); - builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); @@ -40,7 +39,6 @@ class Program return repo; })); - builder.Services.AddSingleton(); builder.Services.AddSingleton(p => { var methodRepo = new CollectableMethodRepository(); @@ -52,17 +50,13 @@ class Program builder.Services.AddSingleton(); var host = builder.Build(); - host.Services.GetService(); // -// builder.Build(); - new RemoteTypeBinder(); + new RemoteTypeBinder(); // -// - _ = host.Services.GetService()!.Start(); - var gateway = host.Services.GetService(); - gateway.Run(); - - Console.ReadLine(); + _ = host.Services.GetService()!.Start(); + var gateway = host.Services.GetService(); + gateway.Run(); + Console.ReadLine(); } } \ No newline at end of file diff --git a/Example.Frontend/Program.cs b/Example.Frontend/Program.cs index f374205..b21b094 100644 --- a/Example.Frontend/Program.cs +++ b/Example.Frontend/Program.cs @@ -39,8 +39,11 @@ class Program builder.Services.AddSingleton(); var serverEndPoint = new IPEndPoint(IPAddress.Loopback, 4567); builder.Services.AddSingleton(); + builder.Services.AddOptions(); builder.Services.Configure(options => options.Endpoint = serverEndPoint); + builder.Services.Configure(o => o.DistributionType = EDistributionType.Channeled); + builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); diff --git a/mROA/Abstract/IPrimaryMessageDistributior.cs b/mROA/Abstract/IPrimaryMessageDistributior.cs deleted file mode 100644 index aa0295c..0000000 --- a/mROA/Abstract/IPrimaryMessageDistributior.cs +++ /dev/null @@ -1,15 +0,0 @@ -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/Implementation/Backend/ConnectionHub.cs b/mROA/Implementation/Backend/ConnectionHub.cs index 54dc096..bfc72e9 100644 --- a/mROA/Implementation/Backend/ConnectionHub.cs +++ b/mROA/Implementation/Backend/ConnectionHub.cs @@ -7,18 +7,10 @@ namespace mROA.Implementation.Backend public class ConnectionHub : IConnectionHub { private readonly Dictionary _connections = new(); - private readonly IContextualSerializationToolKit _serializationToolkit; - - public ConnectionHub(IContextualSerializationToolKit serializationToolkit) - { - _serializationToolkit = serializationToolkit; - } public void RegisterInteraction(IChannelInteractionModule interaction) { _connections.Add(interaction.ConnectionId, interaction); - var module = new RepresentationModule(interaction, _serializationToolkit); - OnConnected?.Invoke(module); } public IChannelInteractionModule GetInteraction(int id) diff --git a/mROA/Implementation/Backend/HubRequestExtractor.cs b/mROA/Implementation/Backend/HubRequestExtractor.cs index 7d5a2cb..362db29 100644 --- a/mROA/Implementation/Backend/HubRequestExtractor.cs +++ b/mROA/Implementation/Backend/HubRequestExtractor.cs @@ -1,3 +1,4 @@ +using Microsoft.Extensions.Options; using mROA.Abstract; using mROA.Implementation.Frontend; @@ -5,28 +6,30 @@ namespace mROA.Implementation.Backend { public class HubRequestExtractor { - private IRealStoreInstanceRepository _contextRepository; - private IInstanceRepository _remoteContextRepository; - private IMethodRepository _methodRepository; - private IContextualSerializationToolKit _serializationToolkit; - private IExecuteModule _executeModule; + private readonly IRealStoreInstanceRepository _contextRepository; + private readonly IInstanceRepository _remoteContextRepository; + private readonly IExecuteModule _executeModule; + private readonly DistributionOptions _mode; - public HubRequestExtractor(IConnectionHub hub, IRealStoreInstanceRepository contextRepository, - IInstanceRepository remoteContextRepository, IMethodRepository methodRepository, - IContextualSerializationToolKit serializationToolkit, IExecuteModule executeModule) + public HubRequestExtractor(IRealStoreInstanceRepository contextRepository, + IInstanceRepository remoteContextRepository, IExecuteModule executeModule, + IOptions mode) { - hub.OnConnected += HubOnOnConnected; _contextRepository = contextRepository; _remoteContextRepository = remoteContextRepository; - _methodRepository = methodRepository; - _serializationToolkit = serializationToolkit; _executeModule = executeModule; + _mode = mode.Value; } - private void HubOnOnConnected(IRepresentationModule interaction) + public IRequestExtractor HubOnOnConnected(IRepresentationModule interaction) { var extractor = CreateExtractor(interaction); - extractor.StartExtraction().ContinueWith(_ => OnDisconnected(interaction)); + if (_mode.DistributionType == EDistributionType.Channeled) + { + extractor.StartExtraction().ContinueWith(_ => OnDisconnected(interaction)); + } + + return extractor; } private void OnDisconnected(IRepresentationModule representationModule) @@ -45,7 +48,7 @@ namespace mROA.Implementation.Backend context.RealRepository = _contextRepository; context.RemoteRepository = _remoteContextRepository; - + return extractor; } } diff --git a/mROA/Implementation/Backend/InstanceRepository.cs b/mROA/Implementation/Backend/InstanceRepository.cs index f61f7f5..cfd42ec 100644 --- a/mROA/Implementation/Backend/InstanceRepository.cs +++ b/mROA/Implementation/Backend/InstanceRepository.cs @@ -11,7 +11,7 @@ namespace mROA.Implementation.Backend { public static object[] EventBinders = { }; - private IRepresentationModuleProducer _representationModuleProducer; + private readonly IRepresentationModuleProducer _representationModuleProducer; private Dictionary _singletons = new(); private readonly IStorage _storage; diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index a8f485e..8c5b30d 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -13,19 +13,21 @@ namespace mROA.Implementation.Backend { private readonly TcpListener _tcpListener; private readonly IConnectionHub _hub; + private readonly HubRequestExtractor _hre; + private readonly DistributionOptions _distribution; private readonly IContextualSerializationToolKit _serialization; private readonly Dictionary _extractorsCTS = new(); - private ICallIndexProvider _callIndexProvider; + private readonly ICallIndexProvider _callIndexProvider; private readonly IIdentityGenerator _identityGenerator; - private readonly IMessageDistributorFactory _distributorFactory; - public NetworkGatewayModule(IOptions options, IIdentityGenerator identityGenerator, IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, IMessageDistributorFactory distributorFactory) + public NetworkGatewayModule(IOptions options, IIdentityGenerator identityGenerator, IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, IOptions distribution, HubRequestExtractor hre) { _tcpListener = new(options.Value.Endpoint); _identityGenerator = identityGenerator; _serialization = serialization; _callIndexProvider = callIndexProvider; _hub = hub; - _distributorFactory = distributorFactory; + _hre = hre; + _distribution = distribution.Value; } public void Run() @@ -57,7 +59,7 @@ namespace mROA.Implementation.Backend interaction.IsConnected = () => streamExtractor.IsConnected; streamExtractor.MessageReceived = async message => { - await interaction.ReceiveChanel.Writer.WriteAsync(message); + await interaction.ReceiveChanel.Writer.WriteAsync(message).ConfigureAwait(false); }; _ = Task.Run(() => streamExtractor.SingleReceive()); var connectionRequest = await interaction.ReceiveChanel.Reader.ReadAsync(); @@ -71,15 +73,21 @@ namespace mROA.Implementation.Backend interaction.Context = context; Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token)); _ = streamExtractor.SendFromChannel(interaction.TrustedPostChanel, cts.Token); - interaction.PostMessageAsync(new NetworkMessageHeader(_serialization!, + interaction.PostMessageAsync(new NetworkMessageHeader(_serialization, new IdAssignment { Id = interaction.ConnectionId }, null)); _extractorsCTS[interaction.ConnectionId] = cts; + _hub.RegisterInteraction(interaction); + _hre.HubOnOnConnected(new RepresentationModule(interaction, _serialization)); + if (_distribution.DistributionType != EDistributionType.Channeled) + { + + } Console.WriteLine("Client registered"); break; case EMessageType.ClientRecovery: { - var recoveryRequest = _serialization!.Deserialize(connectionRequest.Data, null); + var recoveryRequest = _serialization.Deserialize(connectionRequest.Data, null); var recoveryInteraction = _hub.GetInteraction(recoveryRequest.Id); _extractorsCTS[-recoveryRequest.Id].Cancel(); @@ -87,11 +95,11 @@ namespace mROA.Implementation.Backend recoveryInteraction.IsConnected = () => streamExtractor.IsConnected; streamExtractor.MessageReceived = message => { - recoveryInteraction.ReceiveChanel.Writer.WriteAsync(message); + recoveryInteraction.ReceiveChanel.Writer.WriteAsync(message).ConfigureAwait(false); }; _ = streamExtractor.SendFromChannel(recoveryInteraction.TrustedPostChanel, cts.Token); - Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token)); + Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token).ConfigureAwait(false)); recoveryInteraction.Restart(false); diff --git a/mROA/Implementation/Backend/UdpGateway.cs b/mROA/Implementation/Backend/UdpGateway.cs index c1d30a7..20b56b2 100644 --- a/mROA/Implementation/Backend/UdpGateway.cs +++ b/mROA/Implementation/Backend/UdpGateway.cs @@ -12,11 +12,11 @@ namespace mROA.Implementation.Backend { public class UdpGateway : IUntrustedGateway { - private IConnectionHub _hub; - private UdpClient _client; - private Dictionary _reservedPorts = new(); - private CancellationTokenSource _tokenSource = new(); - private IContextualSerializationToolKit _serializationToolkit; + private readonly IConnectionHub _hub; + private readonly UdpClient _client; + private readonly Dictionary _reservedPorts = new(); + private readonly CancellationTokenSource _tokenSource = new(); + private readonly IContextualSerializationToolKit _serializationToolkit; public UdpGateway(IOptions options, IConnectionHub hub, IContextualSerializationToolKit serializationToolkit) { diff --git a/mROA/Implementation/ChannelDistributorFactory.cs b/mROA/Implementation/ChannelDistributorFactory.cs deleted file mode 100644 index 5998694..0000000 --- a/mROA/Implementation/ChannelDistributorFactory.cs +++ /dev/null @@ -1,35 +0,0 @@ -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 0c6b827..53bffc8 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -62,11 +62,7 @@ namespace mROA.Implementation public ValueTask GetNextMessageReceiving(bool infinite = true) { return _receiveReader.ReadAsync(); - // if (_currentReceiving != null) return _currentReceiving; - // _currentReceiving = Task.Run(async () => await GetNextMessage()); - // return _currentReceiving; } -#pragma warning disable CS8602 // Dereference of a possibly null reference. private async ValueTask PostMessageInternal(NetworkMessageHeader messageHeader) { if (!IsConnected()) @@ -77,8 +73,6 @@ namespace mROA.Implementation await _trustedWriter.WriteAsync(messageHeader); return true; } -#pragma warning restore CS8602 // Dereference of a possibly null reference. - public async Task PostMessageAsync(NetworkMessageHeader messageHeader) { diff --git a/mROA/Implementation/CollectableMethodRepository.cs b/mROA/Implementation/CollectableMethodRepository.cs index 8fb9d7e..135380f 100644 --- a/mROA/Implementation/CollectableMethodRepository.cs +++ b/mROA/Implementation/CollectableMethodRepository.cs @@ -5,7 +5,7 @@ namespace mROA.Implementation { public class CollectableMethodRepository : IMethodRepository { - private List _methods = new(); + private readonly List _methods = new(); public void AppendInvokers(IEnumerable methodInvokers) { diff --git a/mROA/Implementation/CreativeRepresentationModuleProducer.cs b/mROA/Implementation/CreativeRepresentationModuleProducer.cs index 66356a1..c5eaf6d 100644 --- a/mROA/Implementation/CreativeRepresentationModuleProducer.cs +++ b/mROA/Implementation/CreativeRepresentationModuleProducer.cs @@ -1,16 +1,13 @@ -using System; -using mROA.Abstract; +using mROA.Abstract; namespace mROA.Implementation { public class CreativeRepresentationModuleProducer : IRepresentationModuleProducer { - private IServiceProvider _creationModules; - private IConnectionHub _hub; - private IContextualSerializationToolKit _serialization; - public CreativeRepresentationModuleProducer(IServiceProvider creationModules, IConnectionHub hub, IContextualSerializationToolKit serialization) + private readonly IConnectionHub _hub; + private readonly IContextualSerializationToolKit _serialization; + public CreativeRepresentationModuleProducer(IConnectionHub hub, IContextualSerializationToolKit serialization) { - _creationModules = creationModules; _hub = hub; _serialization = serialization; } diff --git a/mROA/Implementation/DistributionOptions.cs b/mROA/Implementation/DistributionOptions.cs new file mode 100644 index 0000000..239dc23 --- /dev/null +++ b/mROA/Implementation/DistributionOptions.cs @@ -0,0 +1,13 @@ +namespace mROA.Implementation +{ + public class DistributionOptions + { + public EDistributionType DistributionType { get; set; } + } + + public enum EDistributionType + { + Channeled, + ExtractorFirst + } +} \ No newline at end of file diff --git a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs index e3f3d8e..7e25985 100644 --- a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs +++ b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs @@ -14,11 +14,11 @@ namespace mROA.Implementation.Frontend { private readonly IPEndPoint _serverEndPoint; private TcpClient _tcpClient = new(); - private IChannelInteractionModule _interactionModule; - private IContextualSerializationToolKit _serialization; + private readonly IChannelInteractionModule _interactionModule; + private readonly IContextualSerializationToolKit _serialization; private ChannelInteractionModule.StreamExtractor? _currentExtractor; private CancellationTokenSource _rawExtractorCancellation; - private IEndPointContext _context; + private readonly IEndPointContext _context; public NetworkFrontendBridge(IOptions options, IEndPointContext context, IContextualSerializationToolKit serialization, IChannelInteractionModule interactionModule) { @@ -63,7 +63,7 @@ namespace mROA.Implementation.Frontend _currentExtractor = new ChannelInteractionModule.StreamExtractor(_tcpClient.GetStream(), _serialization, _context); - _ = _currentExtractor.SendFromChannel(_interactionModule!.TrustedPostChanel, + _ = _currentExtractor.SendFromChannel(_interactionModule.TrustedPostChanel, _rawExtractorCancellation.Token); _currentExtractor.MessageReceived = message => { @@ -93,7 +93,7 @@ namespace mROA.Implementation.Frontend public void Disconnect() { - _ = _interactionModule!.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientDisconnect(), + _ = _interactionModule.PostMessageAsync(new NetworkMessageHeader(_serialization, new ClientDisconnect(), _context)); _interactionModule.Dispose(); _tcpClient.Dispose(); diff --git a/mROA/Implementation/Frontend/RequestExtractor.cs b/mROA/Implementation/Frontend/RequestExtractor.cs index 838f4e2..0ebb92a 100644 --- a/mROA/Implementation/Frontend/RequestExtractor.cs +++ b/mROA/Implementation/Frontend/RequestExtractor.cs @@ -10,10 +10,10 @@ namespace mROA.Implementation.Frontend public class RequestExtractor : IRequestExtractor { - private IExecuteModule _executeModule; + private readonly IExecuteModule _executeModule; - private IRepresentationModule _representationModule; - private IEndPointContext _context; + private readonly IRepresentationModule _representationModule; + private readonly IEndPointContext _context; public RequestExtractor(IExecuteModule executeModule, IRepresentationModule representationModule, IEndPointContext context) { @@ -84,7 +84,7 @@ namespace mROA.Implementation.Frontend return; } - _representationModule.PostCallMessage(request.Id, resultType, result, _context); + _representationModule.PostCallMessageAsync(request.Id, resultType, result, _context).ConfigureAwait(false); } private void HandleEventRequest(DefaultCallRequest request) diff --git a/mROA/Implementation/RemoteInstanceRepository.cs b/mROA/Implementation/RemoteInstanceRepository.cs index 451058c..76ad9ab 100644 --- a/mROA/Implementation/RemoteInstanceRepository.cs +++ b/mROA/Implementation/RemoteInstanceRepository.cs @@ -7,10 +7,10 @@ namespace mROA.Implementation { public class RemoteInstanceRepository : IInstanceRepository { - private List _producedProxys = new(); - private ICallIndexProvider _callIndexProvider; + private readonly List _producedProxies = new(); + private readonly ICallIndexProvider _callIndexProvider; - private IRepresentationModuleProducer _representationProducer; + private readonly IRepresentationModuleProducer _representationProducer; public RemoteInstanceRepository(ICallIndexProvider callIndexProvider, IRepresentationModuleProducer representationProducer) { @@ -31,7 +31,7 @@ namespace mROA.Implementation public T GetObject(ComplexObjectIdentifier id, IEndPointContext context) where T : class { - var index = _producedProxys.Find(i => i.Identifier.Equals(id)); + var index = _producedProxies.Find(i => i.Identifier.Equals(id)); if (index is not null) return (T)(index as object); @@ -42,7 +42,7 @@ namespace mROA.Implementation var remote = remoteType(id.ContextId, representationModule, context, _callIndexProvider.GetIndices(typeof(T))); - _producedProxys.Add(remote!); + _producedProxies.Add(remote!); return (remote as T)!; } @@ -62,9 +62,9 @@ namespace mROA.Implementation var remoteObjectBase = instance; - _producedProxys.Add(remoteObjectBase); + _producedProxies.Add(remoteObjectBase); - return _producedProxys.Last(); + return _producedProxies.Last(); } public int GetObjectIndex(object o, IEndPointContext context) diff --git a/mROA/Implementation/RepresentationModule.cs b/mROA/Implementation/RepresentationModule.cs index 9ecb7a5..735ee27 100644 --- a/mROA/Implementation/RepresentationModule.cs +++ b/mROA/Implementation/RepresentationModule.cs @@ -12,8 +12,8 @@ namespace mROA.Implementation { public class RepresentationModule : IRepresentationModule { - private IChannelInteractionModule _interaction; - private IContextualSerializationToolKit _serialization; + private readonly IChannelInteractionModule _interaction; + private readonly IContextualSerializationToolKit _serialization; public RepresentationModule(IChannelInteractionModule interaction, IContextualSerializationToolKit serialization) { From 0c648f0af06e3af4aa924105d8b3d00f363e9dcd Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Fri, 18 Jul 2025 14:08:29 +0300 Subject: [PATCH 2/6] Refactor Network gateway module --- .../Backend/NetworkGatewayModule.cs | 140 ++++++++++-------- 1 file changed, 79 insertions(+), 61 deletions(-) diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index 8c5b30d..a51f4a6 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -16,7 +16,7 @@ namespace mROA.Implementation.Backend private readonly HubRequestExtractor _hre; private readonly DistributionOptions _distribution; private readonly IContextualSerializationToolKit _serialization; - private readonly Dictionary _extractorsCTS = new(); + 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) @@ -49,68 +49,86 @@ namespace mROA.Implementation.Backend while (true) { var client = await _tcpListener.AcceptTcpClientAsync(); - Console.WriteLine($"Client connected from {client.Client.RemoteEndPoint}"); - var interaction = new ChannelInteractionModule(_serialization, _identityGenerator); - - var context = new EndPointContext(null, null); - context.CallIndexProvider = _callIndexProvider; - var streamExtractor = - new ChannelInteractionModule.StreamExtractor(client.GetStream(), _serialization, context); - interaction.IsConnected = () => streamExtractor.IsConnected; - streamExtractor.MessageReceived = async message => - { - await interaction.ReceiveChanel.Writer.WriteAsync(message).ConfigureAwait(false); - }; - _ = Task.Run(() => streamExtractor.SingleReceive()); - var connectionRequest = await interaction.ReceiveChanel.Reader.ReadAsync(); - var cts = new CancellationTokenSource(); - - switch (connectionRequest.MessageType) - { - case EMessageType.ClientConnect: - context.HostId = 0; - context.OwnerId = -interaction.ConnectionId; - interaction.Context = context; - Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token)); - _ = streamExtractor.SendFromChannel(interaction.TrustedPostChanel, cts.Token); - interaction.PostMessageAsync(new NetworkMessageHeader(_serialization, - new IdAssignment { Id = interaction.ConnectionId }, null)); - _extractorsCTS[interaction.ConnectionId] = cts; - - _hub.RegisterInteraction(interaction); - _hre.HubOnOnConnected(new RepresentationModule(interaction, _serialization)); - if (_distribution.DistributionType != EDistributionType.Channeled) - { - - } - Console.WriteLine("Client registered"); - break; - case EMessageType.ClientRecovery: - { - var recoveryRequest = _serialization.Deserialize(connectionRequest.Data, null); - var recoveryInteraction = _hub.GetInteraction(recoveryRequest.Id); - - _extractorsCTS[-recoveryRequest.Id].Cancel(); - - recoveryInteraction.IsConnected = () => streamExtractor.IsConnected; - streamExtractor.MessageReceived = message => - { - recoveryInteraction.ReceiveChanel.Writer.WriteAsync(message).ConfigureAwait(false); - }; - _ = streamExtractor.SendFromChannel(recoveryInteraction.TrustedPostChanel, cts.Token); - - Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token).ConfigureAwait(false)); - - - recoveryInteraction.Restart(false); - break; - } - default: - client.Close(); - break; - } + _ = HandleConnection(client).ConfigureAwait(false); } } + + private async Task HandleConnection(TcpClient client) + { + Console.WriteLine($"Client connected from {client.Client.RemoteEndPoint}"); + var interaction = new ChannelInteractionModule(_serialization, _identityGenerator); + + var context = new EndPointContext(null, null) + { + CallIndexProvider = _callIndexProvider + }; + var streamExtractor = + new ChannelInteractionModule.StreamExtractor(client.GetStream(), _serialization, context); + interaction.IsConnected = () => streamExtractor.IsConnected; + streamExtractor.MessageReceived = async message => + { + await interaction.ReceiveChanel.Writer.WriteAsync(message).ConfigureAwait(false); + }; + _ = Task.Run(() => streamExtractor.SingleReceive()); + var connectionRequest = await interaction.ReceiveChanel.Reader.ReadAsync(); + var cts = new CancellationTokenSource(); + + switch (connectionRequest.MessageType) + { + case EMessageType.ClientConnect: + HandleNewClient(context, interaction, streamExtractor, cts); + break; + case EMessageType.ClientRecovery: + { + RecoverDisconnectedClient(connectionRequest, streamExtractor, cts); + break; + } + default: + client.Close(); + break; + } + } + + private void HandleNewClient(EndPointContext context, ChannelInteractionModule interaction, + ChannelInteractionModule.StreamExtractor streamExtractor, CancellationTokenSource cts) + { + context.HostId = 0; + context.OwnerId = -interaction.ConnectionId; + interaction.Context = context; + Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token)); + _ = streamExtractor.SendFromChannel(interaction.TrustedPostChanel, cts.Token); + interaction.PostMessageAsync(new NetworkMessageHeader(_serialization, + new IdAssignment { Id = interaction.ConnectionId }, null)); + _extractorsTokenSources[interaction.ConnectionId] = cts; + + _hub.RegisterInteraction(interaction); + _hre.HubOnOnConnected(new RepresentationModule(interaction, _serialization)); + + if (_distribution.DistributionType != EDistributionType.Channeled) + { + + } + } + + private void RecoverDisconnectedClient(NetworkMessageHeader connectionRequest, ChannelInteractionModule.StreamExtractor streamExtractor, + CancellationTokenSource cts) + { + var recoveryRequest = _serialization.Deserialize(connectionRequest.Data, null); + var recoveryInteraction = _hub.GetInteraction(recoveryRequest.Id); + + _extractorsTokenSources[-recoveryRequest.Id].Cancel(); + + recoveryInteraction.IsConnected = () => streamExtractor.IsConnected; + streamExtractor.MessageReceived = message => + { + recoveryInteraction.ReceiveChanel.Writer.WriteAsync(message).ConfigureAwait(false); + }; + _ = streamExtractor.SendFromChannel(recoveryInteraction.TrustedPostChanel, cts.Token); + + Task.Run(async () => await streamExtractor.LoopedReceive(cts.Token).ConfigureAwait(false)); + + recoveryInteraction.Restart(false); + } } public class GatewayOptions From 79d0e4b5b1a2be30b895ac73635a600df019c18f Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Sat, 19 Jul 2025 16:01:41 +0300 Subject: [PATCH 3/6] Base of synced distribution model written --- Example.Backend/Program.cs | 2 +- .../Backend/HubRequestExtractor.cs | 5 ++ .../Backend/NetworkGatewayModule.cs | 49 ++++++++++++++++--- .../ChannelInteractionModule.cs | 2 +- 4 files changed, 49 insertions(+), 9 deletions(-) diff --git a/Example.Backend/Program.cs b/Example.Backend/Program.cs index 9763db9..fdf34c8 100644 --- a/Example.Backend/Program.cs +++ b/Example.Backend/Program.cs @@ -24,7 +24,7 @@ class Program builder.Services.AddOptions(); var listening = new IPEndPoint(IPAddress.Any, 4567); builder.Services.Configure(options => options.Endpoint = listening); - builder.Services.Configure(o => o.DistributionType = EDistributionType.Channeled); + builder.Services.Configure(o => o.DistributionType = EDistributionType.ExtractorFirst); builder.Services.AddSingleton(); builder.Services.AddSingleton(); diff --git a/mROA/Implementation/Backend/HubRequestExtractor.cs b/mROA/Implementation/Backend/HubRequestExtractor.cs index 362db29..4a3c48a 100644 --- a/mROA/Implementation/Backend/HubRequestExtractor.cs +++ b/mROA/Implementation/Backend/HubRequestExtractor.cs @@ -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 _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; } diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index a51f4a6..ecdc97b 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -19,7 +19,10 @@ namespace mROA.Implementation.Backend 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) + + public NetworkGatewayModule(IOptions options, IIdentityGenerator identityGenerator, + IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, + IOptions distribution, HubRequestExtractor hre) { _tcpListener = new(options.Value.Endpoint); _identityGenerator = identityGenerator; @@ -57,7 +60,7 @@ namespace mROA.Implementation.Backend { Console.WriteLine($"Client connected from {client.Client.RemoteEndPoint}"); var interaction = new ChannelInteractionModule(_serialization, _identityGenerator); - + var context = new EndPointContext(null, null) { CallIndexProvider = _callIndexProvider @@ -100,17 +103,44 @@ namespace mROA.Implementation.Backend interaction.PostMessageAsync(new NetworkMessageHeader(_serialization, new IdAssignment { Id = interaction.ConnectionId }, null)); _extractorsTokenSources[interaction.ConnectionId] = cts; - - _hub.RegisterInteraction(interaction); - _hre.HubOnOnConnected(new RepresentationModule(interaction, _serialization)); + _hub.RegisterInteraction(interaction); + 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(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); diff --git a/mROA/Implementation/ChannelInteractionModule.cs b/mROA/Implementation/ChannelInteractionModule.cs index 53bffc8..d7b85bf 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -165,7 +165,7 @@ namespace mROA.Implementation await _ioStream.ReadAsync(_lenBuffer); var len = BitConverter.ToUInt16(_lenBuffer); - + return len; } From c7c5b1637c4b0d97e449a9cc96659c783a6fd3df Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Sat, 19 Jul 2025 17:18:46 +0300 Subject: [PATCH 4/6] Direct to exe mode works faster, but has restrictions --- Example.Backend/Program.cs | 2 ++ Example.Frontend/Program.cs | 3 +++ Example.Load/Program.cs | 2 +- mROA/Implementation/Backend/BasicExecutionModule.cs | 6 +++++- mROA/Implementation/Backend/NetworkGatewayModule.cs | 7 +++++-- mROA/Implementation/ChannelInteractionModule.cs | 8 +++++++- mROA/Implementation/Frontend/NetworkFrontendBridge.cs | 7 +++++-- mROA/mROA.csproj | 6 ++++++ 8 files changed, 34 insertions(+), 7 deletions(-) diff --git a/Example.Backend/Program.cs b/Example.Backend/Program.cs index fdf34c8..e80c1d7 100644 --- a/Example.Backend/Program.cs +++ b/Example.Backend/Program.cs @@ -9,12 +9,14 @@ using mROA.Implementation; using mROA.Implementation.Backend; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; class Program { public static void Main(string[] args) { var builder = Host.CreateApplicationBuilder(); + builder.Services.AddLogging(l => l.AddConsole()); builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); diff --git a/Example.Frontend/Program.cs b/Example.Frontend/Program.cs index b21b094..de23bb0 100644 --- a/Example.Frontend/Program.cs +++ b/Example.Frontend/Program.cs @@ -8,6 +8,7 @@ using Example.Frontend; using Example.Shared; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; using mROA.Abstract; using mROA.Cbor; using mROA.Codegen; @@ -23,6 +24,8 @@ class Program new RemoteTypeBinder(); var builder = Host.CreateApplicationBuilder(new HostApplicationBuilderSettings { DisableDefaults = true }); + builder.Services.AddLogging(l => l.SetMinimumLevel(LogLevel.Trace).AddConsole()); + builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(provider => diff --git a/Example.Load/Program.cs b/Example.Load/Program.cs index 36be5f4..5ea37cd 100644 --- a/Example.Load/Program.cs +++ b/Example.Load/Program.cs @@ -32,7 +32,7 @@ Console.WriteLine("End waiting"); var totalRequests = tasks.Sum(i => i.Result); Console.WriteLine($"Total requests: {totalRequests:N0}"); Console.WriteLine($"Results: {totalRequests / time.TotalSeconds:N} RPS"); - +File.AppendAllText("results.txt", $"[DIRECT TO EXE] {totalRequests}\r\n"); async Task> GetLoadEndpoints(int count) { diff --git a/mROA/Implementation/Backend/BasicExecutionModule.cs b/mROA/Implementation/Backend/BasicExecutionModule.cs index b32cdfd..d8ce917 100644 --- a/mROA/Implementation/Backend/BasicExecutionModule.cs +++ b/mROA/Implementation/Backend/BasicExecutionModule.cs @@ -1,5 +1,6 @@ using System; using System.Threading; +using Microsoft.Extensions.Logging; using mROA.Abstract; using mROA.Implementation.CommandExecution; @@ -10,17 +11,20 @@ namespace mROA.Implementation.Backend private readonly ICancellationRepository _cancellationRepo; private readonly IMethodRepository _methodRepo; private readonly IContextualSerializationToolKit _serialization; + private readonly ILogger _logger; - public BasicExecutionModule(ICancellationRepository cancellationRepo, IMethodRepository methodRepo, IContextualSerializationToolKit serialization) + public BasicExecutionModule(ICancellationRepository cancellationRepo, IMethodRepository methodRepo, IContextualSerializationToolKit serialization, ILogger logger) { _cancellationRepo = cancellationRepo; _methodRepo = methodRepo; _serialization = serialization; + _logger = logger; } public ICommandExecution Execute(ICallRequest command, IInstanceRepository instanceRepository, IRepresentationModule representationModule, IEndPointContext endPointContext) { + // _logger.LogInformation("Executing {0}", command.Id); try { if (command is CancelRequest) diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index ecdc97b..ce32fbf 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -4,6 +4,7 @@ using System.Net; using System.Net.Sockets; using System.Threading; using System.Threading.Tasks; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using mROA.Abstract; @@ -14,6 +15,7 @@ namespace mROA.Implementation.Backend private readonly TcpListener _tcpListener; private readonly IConnectionHub _hub; private readonly HubRequestExtractor _hre; + private readonly ILogger _logger; private readonly DistributionOptions _distribution; private readonly IContextualSerializationToolKit _serialization; private readonly Dictionary _extractorsTokenSources = new(); @@ -22,7 +24,7 @@ namespace mROA.Implementation.Backend public NetworkGatewayModule(IOptions options, IIdentityGenerator identityGenerator, IContextualSerializationToolKit serialization, ICallIndexProvider callIndexProvider, IConnectionHub hub, - IOptions distribution, HubRequestExtractor hre) + IOptions distribution, HubRequestExtractor hre, ILogger logger) { _tcpListener = new(options.Value.Endpoint); _identityGenerator = identityGenerator; @@ -30,6 +32,7 @@ namespace mROA.Implementation.Backend _callIndexProvider = callIndexProvider; _hub = hub; _hre = hre; + _logger = logger; _distribution = distribution.Value; } @@ -66,7 +69,7 @@ namespace mROA.Implementation.Backend CallIndexProvider = _callIndexProvider }; var streamExtractor = - new ChannelInteractionModule.StreamExtractor(client.GetStream(), _serialization, context); + new ChannelInteractionModule.StreamExtractor(client.GetStream(), _serialization, context, _logger); interaction.IsConnected = () => streamExtractor.IsConnected; streamExtractor.MessageReceived = async message => { diff --git a/mROA/Implementation/ChannelInteractionModule.cs b/mROA/Implementation/ChannelInteractionModule.cs index d7b85bf..6fd3b38 100644 --- a/mROA/Implementation/ChannelInteractionModule.cs +++ b/mROA/Implementation/ChannelInteractionModule.cs @@ -4,6 +4,7 @@ using System.IO; using System.Threading; using System.Threading.Channels; using System.Threading.Tasks; +using Microsoft.Extensions.Logging; using mROA.Abstract; namespace mROA.Implementation @@ -147,14 +148,16 @@ namespace mROA.Implementation private readonly IContextualSerializationToolKit _serializationToolkit; private readonly Memory _buffer = new byte[BufferSize]; private readonly IEndPointContext _context; + private readonly ILogger _logger; private readonly byte[] _lenBuffer; public StreamExtractor(Stream ioStream, IContextualSerializationToolKit serializationToolkit, - IEndPointContext context) + IEndPointContext context, ILogger logger) { _ioStream = ioStream; _serializationToolkit = serializationToolkit; _context = context; + _logger = logger; _lenBuffer = new byte[2]; } @@ -176,6 +179,7 @@ namespace mROA.Implementation await _ioStream.ReadExactlyAsync(localSpan, cancellationToken: token); var message = _serializationToolkit.Deserialize(localSpan, _context); + // _logger.LogTrace("RECV {0}", message.ToString()); MessageReceived(message); } @@ -195,6 +199,8 @@ namespace mROA.Implementation header.CopyTo(_buffer); var sendingSpan = _buffer[..(len + 2)]; await _ioStream.WriteAsync(sendingSpan, token); + // _logger.LogTrace("SEND {0}", message.ToString()); + } public async Task SendFromChannel(ChannelReader channel, diff --git a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs index 7e25985..40159ee 100644 --- a/mROA/Implementation/Frontend/NetworkFrontendBridge.cs +++ b/mROA/Implementation/Frontend/NetworkFrontendBridge.cs @@ -3,6 +3,7 @@ using System.Net; using System.Net.Sockets; using System.Threading; using System.Threading.Tasks; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using mROA.Abstract; using mROA.Implementation.Backend; @@ -15,17 +16,19 @@ namespace mROA.Implementation.Frontend private readonly IPEndPoint _serverEndPoint; private TcpClient _tcpClient = new(); private readonly IChannelInteractionModule _interactionModule; + private readonly ILogger _logger; private readonly IContextualSerializationToolKit _serialization; private ChannelInteractionModule.StreamExtractor? _currentExtractor; private CancellationTokenSource _rawExtractorCancellation; private readonly IEndPointContext _context; - public NetworkFrontendBridge(IOptions options, IEndPointContext context, IContextualSerializationToolKit serialization, IChannelInteractionModule interactionModule) + public NetworkFrontendBridge(IOptions options, IEndPointContext context, IContextualSerializationToolKit serialization, IChannelInteractionModule interactionModule, ILogger logger) { _serverEndPoint = options.Value.Endpoint; _context = context; _serialization = serialization; _interactionModule = interactionModule; + _logger = logger; _rawExtractorCancellation = new CancellationTokenSource(); } @@ -61,7 +64,7 @@ namespace mROA.Implementation.Frontend private void PrepareExtractor() { _currentExtractor = - new ChannelInteractionModule.StreamExtractor(_tcpClient.GetStream(), _serialization, _context); + new ChannelInteractionModule.StreamExtractor(_tcpClient.GetStream(), _serialization, _context, _logger); _ = _currentExtractor.SendFromChannel(_interactionModule.TrustedPostChanel, _rawExtractorCancellation.Token); diff --git a/mROA/mROA.csproj b/mROA/mROA.csproj index cc70622..f4a4354 100644 --- a/mROA/mROA.csproj +++ b/mROA/mROA.csproj @@ -39,4 +39,10 @@ + + + ..\..\..\..\.nuget\packages\microsoft.extensions.logging.abstractions\9.0.7\lib\netstandard2.0\Microsoft.Extensions.Logging.Abstractions.dll + + + From e4d6709a7761ecdfc81e1a6865ff24db85fd8771 Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Sat, 19 Jul 2025 21:08:25 +0300 Subject: [PATCH 5/6] Backward call with ExtractorFirst option works --- mROA/Implementation/Backend/NetworkGatewayModule.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/mROA/Implementation/Backend/NetworkGatewayModule.cs b/mROA/Implementation/Backend/NetworkGatewayModule.cs index ce32fbf..701cc79 100644 --- a/mROA/Implementation/Backend/NetworkGatewayModule.cs +++ b/mROA/Implementation/Backend/NetworkGatewayModule.cs @@ -130,7 +130,7 @@ namespace mROA.Implementation.Backend if (func(message) is { } t) { var deserialized = _serialization.Deserialize(message.Data, t, context); - requestExtractor.PushMessage(deserialized, message.MessageType); + Task.Run(() => requestExtractor.PushMessage(deserialized, message.MessageType)); break; } } From 83391bd55ee5edb43651e45fbe7322e58fc1bc8e Mon Sep 17 00:00:00 2001 From: Mitrofanov Mikhail Date: Sat, 19 Jul 2025 21:15:03 +0300 Subject: [PATCH 6/6] Load test header update --- Example.Load/Program.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Example.Load/Program.cs b/Example.Load/Program.cs index 5ea37cd..7188833 100644 --- a/Example.Load/Program.cs +++ b/Example.Load/Program.cs @@ -32,7 +32,7 @@ Console.WriteLine("End waiting"); var totalRequests = tasks.Sum(i => i.Result); Console.WriteLine($"Total requests: {totalRequests:N0}"); Console.WriteLine($"Results: {totalRequests / time.TotalSeconds:N} RPS"); -File.AppendAllText("results.txt", $"[DIRECT TO EXE] {totalRequests}\r\n"); +File.AppendAllText("results.txt", $"[DIRECT TO EXE RUN] {totalRequests}\r\n"); async Task> GetLoadEndpoints(int count) {