LumaTunnel
【WinUI】LumaTunnel 是一个面向个人多设备的 Windows 代理系统。客户端在本机提供 HTTP/HTTPS CONNECT 与 SOCKS5 TCP 代理,并通过一个受信任 TLS 证书保护的 WSS 会话,将多个 TCP 流复用到自建 Windows Server 节点。
关注
0
Fork
0
Star
0
using System.Net.Sockets;
using System.Diagnostics;
using System.Threading.Channels;
using LumaTunnel.Shared.Protocol;
namespace LumaTunnel.Server.Core.Runtime;
internal sealed class ServerTunnelStream : IAsyncDisposable
{
private readonly uint _streamId;
private readonly TcpClient _client;
private readonly NetworkStream _networkStream;
private readonly TunnelConnection _connection;
private readonly ServerRuntime _runtime;
private readonly CancellationTokenSource _lifetime = new();
private readonly Channel<ReadOnlyMemory<byte>> _clientData;
private Task? _readTask;
private Task? _writeTask;
private Task? _idleTask;
private long _lastActivity = Stopwatch.GetTimestamp();
private int _disposed;
private int _clientWriteClosed;
private int _targetReadClosed;
public ServerTunnelStream(uint streamId, TcpClient client, TunnelConnection connection, ServerRuntime runtime)
{
_streamId = streamId;
_client = client;
_networkStream = client.GetStream();
_connection = connection;
_runtime = runtime;
_clientData = Channel.CreateBounded<ReadOnlyMemory<byte>>(new BoundedChannelOptions(runtime.Settings.MaxBufferedFramesPerStream)
{
FullMode = BoundedChannelFullMode.Wait,
SingleReader = true,
SingleWriter = true
});
}
public void Start()
{
_readTask = ReadTargetAsync(_lifetime.Token);
_writeTask = WriteTargetAsync(_lifetime.Token);
_idleTask = IdleLoopAsync(_lifetime.Token);
}
public async ValueTask WriteFromClientAsync(ReadOnlyMemory<byte> data, CancellationToken cancellationToken)
{
Interlocked.Exchange(ref _lastActivity, Stopwatch.GetTimestamp());
await _clientData.Writer.WriteAsync(data, cancellationToken).ConfigureAwait(false);
}
public async Task ClientFinishedWritingAsync()
{
if (Interlocked.Exchange(ref _clientWriteClosed, 1) == 0)
{
_clientData.Writer.TryComplete();
try
{
_client.Client.Shutdown(SocketShutdown.Send);
}
catch (SocketException)
{
}
}
if (Volatile.Read(ref _targetReadClosed) != 0)
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
}
public ValueTask DisposeAsync()
{
if (Interlocked.Exchange(ref _disposed, 1) != 0)
return ValueTask.CompletedTask;
_lifetime.Cancel();
_clientData.Writer.TryComplete();
_networkStream.Dispose();
_client.Dispose();
_lifetime.Dispose();
return ValueTask.CompletedTask;
}
private async Task ReadTargetAsync(CancellationToken cancellationToken)
{
var buffer = new byte[TunnelProtocol.MaxPayloadLength];
try
{
while (true)
{
var read = await _networkStream.ReadAsync(buffer, cancellationToken).ConfigureAwait(false);
if (read == 0)
break;
Interlocked.Exchange(ref _lastActivity, Stopwatch.GetTimestamp());
_runtime.Traffic.AddDownload(_connection.DeviceId, read);
await _connection.SendAsync(new TunnelFrame(TunnelFrameType.Data, _streamId, buffer.AsMemory(0, read).ToArray()), cancellationToken).ConfigureAwait(false);
}
Interlocked.Exchange(ref _targetReadClosed, 1);
await _connection.SendAsync(TunnelFrame.Empty(TunnelFrameType.HalfClose, _streamId), cancellationToken).ConfigureAwait(false);
if (Volatile.Read(ref _clientWriteClosed) != 0)
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
}
catch (Exception)
{
await _connection.SendAsync(TunnelFrame.Empty(TunnelFrameType.Reset, _streamId), CancellationToken.None).ConfigureAwait(false);
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
}
}
private async Task WriteTargetAsync(CancellationToken cancellationToken)
{
try
{
await foreach (var data in _clientData.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
{
await _networkStream.WriteAsync(data, cancellationToken).ConfigureAwait(false);
_runtime.Traffic.AddUpload(_connection.DeviceId, data.Length);
}
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
}
catch (Exception)
{
await _connection.SendAsync(TunnelFrame.Empty(TunnelFrameType.Reset, _streamId), CancellationToken.None).ConfigureAwait(false);
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
}
}
private async Task IdleLoopAsync(CancellationToken cancellationToken)
{
var idleTimeout = TimeSpan.FromMinutes(_runtime.Settings.IdleTimeoutMinutes);
var pollInterval = TimeSpan.FromSeconds(Math.Min(60, Math.Max(5, idleTimeout.TotalSeconds / 4)));
try
{
while (!cancellationToken.IsCancellationRequested)
{
await Task.Delay(pollInterval, cancellationToken).ConfigureAwait(false);
if (Stopwatch.GetElapsedTime(Interlocked.Read(ref _lastActivity)) < idleTimeout)
continue;
await _connection.SendAsync(TunnelFrame.Empty(TunnelFrameType.Reset, _streamId), cancellationToken).ConfigureAwait(false);
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
return;
}
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
}
}
}
using System.Net.Sockets;
using System.Diagnostics;
using System.Threading.Channels;
using LumaTunnel.Shared.Protocol;
namespace LumaTunnel.Server.Core.Runtime;
internal sealed class ServerTunnelStream : IAsyncDisposable
{
private readonly uint _streamId;
private readonly TcpClient _client;
private readonly NetworkStream _networkStream;
private readonly TunnelConnection _connection;
private readonly ServerRuntime _runtime;
private readonly CancellationTokenSource _lifetime = new();
private readonly Channel<ReadOnlyMemory<byte>> _clientData;
private Task? _readTask;
private Task? _writeTask;
private Task? _idleTask;
private long _lastActivity = Stopwatch.GetTimestamp();
private int _disposed;
private int _clientWriteClosed;
private int _targetReadClosed;
public ServerTunnelStream(uint streamId, TcpClient client, TunnelConnection connection, ServerRuntime runtime)
{
_streamId = streamId;
_client = client;
_networkStream = client.GetStream();
_connection = connection;
_runtime = runtime;
_clientData = Channel.CreateBounded<ReadOnlyMemory<byte>>(new BoundedChannelOptions(runtime.Settings.MaxBufferedFramesPerStream)
{
FullMode = BoundedChannelFullMode.Wait,
SingleReader = true,
SingleWriter = true
});
}
public void Start()
{
_readTask = ReadTargetAsync(_lifetime.Token);
_writeTask = WriteTargetAsync(_lifetime.Token);
_idleTask = IdleLoopAsync(_lifetime.Token);
}
public async ValueTask WriteFromClientAsync(ReadOnlyMemory<byte> data, CancellationToken cancellationToken)
{
Interlocked.Exchange(ref _lastActivity, Stopwatch.GetTimestamp());
await _clientData.Writer.WriteAsync(data, cancellationToken).ConfigureAwait(false);
}
public async Task ClientFinishedWritingAsync()
{
if (Interlocked.Exchange(ref _clientWriteClosed, 1) == 0)
{
_clientData.Writer.TryComplete();
try
{
_client.Client.Shutdown(SocketShutdown.Send);
}
catch (SocketException)
{
}
}
if (Volatile.Read(ref _targetReadClosed) != 0)
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
}
public ValueTask DisposeAsync()
{
if (Interlocked.Exchange(ref _disposed, 1) != 0)
return ValueTask.CompletedTask;
_lifetime.Cancel();
_clientData.Writer.TryComplete();
_networkStream.Dispose();
_client.Dispose();
_lifetime.Dispose();
return ValueTask.CompletedTask;
}
private async Task ReadTargetAsync(CancellationToken cancellationToken)
{
var buffer = new byte[TunnelProtocol.MaxPayloadLength];
try
{
while (true)
{
var read = await _networkStream.ReadAsync(buffer, cancellationToken).ConfigureAwait(false);
if (read == 0)
break;
Interlocked.Exchange(ref _lastActivity, Stopwatch.GetTimestamp());
_runtime.Traffic.AddDownload(_connection.DeviceId, read);
await _connection.SendAsync(new TunnelFrame(TunnelFrameType.Data, _streamId, buffer.AsMemory(0, read).ToArray()), cancellationToken).ConfigureAwait(false);
}
Interlocked.Exchange(ref _targetReadClosed, 1);
await _connection.SendAsync(TunnelFrame.Empty(TunnelFrameType.HalfClose, _streamId), cancellationToken).ConfigureAwait(false);
if (Volatile.Read(ref _clientWriteClosed) != 0)
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
}
catch (Exception)
{
await _connection.SendAsync(TunnelFrame.Empty(TunnelFrameType.Reset, _streamId), CancellationToken.None).ConfigureAwait(false);
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
}
}
private async Task WriteTargetAsync(CancellationToken cancellationToken)
{
try
{
await foreach (var data in _clientData.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
{
await _networkStream.WriteAsync(data, cancellationToken).ConfigureAwait(false);
_runtime.Traffic.AddUpload(_connection.DeviceId, data.Length);
}
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
}
catch (Exception)
{
await _connection.SendAsync(TunnelFrame.Empty(TunnelFrameType.Reset, _streamId), CancellationToken.None).ConfigureAwait(false);
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
}
}
private async Task IdleLoopAsync(CancellationToken cancellationToken)
{
var idleTimeout = TimeSpan.FromMinutes(_runtime.Settings.IdleTimeoutMinutes);
var pollInterval = TimeSpan.FromSeconds(Math.Min(60, Math.Max(5, idleTimeout.TotalSeconds / 4)));
try
{
while (!cancellationToken.IsCancellationRequested)
{
await Task.Delay(pollInterval, cancellationToken).ConfigureAwait(false);
if (Stopwatch.GetElapsedTime(Interlocked.Read(ref _lastActivity)) < idleTimeout)
continue;
await _connection.SendAsync(TunnelFrame.Empty(TunnelFrameType.Reset, _streamId), cancellationToken).ConfigureAwait(false);
await _connection.RemoveStreamAsync(_streamId).ConfigureAwait(false);
return;
}
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
}
}
}