客户端代码优化

This commit is contained in:
aoshiguchen
2023-09-19 12:09:32 +08:00
parent b026b46647
commit 67c362b423
5 changed files with 146 additions and 135 deletions
+1
View File
@@ -57,3 +57,4 @@ hs_err_pid*
neutrino-proxy-vuepress/deploy.sh
.NEUTRINO_PROXY_CLIENT_ID
logs
**/ixxxk.com/**
@@ -1,8 +1,16 @@
package org.dromara.neutrinoproxy.client.config;
import org.dromara.neutrinoproxy.core.ProxyDataTypeEnum;
import org.dromara.neutrinoproxy.core.ProxyMessage;
import org.dromara.neutrinoproxy.core.ProxyMessageHandler;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.handler.logging.LoggingHandler;
import io.netty.handler.timeout.IdleStateHandler;
import org.dromara.neutrinoproxy.client.core.CmdChannelHandler;
import org.dromara.neutrinoproxy.client.core.ProxyChannelHandler;
import org.dromara.neutrinoproxy.client.core.RealServerChannelHandler;
import org.dromara.neutrinoproxy.client.util.ProxyUtil;
import org.dromara.neutrinoproxy.core.*;
import org.dromara.neutrinoproxy.core.dispatcher.DefaultDispatcher;
import org.dromara.neutrinoproxy.core.dispatcher.Dispatcher;
import io.netty.bootstrap.Bootstrap;
@@ -10,8 +18,10 @@ import io.netty.channel.ChannelHandlerContext;
import org.noear.solon.Solon;
import org.noear.solon.annotation.Bean;
import org.noear.solon.annotation.Configuration;
import org.noear.solon.annotation.Inject;
import org.noear.solon.core.bean.LifecycleBean;
import java.net.InetSocketAddress;
import java.util.List;
/**
@@ -31,19 +41,107 @@ public class ProxyConfiguration implements LifecycleBean {
Solon.context().wrapAndPut(Dispatcher.class, dispatcher);
}
@Bean("cmdTunnelBootstrap")
public Bootstrap cmdTunnelBootstrap() {
return new Bootstrap();
@Bean("tunnelWorkGroup")
public NioEventLoopGroup tunnelWorkGroup(@Inject ProxyConfig proxyConfig) {
return new NioEventLoopGroup(proxyConfig.getTunnel().getThreadCount());
}
@Bean("proxyTunnelBootstrap")
public Bootstrap proxyTunnelBootstrap() {
return new Bootstrap();
@Bean("tcpRealServerWorkGroup")
public NioEventLoopGroup tcpRealServerWorkGroup(@Inject ProxyConfig proxyConfig) {
// 暂时先公用此配置
return new NioEventLoopGroup(proxyConfig.getTunnel().getThreadCount());
}
@Bean("cmdTunnelBootstrap")
public Bootstrap cmdTunnelBootstrap(@Inject ProxyConfig proxyConfig,
@Inject("tunnelWorkGroup") NioEventLoopGroup tunnelWorkGroup) {
Bootstrap bootstrap = new Bootstrap();
bootstrap.group(tunnelWorkGroup);
bootstrap.channel(NioSocketChannel.class);
// bootstrap.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 3000);
// bootstrap.option(ChannelOption.SO_KEEPALIVE, true);
// /**
// * TCP/IP协议中,无论发送多少数据,总是要在数据前面加上协议头,同时,对方接收到数据,也需要发送ACK表示确认。为了尽可能的利用网络带宽,TCP总是希望尽可能的发送足够大的数据。(一个连接会设置MSS参数,因此,TCP/IP希望每次都能够以MSS尺寸的数据块来发送数据)。
// * Nagle算法就是为了尽可能发送大块数据,避免网络中充斥着许多小数据块。
// */
// bootstrap.option(ChannelOption.TCP_NODELAY, true);
bootstrap.remoteAddress(InetSocketAddress.createUnresolved(proxyConfig.getTunnel().getServerIp(), proxyConfig.getTunnel().getServerPort()));
bootstrap.handler(new ChannelInitializer<SocketChannel>() {
@Override
public void initChannel(SocketChannel ch) throws Exception {
if (proxyConfig.getTunnel().getSslEnable()) {
ch.pipeline().addLast(ProxyUtil.createSslHandler(proxyConfig));
}
if (null != proxyConfig.getTunnel().getTransferLogEnable() && proxyConfig.getTunnel().getTransferLogEnable()) {
ch.pipeline().addFirst(new LoggingHandler(CmdChannelHandler.class));
}
ch.pipeline().addLast(new ProxyMessageDecoder(proxyConfig.getProtocol().getMaxFrameLength(),
proxyConfig.getProtocol().getLengthFieldOffset(), proxyConfig.getProtocol().getLengthFieldLength(),
proxyConfig.getProtocol().getLengthAdjustment(), proxyConfig.getProtocol().getInitialBytesToStrip()));
ch.pipeline().addLast(new ProxyMessageEncoder());
ch.pipeline().addLast(new IdleStateHandler(proxyConfig.getProtocol().getReadIdleTime(), proxyConfig.getProtocol().getWriteIdleTime(), proxyConfig.getProtocol().getAllIdleTimeSeconds()));
ch.pipeline().addLast(new CmdChannelHandler());
}
});
return bootstrap;
}
@Bean("tcpProxyTunnelBootstrap")
public Bootstrap tcpProxyTunnelBootstrap(@Inject ProxyConfig proxyConfig,
@Inject("tunnelWorkGroup") NioEventLoopGroup tunnelWorkGroup) {
Bootstrap bootstrap = new Bootstrap();
bootstrap.group(tunnelWorkGroup);
bootstrap.channel(NioSocketChannel.class);
bootstrap.remoteAddress(InetSocketAddress.createUnresolved(proxyConfig.getTunnel().getServerIp(), proxyConfig.getTunnel().getServerPort()));
bootstrap.handler(new ChannelInitializer<SocketChannel>() {
@Override
public void initChannel(SocketChannel ch) throws Exception {
if (proxyConfig.getTunnel().getSslEnable()) {
ch.pipeline().addLast(ProxyUtil.createSslHandler(proxyConfig));
}
if (null != proxyConfig.getTunnel().getTransferLogEnable() && proxyConfig.getTunnel().getTransferLogEnable()) {
ch.pipeline().addFirst(new LoggingHandler(ProxyChannelHandler.class));
}
ch.pipeline().addLast(new ProxyMessageDecoder(proxyConfig.getProtocol().getMaxFrameLength(),
proxyConfig.getProtocol().getLengthFieldOffset(), proxyConfig.getProtocol().getLengthFieldLength(),
proxyConfig.getProtocol().getLengthAdjustment(), proxyConfig.getProtocol().getInitialBytesToStrip()));
ch.pipeline().addLast(new ProxyMessageEncoder());
ch.pipeline().addLast(new IdleStateHandler(proxyConfig.getProtocol().getReadIdleTime(), proxyConfig.getProtocol().getWriteIdleTime(), proxyConfig.getProtocol().getAllIdleTimeSeconds()));
ch.pipeline().addLast(new ProxyChannelHandler());
}
});
return bootstrap;
}
@Bean("udpProxyTunnelBootstrap")
private Bootstrap udpProxyTunnelBootstrap() {
Bootstrap bootstrap = new Bootstrap();
return bootstrap;
}
@Bean("realServerBootstrap")
public Bootstrap realServerBootstrap() {
return new Bootstrap();
public Bootstrap realServerBootstrap(@Inject ProxyConfig proxyConfig,
@Inject("tcpRealServerWorkGroup") NioEventLoopGroup tcpRealServerWorkGroup
) {
Bootstrap bootstrap = new Bootstrap();
bootstrap.group(tcpRealServerWorkGroup);
bootstrap.channel(NioSocketChannel.class);
bootstrap.handler(new ChannelInitializer<SocketChannel>() {
@Override
public void initChannel(SocketChannel ch) throws Exception {
if (null != proxyConfig.getTunnel().getTransferLogEnable() && proxyConfig.getTunnel().getTransferLogEnable()) {
ch.pipeline().addFirst(new LoggingHandler(RealServerChannelHandler.class));
}
ch.pipeline().addLast(new RealServerChannelHandler());
}
});
return bootstrap;
}
}
@@ -1,33 +1,17 @@
package org.dromara.neutrinoproxy.client.core;
import cn.hutool.core.util.StrUtil;
import io.netty.handler.logging.LoggingHandler;
import org.dromara.neutrinoproxy.client.config.ProxyConfig;
import org.dromara.neutrinoproxy.client.util.ProxyUtil;
import org.dromara.neutrinoproxy.core.ProxyMessage;
import org.dromara.neutrinoproxy.core.ProxyMessageDecoder;
import org.dromara.neutrinoproxy.core.ProxyMessageEncoder;
import org.dromara.neutrinoproxy.core.util.FileUtil;
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.handler.ssl.SslHandler;
import io.netty.handler.timeout.IdleStateHandler;
import lombok.extern.slf4j.Slf4j;
import org.noear.solon.Solon;
import org.noear.solon.annotation.Component;
import org.noear.solon.annotation.Init;
import org.noear.solon.annotation.Inject;
import javax.net.ssl.SSLContext;
import javax.net.ssl.SSLEngine;
import javax.net.ssl.TrustManager;
import javax.net.ssl.TrustManagerFactory;
import java.io.InputStream;
import java.net.InetSocketAddress;
import java.security.KeyStore;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@@ -44,24 +28,11 @@ public class ProxyClientService {
private ProxyConfig proxyConfig;
@Inject("cmdTunnelBootstrap")
private Bootstrap cmdTunnelBootstrap;
@Inject("proxyTunnelBootstrap")
private Bootstrap proxyTunnelBootstrap;
@Inject("realServerBootstrap")
private Bootstrap realServerBootstrap;
private volatile Channel channel;
// /**
// * 重连间隔(秒)
// */
// private static final long RECONNECT_INTERVAL_SECONDS = 5;
/**
* 重连次数
*/
private volatile int reconnectCount = 0;
// /**
// * 启用重连服务
// */
// private volatile boolean reconnectServiceEnable = false;
private NioEventLoopGroup workerGroup;
/**
* 重连服务执行器
*/
@@ -70,72 +41,6 @@ public class ProxyClientService {
@Init
public void init() {
this.reconnectExecutor.scheduleWithFixedDelay(this::reconnect, 10, proxyConfig.getTunnel().getReconnection().getIntervalSeconds(), TimeUnit.SECONDS);
this.workerGroup = new NioEventLoopGroup(proxyConfig.getTunnel().getThreadCount());
realServerBootstrap.group(workerGroup);
realServerBootstrap.channel(NioSocketChannel.class);
realServerBootstrap.handler(new ChannelInitializer<SocketChannel>() {
@Override
public void initChannel(SocketChannel ch) throws Exception {
if (null != proxyConfig.getTunnel().getTransferLogEnable() && proxyConfig.getTunnel().getTransferLogEnable()) {
ch.pipeline().addFirst(new LoggingHandler(RealServerChannelHandler.class));
}
ch.pipeline().addLast(new RealServerChannelHandler());
}
});
proxyTunnelBootstrap.group(workerGroup);
proxyTunnelBootstrap.channel(NioSocketChannel.class);
proxyTunnelBootstrap.remoteAddress(InetSocketAddress.createUnresolved(proxyConfig.getTunnel().getServerIp(), proxyConfig.getTunnel().getServerPort()));
proxyTunnelBootstrap.handler(new ChannelInitializer<SocketChannel>() {
@Override
public void initChannel(SocketChannel ch) throws Exception {
if (proxyConfig.getTunnel().getSslEnable()) {
ch.pipeline().addLast(createSslHandler());
}
if (null != proxyConfig.getTunnel().getTransferLogEnable() && proxyConfig.getTunnel().getTransferLogEnable()) {
ch.pipeline().addFirst(new LoggingHandler(ProxyChannelHandler.class));
}
ch.pipeline().addLast(new ProxyMessageDecoder(proxyConfig.getProtocol().getMaxFrameLength(),
proxyConfig.getProtocol().getLengthFieldOffset(), proxyConfig.getProtocol().getLengthFieldLength(),
proxyConfig.getProtocol().getLengthAdjustment(), proxyConfig.getProtocol().getInitialBytesToStrip()));
ch.pipeline().addLast(new ProxyMessageEncoder());
ch.pipeline().addLast(new IdleStateHandler(proxyConfig.getProtocol().getReadIdleTime(), proxyConfig.getProtocol().getWriteIdleTime(), proxyConfig.getProtocol().getAllIdleTimeSeconds()));
ch.pipeline().addLast(new ProxyChannelHandler());
}
});
cmdTunnelBootstrap.group(workerGroup);
cmdTunnelBootstrap.channel(NioSocketChannel.class);
// cmdTunnelBootstrap.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 3000);
// cmdTunnelBootstrap.option(ChannelOption.SO_KEEPALIVE, true);
// /**
// * TCP/IP协议中,无论发送多少数据,总是要在数据前面加上协议头,同时,对方接收到数据,也需要发送ACK表示确认。为了尽可能的利用网络带宽,TCP总是希望尽可能的发送足够大的数据。(一个连接会设置MSS参数,因此,TCP/IP希望每次都能够以MSS尺寸的数据块来发送数据)。
// * Nagle算法就是为了尽可能发送大块数据,避免网络中充斥着许多小数据块。
// */
// cmdTunnelBootstrap.option(ChannelOption.TCP_NODELAY, true);
cmdTunnelBootstrap.remoteAddress(InetSocketAddress.createUnresolved(proxyConfig.getTunnel().getServerIp(), proxyConfig.getTunnel().getServerPort()));
cmdTunnelBootstrap.handler(new ChannelInitializer<SocketChannel>() {
@Override
public void initChannel(SocketChannel ch) throws Exception {
if (proxyConfig.getTunnel().getSslEnable()) {
ch.pipeline().addLast(createSslHandler());
}
if (null != proxyConfig.getTunnel().getTransferLogEnable() && proxyConfig.getTunnel().getTransferLogEnable()) {
ch.pipeline().addFirst(new LoggingHandler(CmdChannelHandler.class));
}
ch.pipeline().addLast(new ProxyMessageDecoder(proxyConfig.getProtocol().getMaxFrameLength(),
proxyConfig.getProtocol().getLengthFieldOffset(), proxyConfig.getProtocol().getLengthFieldLength(),
proxyConfig.getProtocol().getLengthAdjustment(), proxyConfig.getProtocol().getInitialBytesToStrip()));
ch.pipeline().addLast(new ProxyMessageEncoder());
ch.pipeline().addLast(new IdleStateHandler(proxyConfig.getProtocol().getReadIdleTime(), proxyConfig.getProtocol().getWriteIdleTime(), proxyConfig.getProtocol().getAllIdleTimeSeconds()));
ch.pipeline().addLast(new CmdChannelHandler());
}
});
try {
this.start();
@@ -203,33 +108,7 @@ public class ProxyClientService {
}).sync();
}
private ChannelHandler createSslHandler() {
try {
InputStream jksInputStream = FileUtil.getInputStream(proxyConfig.getTunnel().getJksPath());
SSLContext clientContext = SSLContext.getInstance("TLS");
final KeyStore ks = KeyStore.getInstance("JKS");
ks.load(jksInputStream, proxyConfig.getTunnel().getKeyStorePassword().toCharArray());
TrustManagerFactory tmf = TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm());
tmf.init(ks);
TrustManager[] trustManagers = tmf.getTrustManagers();
clientContext.init(null, trustManagers, null);
SSLEngine sslEngine = clientContext.createSSLEngine();
sslEngine.setUseClientMode(true);
return new SslHandler(sslEngine);
} catch (Exception e) {
log.error("创建SSL处理器失败", e);
e.printStackTrace();
}
return null;
}
protected synchronized void reconnect() {
// if (!reconnectServiceEnable) {
// return;
// }
if (null != channel) {
if (channel.isActive()) {
return;
@@ -21,8 +21,8 @@ import org.noear.solon.annotation.Inject;
@Match(type = Constants.ProxyDataTypeName.CONNECT)
@Component
public class ProxyMessageConnectHandler implements ProxyMessageHandler {
@Inject("proxyTunnelBootstrap")
private Bootstrap proxyTunnelBootstrap;
@Inject("tcpProxyTunnelBootstrap")
private Bootstrap tcpProxyTunnelBootstrap;
@Inject("realServerBootstrap")
private Bootstrap realServerBootstrap;
@Inject
@@ -48,7 +48,7 @@ public class ProxyMessageConnectHandler implements ProxyMessageHandler {
realServerChannel.config().setOption(ChannelOption.AUTO_READ, false);
// 获取连接
ProxyUtil.borrowProxyChanel(proxyTunnelBootstrap, new ProxyChannelBorrowListener() {
ProxyUtil.borrowProxyChanel(tcpProxyTunnelBootstrap, new ProxyChannelBorrowListener() {
@Override
public void success(Channel channel) {
@@ -21,6 +21,9 @@
*/
package org.dromara.neutrinoproxy.client.util;
import io.netty.channel.ChannelHandler;
import io.netty.handler.ssl.SslHandler;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.dromara.neutrinoproxy.client.config.ProxyConfig;
import org.dromara.neutrinoproxy.client.core.ProxyChannelBorrowListener;
@@ -34,6 +37,12 @@ import io.netty.util.AttributeKey;
import org.dromara.neutrinoproxy.core.util.FileUtil;
import org.noear.solon.Solon;
import javax.net.ssl.SSLContext;
import javax.net.ssl.SSLEngine;
import javax.net.ssl.TrustManager;
import javax.net.ssl.TrustManagerFactory;
import java.io.InputStream;
import java.security.KeyStore;
import java.util.Iterator;
import java.util.Map;
import java.util.UUID;
@@ -45,6 +54,7 @@ import java.util.concurrent.ConcurrentLinkedQueue;
* @author: aoshiguchen
* @date: 2022/8/31
*/
@Slf4j
public class ProxyUtil {
private static final AttributeKey<Boolean> USER_CHANNEL_WRITEABLE = AttributeKey.newInstance("user_channel_writeable");
@@ -154,4 +164,27 @@ public class ProxyUtil {
clientId = id;
return id;
}
public static ChannelHandler createSslHandler(ProxyConfig proxyConfig) {
try {
InputStream jksInputStream = FileUtil.getInputStream(proxyConfig.getTunnel().getJksPath());
SSLContext clientContext = SSLContext.getInstance("TLS");
final KeyStore ks = KeyStore.getInstance("JKS");
ks.load(jksInputStream, proxyConfig.getTunnel().getKeyStorePassword().toCharArray());
TrustManagerFactory tmf = TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm());
tmf.init(ks);
TrustManager[] trustManagers = tmf.getTrustManagers();
clientContext.init(null, trustManagers, null);
SSLEngine sslEngine = clientContext.createSSLEngine();
sslEngine.setUseClientMode(true);
return new SslHandler(sslEngine);
} catch (Exception e) {
log.error("创建SSL处理器失败", e);
e.printStackTrace();
}
return null;
}
}