From 0dd71195cd89026d95826a365787d761bc15394b Mon Sep 17 00:00:00 2001 From: xgc Date: Sat, 20 Jan 2024 22:56:47 +0800 Subject: [PATCH] =?UTF-8?q?=E5=88=A0=E9=99=A4client=E5=A4=9A=E4=BD=99?= =?UTF-8?q?=E9=83=A8=E5=88=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../client/sdk/ProxyClientSdk.java | 48 ---- .../src/main/resources/app.yml | 68 ------ .../src/main/resources/test.jks | Bin 1388 -> 0 bytes neutrino-proxy-client/pom.xml | 6 +- .../NeutrinoClientRuntimeNativeRegistrar.java | 26 -- .../client/config/ProxyConfig.java | 74 ------ .../client/config/ProxyConfiguration.java | 181 ++++---------- .../client/constant/Constants.java | 13 - .../client/core/CmdChannelHandler.java | 85 ------- .../client/core/CustomThreadFactory.java | 55 ----- .../core/ProxyChannelBorrowListener.java | 38 --- .../client/core/ProxyClientService.java | 130 ---------- .../client/core/RealServerChannelHandler.java | 103 -------- .../client/core/TcpProxyChannelHandler.java | 82 ------- .../client/core/UdpProxyChannelHandler.java | 80 ------- .../client/core/UdpRealServerHandler.java | 41 ---- .../client/handler/BeanHandler.java | 19 ++ .../handler/ProxyMessageAuthHandler.java | 46 ---- .../handler/ProxyMessageConnectHandler.java | 87 ------- .../ProxyMessageDisconnectHandler.java | 40 ---- .../handler/ProxyMessageErrorHandler.java | 38 --- .../handler/ProxyMessageTransferHandler.java | 41 ---- .../UdpProxyMessageConnectHandler.java | 64 ----- .../UdpProxyMessageTransferHandler.java | 48 ---- .../client/util/LockChannel.java | 28 --- .../neutrinoproxy/client/util/ProxyUtil.java | 226 ------------------ .../client/util/UdpChannelBindInfo.java | 22 -- .../client/util/UdpServerUtil.java | 201 ---------------- 28 files changed, 64 insertions(+), 1826 deletions(-) delete mode 100644 neutrino-proxy-client-sdk/src/main/java/org/dromara/neutrinoproxy/client/sdk/ProxyClientSdk.java delete mode 100644 neutrino-proxy-client-sdk/src/main/resources/app.yml delete mode 100644 neutrino-proxy-client-sdk/src/main/resources/test.jks delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/NeutrinoClientRuntimeNativeRegistrar.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/ProxyConfig.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/constant/Constants.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CmdChannelHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CustomThreadFactory.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyChannelBorrowListener.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyClientService.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/RealServerChannelHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/TcpProxyChannelHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpProxyChannelHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpRealServerHandler.java create mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/BeanHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageAuthHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageConnectHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageDisconnectHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageErrorHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageTransferHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageConnectHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageTransferHandler.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/LockChannel.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/ProxyUtil.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpChannelBindInfo.java delete mode 100644 neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpServerUtil.java diff --git a/neutrino-proxy-client-sdk/src/main/java/org/dromara/neutrinoproxy/client/sdk/ProxyClientSdk.java b/neutrino-proxy-client-sdk/src/main/java/org/dromara/neutrinoproxy/client/sdk/ProxyClientSdk.java deleted file mode 100644 index c0570f13..00000000 --- a/neutrino-proxy-client-sdk/src/main/java/org/dromara/neutrinoproxy/client/sdk/ProxyClientSdk.java +++ /dev/null @@ -1,48 +0,0 @@ -package org.dromara.neutrinoproxy.client.sdk; - -import cn.hutool.core.util.StrUtil; -import lombok.extern.slf4j.Slf4j; -import org.noear.solon.Solon; -import org.noear.solon.Utils; -import org.noear.solon.annotation.SolonMain; - -/** - * - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Slf4j -@SolonMain -public class ProxyClientSdk { - - public static void main(String[] args) { - Solon.start(ProxyClientSdk.class, args, app -> { - String loglevel = System.getenv("LOG_LEVEL"); - if (Utils.isNotEmpty(loglevel)) { - app.cfg().put("solon.logging.logger.root.level", loglevel); - } - - setAlias("neutrino.proxy.tunnel.server-ip", "serverIp"); - setAlias("neutrino.proxy.tunnel.server-port", "serverPort"); - setAlias("neutrino.proxy.tunnel.ssl-enable", "sslEnable"); - setAlias("neutrino.proxy.tunnel.jks-path", "jksPath"); - setAlias("neutrino.proxy.tunnel.key-store-password", "keyStorePassword"); - setAlias("neutrino.proxy.tunnel.license-key", "licenseKey"); - - log.info("NeutrinoProxy Client :{}", app.cfg().get("solon.app.version")); - }); - } - - /** - * 别名处理,支持较短的启动参数名 - * @param key - * @param alias - */ - private static void setAlias(String key, String alias) { - String val = Solon.cfg().argx().get(alias); - if (StrUtil.isNotBlank(val)) { - Solon.cfg().put(key, val); - } - } - -} diff --git a/neutrino-proxy-client-sdk/src/main/resources/app.yml b/neutrino-proxy-client-sdk/src/main/resources/app.yml deleted file mode 100644 index a122416c..00000000 --- a/neutrino-proxy-client-sdk/src/main/resources/app.yml +++ /dev/null @@ -1,68 +0,0 @@ -solon: - config: - add: ./app.yml - app: - name: neutrino-proxy-client - version: 2.0.1 -# 日志级别 -solon.logging.appender: - console: - pattern: "%d{yyyy-MM-dd HH:mm:ss.SSS} %highlight(%-5level) %magenta(${PID:-}) --- %-15([%15.15thread]) %-56(%cyan(%-40.40logger{39}%L)) : %msg%n" - file: - enable: true - pattern: "%d{yyyy-MM-dd HH:mm:ss.SSS} %-5level ${PID:-} --- %-15([%15.15thread]) %-56(%-40.40logger{39}%L) : %msg%n" - name: "logs/${neutrino.application.name}" - rolling: "logs/${neutrino.application.name}_%d{yyyy-MM-dd}_%i.log.gz" -solon.logging.logger: - "root": - level: info - -neutrino: - application: - name: neutrino-proxy-client - - proxy: - protocol: - max-frame-length: 2097152 - length-field-offset: 0 - length-field-length: 4 - initial-bytes-to-strip: 0 - length-adjustment: 0 - read-idle-time: 120 - write-idle-time: 20 - all-idle-time-seconds: 0 - tunnel: - # 线程池相关配置,用于技术调优,可忽略 - thread-count: 50 - # 隧道SSL证书配置 - key-store-password: ${STORE_PASS:123456} - jks-path: ${JKS_PATH:classpath:/test.jks} - # 服务端IP - server-ip: ${SERVER_IP:localhost} - # 服务端端口(对应服务端app.yml中的tunnel.port、tunnel.ssl-port) - server-port: ${SERVER_PORT:9002} - # 是否启用SSL(注意:该配置必须和server-port对应上) - ssl-enable: ${SSL_ENABLE:true} - # 客户端连接唯一凭证 - license-key: ${LICENSE_KEY:} - # 客户端唯一身份标识(可忽略,若不设置首次启动会自动生成) - client-id: ${CLIENT_ID:} - # 是否开启隧道传输报文日志(日志级别为debug时开启才有效) - transfer-log-enable: ${CLIENT_LOG:false} - # 是否开启心跳日志 - heartbeat-log-enable: ${HEARTBEAT_LOG:false} - # 重连设置 - reconnection: - # 重连间隔(秒) - interval-seconds: 10 - # 是否开启无限重连(未开启时,客户端license不合法会自动停止应用,开启了则不会,请谨慎开启) - unlimited: false - client: - udp: - # 线程池相关配置,用于技术调优,可忽略 - boss-thread-count: 5 - work-thread-count: 20 - # udp傀儡端口范围 - puppet-port-range: 10000-10500 - # 是否开启隧道传输报文日志(日志级别为debug时开启才有效) - transfer-log-enable: ${CLIENT_LOG:false} diff --git a/neutrino-proxy-client-sdk/src/main/resources/test.jks b/neutrino-proxy-client-sdk/src/main/resources/test.jks deleted file mode 100644 index 3a4d1ad89e1a538f35c3b2996b2241d465ddedad..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 1388 zcmezO_TO6u1_mY|W&~rN#Nzbi;^Lf=)Z!9DpkSlVC1(Smg53s9Oxq3k*toRW7+Dy# zm;@OaSs7TGm=?c3uQzq2R{u|yvt1m9yKNY51-F}KEJ}Vg&o%pt^a}sGfvLwnY29Gd zTK;ObR&0Uur4M;;GlS;#-gw>e_e9sjrFAU2jh#`x$0uYu?)VbdCRlQ;cbop#2hGK= ziyy2KJ)QF9`=T?|dPT+u;&T-HqqpXzeZ1Z;XvAW7?@w^7din<*!I_J#87{mRYhpFg z@10k1i)D8CPu+uI?f*lb3r_A`$5e44A#U5J9|1W>CN7u}_~lBo{fzJ1IL%&NU}<<@ z)V5RUV4@=b#-E8d!w;I?vHr&+aeRy5o2%mfvpOZ$tUL2=>Dq(Id<=a@GJI<6S5H4u znXcaTzss0K{no9IXN%k}NJdS5VmIZ|CVq(%7xkDWSNI2s*L3>c6XI2Qnk3%4xN7-Z z58Kt}*`GbsV867(r#9b5I?Mg%X4$%9wtrgW+e`}+Z6>$Y%)7nFiZNR81w$Cu_AqJb zzqyy!eHU5z`gxZ54>#VWPmcX@Kf_vD9oC%ove36y^O4cD+_1&y;RvIo@|$6)$zqXeQT_l?H8kLAw$r zW&M}^;8?ul=EW1DugazGOx|!dyk>gn(nEVXFF&u44!K+@p_`SSd3O6zjd|~vuP@Eo zy6e1E!t3PCp*e@H-@KvjCZD_N^S!@nri)hSm3uGuJ3nz&&gnN_*k`dPIdv?mE2<0) z56!D}6o2z@sq%h~hv(ulYQmY{DQ6#_`&EQZq|NmC&$fu4C+AjmZ+#GX^7=XZ1V5*` zm6c};Rx>@_RU~_-Y|8n6#)379wv)A$HviU1z3}*UbKtfqO~zM`OX=UypC-e+wJ?>- z;AcTcx$N;XYo=LxTh7|MHchRL>>MUXhuXmYSCi6X8aP@EUM~ z)G!OPR3_%78_0?C8W|aw85$T^n3@=xM~U+q1Gxs~P%eEO(KsL3@4!6G+}O)t(Ade; z*vN3L^6OR^0Yd^Q#lYBs__SiYmr~JCS z&z|e7)OnIFx8r=dQti~DZH^%(Dwy0jzrma7_ z+b`we=F3O6^6gRblG-}2vnTF(*1Cxb^(}4nPC0Sk=Q1%fGB7SyG>|ut1%{g}ABz}^ zNV)8t_!eH_w&B5RX^O|dbLmff|Ae|f#!{yG!~gS3T7CXD&12I zwcEQW@#FDqN5Xx ioku-b)|M>LROU)QG0ot&TE@jEVVOmyMi>7@Z~y={B|CNi diff --git a/neutrino-proxy-client/pom.xml b/neutrino-proxy-client/pom.xml index 1e5d3a8a..a15e1195 100644 --- a/neutrino-proxy-client/pom.xml +++ b/neutrino-proxy-client/pom.xml @@ -16,13 +16,9 @@ org.dromara.neutrino-proxy - neutrino-proxy-core + neutrino-proxy-client-sdk ${revision} - - org.noear - solon-lib - diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/NeutrinoClientRuntimeNativeRegistrar.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/NeutrinoClientRuntimeNativeRegistrar.java deleted file mode 100644 index c9f84ef8..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/NeutrinoClientRuntimeNativeRegistrar.java +++ /dev/null @@ -1,26 +0,0 @@ -package org.dromara.neutrinoproxy.client.config; - -import org.noear.solon.annotation.Component; -import org.noear.solon.aot.RuntimeNativeMetadata; -import org.noear.solon.aot.RuntimeNativeRegistrar; -import org.noear.solon.aot.hint.MemberCategory; -import org.noear.solon.core.AppContext; - -/** - * @author songyinyin - * @since 2023/10/21 22:50 - */ -@Component -public class NeutrinoClientRuntimeNativeRegistrar implements RuntimeNativeRegistrar { - @Override - public void register(AppContext context, RuntimeNativeMetadata metadata) { - metadata.registerResourceInclude("test.jks"); - - metadata.registerReflection(ProxyConfig.Protocol.class, MemberCategory.DECLARED_FIELDS, MemberCategory.INVOKE_PUBLIC_CONSTRUCTORS, MemberCategory.INVOKE_DECLARED_METHODS); - metadata.registerReflection(ProxyConfig.Client.class, MemberCategory.DECLARED_FIELDS, MemberCategory.INVOKE_PUBLIC_CONSTRUCTORS, MemberCategory.INVOKE_DECLARED_METHODS); - metadata.registerReflection(ProxyConfig.Tunnel.class, MemberCategory.DECLARED_FIELDS, MemberCategory.INVOKE_PUBLIC_CONSTRUCTORS, MemberCategory.INVOKE_DECLARED_METHODS); - metadata.registerReflection(ProxyConfig.Tcp.class, MemberCategory.DECLARED_FIELDS, MemberCategory.INVOKE_PUBLIC_CONSTRUCTORS, MemberCategory.INVOKE_DECLARED_METHODS); - metadata.registerReflection(ProxyConfig.Udp.class, MemberCategory.DECLARED_FIELDS, MemberCategory.INVOKE_PUBLIC_CONSTRUCTORS, MemberCategory.INVOKE_DECLARED_METHODS); - metadata.registerReflection(ProxyConfig.Reconnection.class, MemberCategory.DECLARED_FIELDS, MemberCategory.INVOKE_PUBLIC_CONSTRUCTORS, MemberCategory.INVOKE_DECLARED_METHODS); - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/ProxyConfig.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/ProxyConfig.java deleted file mode 100644 index a33f3520..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/config/ProxyConfig.java +++ /dev/null @@ -1,74 +0,0 @@ -package org.dromara.neutrinoproxy.client.config; - -import lombok.Data; -import org.noear.solon.annotation.Component; -import org.noear.solon.annotation.Inject; - -/** - * - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Data -@Component -public class ProxyConfig { - @Inject("${neutrino.proxy.protocol}") - private Protocol protocol; - @Inject("${neutrino.proxy.tunnel}") - private Tunnel tunnel; - @Inject("${neutrino.proxy.client}") - private Client client; - - @Data - public static class Protocol { - private Integer maxFrameLength; - private Integer lengthFieldOffset; - private Integer lengthFieldLength; - private Integer initialBytesToStrip; - private Integer lengthAdjustment; - private Integer readIdleTime; - private Integer writeIdleTime; - private Integer allIdleTimeSeconds; - } - - @Data - public static class Tunnel { - private String keyStorePassword; - private String jksPath; - private String serverIp; - private Integer serverPort; - private Boolean sslEnable; - private Integer obtainLicenseInterval; - private String licenseKey; - private Integer threadCount; - private String clientId; - private Boolean transferLogEnable; - private Boolean heartbeatLogEnable; - private Reconnection reconnection; - } - - @Data - public static class Client { -// private Tcp tcp; - private Udp udp; - } - - @Data - public static class Reconnection { - private Integer intervalSeconds; - private Boolean unlimited; - } - - @Data - public static class Tcp { - - } - - @Data - public static class Udp { - private Integer bossThreadCount; - private Integer workThreadCount; - private String puppetPortRange; - private Boolean transferLogEnable; - } -} 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 87dda221..7ef012e1 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 @@ -1,29 +1,26 @@ package org.dromara.neutrinoproxy.client.config; -import io.netty.channel.ChannelInitializer; -import io.netty.channel.ChannelOption; -import io.netty.channel.ChannelPipeline; +import com.google.common.collect.Lists; +import io.netty.bootstrap.Bootstrap; +import io.netty.channel.ChannelHandlerContext; import io.netty.channel.nio.NioEventLoopGroup; -import io.netty.channel.socket.SocketChannel; -import io.netty.channel.socket.nio.NioDatagramChannel; -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.*; -import org.dromara.neutrinoproxy.client.util.ProxyUtil; -import org.dromara.neutrinoproxy.core.*; +import org.dromara.neutrinoproxy.client.sdk.config.IBeanHandler; +import org.dromara.neutrinoproxy.client.sdk.config.IProxyConfiguration; +import org.dromara.neutrinoproxy.client.sdk.config.ProxyConfig; +import org.dromara.neutrinoproxy.client.sdk.handler.*; +import org.dromara.neutrinoproxy.client.sdk.solon.BeanHandler; +import org.dromara.neutrinoproxy.core.ProxyDataTypeEnum; +import org.dromara.neutrinoproxy.core.ProxyMessage; +import org.dromara.neutrinoproxy.core.ProxyMessageHandler; import org.dromara.neutrinoproxy.core.aot.NeutrinoCoreRuntimeNativeRegistrar; import org.dromara.neutrinoproxy.core.dispatcher.DefaultDispatcher; import org.dromara.neutrinoproxy.core.dispatcher.Dispatcher; -import io.netty.bootstrap.Bootstrap; -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; /** @@ -32,182 +29,92 @@ import java.util.List; * @date: 2022/10/8 */ @Configuration -public class ProxyConfiguration implements LifecycleBean { +public class ProxyConfiguration extends IProxyConfiguration implements LifecycleBean { + @Inject + private ProxyConfig proxyConfig; + @Inject("tcpProxyTunnelBootstrap") + private Bootstrap tcpProxyTunnelBootstrap; + @Inject("realServerBootstrap") + private Bootstrap realServerBootstrap; @Override public void start() throws Throwable { - List list = Solon.context().getBeansOfType(ProxyMessageHandler.class); + List list = Lists.newArrayList( + new ProxyMessageAuthHandler(proxyConfig), + new ProxyMessageConnectHandler(tcpProxyTunnelBootstrap,realServerBootstrap,proxyConfig), + new ProxyMessageDisconnectHandler(), + new ProxyMessageErrorHandler(), + new ProxyMessageTransferHandler(), + new UdpProxyMessageConnectHandler(proxyConfig,tcpProxyTunnelBootstrap), + new UdpProxyMessageTransferHandler() + ); Dispatcher dispatcher = new DefaultDispatcher<>("MessageDispatcher", list, proxyMessage -> ProxyDataTypeEnum.of((int)proxyMessage.getType()) == null ? null : ProxyDataTypeEnum.of((int)proxyMessage.getType()).getName()); Solon.context().wrapAndPut(Dispatcher.class, dispatcher); } + @Override + public IBeanHandler getBeanHandler() { + return new BeanHandler(); + } @Bean("tunnelWorkGroup") public NioEventLoopGroup tunnelWorkGroup(@Inject ProxyConfig proxyConfig) { - return new NioEventLoopGroup(proxyConfig.getTunnel().getThreadCount()); + return super.tunnelWorkGroup(proxyConfig); } @Bean("tcpRealServerWorkGroup") public NioEventLoopGroup tcpRealServerWorkGroup(@Inject ProxyConfig proxyConfig) { // 暂时先公用此配置 - return new NioEventLoopGroup(proxyConfig.getTunnel().getThreadCount()); + return super.tcpRealServerWorkGroup(proxyConfig); } @Bean("udpServerGroup") public NioEventLoopGroup udpServerGroup(@Inject ProxyConfig proxyConfig) { // 暂时先公用此配置 - return new NioEventLoopGroup(proxyConfig.getClient().getUdp().getBossThreadCount()); + return super.udpServerGroup(proxyConfig); } @Bean("udpWorkGroup") public NioEventLoopGroup udpWorkGroup(@Inject ProxyConfig proxyConfig) { // 暂时先公用此配置 - return new NioEventLoopGroup(proxyConfig.getClient().getUdp().getWorkThreadCount()); + return super.udpWorkGroup(proxyConfig); } @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() { - - @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; + return super.cmdTunnelBootstrap(proxyConfig,tunnelWorkGroup); } @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() { - - @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(TcpProxyChannelHandler.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 TcpProxyChannelHandler()); - } - }); - return bootstrap; + return super.tcpProxyTunnelBootstrap(proxyConfig,tunnelWorkGroup); } @Bean("udpProxyTunnelBootstrap") public Bootstrap udpProxyTunnelBootstrap(@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() { - - @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(TcpProxyChannelHandler.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 UdpProxyChannelHandler()); - } - }); - return bootstrap; + @Inject("tunnelWorkGroup") NioEventLoopGroup tunnelWorkGroup) { + return super.udpProxyTunnelBootstrap(proxyConfig,tunnelWorkGroup); } @Bean("realServerBootstrap") 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() { - - @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; + @Inject("tcpRealServerWorkGroup") NioEventLoopGroup tcpRealServerWorkGroup + ) { + return super.realServerBootstrap(proxyConfig,tcpRealServerWorkGroup); } @Bean("udpServerBootstrap") public Bootstrap udpServerBootstrap(@Inject ProxyConfig proxyConfig, @Inject("udpServerGroup") NioEventLoopGroup udpServerGroup, @Inject("udpWorkGroup") NioEventLoopGroup udpWorkGroup) { - Bootstrap bootstrap = new Bootstrap(); - bootstrap.group(udpServerGroup) - // 主线程处理 - .channel(NioDatagramChannel.class) - // 广播 - .option(ChannelOption.SO_BROADCAST, true) - // 设置读缓冲区为2M - .option(ChannelOption.SO_RCVBUF, 2048 * 1024) - // 设置写缓冲区为1M - .option(ChannelOption.SO_SNDBUF, 1024 * 1024) - .handler(new ChannelInitializer() { - @Override - protected void initChannel(NioDatagramChannel ch) { - ChannelPipeline pipeline = ch.pipeline(); - if (null != proxyConfig.getClient().getUdp().getTransferLogEnable() && proxyConfig.getClient().getUdp().getTransferLogEnable()) { - ch.pipeline().addFirst(new LoggingHandler(UdpRealServerHandler.class)); - } - pipeline.addLast(udpWorkGroup, new UdpRealServerHandler()); - } - }); - return bootstrap; + return super.udpServerBootstrap(proxyConfig,udpServerGroup,udpWorkGroup); } @Bean public NeutrinoCoreRuntimeNativeRegistrar neutrinoCoreRuntimeNativeRegistrar() { - return new NeutrinoCoreRuntimeNativeRegistrar(); + return super.neutrinoCoreRuntimeNativeRegistrar(); } - } diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/constant/Constants.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/constant/Constants.java deleted file mode 100644 index f9ab40c3..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/constant/Constants.java +++ /dev/null @@ -1,13 +0,0 @@ -package org.dromara.neutrinoproxy.client.constant; - -import io.netty.util.AttributeKey; -import org.dromara.neutrinoproxy.client.util.UdpChannelBindInfo; - -/** - * @author: aoshiguchen - * @date: 2023/9/21 - */ -public interface Constants { - AttributeKey UDP_CHANNEL_BIND_KEY = AttributeKey.newInstance("udpChannelBindKey"); - -} 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 deleted file mode 100644 index 68db2a0e..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CmdChannelHandler.java +++ /dev/null @@ -1,85 +0,0 @@ -package org.dromara.neutrinoproxy.client.core; - -import org.dromara.neutrinoproxy.client.config.ProxyConfig; -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 { - private static volatile Boolean transferLogEnable = Boolean.FALSE; - - public CmdChannelHandler() { - ProxyConfig proxyConfig = Solon.context().getBean(ProxyConfig.class); - if (null != proxyConfig.getClient() && null != proxyConfig.getTunnel().getHeartbeatLogEnable()) { - transferLogEnable = proxyConfig.getTunnel().getHeartbeatLogEnable(); - } - } - - @Override - protected void channelRead0(ChannelHandlerContext ctx, ProxyMessage proxyMessage) throws Exception { - if (ProxyMessage.TYPE_HEARTBEAT != proxyMessage.getType() || transferLogEnable) { - log.debug("[CMD Channel]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("[CMD Channel]Client CmdChannel disconnect"); - ProxyUtil.setCmdChannel(null); - ProxyUtil.clearRealServerChannels(); - - super.channelInactive(ctx); - } - - @Override - public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { - log.error("[CMD Channel]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.error("[CMD Channel] Read timeout disconnect"); - ctx.channel().close(); - break; - case WRITER_IDLE: - ctx.channel().writeAndFlush(ProxyMessage.buildHeartbeatMessage()); - break; - case ALL_IDLE: - log.error("[CMD Channel] ReadWrite timeout disconnect"); - ctx.close(); - break; - } - } - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CustomThreadFactory.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CustomThreadFactory.java deleted file mode 100644 index b6a3e191..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/CustomThreadFactory.java +++ /dev/null @@ -1,55 +0,0 @@ -/** - * Copyright (c) 2022 aoshiguchen - * - * Permission is hereby granted, free of charge, to any person obtaining a copy - * of this software and associated documentation files (the "Software"), to deal - * in the Software without restriction, including without limitation the rights - * to use, copy, modify, merge, publish, distribute, sublicense, and/or sell - * copies of the Software, and to permit persons to whom the Software is - * furnished to do so, subject to the following conditions: - * - * The above copyright notice and this permission notice shall be included in all - * copies or substantial portions of the Software. - * - * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR - * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, - * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE - * AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER - * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, - * OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE - * SOFTWARE. - */ -package org.dromara.neutrinoproxy.client.core; - -import java.util.concurrent.ThreadFactory; -import java.util.concurrent.atomic.AtomicInteger; - -/** - * - * @author: aoshiguchen - * @date: 2022/9/4 - */ -public class CustomThreadFactory implements ThreadFactory { - private final ThreadGroup group; - private final AtomicInteger threadNumber = new AtomicInteger(1); - private final String namePrefix; - - public CustomThreadFactory(String prefix) { - SecurityManager s = System.getSecurityManager(); - group = (s != null) ? s.getThreadGroup() : - Thread.currentThread().getThreadGroup(); - namePrefix = prefix + "-thread-"; - } - - @Override - public Thread newThread(Runnable r) { - Thread t = new Thread(group, r, namePrefix + threadNumber.getAndIncrement(), 0); - if (t.isDaemon()) { - t.setDaemon(false); - } - if (t.getPriority() != Thread.NORM_PRIORITY) { - t.setPriority(Thread.NORM_PRIORITY); - } - return t; - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyChannelBorrowListener.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyChannelBorrowListener.java deleted file mode 100644 index 21e014fa..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyChannelBorrowListener.java +++ /dev/null @@ -1,38 +0,0 @@ -/** - * Copyright (c) 2022 aoshiguchen - * - * Permission is hereby granted, free of charge, to any person obtaining a copy - * of this software and associated documentation files (the "Software"), to deal - * in the Software without restriction, including without limitation the rights - * to use, copy, modify, merge, publish, distribute, sublicense, and/or sell - * copies of the Software, and to permit persons to whom the Software is - * furnished to do so, subject to the following conditions: - * - * The above copyright notice and this permission notice shall be included in all - * copies or substantial portions of the Software. - * - * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR - * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, - * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE - * AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER - * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, - * OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE - * SOFTWARE. - */ - -package org.dromara.neutrinoproxy.client.core; - -import io.netty.channel.Channel; - -/** - * - * @author: aoshiguchen - * @date: 2022/6/16 - */ -public interface ProxyChannelBorrowListener { - - void success(Channel channel); - - void error(Throwable cause); - -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyClientService.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyClientService.java deleted file mode 100644 index e7ff1e3f..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/ProxyClientService.java +++ /dev/null @@ -1,130 +0,0 @@ -package org.dromara.neutrinoproxy.client.core; - -import cn.hutool.core.util.StrUtil; -import org.dromara.neutrinoproxy.client.config.ProxyConfig; -import org.dromara.neutrinoproxy.client.util.ProxyUtil; -import org.dromara.neutrinoproxy.client.util.UdpServerUtil; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import io.netty.bootstrap.Bootstrap; -import io.netty.channel.*; -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 java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; - -/** - * 代理客户端服务 - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Slf4j -@Component -public class ProxyClientService { - @Inject - private ProxyConfig proxyConfig; - @Inject("cmdTunnelBootstrap") - private Bootstrap cmdTunnelBootstrap; - @Inject("udpServerBootstrap") - private Bootstrap udpServerBootstrap; - private volatile Channel channel; - /** - * 重连次数 - */ - private volatile int reconnectCount = 0; - /** - * 重连服务执行器 - */ - private static final ScheduledExecutorService reconnectExecutor = Executors.newSingleThreadScheduledExecutor(new CustomThreadFactory("ClientReconnect")); - - @Init - public void init() { - this.reconnectExecutor.scheduleWithFixedDelay(this::reconnect, 10, proxyConfig.getTunnel().getReconnection().getIntervalSeconds(), TimeUnit.SECONDS); - - try { - this.start(); - UdpServerUtil.initCache(proxyConfig, udpServerBootstrap); - } catch (Exception e) { - // 启动连不上也做一下重连,因此先catch异常 - log.error("[CmdChannel] start error", e); - } - } - - public void start() { - if (StrUtil.isEmpty(proxyConfig.getTunnel().getServerIp())) { - log.error("not found server-ip config."); - Solon.stop(); - return; - } - if (null == proxyConfig.getTunnel().getServerPort()) { - log.error("not found server-port config."); - Solon.stop(); - return; - } - if (null != proxyConfig.getTunnel().getSslEnable() && proxyConfig.getTunnel().getSslEnable() - && StrUtil.isEmpty(proxyConfig.getTunnel().getJksPath())) { - log.error("not found jks-path config."); - Solon.stop(); - return; - } - if (StrUtil.isEmpty(proxyConfig.getTunnel().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); - } - } else { - channel.writeAndFlush(ProxyMessage.buildAuthMessage(proxyConfig.getTunnel().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.getTunnel().getLicenseKey(), ProxyUtil.getClientId())); - log.info("[CmdChannel] connect proxy server success. channelId:{}", future.channel().id().asLongText()); - -// reconnectServiceEnable = true; - reconnectCount = 0; - } else { - log.info("[CmdChannel] connect proxy server failed!"); - } - } - }).sync(); - } - - protected synchronized void reconnect() { - if (null != channel) { - if (channel.isActive()) { - return; - } - channel.close(); - } - - log.info("[CmdChannel] client reconnect seq:{}", ++reconnectCount); - try { - connectProxyServer(); - } catch (Exception e) { - log.error("[CmdChannel] reconnect 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 deleted file mode 100644 index aafec97a..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/RealServerChannelHandler.java +++ /dev/null @@ -1,103 +0,0 @@ -/** - * Copyright (c) 2022 aoshiguchen - * - * Permission is hereby granted, free of charge, to any person obtaining a copy - * of this software and associated documentation files (the "Software"), to deal - * in the Software without restriction, including without limitation the rights - * to use, copy, modify, merge, publish, distribute, sublicense, and/or sell - * copies of the Software, and to permit persons to whom the Software is - * furnished to do so, subject to the following conditions: - * - * The above copyright notice and this permission notice shall be included in all - * copies or substantial portions of the Software. - * - * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR - * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, - * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE - * AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER - * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, - * OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE - * SOFTWARE. - */ - -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; -import io.netty.buffer.ByteBuf; -import io.netty.channel.Channel; -import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.ChannelOption; -import io.netty.channel.SimpleChannelInboundHandler; - -/** - * 处理与被代理客户端的数据传输 - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Slf4j -public class RealServerChannelHandler extends SimpleChannelInboundHandler { - - - @Override - protected void channelRead0(ChannelHandlerContext ctx, ByteBuf buf) throws Exception { - Channel realServerChannel = ctx.channel(); - Channel proxyChannel = realServerChannel.attr(Constants.NEXT_CHANNEL).get(); - if (null == proxyChannel) { - // 代理客户端连接断开 - ctx.channel().close(); - } else { - - if (proxyChannel.isWritable()) { - if (!realServerChannel.config().isAutoRead()) { - realServerChannel.config().setAutoRead(true); - } - } else { - if (realServerChannel.config().isAutoRead()) { - realServerChannel.config().setAutoRead(false); - } - } - - byte[] bytes = new byte[buf.readableBytes()]; - buf.readBytes(bytes); - String visitorId = ProxyUtil.getVisitorIdByRealServerChannel(realServerChannel); - proxyChannel.writeAndFlush(ProxyMessage.buildTransferMessage(visitorId, bytes)); - } - } - - @Override - public void channelActive(ChannelHandlerContext ctx) throws Exception { - super.channelActive(ctx); - } - - @Override - public void channelInactive(ChannelHandlerContext ctx) throws Exception { - Channel realServerChannel = ctx.channel(); - String visitorId = ProxyUtil.getVisitorIdByRealServerChannel(realServerChannel); - ProxyUtil.removeRealServerChannel(visitorId); - Channel channel = realServerChannel.attr(Constants.NEXT_CHANNEL).get(); - if (channel != null) { - channel.writeAndFlush(ProxyMessage.buildDisconnectMessage(visitorId)); - } - - super.channelInactive(ctx); - } - - @Override - public void channelWritabilityChanged(ChannelHandlerContext ctx) throws Exception { - Channel realServerChannel = ctx.channel(); - Channel proxyChannel = realServerChannel.attr(Constants.NEXT_CHANNEL).get(); - if (proxyChannel != null) { - proxyChannel.config().setOption(ChannelOption.AUTO_READ, realServerChannel.isWritable()); - } - - super.channelWritabilityChanged(ctx); - } - - @Override - public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { - log.error("Client ProxyChannel Error", cause); - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/TcpProxyChannelHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/TcpProxyChannelHandler.java deleted file mode 100644 index c3710e83..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/TcpProxyChannelHandler.java +++ /dev/null @@ -1,82 +0,0 @@ -package org.dromara.neutrinoproxy.client.core; - -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; - -/** - * 处理与服务端之间的数据传输 - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Slf4j -public class TcpProxyChannelHandler extends SimpleChannelInboundHandler { - - - @Override - protected void channelRead0(ChannelHandlerContext ctx, ProxyMessage proxyMessage) throws Exception { - if (ProxyMessage.TYPE_HEARTBEAT != proxyMessage.getType()) { - log.debug("[TCP Proxy Channel]Client ProxyChannel 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 { - // 数据传输连接 - Channel realServerChannel = ctx.channel().attr(Constants.NEXT_CHANNEL).get(); - if (realServerChannel != null && realServerChannel.isActive()) { - realServerChannel.close(); - } - - ProxyUtil.removeTcpProxyChanel(ctx.channel()); - super.channelInactive(ctx); - } - - @Override - public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { - log.error("[TCP Proxy Channel]Client ProxyChannel 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: - if (ctx.channel().isWritable()) { - // 读超时,断开连接 - log.info("[TCP Proxy Channel]Read timeout"); - ctx.channel().close(); - } - break; - case WRITER_IDLE: - ctx.channel().writeAndFlush(ProxyMessage.buildHeartbeatMessage()); - break; - case ALL_IDLE: -// log.debug("[TCP Proxy Channel]ReadWrite timeout"); -// ctx.close(); - break; - } - } - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpProxyChannelHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpProxyChannelHandler.java deleted file mode 100644 index 8504e20f..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpProxyChannelHandler.java +++ /dev/null @@ -1,80 +0,0 @@ -package org.dromara.neutrinoproxy.client.core; - -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; - -/** - * 处理与服务端之间的数据传输 - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Slf4j -public class UdpProxyChannelHandler extends SimpleChannelInboundHandler { - - - @Override - protected void channelRead0(ChannelHandlerContext ctx, ProxyMessage proxyMessage) throws Exception { - if (ProxyMessage.TYPE_HEARTBEAT != proxyMessage.getType()) { - log.debug("[UDP Proxy Channel]Client ProxyChannel 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 { - // 数据传输连接 - Channel realServerChannel = ctx.channel().attr(Constants.NEXT_CHANNEL).get(); - if (realServerChannel != null && realServerChannel.isActive()) { - realServerChannel.close(); - } - - ProxyUtil.removeTcpProxyChanel(ctx.channel()); - super.channelInactive(ctx); - } - - @Override - public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { - log.error("[UDP Proxy Channel]Client ProxyChannel 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("[UDP Proxy Channel]Read timeout"); - ctx.channel().close(); - break; - case WRITER_IDLE: - ctx.channel().writeAndFlush(ProxyMessage.buildHeartbeatMessage()); - break; - case ALL_IDLE: - log.debug("[UDP Proxy Channel]ReadWrite timeout"); - ctx.close(); - break; - } - } - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpRealServerHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpRealServerHandler.java deleted file mode 100644 index 1a288e88..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpRealServerHandler.java +++ /dev/null @@ -1,41 +0,0 @@ -package org.dromara.neutrinoproxy.client.core; - -import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.SimpleChannelInboundHandler; -import io.netty.channel.socket.DatagramPacket; -import lombok.extern.slf4j.Slf4j; -import org.dromara.neutrinoproxy.client.constant.Constants; -import org.dromara.neutrinoproxy.client.util.UdpChannelBindInfo; -import org.dromara.neutrinoproxy.core.ProxyMessage; - -import java.net.InetSocketAddress; - -/** - * @author: aoshiguchen - * @date: 2023/9/21 - */ -@Slf4j -public class UdpRealServerHandler extends SimpleChannelInboundHandler { - - @Override - protected void channelRead0(ChannelHandlerContext ctx, DatagramPacket datagramPacket) throws Exception { - log.debug("chid---<:{} port:{}", ctx.channel().id().asLongText(), ((InetSocketAddress)ctx.channel().localAddress()).getPort()); - UdpChannelBindInfo udpChannelBindInfo = ctx.channel().attr(Constants.UDP_CHANNEL_BIND_KEY).get(); - if (null != udpChannelBindInfo) { - byte[] bytes = new byte[datagramPacket.content().readableBytes()]; - datagramPacket.content().readBytes(bytes); - - udpChannelBindInfo.getTunnelChannel().writeAndFlush(ProxyMessage.buildUdpTransferMessage(new ProxyMessage.UdpBaseInfo() - .setVisitorId(udpChannelBindInfo.getVisitorId()) - .setVisitorIp(udpChannelBindInfo.getVisitorIp()) - .setVisitorPort(udpChannelBindInfo.getVisitorPort()) - .setServerPort(udpChannelBindInfo.getServerPort()) - .setTargetIp(udpChannelBindInfo.getTargetIp()) - .setTargetPort(udpChannelBindInfo.getTargetPort())) - .setData(bytes) - ); - - udpChannelBindInfo.getLockChannel().setResponseCount(udpChannelBindInfo.getLockChannel().getResponseCount() + 1); - } - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/BeanHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/BeanHandler.java new file mode 100644 index 00000000..e6c2debd --- /dev/null +++ b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/BeanHandler.java @@ -0,0 +1,19 @@ +package org.dromara.neutrinoproxy.client.handler; + +import org.dromara.neutrinoproxy.client.sdk.config.IBeanHandler; +import org.dromara.neutrinoproxy.client.sdk.config.ProxyConfig; +import org.dromara.neutrinoproxy.core.dispatcher.Dispatcher; +import org.noear.solon.Solon; + + +public class BeanHandler implements IBeanHandler { + @Override + public Dispatcher getDispatcher(){ + return Solon.context().getBean(Dispatcher.class); + } + + @Override + public ProxyConfig getProxyConfig() { + return Solon.context().getBean(ProxyConfig.class); + } +} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageAuthHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageAuthHandler.java deleted file mode 100644 index ff41bc10..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageAuthHandler.java +++ /dev/null @@ -1,46 +0,0 @@ -package org.dromara.neutrinoproxy.client.handler; - -import io.netty.channel.ChannelHandlerContext; -import lombok.extern.slf4j.Slf4j; -import org.dromara.neutrinoproxy.client.config.ProxyConfig; -import org.dromara.neutrinoproxy.core.Constants; -import org.dromara.neutrinoproxy.core.ExceptionEnum; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.dromara.neutrinoproxy.core.ProxyMessageHandler; -import org.dromara.neutrinoproxy.core.dispatcher.Match; -import org.noear.snack.ONode; -import org.noear.solon.Solon; -import org.noear.solon.annotation.Component; -import org.noear.solon.annotation.Inject; - -/** - * 认证信息处理器 - * @author: aoshiguchen - * @date: 2022/9/4 - */ -@Slf4j -@Match(type = Constants.ProxyDataTypeName.AUTH) -@Component -public class ProxyMessageAuthHandler implements ProxyMessageHandler { - @Inject - private ProxyConfig proxyConfig; - @Override - public void handle(ChannelHandlerContext context, ProxyMessage proxyMessage) { - String info = proxyMessage.getInfo(); - ONode load = ONode.load(info); - Integer code = load.get("code").getInt(); - log.info("Auth result:{}", info); - if (ExceptionEnum.AUTH_FAILED.getCode().equals(code)) { - // 客户端认证失败,直接停止服务 - log.info("client auth failed , client stop."); - context.channel().close(); - if (!proxyConfig.getTunnel().getReconnection().getUnlimited()) { - Solon.stop(); - } - } else if (ExceptionEnum.CONNECT_FAILED.getCode().equals(code) || - ExceptionEnum.LICENSE_CANNOT_REPEAT_CONNECT.getCode().equals(code) - ){ - context.channel().close(); - } - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageConnectHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageConnectHandler.java deleted file mode 100644 index 2bca8667..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageConnectHandler.java +++ /dev/null @@ -1,87 +0,0 @@ -package org.dromara.neutrinoproxy.client.handler; - -import org.dromara.neutrinoproxy.client.config.ProxyConfig; -import org.dromara.neutrinoproxy.client.core.ProxyChannelBorrowListener; -import org.dromara.neutrinoproxy.client.util.ProxyUtil; -import org.dromara.neutrinoproxy.core.Constants; -import org.dromara.neutrinoproxy.core.ProxyDataTypeEnum; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.dromara.neutrinoproxy.core.ProxyMessageHandler; -import org.dromara.neutrinoproxy.core.dispatcher.Match; -import io.netty.bootstrap.Bootstrap; -import io.netty.channel.*; -import org.noear.solon.annotation.Component; -import org.noear.solon.annotation.Inject; - -/** - * 连接信息处理器 - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Match(type = Constants.ProxyDataTypeName.CONNECT) -@Component -public class ProxyMessageConnectHandler implements ProxyMessageHandler { - @Inject("tcpProxyTunnelBootstrap") - private Bootstrap tcpProxyTunnelBootstrap; - @Inject("realServerBootstrap") - private Bootstrap realServerBootstrap; - @Inject - private ProxyConfig proxyConfig; - - @Override - public void handle(ChannelHandlerContext ctx, ProxyMessage proxyMessage) { - final Channel cmdChannel = ctx.channel(); - final String visitorId = proxyMessage.getInfo(); - String[] serverInfo = new String(proxyMessage.getData()).split(":"); - String ip = serverInfo[0]; - int port = Integer.parseInt(serverInfo[1]); - // 连接真实的、被代理的服务 - realServerBootstrap.connect(ip, port).addListener(new ChannelFutureListener() { - - @Override - public void operationComplete(ChannelFuture future) throws Exception { - - // 连接后端服务器成功 - if (future.isSuccess()) { - final Channel realServerChannel = future.channel(); - - realServerChannel.config().setOption(ChannelOption.AUTO_READ, false); - - // 获取连接 - ProxyUtil.borrowTcpProxyChanel(tcpProxyTunnelBootstrap, new ProxyChannelBorrowListener() { - - @Override - public void success(Channel channel) { - // 连接绑定 - channel.attr(Constants.NEXT_CHANNEL).set(realServerChannel); - realServerChannel.attr(Constants.NEXT_CHANNEL).set(channel); - - // 远程绑定 - channel.writeAndFlush(ProxyMessage.buildConnectMessage(visitorId + "@" + proxyConfig.getTunnel().getLicenseKey())); - - realServerChannel.config().setOption(ChannelOption.AUTO_READ, true); - ProxyUtil.addRealServerChannel(visitorId, realServerChannel); - ProxyUtil.setRealServerChannelVisitorId(realServerChannel, visitorId); - } - - @Override - public void error(Throwable cause) { - ProxyMessage proxyMessage = new ProxyMessage(); - proxyMessage.setType(ProxyMessage.TYPE_DISCONNECT); - proxyMessage.setInfo(visitorId); - cmdChannel.writeAndFlush(proxyMessage); - } - }); - - } else { - cmdChannel.writeAndFlush(ProxyMessage.buildDisconnectMessage(visitorId)); - } - } - }); - } - - @Override - public String name() { - return ProxyDataTypeEnum.CONNECT.getDesc(); - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageDisconnectHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageDisconnectHandler.java deleted file mode 100644 index d1ada2d7..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageDisconnectHandler.java +++ /dev/null @@ -1,40 +0,0 @@ -package org.dromara.neutrinoproxy.client.handler; - -import org.dromara.neutrinoproxy.client.util.ProxyUtil; -import org.dromara.neutrinoproxy.core.Constants; -import org.dromara.neutrinoproxy.core.ProxyDataTypeEnum; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.dromara.neutrinoproxy.core.ProxyMessageHandler; -import org.dromara.neutrinoproxy.core.dispatcher.Match; -import io.netty.buffer.Unpooled; -import io.netty.channel.Channel; -import io.netty.channel.ChannelFutureListener; -import io.netty.channel.ChannelHandlerContext; -import org.noear.solon.annotation.Component; - -/** - * 断开连接信息处理器 - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Match(type = Constants.ProxyDataTypeName.DISCONNECT) -@Component -public class ProxyMessageDisconnectHandler implements ProxyMessageHandler { - - @Override - public void handle(ChannelHandlerContext ctx, ProxyMessage proxyMessage) { - Channel realServerChannel = ctx.channel().attr(Constants.NEXT_CHANNEL).get(); - if (null != realServerChannel) { - ctx.channel().attr(Constants.NEXT_CHANNEL).remove(); - ProxyUtil.returnTcpProxyChanel(ctx.channel()); - realServerChannel.writeAndFlush(Unpooled.EMPTY_BUFFER).addListener(ChannelFutureListener.CLOSE); - } - ctx.close(); - } - - @Override - public String name() { - return ProxyDataTypeEnum.DISCONNECT.getDesc(); - } - -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageErrorHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageErrorHandler.java deleted file mode 100644 index 8d986f8e..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageErrorHandler.java +++ /dev/null @@ -1,38 +0,0 @@ -package org.dromara.neutrinoproxy.client.handler; - -import io.netty.channel.ChannelHandlerContext; -import lombok.extern.slf4j.Slf4j; -import org.dromara.neutrinoproxy.core.Constants; -import org.dromara.neutrinoproxy.core.ExceptionEnum; -import org.dromara.neutrinoproxy.core.ProxyDataTypeEnum; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.dromara.neutrinoproxy.core.ProxyMessageHandler; -import org.dromara.neutrinoproxy.core.dispatcher.Match; -import org.noear.snack.ONode; -import org.noear.solon.annotation.Component; - -/** - * 异常信息处理器 - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Slf4j -@Match(type = Constants.ProxyDataTypeName.ERROR) -@Component -public class ProxyMessageErrorHandler implements ProxyMessageHandler { - - @Override - public void handle(ChannelHandlerContext ctx, ProxyMessage proxyMessage) { - log.info("error: {}", proxyMessage.getInfo()); - ONode load = ONode.load(proxyMessage.getInfo()); - Integer code = load.get("code").getInt(); - if (ExceptionEnum.AUTH_FAILED.getCode().equals(code)) { - System.exit(0); - } - } - - @Override - public String name() { - return ProxyDataTypeEnum.DISCONNECT.getDesc(); - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageTransferHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageTransferHandler.java deleted file mode 100644 index 60d51a95..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/ProxyMessageTransferHandler.java +++ /dev/null @@ -1,41 +0,0 @@ -package org.dromara.neutrinoproxy.client.handler; - -import org.dromara.neutrinoproxy.core.Constants; -import org.dromara.neutrinoproxy.core.ProxyDataTypeEnum; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.dromara.neutrinoproxy.core.ProxyMessageHandler; -import org.dromara.neutrinoproxy.core.dispatcher.Match; -import io.netty.buffer.ByteBuf; -import io.netty.channel.Channel; -import io.netty.channel.ChannelHandlerContext; -import org.noear.solon.annotation.Component; - -/** - * 传输信息处理器 - * @author: aoshiguchen - * @date: 2022/6/16 - */ -@Match(type = Constants.ProxyDataTypeName.TRANSFER) -@Component -public class ProxyMessageTransferHandler implements ProxyMessageHandler { - - @Override - public void handle(ChannelHandlerContext ctx, ProxyMessage proxyMessage) { - Channel realServerChannel = ctx.channel().attr(Constants.NEXT_CHANNEL).get(); - if (realServerChannel != null) { - - // 自己可写,则设置来源可读。自己不可写,则设置来源不可读 - ctx.channel().config().setAutoRead(realServerChannel.isWritable()); - - ByteBuf buf = ctx.alloc().buffer(proxyMessage.getData().length); - buf.writeBytes(proxyMessage.getData()); - realServerChannel.writeAndFlush(buf); - } - } - - @Override - public String name() { - return ProxyDataTypeEnum.TRANSFER.getDesc(); - } - -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageConnectHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageConnectHandler.java deleted file mode 100644 index bc0e9bb0..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageConnectHandler.java +++ /dev/null @@ -1,64 +0,0 @@ -package org.dromara.neutrinoproxy.client.handler; - -import io.netty.bootstrap.Bootstrap; -import io.netty.channel.Channel; -import io.netty.channel.ChannelHandlerContext; -import lombok.extern.slf4j.Slf4j; -import org.dromara.neutrinoproxy.client.config.ProxyConfig; -import org.dromara.neutrinoproxy.client.core.ProxyChannelBorrowListener; -import org.dromara.neutrinoproxy.client.util.ProxyUtil; -import org.dromara.neutrinoproxy.core.Constants; -import org.dromara.neutrinoproxy.core.ProxyDataTypeEnum; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.dromara.neutrinoproxy.core.ProxyMessageHandler; -import org.dromara.neutrinoproxy.core.dispatcher.Match; -import org.noear.snack.ONode; -import org.noear.solon.annotation.Component; -import org.noear.solon.annotation.Inject; - -/** - * @author: aoshiguchen - * @date: 2023/9/19 - */ -@Slf4j -@Match(type = Constants.ProxyDataTypeName.UDP_CONNECT) -@Component -public class UdpProxyMessageConnectHandler implements ProxyMessageHandler { - @Inject - private ProxyConfig proxyConfig; - @Inject("udpProxyTunnelBootstrap") - private Bootstrap udpProxyTunnelBootstrap; - - @Override - public void handle(ChannelHandlerContext ctx, ProxyMessage proxyMessage) { - final Channel cmdChannel = ctx.channel(); - final ProxyMessage.UdpBaseInfo udpBaseInfo = ONode.deserialize(proxyMessage.getInfo(), ProxyMessage.UdpBaseInfo.class); - log.info("[UDP connect]info:{}", proxyMessage.getInfo()); - - // 获取连接 - ProxyUtil.borrowUdpProxyChanel(udpProxyTunnelBootstrap, new ProxyChannelBorrowListener() { - - @Override - public void success(Channel channel) { - channel.writeAndFlush(ProxyMessage.buildUdpConnectMessage(new ProxyMessage.UdpBaseInfo() - .setVisitorId(udpBaseInfo.getVisitorId()) - .setServerPort(udpBaseInfo.getServerPort()) - .setTargetIp(udpBaseInfo.getTargetIp()) - .setTargetPort(udpBaseInfo.getTargetPort()) - ).setData(proxyConfig.getTunnel().getLicenseKey().getBytes())); - } - - @Override - public void error(Throwable cause) { - cmdChannel.writeAndFlush(ProxyMessage.buildDisconnectMessage(udpBaseInfo.toJsonString())); - } - }); - - - } - - @Override - public String name() { - return ProxyDataTypeEnum.UDP_CONNECT.getDesc(); - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageTransferHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageTransferHandler.java deleted file mode 100644 index a3d111d5..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageTransferHandler.java +++ /dev/null @@ -1,48 +0,0 @@ -package org.dromara.neutrinoproxy.client.handler; - -import io.netty.buffer.ByteBuf; -import io.netty.buffer.Unpooled; -import io.netty.channel.Channel; -import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.socket.DatagramPacket; -import lombok.extern.slf4j.Slf4j; -import org.dromara.neutrinoproxy.client.util.UdpServerUtil; -import org.dromara.neutrinoproxy.core.Constants; -import org.dromara.neutrinoproxy.core.ProxyDataTypeEnum; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.dromara.neutrinoproxy.core.ProxyMessageHandler; -import org.dromara.neutrinoproxy.core.dispatcher.Match; -import org.noear.snack.ONode; -import org.noear.solon.annotation.Component; - -import java.net.InetSocketAddress; - -/** - * @author: aoshiguchen - * @date: 2023/9/20 - */ -@Slf4j -@Match(type = Constants.ProxyDataTypeName.UDP_TRANSFER) -@Component -public class UdpProxyMessageTransferHandler implements ProxyMessageHandler { - - @Override - public void handle(ChannelHandlerContext ctx, ProxyMessage proxyMessage) { - final ProxyMessage.UdpBaseInfo udpBaseInfo = ONode.deserialize(proxyMessage.getInfo(), ProxyMessage.UdpBaseInfo.class); - log.debug("[UDP transfer]info:{} data:{}", proxyMessage.getInfo(), new String(proxyMessage.getData())); - Channel channel = UdpServerUtil.takeChannel(udpBaseInfo, ctx.channel()); - if (null == channel) { - log.error("[UDP transfer] take udp channel failed."); - return; - } - log.debug("chid--->:{} port:{}", ctx.channel().id().asLongText(), ((InetSocketAddress)channel.localAddress()).getPort()); - InetSocketAddress address = new InetSocketAddress(udpBaseInfo.getTargetIp(), udpBaseInfo.getTargetPort()); - ByteBuf byteBuf = Unpooled.copiedBuffer(proxyMessage.getData()); - channel.writeAndFlush(new DatagramPacket(byteBuf, address)); - } - - @Override - public String name() { - return ProxyDataTypeEnum.UDP_TRANSFER.getDesc(); - } -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/LockChannel.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/LockChannel.java deleted file mode 100644 index d694ba83..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/LockChannel.java +++ /dev/null @@ -1,28 +0,0 @@ -package org.dromara.neutrinoproxy.client.util; - -import io.netty.channel.Channel; -import lombok.Data; -import lombok.experimental.Accessors; - -import java.util.Date; - -/** - * @author: aoshiguchen - * @date: 2023/9/21 - */ -@Accessors(chain = true) -@Data -public class LockChannel { - // 端口号 - private int port; - // 通道 - private Channel channel; - // 期望的响应次数 - private int proxyResponses; - // 超时时间(毫秒) - private long proxyTimeoutMs; - // 被获取的时间 - private Date takeTime; - // 已经响应的次数 - private int responseCount; -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/ProxyUtil.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/ProxyUtil.java deleted file mode 100644 index b44bd4a9..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/ProxyUtil.java +++ /dev/null @@ -1,226 +0,0 @@ -/** - * Copyright (c) 2022 aoshiguchen - * - * Permission is hereby granted, free of charge, to any person obtaining a copy - * of this software and associated documentation files (the "Software"), to deal - * in the Software without restriction, including without limitation the rights - * to use, copy, modify, merge, publish, distribute, sublicense, and/or sell - * copies of the Software, and to permit persons to whom the Software is - * furnished to do so, subject to the following conditions: - * - * The above copyright notice and this permission notice shall be included in all - * copies or substantial portions of the Software. - * - * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR - * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, - * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE - * AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER - * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, - * OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE - * SOFTWARE. - */ -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; -import org.dromara.neutrinoproxy.core.Constants; -import io.netty.bootstrap.Bootstrap; -import io.netty.buffer.Unpooled; -import io.netty.channel.Channel; -import io.netty.channel.ChannelFutureListener; -import io.netty.channel.ChannelOption; -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; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentLinkedQueue; - -/** - * 代理工具 - * @author: aoshiguchen - * @date: 2022/8/31 - */ -@Slf4j -public class ProxyUtil { - private static final AttributeKey USER_CHANNEL_WRITEABLE = AttributeKey.newInstance("user_channel_writeable"); - - private static final AttributeKey CLIENT_CHANNEL_WRITEABLE = AttributeKey.newInstance("client_channel_writeable"); - - private static final int MAX_POOL_SIZE = 100; - - private static Map realServerChannels = new ConcurrentHashMap(); - - private static ConcurrentLinkedQueue tcpProxyChannelPool = new ConcurrentLinkedQueue(); - private static ConcurrentLinkedQueue udpProxyChannelPool = new ConcurrentLinkedQueue<>(); - - private static volatile Channel cmdChannel; - - private static String clientId; - private static final String CLIENT_ID_FILE = ".NEUTRINO_PROXY_CLIENT_ID"; - - public static void borrowTcpProxyChanel(Bootstrap tcpProxyTunnelBootstrap, final ProxyChannelBorrowListener borrowListener) { - Channel channel = tcpProxyChannelPool.poll(); - if (null != channel) { - borrowListener.success(channel); - return; - } - - tcpProxyTunnelBootstrap.connect().addListener((ChannelFutureListener) future -> { - if (future.isSuccess()) { - borrowListener.success(future.channel()); - } else { - borrowListener.error(future.cause()); - } - }); - } - - public static void returnTcpProxyChanel(Channel proxyChanel) { - if (tcpProxyChannelPool.size() > MAX_POOL_SIZE) { - proxyChanel.close(); - } else { - proxyChanel.config().setOption(ChannelOption.AUTO_READ, true); - proxyChanel.attr(Constants.NEXT_CHANNEL).remove(); - tcpProxyChannelPool.offer(proxyChanel); - } - } - - - - public static void removeTcpProxyChanel(Channel proxyChanel) { - tcpProxyChannelPool.remove(proxyChanel); - } - - public static void borrowUdpProxyChanel(Bootstrap tcpProxyTunnelBootstrap, final ProxyChannelBorrowListener borrowListener) { - Channel channel = udpProxyChannelPool.poll(); - if (null != channel) { - borrowListener.success(channel); - return; - } - - tcpProxyTunnelBootstrap.connect().addListener((ChannelFutureListener) future -> { - if (future.isSuccess()) { - borrowListener.success(future.channel()); - } else { - borrowListener.error(future.cause()); - } - }); - } - - public static void returnUdpProxyChanel(Channel proxyChanel) { - if (udpProxyChannelPool.size() > MAX_POOL_SIZE) { - proxyChanel.close(); - } else { - proxyChanel.config().setOption(ChannelOption.AUTO_READ, true); - proxyChanel.attr(Constants.NEXT_CHANNEL).remove(); - udpProxyChannelPool.offer(proxyChanel); - } - } - - - - public static void removeUdpProxyChanel(Channel proxyChanel) { - udpProxyChannelPool.remove(proxyChanel); - } - - public static void setCmdChannel(Channel cmdChannel) { - ProxyUtil.cmdChannel = cmdChannel; - } - - public static Channel getCmdChannel() { - return cmdChannel; - } - - public static void setRealServerChannelVisitorId(Channel realServerChannel, String visitorId) { - realServerChannel.attr(Constants.VISITOR_ID).set(visitorId); - } - - public static String getVisitorIdByRealServerChannel(Channel realServerChannel) { - return realServerChannel.attr(Constants.VISITOR_ID).get(); - } - - public static Channel getRealServerChannel(String userId) { - return realServerChannels.get(userId); - } - - public static void addRealServerChannel(String userId, Channel realServerChannel) { - realServerChannels.put(userId, realServerChannel); - } - - public static Channel removeRealServerChannel(String userId) { - return realServerChannels.remove(userId); - } - - public static boolean isRealServerReadable(Channel realServerChannel) { - return realServerChannel.attr(CLIENT_CHANNEL_WRITEABLE).get() && realServerChannel.attr(USER_CHANNEL_WRITEABLE).get(); - } - - public static void clearRealServerChannels() { - Iterator> ite = realServerChannels.entrySet().iterator(); - while (ite.hasNext()) { - Channel realServerChannel = ite.next().getValue(); - if (realServerChannel.isActive()) { - realServerChannel.writeAndFlush(Unpooled.EMPTY_BUFFER).addListener(ChannelFutureListener.CLOSE); - } - } - - realServerChannels.clear(); - } - - public static String getClientId() { - if (StringUtils.isNotBlank(clientId)) { - return clientId; - } - ProxyConfig proxyConfig = Solon.context().getBean(ProxyConfig.class); - if (StringUtils.isNotBlank(proxyConfig.getTunnel().getClientId())) { - clientId = proxyConfig.getTunnel().getClientId(); - return clientId; - } - String id = FileUtil.readContentAsString(CLIENT_ID_FILE); - if (StringUtils.isNotBlank(id)) { - clientId = id; - return id; - } - id = UUID.randomUUID().toString().replace("-", ""); - FileUtil.write(CLIENT_ID_FILE, id); - 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("create SSL handler failed", e); - e.printStackTrace(); - } - return null; - } - -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpChannelBindInfo.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpChannelBindInfo.java deleted file mode 100644 index fea6aa86..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpChannelBindInfo.java +++ /dev/null @@ -1,22 +0,0 @@ -package org.dromara.neutrinoproxy.client.util; - -import io.netty.channel.Channel; -import lombok.Data; -import lombok.experimental.Accessors; - -/** - * @author: aoshiguchen - * @date: 2023/9/21 - */ -@Accessors(chain = true) -@Data -public class UdpChannelBindInfo { - private Channel tunnelChannel; - private LockChannel lockChannel; - private String visitorId; - private String visitorIp; - private int visitorPort; - private int serverPort; - private String targetIp; - private int targetPort; -} diff --git a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpServerUtil.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpServerUtil.java deleted file mode 100644 index 02f47437..00000000 --- a/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpServerUtil.java +++ /dev/null @@ -1,201 +0,0 @@ -package org.dromara.neutrinoproxy.client.util; - -import io.netty.bootstrap.Bootstrap; -import io.netty.buffer.Unpooled; -import io.netty.channel.Channel; -import io.netty.channel.ChannelFuture; -import io.netty.channel.ChannelFutureListener; -import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.StringUtils; -import org.dromara.neutrinoproxy.client.config.ProxyConfig; -import org.dromara.neutrinoproxy.client.constant.Constants; -import org.dromara.neutrinoproxy.client.core.CustomThreadFactory; -import org.dromara.neutrinoproxy.core.ProxyMessage; -import org.noear.solon.core.runtime.NativeDetector; - -import java.util.*; -import java.util.concurrent.*; - -/** - * @author: aoshiguchen - * @date: 2023/9/21 - */ -@Slf4j -public class UdpServerUtil { - private static final Boolean isSupportUdp = Boolean.FALSE; - private static int udpServerPortMin = 0; - private static int udpServerPortMax = 0; - private static int nextUdpServerPort = 0; - private static Bootstrap udpServerBootstrap; - private static int defaultUdpServerPort; - private static Channel defaultUdpServerChannel; - private static Map portToChannelMap = new ConcurrentHashMap<>(); - /** - * udp服务空闲端口池 - */ - private static ConcurrentLinkedQueue udpServerFreePortPool = new ConcurrentLinkedQueue<>(); - private static List lockChannelList = new ArrayList<>(); - /** - * lockChannel扫描器 - */ - private static final ScheduledExecutorService lockChannelScanner = Executors.newSingleThreadScheduledExecutor(new CustomThreadFactory("lockChannelScanner")); - - /** - * 初始化UDP缓存 - * 1、初始化一个基础UDP服务,用于不需要响应的UDP转发 - * 2、维护一个UDP服务池,用于需要响应的UDP转发 - * @param proxyConfig - */ - public static void initCache(ProxyConfig proxyConfig, Bootstrap udpServerBootstrap) { - // aot 阶段,不初始化UDP服务 - if (NativeDetector.isAotRuntime()) { - return; - } - if (null == proxyConfig.getClient().getUdp() || StringUtils.isEmpty(proxyConfig.getClient().getUdp().getPuppetPortRange())) { - return; - } - ProxyConfig.Udp udpConfig = proxyConfig.getClient().getUdp(); - if (StringUtils.isEmpty(udpConfig.getPuppetPortRange())) { - return; - } - String[] tmp = udpConfig.getPuppetPortRange().split("-"); - if (null == tmp || tmp.length != 2) { - log.error("client udp config error!"); - return; - } - try { - udpServerPortMin = Integer.parseInt(tmp[0]); - udpServerPortMax = Integer.parseInt(tmp[1]); - if (udpServerPortMax <= udpServerPortMin) { - // 至少得给1个udp端口,一个用于基础无响应UDP转发 - throw new RuntimeException("client udp config error!"); - } - nextUdpServerPort = udpServerPortMin; - UdpServerUtil.udpServerBootstrap = udpServerBootstrap; - log.info("udp proxy server port: {} ~ {}", udpServerPortMin, udpServerPortMax); - // 初始化udp服务 - initUdpServer(); - // 初始化lockChannel扫描器 - lockChannelScanner.scheduleWithFixedDelay(UdpServerUtil::lockChannelScan, 5, 3, TimeUnit.SECONDS); - } catch (Exception e) { - log.error("client udp config error!", e); - return; - } - } - - /** - * 初始化udp服务 - */ - private static void initUdpServer() { - defaultUdpServerPort = nextUdpServerPort(); - defaultUdpServerChannel = bindPort(defaultUdpServerPort); - // 初始化默认最多额外开启5个udp服务,其他的需要时再启动 - for (int i = 0; i < 5; i++) { - if (!hasNextUdpServerPort()) { - return; - } - int port = nextUdpServerPort(); - Channel ch = bindPort(port); - portToChannelMap.put(port, ch); - udpServerFreePortPool.offer(port); - } - while (hasNextUdpServerPort()) { - udpServerFreePortPool.offer(nextUdpServerPort()); - } - } - - private static Channel bindPort(int port) { - try { - ChannelFuture channelFuture = udpServerBootstrap.bind(port).sync(); - log.info("[udp server] bind port:{} success!", port); - return channelFuture.channel(); - } catch (InterruptedException e) { - log.error("[udp server] bind port:{} error!", port); - throw new RuntimeException(e); - } - } - - public static Boolean hasNextUdpServerPort() { - return nextUdpServerPort <= udpServerPortMax; - } - - public static synchronized int nextUdpServerPort() { - return nextUdpServerPort++; - } - - /** - * 获取一个可用的udp通道 - * 1、如果期待的响应为0,或者超时时间<=0,则认为不需要响应,直接返回默认的udp服务,否则继续下一步 - * 2、从可用端口队列中找到一个可用端口,若不存在可用端口,则降级为不需要响应,返回默认的udp服务。否则继续下一步 - * 3、根据该端口找到udp服务通道,找不到则绑定端口开启一个通道并返回。将该端口添加到锁定列表 - * 4、维护一个定时器的,定时扫描锁定列表,及时释放锁定的端口 - * @param info - * @return - */ - public static synchronized Channel takeChannel(ProxyMessage.UdpBaseInfo info, Channel tunnelChannel) { - if (info.getProxyResponses() <= 0 || info.getProxyTimeoutMs() <= 0) { - return defaultUdpServerChannel; - } - Integer port = udpServerFreePortPool.poll(); - if (null == port) { - return defaultUdpServerChannel; - } - Channel channel = portToChannelMap.get(port); - if (null == channel) { - channel = bindPort(port); - portToChannelMap.put(port, channel); - } - // 添加到锁定队列 - LockChannel lockChannel = new LockChannel() - .setPort(port) - .setChannel(channel) - .setProxyResponses(info.getProxyResponses()) - .setProxyTimeoutMs(info.getProxyTimeoutMs()) - .setTakeTime(new Date()) - .setResponseCount(0); - lockChannelList.add(lockChannel); - channel.attr(Constants.UDP_CHANNEL_BIND_KEY).set(new UdpChannelBindInfo() - .setTunnelChannel(tunnelChannel) - .setVisitorId(info.getVisitorId()) - .setVisitorIp(info.getVisitorIp()) - .setVisitorPort(info.getVisitorPort()) - .setServerPort(info.getServerPort()) - .setTargetIp(info.getTargetIp()) - .setTargetPort(info.getTargetPort()) - .setLockChannel(lockChannel) - ); - return channel; - } - - /** - * lockChannel扫描 - */ - public static synchronized void lockChannelScan() { - if (lockChannelList.isEmpty()) { - return; - } - Iterator iter = lockChannelList.iterator(); - if (iter.hasNext()) { - LockChannel lockChannel = iter.next(); - if (lockChannel.getResponseCount() >= lockChannel.getProxyResponses() || - System.currentTimeMillis() - lockChannel.getTakeTime().getTime() >= lockChannel.getProxyTimeoutMs() - ) { - iter.remove(); - UdpChannelBindInfo udpChannelBindInfo = lockChannel.getChannel().attr(Constants.UDP_CHANNEL_BIND_KEY).get(); - // 此处必须释放代理隧道 - closeChannel(udpChannelBindInfo.getTunnelChannel()); - lockChannel.getChannel().attr(Constants.UDP_CHANNEL_BIND_KEY).set(null); - udpServerFreePortPool.offer(lockChannel.getPort()); - log.debug("[udp channel]release udp channel port:{}", lockChannel.getPort()); - } - } - } - - private static void closeChannel(Channel channel) { - try { - channel.writeAndFlush(Unpooled.EMPTY_BUFFER).addListener(ChannelFutureListener.CLOSE); - } catch (Exception e) { - // ignore - } - } -}