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 new file mode 100644 index 00000000..f9ab40c3 --- /dev/null +++ b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/constant/Constants.java @@ -0,0 +1,13 @@ +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/UdpRealServerHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/core/UdpRealServerHandler.java index 33fb34cb..1a288e88 100644 --- 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 @@ -1,17 +1,41 @@ package org.dromara.neutrinoproxy.client.core; -import io.netty.buffer.ByteBuf; 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 */ -public class UdpRealServerHandler extends SimpleChannelInboundHandler { +@Slf4j +public class UdpRealServerHandler extends SimpleChannelInboundHandler { @Override - protected void channelRead0(ChannelHandlerContext channelHandlerContext, ByteBuf byteBuf) throws Exception { + 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/UdpProxyMessageTransferHandler.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/handler/UdpProxyMessageTransferHandler.java index 308e4551..2ff61e4e 100644 --- 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 @@ -1,8 +1,13 @@ package org.dromara.neutrinoproxy.client.handler; import com.alibaba.fastjson.JSONObject; +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; @@ -10,6 +15,9 @@ import org.dromara.neutrinoproxy.core.ProxyMessageHandler; import org.dromara.neutrinoproxy.core.dispatcher.Match; import org.noear.solon.annotation.Component; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; + /** * @author: aoshiguchen * @date: 2023/9/20 @@ -19,11 +27,19 @@ import org.noear.solon.annotation.Component; @Component public class UdpProxyMessageTransferHandler implements ProxyMessageHandler { - @Override public void handle(ChannelHandlerContext ctx, ProxyMessage proxyMessage) { final ProxyMessage.UdpBaseInfo udpBaseInfo = JSONObject.parseObject(proxyMessage.getInfo(), ProxyMessage.UdpBaseInfo.class); - log.info("[UDP transfer]info:{} data:{}", proxyMessage.getInfo(), new String(proxyMessage.getData())); + 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 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 new file mode 100644 index 00000000..d694ba83 --- /dev/null +++ b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/LockChannel.java @@ -0,0 +1,28 @@ +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/UdpChannelBindInfo.java b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpChannelBindInfo.java new file mode 100644 index 00000000..fea6aa86 --- /dev/null +++ b/neutrino-proxy-client/src/main/java/org/dromara/neutrinoproxy/client/util/UdpChannelBindInfo.java @@ -0,0 +1,22 @@ +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 index ad3fc281..48b833c0 100644 --- 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 @@ -1,9 +1,17 @@ package org.dromara.neutrinoproxy.client.util; import io.netty.bootstrap.Bootstrap; +import io.netty.channel.Channel; +import io.netty.channel.ChannelFuture; 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 java.util.*; +import java.util.concurrent.*; /** * @author: aoshiguchen @@ -16,7 +24,18 @@ public class UdpServerUtil { private static int udpServerPortMax = 0; private static int nextUdpServerPort = 0; private static Bootstrap udpServerBootstrap; - private static final String defaultUdpServerKey = "default"; + 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缓存 @@ -41,18 +60,54 @@ public class UdpServerUtil { udpServerPortMin = Integer.parseInt(tmp[0]); udpServerPortMax = Integer.parseInt(tmp[1]); if (udpServerPortMax <= udpServerPortMin) { - // 至少得给2个udp端口,一个用于基础无响应UDP转发,一个用于有响应UDP转发 + // 至少得给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; } @@ -60,4 +115,69 @@ public class UdpServerUtil { 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(); + lockChannel.getChannel().attr(Constants.UDP_CHANNEL_BIND_KEY).set(null); + udpServerFreePortPool.offer(lockChannel.getPort()); + log.debug("[udp channel]release udp channel port:{}", lockChannel.getPort()); + } + } + } } diff --git a/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/ProxyMessage.java b/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/ProxyMessage.java index 52eefdad..b77e3829 100644 --- a/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/ProxyMessage.java +++ b/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/ProxyMessage.java @@ -186,7 +186,7 @@ public class ProxyMessage { /** * 超时时间(<=0时,相当于不需要响应) */ - private int proxyTimeout; + private long proxyTimeoutMs; public String toJsonString() { return JSONObject.toJSONString(this); } diff --git a/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/proxy/core/UdpVisitorChannelHandler.java b/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/proxy/core/UdpVisitorChannelHandler.java index d05c9483..5a478adf 100644 --- a/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/proxy/core/UdpVisitorChannelHandler.java +++ b/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/proxy/core/UdpVisitorChannelHandler.java @@ -29,9 +29,6 @@ public class UdpVisitorChannelHandler extends SimpleChannelInboundHandler