From 4c828f814c1b4dc38aa279502626c0bc1ee29f93 Mon Sep 17 00:00:00 2001 From: aoshiguchen <1052045476@qq.com> Date: Sat, 3 Jun 2023 03:19:45 +0800 Subject: [PATCH] =?UTF-8?q?=E6=9C=8D=E5=8A=A1=E7=AB=AF=E3=80=81=E5=AE=A2?= =?UTF-8?q?=E6=88=B7=E7=AB=AF=E8=AE=A4=E8=AF=81=E9=80=BB=E8=BE=91=E8=B0=83?= =?UTF-8?q?=E6=95=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../client/config/ProxyConfiguration.java | 9 +- .../client/core/CmdChannelHandler.java | 77 ++++++ ...lHandler.java => ProxyChannelHandler.java} | 35 +-- .../client/core/ProxyClientService.java | 250 ++++++++++-------- .../client/core/RealServerChannelHandler.java | 4 +- .../handler/ProxyMessageConnectHandler.java | 6 +- .../neutrinoproxy/client/util/ProxyUtil.java | 6 +- 7 files changed, 242 insertions(+), 145 deletions(-) create mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CmdChannelHandler.java rename neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/{ClientChannelHandler.java => ProxyChannelHandler.java} (73%) diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/ProxyConfiguration.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/ProxyConfiguration.java index ebc761ef..ca206083 100644 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/ProxyConfiguration.java +++ b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/ProxyConfiguration.java @@ -31,8 +31,13 @@ public class ProxyConfiguration implements LifecycleBean { Solon.context().wrapAndPut(Dispatcher.class, dispatcher); } - @Bean("bootstrap") - public Bootstrap bootstrap() { + @Bean("cmdTunnelBootstrap") + public Bootstrap cmdTunnelBootstrap() { + return new Bootstrap(); + } + + @Bean("proxyTunnelBootstrap") + public Bootstrap proxyTunnelBootstrap() { return new Bootstrap(); } diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CmdChannelHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CmdChannelHandler.java new file mode 100644 index 00000000..6db7fdb5 --- /dev/null +++ b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CmdChannelHandler.java @@ -0,0 +1,77 @@ +package org.dromara.neutrinoproxy.client.core; + +import org.dromara.neutrinoproxy.client.util.ProxyUtil; +import org.dromara.neutrinoproxy.core.Constants; +import org.dromara.neutrinoproxy.core.ProxyMessage; +import org.dromara.neutrinoproxy.core.dispatcher.Dispatcher; +import io.netty.channel.Channel; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelOption; +import io.netty.channel.SimpleChannelInboundHandler; +import io.netty.handler.timeout.IdleStateEvent; +import lombok.extern.slf4j.Slf4j; +import org.noear.solon.Solon; + +/** + * 处理与服务端之间的数据传输 + * @author: aoshiguchen + * @date: 2022/6/16 + */ +@Slf4j +public class CmdChannelHandler extends SimpleChannelInboundHandler { + + + @Override + protected void channelRead0(ChannelHandlerContext ctx, ProxyMessage proxyMessage) throws Exception { + if (ProxyMessage.TYPE_HEARTBEAT != proxyMessage.getType()) { + log.info("Client CmdChannel recieved proxy message, type is {}", proxyMessage.getType()); + } + Solon.context().getBean(Dispatcher.class).dispatch(ctx, proxyMessage); + } + + @Override + public void channelWritabilityChanged(ChannelHandlerContext ctx) throws Exception { + Channel realServerChannel = ctx.channel().attr(Constants.NEXT_CHANNEL).get(); + if (realServerChannel != null) { + realServerChannel.config().setOption(ChannelOption.AUTO_READ, ctx.channel().isWritable()); + } + + super.channelWritabilityChanged(ctx); + } + + @Override + public void channelInactive(ChannelHandlerContext ctx) throws Exception { + log.info("Client CmdChannel 与服务端断开连接"); + ProxyUtil.setCmdChannel(null); + ProxyUtil.clearRealServerChannels(); + + super.channelInactive(ctx); + } + + @Override + public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { + log.error("Client CmdChannel Error channelId:{}", ctx.channel().id().asLongText(), cause); + ctx.close(); + } + + @Override + public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { + if(evt instanceof IdleStateEvent) { + IdleStateEvent event = (IdleStateEvent)evt; + switch (event.state()) { + case READER_IDLE: + // 读超时,断开连接 +// log.info("读超时"); +// ctx.channel().close(); + break; + case WRITER_IDLE: + ctx.channel().writeAndFlush(ProxyMessage.buildHeartbeatMessage()); + break; + case ALL_IDLE: + log.info("Client CmdChannel 读写超时"); + ctx.close(); + break; + } + } + } +} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ClientChannelHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyChannelHandler.java similarity index 73% rename from neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ClientChannelHandler.java rename to neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyChannelHandler.java index b6d85001..ebe217f6 100644 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ClientChannelHandler.java +++ b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyChannelHandler.java @@ -1,15 +1,15 @@ package org.dromara.neutrinoproxy.client.core; -import org.dromara.neutrinoproxy.client.util.ProxyUtil; -import org.dromara.neutrinoproxy.core.Constants; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.dromara.neutrinoproxy.core.dispatcher.Dispatcher; import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelOption; import io.netty.channel.SimpleChannelInboundHandler; import io.netty.handler.timeout.IdleStateEvent; import lombok.extern.slf4j.Slf4j; +import org.dromara.neutrinoproxy.client.util.ProxyUtil; +import org.dromara.neutrinoproxy.core.Constants; +import org.dromara.neutrinoproxy.core.ProxyMessage; +import org.dromara.neutrinoproxy.core.dispatcher.Dispatcher; import org.noear.solon.Solon; /** @@ -18,13 +18,13 @@ import org.noear.solon.Solon; * @date: 2022/6/16 */ @Slf4j -public class ClientChannelHandler extends SimpleChannelInboundHandler { +public class ProxyChannelHandler extends SimpleChannelInboundHandler { @Override protected void channelRead0(ChannelHandlerContext ctx, ProxyMessage proxyMessage) throws Exception { if (ProxyMessage.TYPE_HEARTBEAT != proxyMessage.getType()) { - log.info("recieved proxy message, type is {}", proxyMessage.getType()); + log.info("Client ProxyChannel recieved proxy message, type is {}", proxyMessage.getType()); } Solon.context().getBean(Dispatcher.class).dispatch(ctx, proxyMessage); } @@ -41,17 +41,10 @@ public class ClientChannelHandler extends SimpleChannelInboundHandler() { @@ -78,17 +81,10 @@ public class ProxyClientService { } }); - bootstrap.group(workerGroup); - 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.handler(new ChannelInitializer() { + proxyTunnelBootstrap.group(workerGroup); + proxyTunnelBootstrap.channel(NioSocketChannel.class); + proxyTunnelBootstrap.remoteAddress(InetSocketAddress.createUnresolved(proxyConfig.getClient().getServerIp(), proxyConfig.getClient().getServerPort())); + proxyTunnelBootstrap.handler(new ChannelInitializer() { @Override public void initChannel(SocketChannel ch) throws Exception { @@ -101,115 +97,143 @@ public class ProxyClientService { 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 ClientChannelHandler()); + ch.pipeline().addLast(new ProxyChannelHandler()); } }); - try { - this.start(); - } catch (Exception e) { - // 启动连不上也做一下重连,因此先catch异常 - log.error("启动异常", e); - } - } + 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.getClient().getServerIp(), proxyConfig.getClient().getServerPort())); - public void start() { - if (StrUtil.isEmpty(proxyConfig.getClient().getServerIp())) { - log.error("not found server-ip config."); - Solon.stop(); - return; - } - if (null == proxyConfig.getClient().getServerPort()) { - log.error("not found server-port config."); - Solon.stop(); - return; - } - if (null != proxyConfig.getClient().getSslEnable() && proxyConfig.getClient().getSslEnable() - && StrUtil.isEmpty(proxyConfig.getClient().getJksPath())) { - log.error("not found jks-path config."); - Solon.stop(); - return; - } - if (StrUtil.isEmpty(proxyConfig.getClient().getLicenseKey())) { - log.error("not found license-key config."); - Solon.stop(); - return; - } - if (null == channel || !channel.isActive()) { - try { - connectProxyServer(); - } catch (Exception e) { - log.error("client start error", e); + cmdTunnelBootstrap.handler(new ChannelInitializer() { + + @Override + public void initChannel(SocketChannel ch) throws Exception { + if (proxyConfig.getClient().getSslEnable()) { + ch.pipeline().addLast(createSslHandler()); } - } else { - channel.writeAndFlush(ProxyMessage.buildAuthMessage(proxyConfig.getClient().getLicenseKey(), ProxyUtil.getClientId())); +// ch.pipeline().addFirst(new LoggingHandler(ProxyClientService.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(); + } catch (Exception e) { + // 启动连不上也做一下重连,因此先catch异常 + log.error("[客户端指令隧道] 启动异常", e); } + } - /** - * 连接代理服务器 - */ - private void connectProxyServer() throws InterruptedException { - bootstrap.connect(proxyConfig.getClient().getServerIp(), proxyConfig.getClient().getServerPort()) - .addListener(new ChannelFutureListener() { - - @Override - public void operationComplete(ChannelFuture future) throws Exception { - if (future.isSuccess()) { - channel = future.channel(); - // 连接成功,向服务器发送客户端认证信息(licenseKey) - ProxyUtil.setCmdChannel(future.channel()); - future.channel().writeAndFlush(ProxyMessage.buildAuthMessage(proxyConfig.getClient().getLicenseKey(), ProxyUtil.getClientId())); - log.info("连接代理服务成功. channelId:{}", future.channel().id().asLongText()); - -// reconnectServiceEnable = true; - reconnectCount = 0; - } else { - log.info("连接代理服务失败!"); - } - } - }).sync(); + public void start() { + if (StrUtil.isEmpty(proxyConfig.getClient().getServerIp())) { + log.error("not found server-ip config."); + Solon.stop(); + return; } - - private ChannelHandler createSslHandler() { - try { - InputStream jksInputStream = FileUtil.getInputStream(proxyConfig.getClient().getJksPath()); - - SSLContext clientContext = SSLContext.getInstance("TLS"); - final KeyStore ks = KeyStore.getInstance("JKS"); - ks.load(jksInputStream, proxyConfig.getClient().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; + if (null == proxyConfig.getClient().getServerPort()) { + log.error("not found server-port config."); + Solon.stop(); + return; } - - protected synchronized void reconnect() { -// if (!reconnectServiceEnable) { -// return; -// } - if (null != channel) { - if (channel.isActive()) { - return; - } - channel.close(); - } - - log.info("客户端重连 seq:{}", ++reconnectCount); + if (null != proxyConfig.getClient().getSslEnable() && proxyConfig.getClient().getSslEnable() + && StrUtil.isEmpty(proxyConfig.getClient().getJksPath())) { + log.error("not found jks-path config."); + Solon.stop(); + return; + } + if (StrUtil.isEmpty(proxyConfig.getClient().getLicenseKey())) { + log.error("not found license-key config."); + Solon.stop(); + return; + } + if (null == channel || !channel.isActive()) { try { connectProxyServer(); } catch (Exception e) { - log.error("重连异常", e); + log.error("client start error", e); } + } else { + channel.writeAndFlush(ProxyMessage.buildAuthMessage(proxyConfig.getClient().getLicenseKey(), ProxyUtil.getClientId())); } + } + + /** + * 连接代理服务器 + */ + private void connectProxyServer() throws InterruptedException { + cmdTunnelBootstrap.connect() + .addListener(new ChannelFutureListener() { + + @Override + public void operationComplete(ChannelFuture future) throws Exception { + if (future.isSuccess()) { + channel = future.channel(); + // 连接成功,向服务器发送客户端认证信息(licenseKey) + ProxyUtil.setCmdChannel(future.channel()); + future.channel().writeAndFlush(ProxyMessage.buildAuthMessage(proxyConfig.getClient().getLicenseKey(), ProxyUtil.getClientId())); + log.info("[客户端指令隧道] 连接代理服务成功. channelId:{}", future.channel().id().asLongText()); + +// reconnectServiceEnable = true; + reconnectCount = 0; + } else { + log.info("[客户端指令隧道] 连接代理服务失败!"); + } + } + }).sync(); + } + + private ChannelHandler createSslHandler() { + try { + InputStream jksInputStream = FileUtil.getInputStream(proxyConfig.getClient().getJksPath()); + + SSLContext clientContext = SSLContext.getInstance("TLS"); + final KeyStore ks = KeyStore.getInstance("JKS"); + ks.load(jksInputStream, proxyConfig.getClient().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; + } + channel.close(); + } + + log.info("[客户端指令隧道] 客户端重连 seq:{}", ++reconnectCount); + try { + connectProxyServer(); + } catch (Exception e) { + log.error("[客户端指令隧道] 重连异常", e); + } + } } diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/RealServerChannelHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/RealServerChannelHandler.java index 0506217b..69c47d45 100644 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/RealServerChannelHandler.java +++ b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/RealServerChannelHandler.java @@ -22,6 +22,7 @@ package org.dromara.neutrinoproxy.client.core; +import lombok.extern.slf4j.Slf4j; import org.dromara.neutrinoproxy.client.util.ProxyUtil; import org.dromara.neutrinoproxy.core.Constants; import org.dromara.neutrinoproxy.core.ProxyMessage; @@ -36,6 +37,7 @@ import io.netty.channel.SimpleChannelInboundHandler; * @author: aoshiguchen * @date: 2022/6/16 */ +@Slf4j public class RealServerChannelHandler extends SimpleChannelInboundHandler { @@ -85,6 +87,6 @@ public class RealServerChannelHandler extends SimpleChannelInboundHandler { + proxyTunnelBootstrap.connect().addListener((ChannelFutureListener) future -> { if (future.isSuccess()) { borrowListener.success(future.channel()); } else {