增加流量不重置功能,ping改为tcpping

This commit is contained in:
qaq
2025-07-22 10:30:16 +08:00
parent 1448f96fa5
commit de4da12f42
17 changed files with 277 additions and 243 deletions
@@ -9,6 +9,7 @@ import org.springframework.web.servlet.HandlerInterceptor;
import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;
/**
* JWT拦截器,验证用户是否登录
*/
@@ -68,9 +68,12 @@ public class ResetFlowAsync {
*/
private void resetUserFlow(int currentDay, int lastDayOfMonth) {
try {
// flowResetTime字段存储的是1-31的数字,表示每月第几号重置
// flowResetTime字段存储的是0-31的数字,0表示不重置,1-31表示每月第几号重置
// 构建查询条件:重置日期等于今天,或者重置日期大于当月最大天数且今天是月末
// 排除flowResetTime为0的记录(不重置)
QueryWrapper<User> queryWrapper = new QueryWrapper<>();
queryWrapper.ne("flow_reset_time", 0); // 排除不重置的用户
if (currentDay == lastDayOfMonth) {
// 如果今天是月末,查询重置日期等于今天或者大于当月最大天数的记录
// 例如:当月30天,但用户设置31号重置,则在30号执行重置
@@ -118,9 +121,12 @@ public class ResetFlowAsync {
*/
private void resetUserTunnelFlow(int currentDay, int lastDayOfMonth) {
try {
// flowResetTime字段存储的是1-31的数字,表示每月第几号重置
// flowResetTime字段存储的是0-31的数字,0表示不重置,1-31表示每月第几号重置
// 构建查询条件:重置日期等于今天,或者重置日期大于当月最大天数且今天是月末
// 排除flowResetTime为0的记录(不重置)
QueryWrapper<UserTunnel> queryWrapper = new QueryWrapper<>();
queryWrapper.ne("flow_reset_time", 0); // 排除不重置的用户隧道
if (currentDay == lastDayOfMonth) {
// 如果今天是月末,查询重置日期等于今天或者大于当月最大天数的记录
// 例如:当月30天,但用户设置31号重置,则在30号执行重置
@@ -21,6 +21,7 @@ import java.util.Objects;
@Configuration
@Slf4j
public class WebSocketInterceptor extends HttpSessionHandshakeInterceptor {
@Resource
NodeService nodeService;
@@ -411,9 +411,10 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
String[] remoteAddresses = forward.getRemoteAddr().split(",");
String targetAddress = remoteAddresses[0].trim();
// 提取IP部分(去掉端口)
// 提取IP和端口
String targetIp = extractIpFromAddress(targetAddress);
if (targetIp == null) {
int targetPort = extractPortFromAddress(targetAddress);
if (targetIp == null || targetPort == -1) {
return R.err("无法解析目标地址: " + targetAddress);
}
@@ -421,22 +422,22 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
// 6. 根据隧道类型执行不同的诊断策略
if (tunnel.getType() == TUNNEL_TYPE_PORT_FORWARD) {
// 端口转发:入口节点直接ping目标地址
DiagnosisResult result = performPingDiagnosis(inNode, targetIp, "转发->目标");
// 端口转发:入口节点直接TCP ping目标地址
DiagnosisResult result = performTcpPingDiagnosis(inNode, targetIp, targetPort, "转发->目标");
results.add(result);
} else {
// 隧道转发:入口ping出口,出口ping目标
// 隧道转发:入口TCP ping出口,出口TCP ping目标
Node outNode = nodeService.getNodeById(tunnel.getOutNodeId());
if (outNode == null) {
return R.err("出口节点不存在");
}
// 入口ping出口
DiagnosisResult inToOutResult = performPingDiagnosis(inNode, outNode.getServerIp(), "入口->出口");
// 入口TCP ping出口(使用转发的出口端口)
DiagnosisResult inToOutResult = performTcpPingDiagnosis(inNode, outNode.getServerIp(), forward.getOutPort(), "入口->出口");
results.add(inToOutResult);
// 出口ping目标
DiagnosisResult outToTargetResult = performPingDiagnosis(outNode, targetIp, "出口->目标");
// 出口TCP ping目标
DiagnosisResult outToTargetResult = performTcpPingDiagnosis(outNode, targetIp, targetPort, "出口->目标");
results.add(outToTargetResult);
}
@@ -481,58 +482,101 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
}
/**
* 执行ping诊断
* 从地址字符串中提取端口号
* 支持格式: ip:port, [ipv6]:port, domain:port
*/
private int extractPortFromAddress(String address) {
if (address == null || address.trim().isEmpty()) {
return -1;
}
address = address.trim();
// IPv6格式: [ipv6]:port
if (address.startsWith("[")) {
int closeBracket = address.indexOf(']');
if (closeBracket > 1 && closeBracket + 1 < address.length() && address.charAt(closeBracket + 1) == ':') {
String portStr = address.substring(closeBracket + 2);
try {
return Integer.parseInt(portStr);
} catch (NumberFormatException e) {
return -1;
}
}
}
// IPv4或域名格式: ip:port 或 domain:port
int lastColon = address.lastIndexOf(':');
if (lastColon > 0 && lastColon + 1 < address.length()) {
String portStr = address.substring(lastColon + 1);
try {
return Integer.parseInt(portStr);
} catch (NumberFormatException e) {
return -1;
}
}
// 如果没有端口,返回-1表示无法解析
return -1;
}
/**
* 执行TCP ping诊断
*
* @param node 执行ping的节点
* @param node 执行TCP ping的节点
* @param targetIp 目标IP地址
* @param port 目标端口
* @param description 诊断描述
* @return 诊断结果
*/
private DiagnosisResult performPingDiagnosis(Node node, String targetIp, String description) {
private DiagnosisResult performTcpPingDiagnosis(Node node, String targetIp, int port, String description) {
try {
// 构建ping请求数据
JSONObject pingData = new JSONObject();
pingData.put("ip", targetIp);
pingData.put("count", 4);
// 构建TCP ping请求数据
JSONObject tcpPingData = new JSONObject();
tcpPingData.put("ip", targetIp);
tcpPingData.put("port", port);
tcpPingData.put("count", 4);
tcpPingData.put("timeout", 5000); // 5秒超时
// 发送ping命令到节点
GostDto gostResult = WebSocketServer.send_msg(node.getId(), pingData, "Ping");
// 发送TCP ping命令到节点
GostDto gostResult = WebSocketServer.send_msg(node.getId(), tcpPingData, "TcpPing");
DiagnosisResult result = new DiagnosisResult();
result.setNodeId(node.getId());
result.setNodeName(node.getName());
result.setTargetIp(targetIp);
result.setTargetPort(port);
result.setDescription(description);
result.setTimestamp(System.currentTimeMillis());
if (gostResult != null && "OK".equals(gostResult.getMsg())) {
// 尝试解析ping响应数据
// 尝试解析TCP ping响应数据
try {
if (gostResult.getData() != null) {
JSONObject pingResponse = (JSONObject) gostResult.getData();
boolean success = pingResponse.getBooleanValue("success");
JSONObject tcpPingResponse = (JSONObject) gostResult.getData();
boolean success = tcpPingResponse.getBooleanValue("success");
result.setSuccess(success);
if (success) {
result.setMessage("ping成功");
result.setAverageTime(pingResponse.getDoubleValue("averageTime"));
result.setPacketLoss(pingResponse.getDoubleValue("packetLoss"));
result.setMessage("TCP连接成功");
result.setAverageTime(tcpPingResponse.getDoubleValue("averageTime"));
result.setPacketLoss(tcpPingResponse.getDoubleValue("packetLoss"));
} else {
result.setMessage(pingResponse.getString("errorMessage"));
result.setMessage(tcpPingResponse.getString("errorMessage"));
result.setAverageTime(-1.0);
result.setPacketLoss(100.0);
}
} else {
// 没有详细数据,使用默认值
result.setSuccess(true);
result.setMessage("ping成功");
result.setMessage("TCP连接成功");
result.setAverageTime(0.0);
result.setPacketLoss(0.0);
}
} catch (Exception e) {
// 解析响应数据失败,但ping命令本身成功了
// 解析响应数据失败,但TCP ping命令本身成功了
result.setSuccess(true);
result.setMessage("ping成功,但无法解析详细数据");
result.setMessage("TCP连接成功,但无法解析详细数据");
result.setAverageTime(0.0);
result.setPacketLoss(0.0);
}
@@ -549,6 +593,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
result.setNodeId(node.getId());
result.setNodeName(node.getName());
result.setTargetIp(targetIp);
result.setTargetPort(port);
result.setDescription(description);
result.setSuccess(false);
result.setMessage("诊断执行异常: " + e.getMessage());
@@ -1326,6 +1371,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
private Long nodeId;
private String nodeName;
private String targetIp;
private Integer targetPort;
private String description;
private boolean success;
private String message;
@@ -643,16 +643,17 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
// 3. 根据隧道类型执行不同的诊断策略
if (tunnel.getType() == TUNNEL_TYPE_PORT_FORWARD) {
// 端口转发:只给入口节点发送诊断指令,ping谷歌DNS
DiagnosisResult inResult = performPingDiagnosisWithConnectionCheck(inNode, "www.google.com", "入口->外网");
// 端口转发:只给入口节点发送诊断指令,TCP ping谷歌443端口
DiagnosisResult inResult = performTcpPingDiagnosisWithConnectionCheck(inNode, "www.google.com", 443, "入口->外网");
results.add(inResult);
} else {
// 隧道转发:入口ping出口,出口ping谷歌DNS
DiagnosisResult inToOutResult = performPingDiagnosisWithConnectionCheck(inNode, outNode.getServerIp(), "入口->出口");
// 隧道转发:入口TCP ping出口,出口TCP ping谷歌443端口
int outNodePort = getOutNodeTcpPort(tunnel.getId());
DiagnosisResult inToOutResult = performTcpPingDiagnosisWithConnectionCheck(inNode, outNode.getServerIp(), outNodePort, "入口->出口");
results.add(inToOutResult);
// 先检查出口节点的真实连接状态,然后再进行诊断
DiagnosisResult outToExternalResult = performPingDiagnosisWithConnectionCheck(outNode, "www.google.com", "出口->外网");
DiagnosisResult outToExternalResult = performTcpPingDiagnosisWithConnectionCheck(outNode, "www.google.com", 443, "出口->外网");
results.add(outToExternalResult);
}
@@ -668,58 +669,78 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
}
/**
* 执行ping诊断
* 获取出口节点的TCP端口
* 通过隧道ID查找转发服务的出口端口,如果没有则使用默认SSH端口22
*
* @param node 执行ping的节点
* @param tunnelId 隧道ID
* @return TCP端口号
*/
private int getOutNodeTcpPort(Long tunnelId) {
List<Forward> forwards = forwardService.list(new QueryWrapper<Forward>().eq("tunnel_id", tunnelId));
if (!forwards.isEmpty()) {
return forwards.get(0).getOutPort();
}
// 如果没有转发服务,使用默认SSH端口22
return 22;
}
/**
* 执行TCP ping诊断
*
* @param node 执行TCP ping的节点
* @param targetIp 目标IP地址
* @param port 目标端口
* @param description 诊断描述
* @return 诊断结果
*/
private DiagnosisResult performPingDiagnosis(Node node, String targetIp, String description) {
private DiagnosisResult performTcpPingDiagnosis(Node node, String targetIp, int port, String description) {
try {
// 构建ping请求数据
JSONObject pingData = new JSONObject();
pingData.put("ip", targetIp);
pingData.put("count", 4);
// 构建TCP ping请求数据
JSONObject tcpPingData = new JSONObject();
tcpPingData.put("ip", targetIp);
tcpPingData.put("port", port);
tcpPingData.put("count", 4);
tcpPingData.put("timeout", 5000); // 5秒超时
// 发送ping命令到节点
GostDto gostResult = WebSocketServer.send_msg(node.getId(), pingData, "Ping");
// 发送TCP ping命令到节点
GostDto gostResult = WebSocketServer.send_msg(node.getId(), tcpPingData, "TcpPing");
DiagnosisResult result = new DiagnosisResult();
result.setNodeId(node.getId());
result.setNodeName(node.getName());
result.setTargetIp(targetIp);
result.setTargetPort(port);
result.setDescription(description);
result.setTimestamp(System.currentTimeMillis());
if (gostResult != null && "OK".equals(gostResult.getMsg())) {
// 尝试解析ping响应数据
// 尝试解析TCP ping响应数据
try {
if (gostResult.getData() != null) {
JSONObject pingResponse = (JSONObject) gostResult.getData();
boolean success = pingResponse.getBooleanValue("success");
JSONObject tcpPingResponse = (JSONObject) gostResult.getData();
boolean success = tcpPingResponse.getBooleanValue("success");
result.setSuccess(success);
if (success) {
result.setMessage("ping成功");
result.setAverageTime(pingResponse.getDoubleValue("averageTime"));
result.setPacketLoss(pingResponse.getDoubleValue("packetLoss"));
result.setMessage("TCP连接成功");
result.setAverageTime(tcpPingResponse.getDoubleValue("averageTime"));
result.setPacketLoss(tcpPingResponse.getDoubleValue("packetLoss"));
} else {
result.setMessage(pingResponse.getString("errorMessage"));
result.setMessage(tcpPingResponse.getString("errorMessage"));
result.setAverageTime(-1.0);
result.setPacketLoss(100.0);
}
} else {
// 没有详细数据,使用默认值
result.setSuccess(true);
result.setMessage("ping成功");
result.setMessage("TCP连接成功");
result.setAverageTime(0.0);
result.setPacketLoss(0.0);
}
} catch (Exception e) {
// 解析响应数据失败,但ping命令本身成功了
// 解析响应数据失败,但TCP ping命令本身成功了
result.setSuccess(true);
result.setMessage("ping成功,但无法解析详细数据");
result.setMessage("TCP连接成功,但无法解析详细数据");
result.setAverageTime(0.0);
result.setPacketLoss(0.0);
}
@@ -736,6 +757,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
result.setNodeId(node.getId());
result.setNodeName(node.getName());
result.setTargetIp(targetIp);
result.setTargetPort(port);
result.setDescription(description);
result.setSuccess(false);
result.setMessage("诊断执行异常: " + e.getMessage());
@@ -747,23 +769,25 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
}
/**
* 执行ping诊断(带连接状态检查)
* 执行TCP ping诊断(带连接状态检查)
*
* @param node 执行ping的节点
* @param node 执行TCP ping的节点
* @param targetIp 目标IP地址
* @param port 目标端口
* @param description 诊断描述
* @return 诊断结果
*/
private DiagnosisResult performPingDiagnosisWithConnectionCheck(Node node, String targetIp, String description) {
private DiagnosisResult performTcpPingDiagnosisWithConnectionCheck(Node node, String targetIp, int port, String description) {
DiagnosisResult result = new DiagnosisResult();
result.setNodeId(node.getId());
result.setNodeName(node.getName());
result.setTargetIp(targetIp);
result.setTargetPort(port);
result.setDescription(description);
result.setTimestamp(System.currentTimeMillis());
try {
return performPingDiagnosis(node, targetIp, description);
return performTcpPingDiagnosis(node, targetIp, port, description);
} catch (Exception e) {
result.setSuccess(false);
result.setMessage("连接检查异常: " + e.getMessage());
@@ -817,6 +841,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
private Long nodeId;
private String nodeName;
private String targetIp;
private Integer targetPort;
private String description;
private boolean success;
private String message;