调通UDP代理

This commit is contained in:
aoshiguchen
2023-09-21 23:40:00 +08:00
parent 33704b1694
commit e416c76ea8
9 changed files with 277 additions and 11 deletions
@@ -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<UdpChannelBindInfo> UDP_CHANNEL_BIND_KEY = AttributeKey.newInstance("udpChannelBindKey");
}
@@ -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<ByteBuf> {
@Slf4j
public class UdpRealServerHandler extends SimpleChannelInboundHandler<DatagramPacket> {
@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);
}
}
}
@@ -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
@@ -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;
}
@@ -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;
}
@@ -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<Integer, Channel> portToChannelMap = new ConcurrentHashMap<>();
/**
* udp服务空闲端口池
*/
private static ConcurrentLinkedQueue<Integer> udpServerFreePortPool = new ConcurrentLinkedQueue<>();
private static List<LockChannel> 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<LockChannel> 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());
}
}
}
}
@@ -186,7 +186,7 @@ public class ProxyMessage {
/**
* 超时时间(<=0时,相当于不需要响应)
*/
private int proxyTimeout;
private long proxyTimeoutMs;
public String toJsonString() {
return JSONObject.toJSONString(this);
}
@@ -29,9 +29,6 @@ public class UdpVisitorChannelHandler extends SimpleChannelInboundHandler<Datagr
@Override
protected void channelRead0(ChannelHandlerContext ctx, DatagramPacket datagramPacket) throws Exception {
System.out.println("channelId:" + ctx.channel().id().asLongText());
System.out.println("服务端接收到消息 \nsender:" + datagramPacket.sender().toString() + "内容\n" + datagramPacket.content().toString(StandardCharsets.UTF_8));
byte[] bytes = new byte[datagramPacket.content().readableBytes()];
datagramPacket.content().readBytes(bytes);
datagramPacket.content().resetReaderIndex();
@@ -54,6 +51,8 @@ public class UdpVisitorChannelHandler extends SimpleChannelInboundHandler<Datagr
.setVisitorPort(datagramPacket.sender().getPort())
.setTargetIp(targetIp)
.setTargetPort(targetPort)
.setProxyTimeoutMs(10000)
.setProxyResponses(3)
).setData(bytes));
// 增加流量计数
@@ -0,0 +1,44 @@
package org.dromara.neutrinoproxy.server.proxy.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.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.solon.annotation.Component;
import java.net.InetSocketAddress;
/**
* @author: aoshiguchen
* @date: 2023/9/21
*/
@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 = JSONObject.parseObject(proxyMessage.getInfo(), ProxyMessage.UdpBaseInfo.class);
log.debug("[UDP transfer]info:{} data:{}", proxyMessage.getInfo(), new String(proxyMessage.getData()));
Channel visitorChannel = ctx.channel().attr(Constants.NEXT_CHANNEL).get();
if (null != visitorChannel) {
InetSocketAddress address = new InetSocketAddress(udpBaseInfo.getVisitorIp(), udpBaseInfo.getVisitorPort());
ByteBuf byteBuf = Unpooled.copiedBuffer(proxyMessage.getData());
visitorChannel.writeAndFlush(new DatagramPacket(byteBuf, address));
}
}
@Override
public String name() {
return ProxyDataTypeEnum.UDP_TRANSFER.getDesc();
}
}