Untrusted channel works, but it breaks trusted channel
This commit is contained in:
@@ -19,9 +19,10 @@ class Program
|
|||||||
builder.Modules.Add(new BackendIdentityGenerator());
|
builder.Modules.Add(new BackendIdentityGenerator());
|
||||||
// builder.UseNetworkGateway(new IPEndPoint(IPAddress.Loopback, 4567), typeof(NextGenerationInteractionModule),
|
// builder.UseNetworkGateway(new IPEndPoint(IPAddress.Loopback, 4567), typeof(NextGenerationInteractionModule),
|
||||||
// builder.GetModule<IIdentityGenerator>()!);
|
// builder.GetModule<IIdentityGenerator>()!);
|
||||||
builder.UseNetworkGateway(new IPEndPoint(IPAddress.Loopback, 4567), typeof(ChannelInteractionModule),
|
var listening = new IPEndPoint(IPAddress.Loopback, 4567);
|
||||||
|
builder.UseNetworkGateway(listening, typeof(ChannelInteractionModule),
|
||||||
builder.GetModule<IIdentityGenerator>()!);
|
builder.GetModule<IIdentityGenerator>()!);
|
||||||
|
builder.Modules.Add(new UdpGateway(listening));
|
||||||
builder.Modules.Add(new ConnectionHub());
|
builder.Modules.Add(new ConnectionHub());
|
||||||
builder.Modules.Add(new HubRequestExtractor(typeof(RequestExtractor)));
|
builder.Modules.Add(new HubRequestExtractor(typeof(RequestExtractor)));
|
||||||
|
|
||||||
@@ -45,12 +46,12 @@ class Program
|
|||||||
builder.Build();
|
builder.Build();
|
||||||
new RemoteTypeBinder();
|
new RemoteTypeBinder();
|
||||||
|
|
||||||
TransmissionConfig.RealContextRepository = builder.GetModule<MultiClientContextRepository>();
|
TransmissionConfig.RealContextRepository = builder.GetModule<MultiClientContextRepository>()!;
|
||||||
TransmissionConfig.RemoteEndpointContextRepository = builder.GetModule<RemoteContextRepository>();
|
TransmissionConfig.RemoteEndpointContextRepository = builder.GetModule<RemoteContextRepository>()!;
|
||||||
TransmissionConfig.OwnershipRepository = new MultiClientOwnershipRepository();
|
TransmissionConfig.OwnershipRepository = new MultiClientOwnershipRepository();
|
||||||
|
|
||||||
|
_ = builder.GetModule<UdpGateway>()!.Start();
|
||||||
var gateway = builder.GetModule<IGatewayModule>();
|
var gateway = builder.GetModule<IGatewayModule>();
|
||||||
|
|
||||||
gateway.Run();
|
gateway.Run();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -25,8 +25,10 @@ class Program
|
|||||||
|
|
||||||
builder.Modules.Add(new RemoteContextRepository());
|
builder.Modules.Add(new RemoteContextRepository());
|
||||||
builder.Modules.Add(new ChannelInteractionModule());
|
builder.Modules.Add(new ChannelInteractionModule());
|
||||||
|
builder.Modules.Add(new UdpUntrustedInteraction());
|
||||||
builder.Modules.Add(new RepresentationModule());
|
builder.Modules.Add(new RepresentationModule());
|
||||||
builder.Modules.Add(new NetworkFrontendBridge(new IPEndPoint(IPAddress.Loopback, 4567)));
|
var serverEndPoint = new IPEndPoint(IPAddress.Loopback, 4567);
|
||||||
|
builder.Modules.Add(new NetworkFrontendBridge(serverEndPoint));
|
||||||
builder.Modules.Add(new StaticRepresentationModuleProducer());
|
builder.Modules.Add(new StaticRepresentationModuleProducer());
|
||||||
builder.Modules.Add(new RequestExtractor());
|
builder.Modules.Add(new RequestExtractor());
|
||||||
builder.Modules.Add(new BasicExecutionModule());
|
builder.Modules.Add(new BasicExecutionModule());
|
||||||
@@ -43,10 +45,13 @@ class Program
|
|||||||
var frontendBridge = builder.GetModule<IFrontendBridge>()!;
|
var frontendBridge = builder.GetModule<IFrontendBridge>()!;
|
||||||
frontendBridge.Connect();
|
frontendBridge.Connect();
|
||||||
_ = builder.GetModule<RequestExtractor>()!.StartExtraction();
|
_ = builder.GetModule<RequestExtractor>()!.StartExtraction();
|
||||||
|
_ = builder.GetModule<UdpUntrustedInteraction>().Start(serverEndPoint);
|
||||||
Console.WriteLine(TransmissionConfig.OwnershipRepository.GetOwnershipId());
|
Console.WriteLine(TransmissionConfig.OwnershipRepository.GetOwnershipId());
|
||||||
var context = builder.GetModule<RemoteContextRepository>();
|
var context = builder.GetModule<RemoteContextRepository>();
|
||||||
|
|
||||||
var factory = context.GetSingleObject(typeof(IPrinterFactory), -TransmissionConfig.OwnershipRepository.GetHostOwnershipId()) as IPrinterFactory;
|
var factory =
|
||||||
|
context.GetSingleObject(typeof(IPrinterFactory),
|
||||||
|
-TransmissionConfig.OwnershipRepository.GetHostOwnershipId()) as IPrinterFactory;
|
||||||
|
|
||||||
using (var disposingPrinter = factory.Create("Test"))
|
using (var disposingPrinter = factory.Create("Test"))
|
||||||
{
|
{
|
||||||
@@ -66,6 +71,9 @@ class Program
|
|||||||
|
|
||||||
Thread.Sleep(100);
|
Thread.Sleep(100);
|
||||||
|
|
||||||
|
disposingPrinter.SomeoneIsApproaching("Mikhail");
|
||||||
|
Console.WriteLine("Approaching detected");
|
||||||
|
|
||||||
factory.Register(new ClientBasedPrinter());
|
factory.Register(new ClientBasedPrinter());
|
||||||
DemoCheck.ClientBasedImplementation = true;
|
DemoCheck.ClientBasedImplementation = true;
|
||||||
Console.WriteLine("Registered printer");
|
Console.WriteLine("Registered printer");
|
||||||
@@ -102,6 +110,7 @@ class Program
|
|||||||
|
|
||||||
Console.WriteLine("Dispose printer");
|
Console.WriteLine("Dispose printer");
|
||||||
}
|
}
|
||||||
|
|
||||||
DemoCheck.Dispose = true;
|
DemoCheck.Dispose = true;
|
||||||
|
|
||||||
|
|
||||||
@@ -116,12 +125,12 @@ class Program
|
|||||||
cts.Cancel();
|
cts.Cancel();
|
||||||
Console.WriteLine($"Token state {cts.Token.IsCancellationRequested}");
|
Console.WriteLine($"Token state {cts.Token.IsCancellationRequested}");
|
||||||
DemoCheck.TaskCancelation = true;
|
DemoCheck.TaskCancelation = true;
|
||||||
|
|
||||||
frontendBridge.Disconnect();
|
frontendBridge.Disconnect();
|
||||||
|
|
||||||
DemoCheck.Show();
|
DemoCheck.Show();
|
||||||
Console.ReadKey();
|
Console.ReadKey();
|
||||||
|
|
||||||
//
|
//
|
||||||
// const int iterations = 10000;
|
// const int iterations = 10000;
|
||||||
// var timer = Stopwatch.StartNew();
|
// var timer = Stopwatch.StartNew();
|
||||||
|
|||||||
@@ -0,0 +1,11 @@
|
|||||||
|
using System;
|
||||||
|
using System.Net;
|
||||||
|
using System.Threading.Tasks;
|
||||||
|
|
||||||
|
namespace mROA.Abstract
|
||||||
|
{
|
||||||
|
public interface IUntrustedGateway : IInjectableModule, IDisposable
|
||||||
|
{
|
||||||
|
Task Start();
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -18,6 +18,8 @@ namespace mROA.Implementation.Backend
|
|||||||
private ISerializationToolkit? _serialization;
|
private ISerializationToolkit? _serialization;
|
||||||
private Dictionary<int, CancellationTokenSource> _extractorsCTS = new();
|
private Dictionary<int, CancellationTokenSource> _extractorsCTS = new();
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
public NetworkGatewayModule(IPEndPoint endpoint, Type interactionModuleType,
|
public NetworkGatewayModule(IPEndPoint endpoint, Type interactionModuleType,
|
||||||
IInjectableModule[] injectableModules)
|
IInjectableModule[] injectableModules)
|
||||||
{
|
{
|
||||||
@@ -124,7 +126,7 @@ namespace mROA.Implementation.Backend
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private void ThrowIfNotInjected()
|
private void ThrowIfNotInjected()
|
||||||
{
|
{
|
||||||
if (_hub is null)
|
if (_hub is null)
|
||||||
|
|||||||
@@ -0,0 +1,92 @@
|
|||||||
|
using System;
|
||||||
|
using System.Collections.Generic;
|
||||||
|
using System.Net;
|
||||||
|
using System.Net.Sockets;
|
||||||
|
using System.Threading;
|
||||||
|
using System.Threading.Channels;
|
||||||
|
using System.Threading.Tasks;
|
||||||
|
using mROA.Abstract;
|
||||||
|
|
||||||
|
namespace mROA.Implementation.Backend
|
||||||
|
{
|
||||||
|
public class UdpGateway : IUntrustedGateway
|
||||||
|
{
|
||||||
|
private IConnectionHub _hub;
|
||||||
|
private UdpClient _client;
|
||||||
|
private Dictionary<IPEndPoint, int> _reservedPorts = new();
|
||||||
|
private CancellationTokenSource _tokenSource = new();
|
||||||
|
private ISerializationToolkit _serializationToolkit;
|
||||||
|
|
||||||
|
public UdpGateway(IPEndPoint listeningEndpoint)
|
||||||
|
{
|
||||||
|
_client = new UdpClient(listeningEndpoint);
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
public void Inject<T>(T dependency)
|
||||||
|
{
|
||||||
|
switch (dependency)
|
||||||
|
{
|
||||||
|
case IConnectionHub hub:
|
||||||
|
_hub = hub;
|
||||||
|
break;
|
||||||
|
case ISerializationToolkit serializationToolkit:
|
||||||
|
_serializationToolkit = serializationToolkit;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public void Dispose()
|
||||||
|
{
|
||||||
|
_tokenSource.Cancel();
|
||||||
|
_client.Close();
|
||||||
|
}
|
||||||
|
|
||||||
|
public Task Start()
|
||||||
|
{
|
||||||
|
var token = _tokenSource.Token;
|
||||||
|
return Task.Run(async () =>
|
||||||
|
{
|
||||||
|
while (token.IsCancellationRequested == false)
|
||||||
|
{
|
||||||
|
var incoming = await _client.ReceiveAsync();
|
||||||
|
var parsed = _serializationToolkit.Deserialize<NetworkMessageHeader>(incoming.Buffer);
|
||||||
|
try
|
||||||
|
{
|
||||||
|
int channelId;
|
||||||
|
switch (parsed.MessageType)
|
||||||
|
{
|
||||||
|
case EMessageType.UntrustedConnect:
|
||||||
|
channelId = BitConverter.ToInt32(parsed.Data);
|
||||||
|
_reservedPorts[incoming.RemoteEndPoint] = channelId;
|
||||||
|
_ = UntrustedSend(_hub.GetInteraction(channelId), incoming.RemoteEndPoint);
|
||||||
|
break;
|
||||||
|
default:
|
||||||
|
channelId = _reservedPorts[incoming.RemoteEndPoint];
|
||||||
|
var interaction = _hub.GetInteraction(channelId);
|
||||||
|
await interaction.ReceiveChanel.Writer.WriteAsync(parsed, token);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
Console.WriteLine(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}, token);
|
||||||
|
}
|
||||||
|
|
||||||
|
private Task UntrustedSend(IChannelInteractionModule interaction, IPEndPoint endpoint)
|
||||||
|
{
|
||||||
|
return Task.Run(async () =>
|
||||||
|
{
|
||||||
|
await foreach (var post in interaction.UntrustedPostChanel.ReadAllAsync())
|
||||||
|
{
|
||||||
|
var parsed = _serializationToolkit.Serialize(post);
|
||||||
|
await _client.SendAsync(parsed, parsed.Length, endpoint);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -12,5 +12,6 @@ namespace mROA.Implementation
|
|||||||
ClientRecovery,
|
ClientRecovery,
|
||||||
ClientConnect,
|
ClientConnect,
|
||||||
ClientDisconnect,
|
ClientDisconnect,
|
||||||
|
UntrustedConnect,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -12,6 +12,7 @@ namespace mROA.Implementation
|
|||||||
private ISerializationToolkit _serializationToolkit;
|
private ISerializationToolkit _serializationToolkit;
|
||||||
private IChannelInteractionModule _channelInteractionModule;
|
private IChannelInteractionModule _channelInteractionModule;
|
||||||
private CancellationTokenSource _tokenSource = new CancellationTokenSource();
|
private CancellationTokenSource _tokenSource = new CancellationTokenSource();
|
||||||
|
|
||||||
public void Dispose()
|
public void Dispose()
|
||||||
{
|
{
|
||||||
_tokenSource.Cancel();
|
_tokenSource.Cancel();
|
||||||
@@ -35,19 +36,36 @@ namespace mROA.Implementation
|
|||||||
{
|
{
|
||||||
var message = new Memory<byte>((await udpClient.ReceiveAsync()).Buffer);
|
var message = new Memory<byte>((await udpClient.ReceiveAsync()).Buffer);
|
||||||
var parsed = _serializationToolkit.Deserialize<NetworkMessageHeader>(message.Span)!;
|
var parsed = _serializationToolkit.Deserialize<NetworkMessageHeader>(message.Span)!;
|
||||||
|
|
||||||
await writer.WriteAsync(parsed, token);
|
await writer.WriteAsync(parsed, token);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private async Task Posting(UdpClient udpClient, CancellationToken token)
|
private async Task Posting(UdpClient udpClient, CancellationToken token)
|
||||||
{
|
{
|
||||||
|
var initMessage = new NetworkMessageHeader
|
||||||
|
{
|
||||||
|
MessageType = EMessageType.UntrustedConnect, Id = Guid.NewGuid(),
|
||||||
|
Data = BitConverter.GetBytes(Math.Abs(_channelInteractionModule.ConnectionId))
|
||||||
|
};
|
||||||
|
|
||||||
|
var initParsed = _serializationToolkit.Serialize(initMessage);
|
||||||
|
|
||||||
|
await udpClient.SendAsync(initParsed, initParsed.Length);
|
||||||
|
|
||||||
await foreach (var post in _channelInteractionModule.UntrustedPostChanel.ReadAllAsync(token))
|
await foreach (var post in _channelInteractionModule.UntrustedPostChanel.ReadAllAsync(token))
|
||||||
{
|
{
|
||||||
var serialized = _serializationToolkit.Serialize(post);
|
var serialized = _serializationToolkit.Serialize(post);
|
||||||
|
#if TRACE
|
||||||
|
Console.WriteLine("Untrusted write start");
|
||||||
|
#endif
|
||||||
await udpClient.SendAsync(serialized, serialized.Length);
|
await udpClient.SendAsync(serialized, serialized.Length);
|
||||||
|
#if TRACE
|
||||||
|
Console.WriteLine("Untrusted write finished");
|
||||||
|
#endif
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public void Inject<T>(T dependency)
|
public void Inject<T>(T dependency)
|
||||||
{
|
{
|
||||||
switch (dependency)
|
switch (dependency)
|
||||||
|
|||||||
+1
-1
@@ -21,7 +21,7 @@
|
|||||||
</PropertyGroup>
|
</PropertyGroup>
|
||||||
|
|
||||||
<PropertyGroup Condition=" '$(Configuration)' == 'Debug' ">
|
<PropertyGroup Condition=" '$(Configuration)' == 'Debug' ">
|
||||||
<DefineConstants>;</DefineConstants>
|
<DefineConstants></DefineConstants>
|
||||||
</PropertyGroup>
|
</PropertyGroup>
|
||||||
|
|
||||||
<ItemGroup>
|
<ItemGroup>
|
||||||
|
|||||||
Reference in New Issue
Block a user