修复v6,重复分配端口,指定网卡失败的问题

This commit is contained in:
qaq
2025-11-20 11:27:58 +08:00
parent 252c16de29
commit f6dd3ad657
16 changed files with 661 additions and 379 deletions
@@ -1,5 +1,6 @@
package com.admin.common.task;
import com.admin.common.dto.GostDto;
import com.admin.common.utils.GostUtil;
import com.admin.entity.*;
import com.admin.service.*;
@@ -1,5 +1,7 @@
package com.admin.common.utils;
import cn.hutool.core.util.StrUtil;
import com.admin.common.dto.GostDto;
import com.admin.entity.*;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
@@ -12,26 +14,34 @@ import java.util.Objects;
public class GostUtil {
public static void AddLimiters(Long node_id, Long name, String speed) {
public static GostDto AddLimiters(Long node_id, Long name, String speed) {
JSONObject data = createLimiterData(name, speed);
WebSocketServer.send_msg(node_id, data, "AddLimiters");
GostDto gostDto = WebSocketServer.send_msg(node_id, data, "AddLimiters");
if (gostDto.getMsg().contains("exists")){
gostDto.setMsg("OK");
}
return gostDto;
}
public static void UpdateLimiters(Long node_id, Long name, String speed) {
public static GostDto UpdateLimiters(Long node_id, Long name, String speed) {
JSONObject data = createLimiterData(name, speed);
JSONObject req = new JSONObject();
req.put("limiter", name + "");
req.put("data", data);
WebSocketServer.send_msg(node_id, req, "UpdateLimiters");
return WebSocketServer.send_msg(node_id, req, "UpdateLimiters");
}
public static void DeleteLimiters(Long node_id, Long name) {
public static GostDto DeleteLimiters(Long node_id, Long name) {
JSONObject req = new JSONObject();
req.put("limiter", name + "");
WebSocketServer.send_msg(node_id, req, "DeleteLimiters");
GostDto gostDto = WebSocketServer.send_msg(node_id, req, "DeleteLimiters");
if (gostDto.getMsg().contains("not found")){
gostDto.setMsg("OK");
}
return gostDto;
}
public static void AddChains(Long node_id, List<ChainTunnel> chainTunnels, Map<Long, Node> node_s) {
public static GostDto AddChains(Long node_id, List<ChainTunnel> chainTunnels, Map<Long, Node> node_s) {
JSONArray nodes = new JSONArray();
for (ChainTunnel chainTunnel : chainTunnels) {
JSONObject dialer = new JSONObject();
@@ -43,19 +53,23 @@ public class GostUtil {
Node node_info = node_s.get(chainTunnel.getNodeId());
JSONObject node = new JSONObject();
node.put("name", "node_" + chainTunnel.getInx());
node.put("addr", node_info.getServerIp() + ":" + chainTunnel.getPort());
node.put("addr", processServerAddress(node_info.getServerIp()) + ":" + chainTunnel.getPort());
node.put("connector", connector);
node.put("dialer", dialer);
if (StringUtils.isNotBlank(node_info.getInterfaceName())) {
node.put("interface", node_info.getInterfaceName());
}
nodes.add(node);
}
JSONObject hop = new JSONObject();
hop.put("name", "hop_" + chainTunnels.getFirst().getTunnelId());
// interface设置在转发链
if (StringUtils.isNotBlank(node_s.get(node_id).getInterfaceName())) {
hop.put("interface", node_s.get(node_id).getInterfaceName());
}
JSONObject selector = new JSONObject();
selector.put("strategy", chainTunnels.getFirst().getStrategy());
selector.put("maxFails", 1);
@@ -72,24 +86,34 @@ public class GostUtil {
data.put("name", "chains_" + chainTunnels.getFirst().getTunnelId());
data.put("hops", hops);
WebSocketServer.send_msg(node_id, data, "AddChains");
GostDto gostDto = WebSocketServer.send_msg(node_id, data, "AddChains");
if (gostDto.getMsg().contains("exists")){
gostDto.setMsg("OK");
}
return gostDto;
}
public static void DeleteChains(Long node_id, String name) {
public static GostDto DeleteChains(Long node_id, String name) {
JSONObject data = new JSONObject();
data.put("chain", name);
WebSocketServer.send_msg(node_id, data, "DeleteChains");
GostDto gostDto = WebSocketServer.send_msg(node_id, data, "DeleteChains");
if (gostDto.getMsg().contains("not found")){
gostDto.setMsg("OK");
}
return gostDto;
}
public static void AddChainService(Long node_id, ChainTunnel chainTunnel, Map<Long, Node> node_s) {
public static GostDto AddChainService(Long node_id, ChainTunnel chainTunnel, Map<Long, Node> node_s) {
JSONArray services = new JSONArray();
Node node_info = node_s.get(chainTunnel.getNodeId());
JSONObject service_item = new JSONObject();
service_item.put("name", chainTunnel.getTunnelId() + "_tls");
service_item.put("addr", node_info.getTcpListenAddr() + ":" + chainTunnel.getPort());
if (StringUtils.isNotBlank(node_info.getInterfaceName())) {
// 只为出口节点(chainType=3)设置 interface
if (chainTunnel.getChainType() == 3 && StringUtils.isNotBlank(node_s.get(node_id).getInterfaceName())) {
JSONObject metadata = new JSONObject();
metadata.put("interface", node_info.getInterfaceName());
metadata.put("interface", node_s.get(node_id).getInterfaceName());
service_item.put("metadata", metadata);
}
@@ -106,16 +130,14 @@ public class GostUtil {
services.add(service_item);
WebSocketServer.send_msg(node_id, services, "AddService");
GostDto gostDto = WebSocketServer.send_msg(node_id, services, "AddService");
if (gostDto.getMsg().contains("exists")){
gostDto.setMsg("OK");
}
return gostDto;
}
public static void DeleteChainService(Long node_id, JSONArray services) {
JSONObject data = new JSONObject();
data.put("services", services);
WebSocketServer.send_msg(node_id, data, "DeleteService");
}
public static void AddAndUpdateService(String name, Integer limiter, Node node, Forward forward, ForwardPort forwardPort, Tunnel tunnel, String meth) {
public static GostDto AddAndUpdateService(String name, Integer limiter, Node node, Forward forward, ForwardPort forwardPort, Tunnel tunnel, String meth) {
JSONArray services = new JSONArray();
String[] protocols = {"tcp", "udp"};
for (String protocol : protocols) {
@@ -127,7 +149,8 @@ public class GostUtil {
service.put("addr", node.getUdpListenAddr() + ":" + forwardPort.getPort());
}
if (StringUtils.isNotBlank(node.getInterfaceName())) {
// 只在端口转发时设置 interface(隧道转发时 interface 在转发链的节点上设置)
if (tunnel.getType() == 1 && StringUtils.isNotBlank(node.getInterfaceName())) {
JSONObject metadata = new JSONObject();
metadata.put("interface", node.getInterfaceName());
service.put("metadata", metadata);
@@ -155,22 +178,30 @@ public class GostUtil {
services.add(service);
}
WebSocketServer.send_msg(node.getId(), services, meth);
GostDto gostDto = WebSocketServer.send_msg(node.getId(), services, meth);
if (gostDto.getMsg().contains("exists")){
gostDto.setMsg("OK");
}
return gostDto;
}
public static void DeleteService(Long node_id, JSONArray services) {
public static GostDto DeleteService(Long node_id, JSONArray services) {
JSONObject data = new JSONObject();
data.put("services", services);
WebSocketServer.send_msg(node_id, data, "DeleteService");
GostDto gostDto = WebSocketServer.send_msg(node_id, data, "DeleteService");
if (gostDto.getMsg().contains("not found")){
gostDto.setMsg("OK");
}
return gostDto;
}
public static void PauseAndResumeService(Long node_id, String name, String meth) {
public static GostDto PauseAndResumeService(Long node_id, String name, String meth) {
JSONObject data = new JSONObject();
JSONArray services = new JSONArray();
services.add(name + "_tcp");
services.add(name + "_udp");
data.put("services", services);
WebSocketServer.send_msg(node_id, data, meth);
return WebSocketServer.send_msg(node_id, data, meth);
}
@@ -208,7 +239,7 @@ public class GostUtil {
num++;
}
if (strategy == null || strategy.equals("")) {
if (strategy == null || strategy.isEmpty()) {
strategy = "fifo";
}
@@ -222,5 +253,42 @@ public class GostUtil {
return forwarder;
}
public static String processServerAddress(String serverAddr) {
if (StrUtil.isBlank(serverAddr)) {
return serverAddr;
}
// 如果已经被方括号包裹,直接返回
if (serverAddr.startsWith("[")) {
return serverAddr;
}
// 查找最后一个冒号,分离主机和端口
int lastColonIndex = serverAddr.lastIndexOf(':');
if (lastColonIndex == -1) {
// 没有端口号,直接检查是否需要包裹
return isIPv6Address(serverAddr) ? "[" + serverAddr + "]" : serverAddr;
}
String host = serverAddr.substring(0, lastColonIndex);
String port = serverAddr.substring(lastColonIndex);
// 检查主机部分是否为IPv6地址
if (isIPv6Address(host)) {
return "[" + host + "]" + port;
}
return serverAddr;
}
private static boolean isIPv6Address(String address) {
// IPv6地址包含多个冒号,至少2个
if (!address.contains(":")) {
return false;
}
// 计算冒号数量,IPv6地址至少有2个冒号
long colonCount = address.chars().filter(ch -> ch == ':').count();
return colonCount >= 2;
}
}
@@ -56,6 +56,4 @@ public interface TunnelService extends IService<Tunnel> {
* @return 诊断结果
*/
R diagnoseTunnel(Long tunnelId);
Integer getNodePort(Long nodeId, Integer type, Integer port);
}
@@ -18,6 +18,7 @@ import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.util.*;
@@ -82,53 +83,60 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
// 判断是否使用隧道的inIp
boolean useTunnelInIp = tunnel.getInIp() != null && !tunnel.getInIp().trim().isEmpty();
// 收集所有的IP列表
List<String> ipList = new ArrayList<>();
// 收集所有的端口列表
List<Integer> portList = new ArrayList<>();
Set<String> ipPortSet = new LinkedHashSet<>();
if (useTunnelInIp) {
// 使用隧道的inIp
// 使用隧道的inIp(求笛卡尔积)
List<String> ipList = new ArrayList<>();
List<Integer> portList = new ArrayList<>();
String[] tunnelInIps = tunnel.getInIp().split(",");
for (String ip : tunnelInIps) {
if (ip != null && !ip.trim().isEmpty()) {
ipList.add(ip.trim());
}
}
} else {
// 使用节点的serverIp
// 收集所有端口
for (ForwardPort forwardPort : forwardPorts) {
Node node = nodeService.getById(forwardPort.getNodeId());
if (node != null && node.getServerIp() != null) {
ipList.add(node.getServerIp());
if (forwardPort.getPort() != null) {
portList.add(forwardPort.getPort());
}
}
}
// 收集所有端口
for (ForwardPort forwardPort : forwardPorts) {
if (forwardPort.getPort() != null) {
portList.add(forwardPort.getPort());
// 去重
List<String> uniqueIps = ipList.stream().distinct().toList();
List<Integer> uniquePorts = portList.stream().distinct().toList();
// 组合 IP:Port(笛卡尔积)
for (String ip : uniqueIps) {
for (Integer port : uniquePorts) {
ipPortSet.add(ip + ":" + port);
}
}
// inPort设置为第一个端口(用于向后兼容)
if (!uniquePorts.isEmpty()) {
forward.setInPort(uniquePorts.getFirst());
}
} else {
// 使用节点的serverIp(一对一,不求笛卡尔积)
for (ForwardPort forwardPort : forwardPorts) {
Node node = nodeService.getById(forwardPort.getNodeId());
if (node != null && node.getServerIp() != null && forwardPort.getPort() != null) {
ipPortSet.add(node.getServerIp() + ":" + forwardPort.getPort());
}
}
// inPort设置为第一个端口(用于向后兼容)
if (!forwardPorts.isEmpty() && forwardPorts.getFirst().getPort() != null) {
forward.setInPort(forwardPorts.getFirst().getPort());
}
}
// 去重
List<String> uniqueIps = ipList.stream().distinct().toList();
List<Integer> uniquePorts = portList.stream().distinct().toList();
// 组合 IP:Port(笛卡尔积)
Set<String> ipPortSet = new LinkedHashSet<>();
for (String ip : uniqueIps) {
for (Integer port : uniquePorts) {
ipPortSet.add(ip + ":" + port);
}
}
// 设置入口IP和端口
// 设置入口IP
if (!ipPortSet.isEmpty()) {
forward.setInIp(String.join(",", ipPortSet));
// inPort设置为第一个端口(用于向后兼容)
forward.setInPort(uniquePorts.getFirst());
}
}
@@ -159,25 +167,43 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
forward.setUserName(currentUser.getUserName());
forward.setCreatedTime(System.currentTimeMillis());
forward.setUpdatedTime(System.currentTimeMillis());
List<JSONObject> success = new ArrayList<>();
List<ChainTunnel> chainTunnels = chainTunnelService.list(new QueryWrapper<ChainTunnel>().eq("tunnel_id", tunnel.getId()).eq("chain_type", 1));
chainTunnels = get_port(chainTunnels, forwardDto.getInPort());
this.save(forward);
List<ChainTunnel> chainTunnels = chainTunnelService.list(new QueryWrapper<ChainTunnel>().eq("tunnel_id", tunnel.getId()).eq("chain_type", 1));
for (ChainTunnel chainTunnel : chainTunnels) {
Integer nodePort = tunnelService.getNodePort(chainTunnel.getNodeId(), 2, forwardDto.getInPort());
ForwardPort forwardPort = new ForwardPort();
forwardPort.setForwardId(forward.getId());
forwardPort.setNodeId(chainTunnel.getNodeId());
forwardPort.setPort(nodePort);
forwardPort.setPort(chainTunnel.getPort());
forwardPortService.save(forwardPort);
String serviceName = buildServiceName(forward.getId(), forward.getUserId(), permissionResult.getUserTunnel());
Integer limiter = permissionResult.getLimiter();
Node node = nodeService.getById(chainTunnel.getNodeId());
if (node == null){
if (node == null) {
return R.err("部分节点不存在");
}
GostUtil.AddAndUpdateService(serviceName, limiter, node, forward, forwardPort, tunnel, "AddService");
GostDto gostDto = GostUtil.AddAndUpdateService(serviceName, limiter, node, forward, forwardPort, tunnel, "AddService");
if (Objects.equals(gostDto.getMsg(), "OK")) {
JSONObject data = new JSONObject();
data.put("node_id", node.getId());
data.put("name", serviceName);
success.add(data);
} else {
this.removeById(forward.getId());
forwardPortService.remove(new QueryWrapper<ForwardPort>().eq("forward_id", forward.getId()));
for (JSONObject jsonObject : success) {
JSONArray se = new JSONArray();
se.add(jsonObject.getString("name") + "_tcp");
se.add(jsonObject.getString("name") + "_udp");
GostUtil.DeleteService(jsonObject.getLong("node_id"), se);
return R.err(gostDto.getMsg());
}
}
}
return R.ok();
}
@@ -225,25 +251,22 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
List<ChainTunnel> chainTunnels = chainTunnelService.list(new QueryWrapper<ChainTunnel>().eq("tunnel_id", tunnel.getId()).eq("chain_type", 1));
chainTunnels = get_port(chainTunnels, forwardUpdateDto.getInPort());
for (ChainTunnel chainTunnel : chainTunnels) {
String serviceName = buildServiceName(existForward.getId(), existForward.getUserId(), userTunnel);
Integer limiter = permissionResult.getLimiter();
Node node = nodeService.getById(chainTunnel.getNodeId());
if (node == null){
if (node == null) {
return R.err("部分节点不存在");
}
ForwardPort forwardPort = forwardPortService.getOne(new QueryWrapper<ForwardPort>().eq("forward_id", existForward.getId()).eq("node_id", node.getId()));
if (forwardPort == null){
if (forwardPort == null) {
return R.err("部分节点不存在1");
}
if (forwardUpdateDto.getInPort() != null && !forwardUpdateDto.getInPort().equals(forwardPort.getPort())) {
Integer nodePort = tunnelService.getNodePort(forwardPort.getNodeId(), 2, forwardUpdateDto.getInPort());
if (Objects.equals(nodePort, forwardUpdateDto.getInPort())) {
forwardPort.setPort(nodePort);
forwardPortService.updateById(forwardPort);
}
}
GostUtil.AddAndUpdateService(serviceName, limiter, node, existForward, forwardPort, tunnel, "UpdateService");
forwardPort.setPort(chainTunnel.getPort());
forwardPortService.updateById(forwardPort);
GostDto gostDto = GostUtil.AddAndUpdateService(serviceName, limiter, node, existForward, forwardPort, tunnel, "UpdateService");
if (!Objects.equals(gostDto.getMsg(), "OK")) return R.err(gostDto.getMsg());
}
return R.ok();
@@ -289,7 +312,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
String serviceName = buildServiceName(forward.getId(), forward.getUserId(), userTunnel);
Node node = nodeService.getById(chainTunnel.getNodeId());
if (node == null){
if (node == null) {
return R.err("部分节点不存在");
}
@@ -305,7 +328,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
@Override
public R pauseForward(Long id) {
return changeForwardStatus(id, 0, "PauseService");
return changeForwardStatus(id, 0, "PauseService");
}
@Override
@@ -400,7 +423,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
// 1. 入口->第一跳(或出口)
for (ChainTunnel inNode : inNodes) {
Node fromNode = nodeService.getById(inNode.getNodeId());
if (fromNode != null) {
if (!chainNodesList.isEmpty()) {
for (ChainTunnel firstChainNode : chainNodesList.getFirst()) {
@@ -436,10 +459,10 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
// 2. 链路测试
for (int i = 0; i < chainNodesList.size(); i++) {
List<ChainTunnel> currentHop = chainNodesList.get(i);
for (ChainTunnel currentNode : currentHop) {
Node fromNode = nodeService.getById(currentNode.getNodeId());
if (fromNode != null) {
if (i + 1 < chainNodesList.size()) {
for (ChainTunnel nextNode : chainNodesList.get(i + 1)) {
@@ -507,65 +530,57 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
}
@Override
@Transactional
public R updateForwardOrder(Map<String, Object> params) {
try {
// 1. 获取当前用户信息
UserInfo currentUser = getCurrentUserInfo();
// 1. 获取当前用户信息
UserInfo currentUser = getCurrentUserInfo();
// 2. 验证参数
if (!params.containsKey("forwards")) {
return R.err("缺少forwards参数");
}
@SuppressWarnings("unchecked")
List<Map<String, Object>> forwardsList = (List<Map<String, Object>>) params.get("forwards");
if (forwardsList == null || forwardsList.isEmpty()) {
return R.err("forwards参数不能为空");
}
// 3. 验证用户权限(只能更新自己的转发)
if (currentUser.getRoleId() != 0) {
// 普通用户只能更新自己的转发
List<Long> forwardIds = forwardsList.stream()
.map(item -> Long.valueOf(item.get("id").toString()))
.collect(Collectors.toList());
// 检查所有转发是否属于当前用户
QueryWrapper<Forward> queryWrapper = new QueryWrapper<>();
queryWrapper.in("id", forwardIds);
queryWrapper.eq("user_id", currentUser.getUserId());
long count = this.count(queryWrapper);
if (count != forwardIds.size()) {
return R.err("只能更新自己的转发排序");
}
}
// 4. 批量更新排序
List<Forward> forwardsToUpdate = new ArrayList<>();
for (Map<String, Object> forwardData : forwardsList) {
Long id = Long.valueOf(forwardData.get("id").toString());
Integer inx = Integer.valueOf(forwardData.get("inx").toString());
Forward forward = new Forward();
forward.setId(id);
forward.setInx(inx);
forwardsToUpdate.add(forward);
}
// 5. 执行批量更新
boolean success = this.updateBatchById(forwardsToUpdate);
if (success) {
log.info("用户 {} 更新了 {} 个转发的排序", currentUser.getUserName(), forwardsToUpdate.size());
return R.ok("排序更新成功");
} else {
return R.err("排序更新失败");
}
} catch (Exception e) {
log.error("更新转发排序失败", e);
return R.err("更新排序时发生错误: " + e.getMessage());
// 2. 验证参数
if (!params.containsKey("forwards")) {
return R.err("缺少forwards参数");
}
@SuppressWarnings("unchecked")
List<Map<String, Object>> forwardsList = (List<Map<String, Object>>) params.get("forwards");
if (forwardsList == null || forwardsList.isEmpty()) {
return R.err("forwards参数不能为空");
}
// 3. 验证用户权限(只能更新自己的转发)
if (currentUser.getRoleId() != 0) {
// 普通用户只能更新自己的转发
List<Long> forwardIds = forwardsList.stream()
.map(item -> Long.valueOf(item.get("id").toString()))
.collect(Collectors.toList());
// 检查所有转发是否属于当前用户
QueryWrapper<Forward> queryWrapper = new QueryWrapper<>();
queryWrapper.in("id", forwardIds);
queryWrapper.eq("user_id", currentUser.getUserId());
long count = this.count(queryWrapper);
if (count != forwardIds.size()) {
return R.err("只能更新自己的转发排序");
}
}
// 4. 批量更新排序
List<Forward> forwardsToUpdate = new ArrayList<>();
for (Map<String, Object> forwardData : forwardsList) {
Long id = Long.valueOf(forwardData.get("id").toString());
Integer inx = Integer.valueOf(forwardData.get("inx").toString());
Forward forward = new Forward();
forward.setId(id);
forward.setInx(inx);
forwardsToUpdate.add(forward);
}
// 5. 执行批量更新
this.updateBatchById(forwardsToUpdate);
return R.ok();
}
@@ -621,10 +636,11 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
for (ChainTunnel chainTunnel : chainTunnels) {
String serviceName = buildServiceName(forward.getId(), forward.getUserId(), userTunnel);
Node node = nodeService.getById(chainTunnel.getNodeId());
if (node == null){
if (node == null) {
return R.err("部分节点不存在");
}
GostUtil.PauseAndResumeService(node.getId(), serviceName, gostMethod);
GostDto gostDto = GostUtil.PauseAndResumeService(node.getId(), serviceName, gostMethod);
if (!Objects.equals(gostDto.getMsg(), "OK")) return R.err(gostDto.getMsg());
}
forward.setStatus(targetStatus);
forward.setUpdatedTime(System.currentTimeMillis());
@@ -923,6 +939,94 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
return forwardId + "_" + userId + "_" + userTunnelId;
}
public List<ChainTunnel> get_port(List<ChainTunnel> chainTunnelList, Integer in_port) {
List<List<Integer>> list = new ArrayList<>();
// 获取每个节点的端口列表
for (ChainTunnel tunnel : chainTunnelList) {
List<Integer> nodePort = getNodePort(tunnel.getNodeId());
if (nodePort.isEmpty()) {
throw new RuntimeException("暂无可用端口");
}
list.add(nodePort);
}
// ========== 如果指定了 in_port,优先检查公有 ==========
if (in_port != null) {
for (List<Integer> ports : list) {
if (!ports.contains(in_port)) {
throw new RuntimeException("指定端口 " + in_port + " 不可用(并非所有节点都有此端口)");
}
}
// 所有节点都有该端口 设置回 ChainTunnel
for (ChainTunnel tunnel : chainTunnelList) {
tunnel.setPort(in_port);
}
return chainTunnelList;
}
// ========== 未指定 in_port 查找最小的共同端口 ==========
Set<Integer> intersection = new HashSet<>(list.get(0));
for (int i = 1; i < list.size(); i++) {
intersection.retainAll(list.get(i));
}
if (!intersection.isEmpty()) {
// 找最小端口
Integer commonMin = intersection.stream().min(Integer::compareTo).orElseThrow();
// 设置到所有节点
for (ChainTunnel tunnel : chainTunnelList) {
tunnel.setPort(commonMin);
}
return chainTunnelList;
}
// ========== 没有共同端口取各自第一个可用端口 ==========
for (int i = 0; i < chainTunnelList.size(); i++) {
List<Integer> ports = list.get(i);
Integer first = ports.getFirst();
chainTunnelList.get(i).setPort(first);
}
return chainTunnelList;
}
public List<Integer> getNodePort(Long nodeId) {
Node node = nodeService.getById(nodeId);
if (node == null) {
throw new RuntimeException("节点不存在");
}
// 1. 查询隧道转发链占用的端口
List<ChainTunnel> chainTunnels = chainTunnelService.list(
new QueryWrapper<ChainTunnel>().eq("node_id", nodeId)
);
Set<Integer> usedPorts = chainTunnels.stream()
.map(ChainTunnel::getPort)
.filter(Objects::nonNull)
.collect(Collectors.toSet());
List<ForwardPort> list = forwardPortService.list(new QueryWrapper<ForwardPort>().eq("node_id", nodeId));
Set<Integer> forwardUsedPorts = new HashSet<>();
for (ForwardPort forwardPort : list) {
forwardUsedPorts.add(forwardPort.getPort());
}
usedPorts.addAll(forwardUsedPorts);
List<Integer> parsedPorts = TunnelServiceImpl.parsePorts(node.getPort());
return parsedPorts.stream()
.filter(p -> !usedPorts.contains(p))
.toList();
}
// ========== 内部数据类 ==========
@Data
@@ -973,7 +1077,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
private double averageTime;
private double packetLoss;
private long timestamp;
// 链路类型相关字段
private Integer fromChainType; // 1: 入口, 2: 链, 3: 出口
private Integer fromInx;
@@ -6,6 +6,7 @@ import com.admin.common.dto.GostDto;
import com.admin.common.dto.NodeDto;
import com.admin.common.dto.NodeUpdateDto;
import com.admin.common.lang.R;
import com.admin.common.utils.GostUtil;
import com.admin.common.utils.WebSocketServer;
import com.admin.entity.*;
import com.admin.mapper.NodeMapper;
@@ -125,9 +126,9 @@ public class NodeServiceImpl extends ServiceImpl<NodeMapper, Node> implements No
ViteConfig viteConfig = viteConfigService.getOne(new QueryWrapper<ViteConfig>().eq("name", "ip"));
if (viteConfig == null) return R.err("请先前往网站配置中设置ip");
StringBuilder command = new StringBuilder();
command.append("curl -L https://github.com/bqlpfy/flux-panel/releases/download/2.0.2-beta/install.sh")
command.append("curl -L https://github.com/bqlpfy/flux-panel/releases/download/2.0.3-beta/install.sh")
.append(" -o ./install.sh && chmod +x ./install.sh && ");
String processedServerAddr = processServerAddress(viteConfig.getValue());
String processedServerAddr = GostUtil.processServerAddress(viteConfig.getValue());
command.append("./install.sh")
.append(" -a ").append(processedServerAddr) // 服务器地址
.append(" -s ").append(node.getSecret()); // 节点密钥
@@ -182,44 +183,6 @@ public class NodeServiceImpl extends ServiceImpl<NodeMapper, Node> implements No
}
private String processServerAddress(String serverAddr) {
if (StrUtil.isBlank(serverAddr)) {
return serverAddr;
}
// 如果已经被方括号包裹,直接返回
if (serverAddr.startsWith("[")) {
return serverAddr;
}
// 查找最后一个冒号,分离主机和端口
int lastColonIndex = serverAddr.lastIndexOf(':');
if (lastColonIndex == -1) {
// 没有端口号,直接检查是否需要包裹
return isIPv6Address(serverAddr) ? "[" + serverAddr + "]" : serverAddr;
}
String host = serverAddr.substring(0, lastColonIndex);
String port = serverAddr.substring(lastColonIndex);
// 检查主机部分是否为IPv6地址
if (isIPv6Address(host)) {
return "[" + host + "]" + port;
}
return serverAddr;
}
private boolean isIPv6Address(String address) {
// IPv6地址包含多个冒号,至少2个
if (!address.contains(":")) {
return false;
}
// 计算冒号数量,IPv6地址至少有2个冒号
long colonCount = address.chars().filter(ch -> ch == ':').count();
return colonCount >= 2;
}
}
@@ -8,6 +8,7 @@ import com.admin.common.utils.GostUtil;
import com.admin.entity.*;
import com.admin.mapper.SpeedLimitMapper;
import com.admin.service.*;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import lombok.Data;
@@ -19,6 +20,7 @@ import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.UUID;
@@ -65,11 +67,23 @@ public class SpeedLimitServiceImpl extends ServiceImpl<SpeedLimitMapper, SpeedLi
String speedInMBps = convertBitsToMBps(speedLimit.getSpeed());
List<Long> limit_success = new ArrayList<>();
List<ChainTunnel> tunnelList = chainTunnelService.list(new QueryWrapper<ChainTunnel>().eq("tunnel_id", speedLimit.getTunnelId()));
for (ChainTunnel chainTunnel : tunnelList) {
Node node = nodeService.getById(chainTunnel.getNodeId());
if (node != null) {
GostUtil.AddLimiters(node.getId(),speedLimit.getId(),speedInMBps);
GostDto gostDto = GostUtil.AddLimiters(node.getId(), speedLimit.getId(), speedInMBps);
if (Objects.equals(gostDto.getMsg(), "OK")){
limit_success.add(node.getId());
}else {
this.removeById(speedLimit.getId());
for (Long node_id : limit_success) {
GostDto deleteLimiters = GostUtil.DeleteLimiters(node_id, speedLimit.getId());
System.out.println(deleteLimiters);
}
return R.err(gostDto.getMsg());
}
}
}
return R.ok();
@@ -94,7 +108,8 @@ public class SpeedLimitServiceImpl extends ServiceImpl<SpeedLimitMapper, SpeedLi
for (ChainTunnel chainTunnel : tunnelList) {
Node node = nodeService.getById(chainTunnel.getNodeId());
if (node != null) {
GostUtil.UpdateLimiters(node.getId(),speedLimit.getId(),speedInMBps);
GostDto gostDto = GostUtil.UpdateLimiters(node.getId(), speedLimit.getId(), speedInMBps);
if (!Objects.equals(gostDto.getMsg(), "OK")) return R.err(gostDto.getMsg());
}
}
this.updateById(speedLimit);
@@ -116,7 +131,8 @@ public class SpeedLimitServiceImpl extends ServiceImpl<SpeedLimitMapper, SpeedLi
for (ChainTunnel chainTunnel : tunnelList) {
Node node = nodeService.getById(chainTunnel.getNodeId());
if (node != null) {
GostUtil.DeleteLimiters(node.getId(),speedLimit.getId());
GostDto gostDto = GostUtil.DeleteLimiters(node.getId(), speedLimit.getId());
if (!Objects.equals(gostDto.getMsg(), "OK"))return R.err(gostDto.getMsg());
}
}
this.removeById(id);
@@ -81,7 +81,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
Node node = nodeService.getById(chain_node.getNodeId());
if (node == null) return R.err("节点不存在");
nodes.put(node.getId(), node);
Integer nodePort = getNodePort(chain_node.getNodeId(), 1, null);
Integer nodePort = getNodePort(chain_node.getNodeId());
chain_node.setPort(nodePort);
chain_node.setInx(inx); // 设置转发链序号
chainTunnels.add(chain_node);
@@ -93,7 +93,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
Node node = nodeService.getById(out_node.getNodeId());
if (node == null) return R.err("节点不存在");
nodes.put(node.getId(), node);
Integer nodePort = getNodePort(out_node.getNodeId(), 1, null);
Integer nodePort = getNodePort(out_node.getNodeId());
out_node.setPort(nodePort);
chainTunnels.add(out_node);
}
@@ -132,14 +132,35 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
}
chainTunnelService.saveBatch(chainTunnels);
List<JSONObject> chain_success = new ArrayList<>();
List<JSONObject> service_success = new ArrayList<>();
if (tunnel.getType() == 2) {
for (ChainTunnel in_node : tunnelDto.getInNodeId()) {
// 创建Chain, 指向chainNode的第一跳。如果chainNode为空就是指向出口
if (tunnelDto.getChainNodes().isEmpty()) { // 指向出口
GostUtil.AddChains(in_node.getNodeId(), tunnelDto.getOutNodeId(), nodes);
GostDto gostDto = GostUtil.AddChains(in_node.getNodeId(), tunnelDto.getOutNodeId(), nodes);
isError(gostDto);
} else {
GostUtil.AddChains(in_node.getNodeId(), tunnelDto.getChainNodes().getFirst(), nodes);// 指向第一跳
GostDto gostDto = GostUtil.AddChains(in_node.getNodeId(), tunnelDto.getChainNodes().getFirst(), nodes);// 指向第一跳
if (Objects.equals(gostDto.getMsg(), "OK")){
JSONObject data = new JSONObject();
data.put("node_id", in_node.getNodeId());
data.put("name", "chains_" + tunnel.getId());
chain_success.add(data);
}else {
this.removeById(tunnel.getId());
chainTunnelService.remove(new QueryWrapper<ChainTunnel>().eq("tunnel_id", tunnel.getId()));
for (JSONObject chainSuccess : chain_success) {
GostDto deleteChains = GostUtil.DeleteChains(chainSuccess.getLong("node_id"), chainSuccess.getString("name"));
System.out.println(deleteChains);
}
return R.err(gostDto.getMsg());
}
}
}
@@ -149,19 +170,79 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
for (ChainTunnel chainTunnel : chainTunnels1) {
int inx = i+1;
if (inx >= tunnelDto.getChainNodes().size()) { // 指向出口
GostUtil.AddChains(chainTunnel.getNodeId(), tunnelDto.getOutNodeId(), nodes);
GostDto gostDto = GostUtil.AddChains(chainTunnel.getNodeId(), tunnelDto.getOutNodeId(), nodes);
if (Objects.equals(gostDto.getMsg(), "OK")){
JSONObject data = new JSONObject();
data.put("node_id", chainTunnel.getNodeId());
data.put("name", "chains_" + tunnel.getId());
chain_success.add(data);
}else {
this.removeById(tunnel.getId());
chainTunnelService.remove(new QueryWrapper<ChainTunnel>().eq("tunnel_id", tunnel.getId()));
for (JSONObject chainSuccess : chain_success) {
GostDto deleteChains = GostUtil.DeleteChains(chainSuccess.getLong("node_id"), chainSuccess.getString("name"));
System.out.println(deleteChains);
}
return R.err(gostDto.getMsg());
}
} else {
GostUtil.AddChains(chainTunnel.getNodeId(), tunnelDto.getChainNodes().get(inx), nodes);
GostDto gostDto = GostUtil.AddChains(chainTunnel.getNodeId(), tunnelDto.getChainNodes().get(inx), nodes);
if (Objects.equals(gostDto.getMsg(), "OK")){
JSONObject data = new JSONObject();
data.put("node_id", chainTunnel.getNodeId());
data.put("name", "chains_" + tunnel.getId());
chain_success.add(data);
}else {
this.removeById(tunnel.getId());
chainTunnelService.remove(new QueryWrapper<ChainTunnel>().eq("tunnel_id", tunnel.getId()));
for (JSONObject chainSuccess : chain_success) {
GostDto deleteChains = GostUtil.DeleteChains(chainSuccess.getLong("node_id"), chainSuccess.getString("name"));
System.out.println(deleteChains);
}
return R.err(gostDto.getMsg());
}
}
GostUtil.AddChainService(chainTunnel.getNodeId(), chainTunnel, nodes);
GostDto gostDto = GostUtil.AddChainService(chainTunnel.getNodeId(), chainTunnel, nodes);
if (Objects.equals(gostDto.getMsg(), "OK")){
JSONObject data = new JSONObject();
data.put("node_id", chainTunnel.getNodeId());
data.put("name", tunnel.getId() + "_tls");
service_success.add(data);
}else {
this.removeById(tunnel.getId());
chainTunnelService.remove(new QueryWrapper<ChainTunnel>().eq("tunnel_id", tunnel.getId()));
for (JSONObject serviceSuccess : service_success) {
JSONArray jsonArray = new JSONArray();
jsonArray.add(serviceSuccess.getString("name"));
GostDto deleteService = GostUtil.DeleteService(serviceSuccess.getLong("node_id"), jsonArray);
System.out.println(deleteService);
}
return R.err(gostDto.getMsg());
}
}
}
for (ChainTunnel out_node : tunnelDto.getOutNodeId()) {
GostUtil.AddChainService(out_node.getNodeId(), out_node, nodes);
GostDto gostDto = GostUtil.AddChainService(out_node.getNodeId(), out_node, nodes);
if (Objects.equals(gostDto.getMsg(), "OK")){
JSONObject data = new JSONObject();
data.put("node_id", out_node.getNodeId());
data.put("name", tunnel.getId() + "_tls");
service_success.add(data);
}else {
this.removeById(tunnel.getId());
chainTunnelService.remove(new QueryWrapper<ChainTunnel>().eq("tunnel_id", tunnel.getId()));
for (JSONObject serviceSuccess : service_success) {
JSONArray jsonArray = new JSONArray();
jsonArray.add(serviceSuccess.getString("name"));
GostDto deleteService = GostUtil.DeleteService(serviceSuccess.getLong("node_id"), jsonArray);
System.out.println(deleteService);
}
return R.err(gostDto.getMsg());
}
}
}
@@ -286,12 +367,12 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
GostUtil.DeleteChains(chainTunnel.getNodeId(), "chains_" + chainTunnel.getTunnelId());
JSONArray services = new JSONArray();
services.add(chainTunnel.getTunnelId() + "_tls");
GostUtil.DeleteChainService(chainTunnel.getNodeId(), services);
GostUtil.DeleteService(chainTunnel.getNodeId(), services);
}
else { // 出口
JSONArray services = new JSONArray();
services.add(chainTunnel.getTunnelId() + "_tls");
GostUtil.DeleteChainService(chainTunnel.getNodeId(), services);
GostUtil.DeleteService(chainTunnel.getNodeId(), services);
}
}
chainTunnelService.remove(new QueryWrapper<ChainTunnel>().eq("tunnel_id", id));
@@ -471,8 +552,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
return R.ok(diagnosisReport);
}
@Override
public Integer getNodePort(Long nodeId,Integer type, Integer port) {
public Integer getNodePort(Long nodeId) {
Node node = nodeService.getById(nodeId);
if (node == null){
@@ -505,14 +585,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
if (availablePorts.isEmpty()) {
throw new RuntimeException("节点端口已满,无可用端口");
}
if (type == 1) {
return availablePorts.getLast();
}else {
if (port != null && availablePorts.contains(port)) {
return port;
}
return availablePorts.getFirst();
}
return availablePorts.getFirst();
}
public static List<Integer> parsePorts(String input) {
@@ -534,6 +607,10 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
return set.stream().sorted().collect(Collectors.toList());
}
private void isError(GostDto gostDto){
}
private DiagnosisResult performTcpPingDiagnosis(Node node, String targetIp, int port, String description) {
try {
// 构建TCP ping请求数据
@@ -264,44 +264,61 @@ public class UserServiceImpl extends ServiceImpl<UserMapper, User> implements Us
for (UserPackageDto.UserForwardDetailDto forward : forwards) {
Tunnel tunnel = tunnelService.getById(forward.getTunnelId());
if (tunnel == null) continue;
List<ForwardPort> forwardPorts = forwardPortService.list(
new QueryWrapper<ForwardPort>().eq("forward_id", forward.getId())
);
if (forwardPorts.isEmpty()) continue;
boolean useTunnelInIp = tunnel.getInIp() != null && !tunnel.getInIp().trim().isEmpty();
List<String> ipList = new ArrayList<>();
List<Integer> portList = new ArrayList<>();
java.util.Set<String> ipPortSet = new java.util.LinkedHashSet<>();
if (useTunnelInIp) {
// 使用隧道的inIp(求笛卡尔积)
List<String> ipList = new ArrayList<>();
List<Integer> portList = new ArrayList<>();
String[] tunnelInIps = tunnel.getInIp().split(",");
for (String ip : tunnelInIps) {
if (ip != null && !ip.trim().isEmpty()) {
ipList.add(ip.trim());
}
}
} else {
for (ForwardPort forwardPort : forwardPorts) {
Node node = nodeService.getById(forwardPort.getNodeId());
if (node != null && node.getServerIp() != null) {
ipList.add(node.getServerIp());
if (forwardPort.getPort() != null) {
portList.add(forwardPort.getPort());
}
}
}
for (ForwardPort forwardPort : forwardPorts) {
if (forwardPort.getPort() != null) {
portList.add(forwardPort.getPort());
}
}
List<String> uniqueIps = ipList.stream().distinct().toList();
List<Integer> uniquePorts = portList.stream().distinct().toList();
java.util.Set<String> ipPortSet = new java.util.LinkedHashSet<>();
for (String ip : uniqueIps) {
for (Integer port : uniquePorts) {
ipPortSet.add(ip + ":" + port);
List<String> uniqueIps = ipList.stream().distinct().toList();
List<Integer> uniquePorts = portList.stream().distinct().toList();
for (String ip : uniqueIps) {
for (Integer port : uniquePorts) {
ipPortSet.add(ip + ":" + port);
}
}
if (!uniquePorts.isEmpty()) {
forward.setInPort(uniquePorts.getFirst());
}
} else {
// 使用节点的serverIp(一对一,不求笛卡尔积)
for (ForwardPort forwardPort : forwardPorts) {
Node node = nodeService.getById(forwardPort.getNodeId());
if (node != null && node.getServerIp() != null && forwardPort.getPort() != null) {
ipPortSet.add(node.getServerIp() + ":" + forwardPort.getPort());
}
}
if (!forwardPorts.isEmpty() && forwardPorts.getFirst().getPort() != null) {
forward.setInPort(forwardPorts.getFirst().getPort());
}
}
if (!ipPortSet.isEmpty()) {
forward.setInIp(String.join(",", ipPortSet));
forward.setInPort(uniquePorts.getFirst());
}
}
}