UDP代理隧道连接逻辑调整
This commit is contained in:
+2
-2
@@ -36,11 +36,11 @@ public class UdpProxyMessageConnectHandler implements ProxyMessageHandler {
|
||||
log.info("[UDP connect]info:{}", proxyMessage.getInfo());
|
||||
|
||||
// 获取连接
|
||||
ProxyUtil.borrowTcpProxyChanel(udpProxyTunnelBootstrap, new ProxyChannelBorrowListener() {
|
||||
ProxyUtil.borrowUdpProxyChanel(udpProxyTunnelBootstrap, new ProxyChannelBorrowListener() {
|
||||
|
||||
@Override
|
||||
public void success(Channel channel) {
|
||||
ctx.channel().writeAndFlush(ProxyMessage.buildUdpConnectMessage(new ProxyMessage.UdpBaseInfo()
|
||||
channel.writeAndFlush(ProxyMessage.buildUdpConnectMessage(new ProxyMessage.UdpBaseInfo()
|
||||
.setVisitorId(udpBaseInfo.getVisitorId())
|
||||
.setServerPort(udpBaseInfo.getServerPort())
|
||||
.setTargetIp(udpBaseInfo.getTargetIp())
|
||||
|
||||
+51
-34
@@ -7,9 +7,11 @@ import io.netty.channel.ChannelOption;
|
||||
import io.netty.channel.SimpleChannelInboundHandler;
|
||||
import io.netty.channel.socket.DatagramPacket;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.dromara.neutrinoproxy.core.Constants;
|
||||
import org.dromara.neutrinoproxy.core.ProxyMessage;
|
||||
import org.dromara.neutrinoproxy.server.constant.NetworkProtocolEnum;
|
||||
import org.dromara.neutrinoproxy.server.proxy.domain.ProxyAttachment;
|
||||
import org.dromara.neutrinoproxy.server.proxy.domain.VisitorChannelAttachInfo;
|
||||
import org.dromara.neutrinoproxy.server.service.FlowReportService;
|
||||
import org.dromara.neutrinoproxy.server.util.ProxyUtil;
|
||||
@@ -30,42 +32,53 @@ public class UdpVisitorChannelHandler extends SimpleChannelInboundHandler<Datagr
|
||||
System.out.println("channelId:" + ctx.channel().id().asLongText());
|
||||
System.out.println("服务端接收到消息 \nsender:" + datagramPacket.sender().toString() + "内容\n" + datagramPacket.content().toString(StandardCharsets.UTF_8));
|
||||
|
||||
// 通知代理客户端
|
||||
Channel visitorChannel = ctx.channel();
|
||||
Channel proxyChannel = visitorChannel.attr(Constants.NEXT_CHANNEL).get();
|
||||
|
||||
if (null == proxyChannel) {
|
||||
// 该端口还没有代理客户端
|
||||
ctx.channel().close();
|
||||
return;
|
||||
}
|
||||
String targetIp = proxyChannel.attr(Constants.TARGET_IP).get();
|
||||
int targetPort = proxyChannel.attr(Constants.TARGET_PORT).get();
|
||||
|
||||
// 转发代理数据
|
||||
byte[] bytes = new byte[datagramPacket.content().readableBytes()];
|
||||
datagramPacket.content().readBytes(bytes);
|
||||
String visitorId = ProxyUtil.getVisitorIdByChannel(visitorChannel);
|
||||
proxyChannel.writeAndFlush(ProxyMessage.buildUdpTransferMessage(new ProxyMessage.UdpBaseInfo()
|
||||
.setVisitorId(visitorId)
|
||||
.setVisitorIp(datagramPacket.sender().getAddress().getHostAddress())
|
||||
.setVisitorPort(datagramPacket.sender().getPort())
|
||||
.setTargetIp(targetIp)
|
||||
.setTargetPort(targetPort)
|
||||
).setData(bytes));
|
||||
datagramPacket.content().resetReaderIndex();
|
||||
ProxyAttachment proxyAttachment = new ProxyAttachment(ctx.channel(), bytes, (channel, buf) -> {
|
||||
Channel proxyChannel = channel.attr(Constants.NEXT_CHANNEL).get();
|
||||
|
||||
// 增加流量计数
|
||||
VisitorChannelAttachInfo visitorChannelAttachInfo = ProxyUtil.getAttachInfo(visitorChannel);
|
||||
Solon.context().getBean(FlowReportService.class).addWriteByte(visitorChannelAttachInfo.getLicenseId(), bytes.length);
|
||||
}
|
||||
if (null == proxyChannel) {
|
||||
// 该端口还没有代理客户端
|
||||
ctx.channel().close();
|
||||
return;
|
||||
}
|
||||
String targetIp = proxyChannel.attr(Constants.TARGET_IP).get();
|
||||
int targetPort = proxyChannel.attr(Constants.TARGET_PORT).get();
|
||||
|
||||
// 转发代理数据
|
||||
String visitorId = ProxyUtil.getVisitorIdByChannel(channel);
|
||||
proxyChannel.writeAndFlush(ProxyMessage.buildUdpTransferMessage(new ProxyMessage.UdpBaseInfo()
|
||||
.setVisitorId(visitorId)
|
||||
.setVisitorIp(datagramPacket.sender().getAddress().getHostAddress())
|
||||
.setVisitorPort(datagramPacket.sender().getPort())
|
||||
.setTargetIp(targetIp)
|
||||
.setTargetPort(targetPort)
|
||||
).setData(bytes));
|
||||
|
||||
// 增加流量计数
|
||||
VisitorChannelAttachInfo visitorChannelAttachInfo = ProxyUtil.getAttachInfo(channel);
|
||||
Solon.context().getBean(FlowReportService.class).addWriteByte(visitorChannelAttachInfo.getLicenseId(), bytes.length);
|
||||
});
|
||||
|
||||
// String visitorId = ProxyUtil.getVisitorIdByChannel(ctx.channel());
|
||||
// if (StringUtils.isNotBlank(visitorId)) {
|
||||
// // UDP代理隧道已就绪,直接转发
|
||||
// proxyAttachment.execute();
|
||||
// return;
|
||||
// }
|
||||
Channel proxyChannel = ctx.channel().attr(Constants.NEXT_CHANNEL).get();
|
||||
if (null != proxyChannel && proxyChannel.isActive()) {
|
||||
// UDP代理隧道已就绪,直接转发
|
||||
proxyAttachment.execute();
|
||||
return;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void channelActive(ChannelHandlerContext ctx) throws Exception {
|
||||
System.out.println("active channelId:" + ctx.channel().id().asLongText());
|
||||
Channel visitorChannel = ctx.channel();
|
||||
InetSocketAddress sa = (InetSocketAddress) visitorChannel.localAddress();
|
||||
Channel cmdChannel = ProxyUtil.getCmdChannelByServerPort(sa.getPort());
|
||||
|
||||
// 没有指令通道,直接结束
|
||||
if (null == cmdChannel) {
|
||||
// 该端口还没有代理客户端
|
||||
ctx.channel().close();
|
||||
@@ -86,17 +99,21 @@ public class UdpVisitorChannelHandler extends SimpleChannelInboundHandler<Datagr
|
||||
// 用户连接到代理服务器时,设置用户连接不可读,等待代理后端服务器连接成功后再改变为可读状态
|
||||
visitorChannel.config().setOption(ChannelOption.AUTO_READ, false);
|
||||
|
||||
// UDP此处叫visitor似有不妥,与TCP不同
|
||||
// TODO UDP此处叫visitor似有不妥,与TCP不同,2.x重构思考
|
||||
String visitorId = ProxyUtil.newVisitorId();
|
||||
// 此处需要和tcp分开
|
||||
ProxyUtil.addVisitorChannelToCmdChannel(NetworkProtocolEnum.UDP, cmdChannel, visitorId, visitorChannel, sa.getPort());
|
||||
ProxyUtil.addProxyConnectAttachment(visitorId, proxyAttachment);
|
||||
cmdChannel.writeAndFlush(ProxyMessage.buildUdpConnectMessage(new ProxyMessage.UdpBaseInfo()
|
||||
.setVisitorId(visitorId)
|
||||
.setServerPort(sa.getPort())
|
||||
.setTargetIp(targetIp)
|
||||
.setTargetPort(targetPort)
|
||||
));
|
||||
.setVisitorId(visitorId)
|
||||
.setServerPort(sa.getPort())
|
||||
.setTargetIp(targetIp)
|
||||
.setTargetPort(targetPort)
|
||||
));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void channelActive(ChannelHandlerContext ctx) throws Exception {
|
||||
super.channelActive(ctx);
|
||||
}
|
||||
|
||||
|
||||
+6
-1
@@ -12,6 +12,7 @@ import org.dromara.neutrinoproxy.server.dal.PortMappingMapper;
|
||||
import org.dromara.neutrinoproxy.server.dal.entity.LicenseDO;
|
||||
import org.dromara.neutrinoproxy.server.dal.entity.PortMappingDO;
|
||||
import org.dromara.neutrinoproxy.server.dal.entity.UserDO;
|
||||
import org.dromara.neutrinoproxy.server.proxy.domain.ProxyAttachment;
|
||||
import org.dromara.neutrinoproxy.server.service.LicenseService;
|
||||
import org.dromara.neutrinoproxy.server.service.UserService;
|
||||
import org.dromara.neutrinoproxy.server.util.ProxyUtil;
|
||||
@@ -35,7 +36,6 @@ public class UdpProxyMessageConnectHandler implements ProxyMessageHandler {
|
||||
|
||||
@Override
|
||||
public void handle(ChannelHandlerContext ctx, ProxyMessage proxyMessage) {
|
||||
final Channel proxyChannel = ctx.channel();
|
||||
final ProxyMessage.UdpBaseInfo udpBaseInfo = JSONObject.parseObject(proxyMessage.getInfo(), ProxyMessage.UdpBaseInfo.class);
|
||||
final String licenseKey = new String(proxyMessage.getData());
|
||||
log.info("[UDP connect]info:{} licenseKey:{}", proxyMessage.getInfo(), licenseKey);
|
||||
@@ -85,6 +85,11 @@ public class UdpProxyMessageConnectHandler implements ProxyMessageHandler {
|
||||
visitorChannel.attr(Constants.NEXT_CHANNEL).set(ctx.channel());
|
||||
// 代理客户端与后端服务器连接成功,修改用户连接为可读状态
|
||||
visitorChannel.config().setOption(ChannelOption.AUTO_READ, true);
|
||||
// 获取代理附加对象
|
||||
ProxyAttachment proxyAttachment = ProxyUtil.getProxyConnectAttachment(udpBaseInfo.getVisitorId());
|
||||
if (null != proxyAttachment) {
|
||||
proxyAttachment.execute();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
Reference in New Issue
Block a user