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 index 45373152..73283be3 100644 --- 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 @@ -47,7 +47,7 @@ public class TcpProxyChannelHandler extends SimpleChannelInboundHandler realServerChannels = new ConcurrentHashMap(); - private static ConcurrentLinkedQueue proxyChannelPool = new ConcurrentLinkedQueue(); + private static ConcurrentLinkedQueue tcpProxyChannelPool = new ConcurrentLinkedQueue(); + private static ConcurrentLinkedQueue udpProxyChannelPool = new ConcurrentLinkedQueue<>(); private static volatile Channel cmdChannel; @@ -72,7 +73,7 @@ public class ProxyUtil { private static final String CLIENT_ID_FILE = ".NEUTRINO_PROXY_CLIENT_ID"; public static void borrowTcpProxyChanel(Bootstrap tcpProxyTunnelBootstrap, final ProxyChannelBorrowListener borrowListener) { - Channel channel = proxyChannelPool.poll(); + Channel channel = tcpProxyChannelPool.poll(); if (null != channel) { borrowListener.success(channel); return; @@ -87,18 +88,52 @@ public class ProxyUtil { }); } - public static void returnProxyChanel(Channel proxyChanel) { - if (proxyChannelPool.size() > MAX_POOL_SIZE) { + 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(); - proxyChannelPool.offer(proxyChanel); + tcpProxyChannelPool.offer(proxyChanel); } } - public static void removeProxyChanel(Channel proxyChanel) { - proxyChannelPool.remove(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) { diff --git a/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/Constants.java b/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/Constants.java index c6d04a16..5a3d916c 100644 --- a/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/Constants.java +++ b/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/Constants.java @@ -38,6 +38,9 @@ public interface Constants { AttributeKey LICENSE_ID = AttributeKey.newInstance("license_id"); + AttributeKey TARGET_IP = AttributeKey.newInstance("targetIp"); + AttributeKey TARGET_PORT = AttributeKey.newInstance("targetPort"); + int HEADER_SIZE = 4; int TYPE_SIZE = 1; int SERIAL_NUMBER_SIZE = 8; @@ -49,6 +52,9 @@ public interface Constants { String CONNECT = "CONNECT"; String DISCONNECT = "DISCONNECT"; String TRANSFER = "TRANSFER"; + String UDP_CONNECT = "UDP_CONNECT"; + String UDP_DISCONNECT = "UDP_DISCONNECT"; + String UDP_TRANSFER = "UDP_TRANSFER"; String ERROR = "ERROR"; String PORT_MAPPING_SYNC = "PORT_MAPPING_SYNC"; } diff --git a/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/ProxyDataTypeEnum.java b/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/ProxyDataTypeEnum.java index 0a62897b..e533bf8e 100644 --- a/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/ProxyDataTypeEnum.java +++ b/neutrino-proxy-core/src/main/java/org/dromara/neutrinoproxy/core/ProxyDataTypeEnum.java @@ -44,7 +44,10 @@ public enum ProxyDataTypeEnum { DISCONNECT(0x04, Constants.ProxyDataTypeName.DISCONNECT,"DISCONNECT"), TRANSFER(0x05, Constants.ProxyDataTypeName.TRANSFER,"TRANSFER"), ERROR(0x06, Constants.ProxyDataTypeName.ERROR,"ERROR"), - PORT_MAPPING_SYNC(0x07, Constants.ProxyDataTypeName.PORT_MAPPING_SYNC, "PORT_MAPPING_SYNC"); + PORT_MAPPING_SYNC(0x07, Constants.ProxyDataTypeName.PORT_MAPPING_SYNC, "PORT_MAPPING_SYNC"), + UDP_CONNECT(0x08, Constants.ProxyDataTypeName.UDP_CONNECT,"UDP_CONNECT"), + UDP_DISCONNECT(0x09, Constants.ProxyDataTypeName.UDP_DISCONNECT,"UDP_DISCONNECT"), + UDP_TRANSFER(0x10, Constants.ProxyDataTypeName.UDP_TRANSFER,"UDP_TRANSFER"); private static Map cache = Stream.of(values()).collect(Collectors.toMap(ProxyDataTypeEnum::getType, Function.identity())); private int type; 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 43d138e8..b349d91a 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 @@ -66,10 +66,18 @@ public class ProxyMessage { * 通用异常信息 */ public static final byte TYPE_ERROR = 0x06; + /** + * UDP代理隧道连接 + */ + public static final byte TYPE_UDP_CONNECT = 0x08; + /** + * UDP代理隧道断开连接 + */ + public static final byte TYPE_UDP_DISCONNECT = 0x09; /** * UDP数据传输 */ - private static final byte TYPE_UDP_TRANSFER = 0x08; + public static final byte TYPE_UDP_TRANSFER = 0x10; /** * 消息类型 @@ -134,15 +142,19 @@ public class ProxyMessage { .setData(data); } - public static ProxyMessage buildUdpTransferMessage(String visitorIp, int visitorPort, String targetIp, int targetPort, byte[] data) { + + public static ProxyMessage buildUdpConnectMessage(UdpBaseInfo info) { + return create().setType(TYPE_UDP_CONNECT) + .setInfo(info.toJsonString()); + } + + public static ProxyMessage buildUdpDisconnectMessage() { + return create().setType(TYPE_UDP_DISCONNECT); + } + + public static ProxyMessage buildUdpTransferMessage(UdpBaseInfo info) { return create().setType(TYPE_UDP_TRANSFER) - .setInfo(JSONObject.toJSONString(new UdpBaseInfo() - .setVisitorIp(visitorIp) - .setVisitorPort(visitorPort) - .setTargetIp(targetIp) - .setTargetPort(targetPort) - )) - .setData(data); + .setInfo(info.toJsonString()); } public static ProxyMessage buildErrMessage(ExceptionEnum exceptionEnum, String info) { @@ -161,9 +173,14 @@ public class ProxyMessage { @Accessors(chain = true) @Data public static class UdpBaseInfo { + private String visitorId; private String visitorIp; private int visitorPort; + private int serverPort; private String targetIp; private int targetPort; + public String toJsonString() { + return JSONObject.toJSONString(this); + } } } diff --git a/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/dal/PortMappingMapper.java b/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/dal/PortMappingMapper.java index cf35a7c0..f8b7a054 100644 --- a/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/dal/PortMappingMapper.java +++ b/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/dal/PortMappingMapper.java @@ -74,6 +74,14 @@ public interface PortMappingMapper extends BaseMapper { ); } + default PortMappingDO findByLicenseIdAndServerPort(Integer licenseId, Integer serverPort) { + return this.selectOne(new LambdaQueryWrapper() + .eq(PortMappingDO::getLicenseId, licenseId) + .eq(PortMappingDO::getServerPort, serverPort) + .last("limit 1") + ); + } + default void updateOnlineStatus(Integer licenseId,Integer serverPort, Integer isOnline, Date updateTime) { this.update(null, new LambdaUpdateWrapper() .eq(PortMappingDO::getLicenseId, licenseId) diff --git a/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/proxy/core/TcpVisitorChannelHandler.java b/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/proxy/core/TcpVisitorChannelHandler.java index 0c52037c..8f0b2186 100644 --- a/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/proxy/core/TcpVisitorChannelHandler.java +++ b/neutrino-proxy-server/src/main/java/org/dromara/neutrinoproxy/server/proxy/core/TcpVisitorChannelHandler.java @@ -4,6 +4,7 @@ import cn.hutool.core.util.StrUtil; import lombok.extern.slf4j.Slf4j; 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.VisitorChannelAttachInfo; import org.dromara.neutrinoproxy.server.service.FlowReportService; import org.dromara.neutrinoproxy.server.util.ProxyUtil; @@ -76,7 +77,7 @@ public class TcpVisitorChannelHandler extends SimpleChannelInboundHandler { @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)); + // 通知代理客户端 + 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)); + + // 增加流量计数 + VisitorChannelAttachInfo visitorChannelAttachInfo = ProxyUtil.getAttachInfo(visitorChannel); + Solon.context().getBean(FlowReportService.class).addWriteByte(visitorChannelAttachInfo.getLicenseId(), bytes.length); + } + + @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()); @@ -43,16 +83,27 @@ public class UdpVisitorChannelHandler extends SimpleChannelInboundHandler