mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-29 16:06:36 +08:00
优化ui 添加转发诊断 移除vue前端
This commit is contained in:
@@ -0,0 +1,52 @@
|
||||
package com.admin.common.task;
|
||||
|
||||
import com.admin.common.utils.WebSocketServer;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* WebSocket资源清理定时任务
|
||||
* 定期清理已完成的请求、无效的session锁等,防止内存泄漏
|
||||
*/
|
||||
@Slf4j
|
||||
@Component
|
||||
public class WebSocketCleanupTask {
|
||||
|
||||
/**
|
||||
* 每5分钟执行一次轻量级清理
|
||||
* 清理已完成的请求和无效的session锁
|
||||
*/
|
||||
@Scheduled(fixedRate = 5 * 60 * 1000) // 5分钟
|
||||
public void lightweightCleanup() {
|
||||
try {
|
||||
log.debug("开始执行WebSocket轻量级清理...");
|
||||
|
||||
// 清理已完成的请求
|
||||
WebSocketServer.cleanupCompletedRequests();
|
||||
|
||||
// 清理无效的session锁
|
||||
WebSocketServer.cleanupInvalidSessionLocks();
|
||||
|
||||
} catch (Exception e) {
|
||||
log.error("WebSocket轻量级清理失败", e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 每30分钟执行一次全面清理
|
||||
* 包括所有类型的资源清理
|
||||
*/
|
||||
@Scheduled(fixedRate = 30 * 60 * 1000) // 30分钟
|
||||
public void fullCleanup() {
|
||||
try {
|
||||
log.info("开始执行WebSocket全面清理...");
|
||||
|
||||
// 执行全面清理
|
||||
WebSocketServer.performFullCleanup();
|
||||
|
||||
} catch (Exception e) {
|
||||
log.error("WebSocket全面清理失败", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -41,6 +41,9 @@ public class WebSocketServer extends TextWebSocketHandler {
|
||||
|
||||
// 存储等待响应的请求,key为requestId,value为CompletableFuture
|
||||
private static final ConcurrentHashMap<String, CompletableFuture<GostDto>> pendingRequests = new ConcurrentHashMap<>();
|
||||
|
||||
// 存储请求ID与节点ID的映射关系,用于清理特定节点的请求
|
||||
private static final ConcurrentHashMap<String, Long> requestNodeMapping = new ConcurrentHashMap<>();
|
||||
|
||||
//接受客户端消息
|
||||
@Override
|
||||
@@ -66,6 +69,9 @@ public class WebSocketServer extends TextWebSocketHandler {
|
||||
|
||||
if (requestId != null) {
|
||||
CompletableFuture<GostDto> future = pendingRequests.remove(requestId);
|
||||
// 同时清理请求-节点映射关系
|
||||
requestNodeMapping.remove(requestId);
|
||||
|
||||
if (future != null) {
|
||||
GostDto result = new GostDto();
|
||||
|
||||
@@ -123,29 +129,60 @@ public class WebSocketServer extends TextWebSocketHandler {
|
||||
if (!Objects.equals(type, "1")) {
|
||||
// 网页管理员连接
|
||||
activeSessions.add(session);
|
||||
log.info("管理员连接建立,sessionId: {}", session.getId());
|
||||
} else {
|
||||
// 客户端节点连接
|
||||
Long nodeId = Long.valueOf(id);
|
||||
String version = (String) session.getAttributes().get("nodeVersion");
|
||||
|
||||
log.info("节点 {} 连接建立,开始更新状态", nodeId);
|
||||
|
||||
// 先添加到会话映射
|
||||
nodeSessions.put(nodeId, session);
|
||||
|
||||
// 更新节点状态为在线
|
||||
Node byId = nodeService.getById(nodeId);
|
||||
if (byId != null) {
|
||||
byId.setStatus(1);
|
||||
nodeService.updateById(byId);
|
||||
Node node = nodeService.getById(nodeId);
|
||||
if (node != null) {
|
||||
// 更新状态和版本信息
|
||||
node.setStatus(1);
|
||||
if (version != null) {
|
||||
node.setVersion(version);
|
||||
}
|
||||
boolean updateResult = nodeService.updateById(node);
|
||||
|
||||
// 广播节点上线状态给所有管理员
|
||||
JSONObject res = new JSONObject();
|
||||
res.put("id", id);
|
||||
res.put("type", "status");
|
||||
res.put("data", 1);
|
||||
broadcastMessage(res.toJSONString());
|
||||
if (updateResult) {
|
||||
log.info("节点 {} 状态更新为在线成功,版本: {}", nodeId, version);
|
||||
|
||||
// 广播节点上线状态给所有管理员
|
||||
JSONObject res = new JSONObject();
|
||||
res.put("id", id);
|
||||
res.put("type", "status");
|
||||
res.put("data", 1);
|
||||
broadcastMessage(res.toJSONString());
|
||||
} else {
|
||||
log.error("节点 {} 状态更新失败", nodeId);
|
||||
}
|
||||
} else {
|
||||
log.error("节点 {} 不存在,无法更新状态", nodeId);
|
||||
// 移除无效的会话
|
||||
nodeSessions.remove(nodeId);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
} catch (Exception e) {
|
||||
log.error("建立连接时发生异常: {}", e.getMessage(), e);
|
||||
// 异常情况下,确保清理会话
|
||||
try {
|
||||
String id = session.getAttributes().get("id").toString();
|
||||
String type = session.getAttributes().get("type").toString();
|
||||
if (Objects.equals(type, "1")) {
|
||||
Long nodeId = Long.valueOf(id);
|
||||
nodeSessions.remove(nodeId);
|
||||
log.warn("由于异常,移除节点 {} 的会话", nodeId);
|
||||
}
|
||||
} catch (Exception cleanupException) {
|
||||
log.error("清理异常会话时出错: {}", cleanupException.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -157,36 +194,46 @@ public class WebSocketServer extends TextWebSocketHandler {
|
||||
String type = session.getAttributes().get("type").toString();
|
||||
String sessionId = session.getId();
|
||||
|
||||
log.info("连接关闭,ID: {}, 类型: {}, 状态: {}", id, type, status);
|
||||
|
||||
if (!Objects.equals(type, "1")) {
|
||||
// 连接关闭
|
||||
activeSessions.remove(session);
|
||||
// 管理员连接关闭
|
||||
boolean removed = activeSessions.remove(session);
|
||||
log.info("管理员连接关闭,sessionId: {}, 移除结果: {}", sessionId, removed);
|
||||
} else {
|
||||
// 客户端节点连接关闭
|
||||
Long nodeId = Long.valueOf(id);
|
||||
nodeSessions.remove(nodeId);
|
||||
WebSocketSession removedSession = nodeSessions.remove(nodeId);
|
||||
|
||||
log.info("节点 {} 连接关闭,开始更新状态为离线", nodeId);
|
||||
|
||||
// 更新节点状态为离线
|
||||
Node byId = nodeService.getById(nodeId);
|
||||
if (byId != null) {
|
||||
byId.setStatus(0);
|
||||
nodeService.updateById(byId);
|
||||
Node node = nodeService.getById(nodeId);
|
||||
if (node != null) {
|
||||
node.setStatus(0);
|
||||
boolean updateResult = nodeService.updateById(node);
|
||||
|
||||
JSONObject res = new JSONObject();
|
||||
res.put("id", id);
|
||||
res.put("type", "status");
|
||||
res.put("data", 0);
|
||||
broadcastMessage(res.toJSONString());
|
||||
if (updateResult) {
|
||||
log.info("节点 {} 状态更新为离线成功", nodeId);
|
||||
|
||||
JSONObject res = new JSONObject();
|
||||
res.put("id", id);
|
||||
res.put("type", "status");
|
||||
res.put("data", 0);
|
||||
broadcastMessage(res.toJSONString());
|
||||
} else {
|
||||
log.error("节点 {} 状态更新为离线失败", nodeId);
|
||||
}
|
||||
} else {
|
||||
log.warn("节点 {} 不存在,无法更新离线状态", nodeId);
|
||||
}
|
||||
|
||||
|
||||
// 清理该节点的待处理请求
|
||||
clearPendingRequestsForNode(nodeId);
|
||||
}
|
||||
|
||||
// 清理session锁对象
|
||||
sessionLocks.remove(sessionId);
|
||||
|
||||
// 清理该节点的待处理请求
|
||||
if (Objects.equals(type, "1")) {
|
||||
clearPendingRequestsForNode(Long.valueOf(id));
|
||||
}
|
||||
|
||||
} catch (Exception e) {
|
||||
log.error("关闭连接时发生异常: {}", e.getMessage(), e);
|
||||
@@ -237,6 +284,50 @@ public class WebSocketServer extends TextWebSocketHandler {
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 清理无效的session锁(定期清理任务)
|
||||
*/
|
||||
public static void cleanupInvalidSessionLocks() {
|
||||
java.util.List<String> invalidSessionIds = new java.util.ArrayList<>();
|
||||
|
||||
sessionLocks.keySet().forEach(sessionId -> {
|
||||
boolean isValidSession = false;
|
||||
|
||||
// 检查是否为有效的管理员session
|
||||
for (WebSocketSession adminSession : activeSessions) {
|
||||
if (adminSession != null && sessionId.equals(adminSession.getId())) {
|
||||
isValidSession = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
// 检查是否为有效的节点session
|
||||
if (!isValidSession) {
|
||||
for (WebSocketSession nodeSession : nodeSessions.values()) {
|
||||
if (nodeSession != null && sessionId.equals(nodeSession.getId())) {
|
||||
isValidSession = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (!isValidSession) {
|
||||
invalidSessionIds.add(sessionId);
|
||||
}
|
||||
});
|
||||
|
||||
int cleanedCount = 0;
|
||||
for (String sessionId : invalidSessionIds) {
|
||||
if (sessionLocks.remove(sessionId) != null) {
|
||||
cleanedCount++;
|
||||
}
|
||||
}
|
||||
|
||||
if (cleanedCount > 0) {
|
||||
log.info("清理了 {} 个无效的session锁", cleanedCount);
|
||||
}
|
||||
}
|
||||
|
||||
// 广播消息
|
||||
public static void broadcastMessage(String message) {
|
||||
@@ -247,31 +338,254 @@ public class WebSocketServer extends TextWebSocketHandler {
|
||||
|
||||
/**
|
||||
* 清理指定节点的待处理请求
|
||||
* 只清理属于该节点的未完成请求,避免影响其他节点的请求
|
||||
*/
|
||||
private static void clearPendingRequestsForNode(Long nodeId) {
|
||||
// 完成所有待处理的请求,设置为连接断开错误
|
||||
pendingRequests.entrySet().removeIf(entry -> {
|
||||
CompletableFuture<GostDto> future = entry.getValue();
|
||||
if (!future.isDone()) {
|
||||
GostDto errorResult = new GostDto();
|
||||
errorResult.setMsg("节点连接已断开");
|
||||
future.complete(errorResult);
|
||||
if (nodeId == null) {
|
||||
log.warn("节点ID为空,无法清理待处理请求");
|
||||
return;
|
||||
}
|
||||
|
||||
java.util.concurrent.atomic.AtomicInteger clearedCount = new java.util.concurrent.atomic.AtomicInteger(0);
|
||||
java.util.List<String> requestIdsToRemove = new java.util.ArrayList<>();
|
||||
|
||||
// 找出属于该节点的请求
|
||||
requestNodeMapping.entrySet().forEach(entry -> {
|
||||
String requestId = entry.getKey();
|
||||
Long mappedNodeId = entry.getValue();
|
||||
|
||||
if (nodeId.equals(mappedNodeId)) {
|
||||
CompletableFuture<GostDto> future = pendingRequests.get(requestId);
|
||||
if (future != null && !future.isDone()) {
|
||||
// 完成该请求并设置错误信息
|
||||
GostDto errorResult = new GostDto();
|
||||
errorResult.setMsg("节点连接已断开");
|
||||
future.complete(errorResult);
|
||||
clearedCount.incrementAndGet();
|
||||
}
|
||||
requestIdsToRemove.add(requestId);
|
||||
}
|
||||
return true; // 移除所有请求
|
||||
});
|
||||
|
||||
// 批量清理映射关系和请求
|
||||
requestIdsToRemove.forEach(requestId -> {
|
||||
pendingRequests.remove(requestId);
|
||||
requestNodeMapping.remove(requestId);
|
||||
});
|
||||
|
||||
if (clearedCount.get() > 0) {
|
||||
log.info("清理了节点 {} 的 {} 个待处理请求", nodeId, clearedCount.get());
|
||||
} else {
|
||||
log.debug("节点 {} 没有待处理的请求需要清理", nodeId);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* 检查节点的实际连接状态,如果状态不一致则修复
|
||||
*/
|
||||
public static boolean checkAndFixNodeStatus(NodeService nodeService, Long nodeId) {
|
||||
try {
|
||||
WebSocketSession session = nodeSessions.get(nodeId);
|
||||
boolean isConnected = session != null && session.isOpen();
|
||||
|
||||
Node node = nodeService.getById(nodeId);
|
||||
if (node != null) {
|
||||
int currentStatus = node.getStatus();
|
||||
int expectedStatus = isConnected ? 1 : 0;
|
||||
|
||||
if (currentStatus != expectedStatus) {
|
||||
log.warn("节点 {} 状态不一致,数据库状态: {}, 实际连接状态: {}, 正在修复...",
|
||||
nodeId, currentStatus, expectedStatus);
|
||||
|
||||
node.setStatus(expectedStatus);
|
||||
boolean updateResult = nodeService.updateById(node);
|
||||
|
||||
if (updateResult) {
|
||||
log.info("节点 {} 状态修复成功,更新为: {}", nodeId, expectedStatus);
|
||||
|
||||
// 广播状态变更
|
||||
JSONObject res = new JSONObject();
|
||||
res.put("id", nodeId.toString());
|
||||
res.put("type", "status");
|
||||
res.put("data", expectedStatus);
|
||||
broadcastMessage(res.toJSONString());
|
||||
|
||||
return true;
|
||||
} else {
|
||||
log.error("节点 {} 状态修复失败", nodeId);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return isConnected;
|
||||
} catch (Exception e) {
|
||||
log.error("检查节点 {} 状态时发生异常: {}", nodeId, e.getMessage(), e);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取节点的实际连接状态
|
||||
*/
|
||||
public static boolean isNodeConnected(Long nodeId) {
|
||||
WebSocketSession session = nodeSessions.get(nodeId);
|
||||
return session != null && session.isOpen();
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取所有在线节点的ID列表
|
||||
*/
|
||||
public static java.util.Set<Long> getConnectedNodeIds() {
|
||||
return nodeSessions.entrySet().stream()
|
||||
.filter(entry -> entry.getValue() != null && entry.getValue().isOpen())
|
||||
.map(java.util.Map.Entry::getKey)
|
||||
.collect(java.util.stream.Collectors.toSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取指定节点的待处理请求数量
|
||||
*/
|
||||
public static int getPendingRequestCount(Long nodeId) {
|
||||
if (nodeId == null) {
|
||||
return 0;
|
||||
}
|
||||
return (int) requestNodeMapping.entrySet().stream()
|
||||
.filter(entry -> nodeId.equals(entry.getValue()))
|
||||
.map(java.util.Map.Entry::getKey)
|
||||
.filter(requestId -> {
|
||||
CompletableFuture<GostDto> future = pendingRequests.get(requestId);
|
||||
return future != null && !future.isDone();
|
||||
})
|
||||
.count();
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取所有待处理请求的总数
|
||||
*/
|
||||
public static int getTotalPendingRequestCount() {
|
||||
return (int) pendingRequests.entrySet().stream()
|
||||
.filter(entry -> entry.getValue() != null && !entry.getValue().isDone())
|
||||
.count();
|
||||
}
|
||||
|
||||
/**
|
||||
* 清理所有已完成但未被移除的请求(定期清理任务)
|
||||
*/
|
||||
public static void cleanupCompletedRequests() {
|
||||
java.util.List<String> completedRequestIds = new java.util.ArrayList<>();
|
||||
|
||||
pendingRequests.entrySet().forEach(entry -> {
|
||||
String requestId = entry.getKey();
|
||||
CompletableFuture<GostDto> future = entry.getValue();
|
||||
|
||||
if (future != null && future.isDone()) {
|
||||
completedRequestIds.add(requestId);
|
||||
}
|
||||
});
|
||||
|
||||
int cleanedCount = 0;
|
||||
for (String requestId : completedRequestIds) {
|
||||
if (pendingRequests.remove(requestId) != null) {
|
||||
requestNodeMapping.remove(requestId);
|
||||
cleanedCount++;
|
||||
}
|
||||
}
|
||||
|
||||
if (cleanedCount > 0) {
|
||||
log.info("清理了 {} 个已完成的请求", cleanedCount);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取WebSocket连接统计信息
|
||||
*/
|
||||
public static java.util.Map<String, Object> getConnectionStats() {
|
||||
java.util.Map<String, Object> stats = new java.util.HashMap<>();
|
||||
|
||||
// 节点连接统计
|
||||
int totalNodes = nodeSessions.size();
|
||||
int onlineNodes = (int) nodeSessions.entrySet().stream()
|
||||
.filter(entry -> entry.getValue() != null && entry.getValue().isOpen())
|
||||
.count();
|
||||
|
||||
// 管理员连接统计
|
||||
int adminConnections = activeSessions.size();
|
||||
|
||||
// 请求统计
|
||||
int totalPendingRequests = getTotalPendingRequestCount();
|
||||
int totalMappings = requestNodeMapping.size();
|
||||
|
||||
stats.put("totalNodes", totalNodes);
|
||||
stats.put("onlineNodes", onlineNodes);
|
||||
stats.put("offlineNodes", totalNodes - onlineNodes);
|
||||
stats.put("adminConnections", adminConnections);
|
||||
stats.put("totalPendingRequests", totalPendingRequests);
|
||||
stats.put("totalRequestMappings", totalMappings);
|
||||
|
||||
return stats;
|
||||
}
|
||||
|
||||
/**
|
||||
* 执行全面的内存清理操作
|
||||
* 建议定期调用以防止内存泄漏
|
||||
*/
|
||||
public static void performFullCleanup() {
|
||||
log.info("开始执行WebSocket全面清理操作...");
|
||||
|
||||
// 获取清理前的统计信息
|
||||
java.util.Map<String, Object> statsBefore = getConnectionStats();
|
||||
|
||||
// 执行各种清理操作
|
||||
cleanupCompletedRequests();
|
||||
cleanupInvalidSessionLocks();
|
||||
|
||||
// 清理失效的节点session
|
||||
java.util.List<Long> invalidNodeIds = new java.util.ArrayList<>();
|
||||
nodeSessions.entrySet().forEach(entry -> {
|
||||
WebSocketSession session = entry.getValue();
|
||||
if (session == null || !session.isOpen()) {
|
||||
invalidNodeIds.add(entry.getKey());
|
||||
}
|
||||
});
|
||||
|
||||
int removedNodeSessions = 0;
|
||||
for (Long nodeId : invalidNodeIds) {
|
||||
if (nodeSessions.remove(nodeId) != null) {
|
||||
removedNodeSessions++;
|
||||
}
|
||||
}
|
||||
|
||||
// 清理失效的管理员session
|
||||
int removedAdminSessions = 0;
|
||||
java.util.Iterator<WebSocketSession> adminIterator = activeSessions.iterator();
|
||||
while (adminIterator.hasNext()) {
|
||||
WebSocketSession session = adminIterator.next();
|
||||
if (session == null || !session.isOpen()) {
|
||||
adminIterator.remove();
|
||||
removedAdminSessions++;
|
||||
}
|
||||
}
|
||||
|
||||
// 获取清理后的统计信息
|
||||
java.util.Map<String, Object> statsAfter = getConnectionStats();
|
||||
|
||||
log.info("WebSocket清理完成 - 清理前: {}, 清理后: {}, 移除节点session: {}, 移除管理员session: {}",
|
||||
statsBefore, statsAfter, removedNodeSessions, removedAdminSessions);
|
||||
}
|
||||
|
||||
public static GostDto send_msg(Long node_id, Object msg, String type) {
|
||||
WebSocketSession nodeSession = nodeSessions.get(node_id);
|
||||
|
||||
if (nodeSession == null) {
|
||||
log.warn("发送消息失败:节点 {} 不在线或会话不存在", node_id);
|
||||
GostDto result = new GostDto();
|
||||
result.setMsg("节点不在线");
|
||||
return result;
|
||||
}
|
||||
|
||||
if (!nodeSession.isOpen()) {
|
||||
log.warn("发送消息失败:节点 {} 连接已断开,清理会话", node_id);
|
||||
nodeSessions.remove(node_id);
|
||||
sessionLocks.remove(nodeSession.getId());
|
||||
GostDto result = new GostDto();
|
||||
@@ -285,6 +599,9 @@ public class WebSocketServer extends TextWebSocketHandler {
|
||||
// 创建CompletableFuture用于等待响应
|
||||
CompletableFuture<GostDto> future = new CompletableFuture<>();
|
||||
pendingRequests.put(requestId, future);
|
||||
|
||||
// 建立请求ID与节点ID的映射关系
|
||||
requestNodeMapping.put(requestId, node_id);
|
||||
|
||||
try {
|
||||
JSONObject data = new JSONObject();
|
||||
@@ -293,17 +610,23 @@ public class WebSocketServer extends TextWebSocketHandler {
|
||||
data.put("requestId", requestId);
|
||||
sendToUser(nodeSession, data.toJSONString());
|
||||
GostDto result = future.get(10, TimeUnit.SECONDS);
|
||||
|
||||
log.debug("成功发送消息到节点 {} 并收到响应: {}", node_id, result.getMsg());
|
||||
return result;
|
||||
|
||||
} catch (Exception e) {
|
||||
// 清理请求和映射关系
|
||||
pendingRequests.remove(requestId);
|
||||
requestNodeMapping.remove(requestId);
|
||||
|
||||
GostDto result = new GostDto();
|
||||
if (e instanceof java.util.concurrent.TimeoutException) {
|
||||
result.setMsg("等待响应超时");
|
||||
log.warn("节点 {} 响应超时,可能存在连接问题", node_id);
|
||||
} else {
|
||||
result.setMsg("发送消息失败: " + e.getMessage());
|
||||
log.error("发送消息到节点 {} 失败: {}", node_id, e.getMessage(), e);
|
||||
}
|
||||
log.error("发送消息到节点{}失败: {}", node_id, e.getMessage(), e);
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -37,11 +37,15 @@ public class WebSocketInterceptor extends HttpSessionHandshakeInterceptor {
|
||||
String version = serverHttpRequest.getServletRequest().getParameter("version");
|
||||
if (Objects.equals(type, "1")) {
|
||||
Node node = nodeService.getOne(new QueryWrapper<Node>().eq("secret", secret));
|
||||
if (node == null) return false;
|
||||
if (node == null) {
|
||||
log.warn("节点验证失败:未找到匹配的secret");
|
||||
return false;
|
||||
}
|
||||
attributes.put("id", node.getId());
|
||||
node.setStatus(1);
|
||||
node.setVersion(version);
|
||||
nodeService.updateById(node);
|
||||
attributes.put("nodeSecret", secret);
|
||||
attributes.put("nodeVersion", version);
|
||||
log.info("节点 {} 通过验证,版本: {}", node.getId(), version);
|
||||
// 不在这里更新状态,等到连接建立后再统一更新
|
||||
}else {
|
||||
boolean b = JwtUtil.validateToken(secret);
|
||||
if (!b) return false;
|
||||
|
||||
@@ -26,6 +26,8 @@ import java.util.Map;
|
||||
@RequestMapping("/api/v1/forward")
|
||||
public class ForwardController extends BaseController {
|
||||
|
||||
@Autowired
|
||||
private ForwardService forwardService;
|
||||
|
||||
@LogAnnotation
|
||||
@PostMapping("/create")
|
||||
@@ -73,5 +75,17 @@ public class ForwardController extends BaseController {
|
||||
return forwardService.resumeForward(id);
|
||||
}
|
||||
|
||||
/**
|
||||
* 转发诊断功能
|
||||
* @param params 包含forwardId的参数
|
||||
* @return 诊断结果
|
||||
*/
|
||||
@LogAnnotation
|
||||
@RequireRole
|
||||
@PostMapping("/diagnose")
|
||||
public R diagnoseForward(@RequestBody Map<String, Object> params) {
|
||||
Long forwardId = Long.valueOf(params.get("forwardId").toString());
|
||||
return forwardService.diagnoseForward(forwardId);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -62,4 +62,16 @@ public class NodeController extends BaseController {
|
||||
Long id = Long.valueOf(params.get("id").toString());
|
||||
return nodeService.getInstallCommand(id);
|
||||
}
|
||||
|
||||
/**
|
||||
* 检查和修复节点状态
|
||||
* @param params 包含节点ID的参数(可选)
|
||||
* @return 检查结果
|
||||
*/
|
||||
@LogAnnotation
|
||||
@RequireRole
|
||||
@PostMapping("/check-status")
|
||||
public R checkNodeStatus(@RequestBody(required = false) Map<String, Object> params) {
|
||||
return nodeService.checkAndFixNodeStatus(params);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -64,4 +64,11 @@ public interface ForwardService extends IService<Forward> {
|
||||
* @return 结果
|
||||
*/
|
||||
R resumeForward(Long id);
|
||||
|
||||
/**
|
||||
* 转发诊断功能
|
||||
* @param id 转发ID
|
||||
* @return 诊断结果
|
||||
*/
|
||||
R diagnoseForward(Long id);
|
||||
}
|
||||
|
||||
@@ -28,4 +28,11 @@ public interface NodeService extends IService<Node> {
|
||||
Node getNodeById(Long id);
|
||||
|
||||
R getInstallCommand(Long id);
|
||||
|
||||
/**
|
||||
* 检查和修复节点状态
|
||||
* @param params 包含节点ID的参数(可选)
|
||||
* @return 检查结果
|
||||
*/
|
||||
R checkAndFixNodeStatus(java.util.Map<String, Object> params);
|
||||
}
|
||||
|
||||
@@ -7,11 +7,13 @@ import com.admin.common.dto.GostDto;
|
||||
import com.admin.common.lang.R;
|
||||
import com.admin.common.utils.GostUtil;
|
||||
import com.admin.common.utils.JwtUtil;
|
||||
import com.admin.common.utils.WebSocketServer;
|
||||
import com.admin.entity.*;
|
||||
import com.admin.mapper.ForwardMapper;
|
||||
import com.admin.service.*;
|
||||
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
|
||||
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
|
||||
import com.alibaba.fastjson.JSONObject;
|
||||
import lombok.Data;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.beans.BeanUtils;
|
||||
@@ -19,10 +21,7 @@ import org.springframework.context.annotation.Lazy;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
import javax.annotation.Resource;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.*;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
@@ -385,6 +384,181 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
|
||||
return result ? R.ok("服务已" + operation) : R.err("更新状态失败");
|
||||
}
|
||||
|
||||
@Override
|
||||
public R diagnoseForward(Long id) {
|
||||
// 1. 获取当前用户信息
|
||||
UserInfo currentUser = getCurrentUserInfo();
|
||||
|
||||
// 2. 检查转发是否存在且用户有权限访问
|
||||
Forward forward = validateForwardExists(id, currentUser);
|
||||
if (forward == null) {
|
||||
return R.err("转发不存在");
|
||||
}
|
||||
|
||||
// 3. 获取隧道信息
|
||||
Tunnel tunnel = validateTunnel(forward.getTunnelId());
|
||||
if (tunnel == null) {
|
||||
return R.err("隧道不存在");
|
||||
}
|
||||
|
||||
// 4. 获取入口节点信息
|
||||
Node inNode = nodeService.getNodeById(tunnel.getInNodeId());
|
||||
if (inNode == null) {
|
||||
return R.err("入口节点不存在");
|
||||
}
|
||||
|
||||
// 5. 解析目标地址,取第一个地址作为诊断目标
|
||||
String[] remoteAddresses = forward.getRemoteAddr().split(",");
|
||||
String targetAddress = remoteAddresses[0].trim();
|
||||
|
||||
// 提取IP部分(去掉端口)
|
||||
String targetIp = extractIpFromAddress(targetAddress);
|
||||
if (targetIp == null) {
|
||||
return R.err("无法解析目标地址: " + targetAddress);
|
||||
}
|
||||
|
||||
List<DiagnosisResult> results = new ArrayList<>();
|
||||
|
||||
// 6. 根据隧道类型执行不同的诊断策略
|
||||
if (tunnel.getType() == TUNNEL_TYPE_PORT_FORWARD) {
|
||||
// 端口转发:入口节点直接ping目标地址
|
||||
DiagnosisResult result = performPingDiagnosis(inNode, targetIp, "转发->目标");
|
||||
results.add(result);
|
||||
} else {
|
||||
// 隧道转发:入口ping出口,出口ping目标
|
||||
Node outNode = nodeService.getNodeById(tunnel.getOutNodeId());
|
||||
if (outNode == null) {
|
||||
return R.err("出口节点不存在");
|
||||
}
|
||||
|
||||
// 入口ping出口
|
||||
DiagnosisResult inToOutResult = performPingDiagnosis(inNode, outNode.getServerIp(), "入口->出口");
|
||||
results.add(inToOutResult);
|
||||
|
||||
// 出口ping目标
|
||||
DiagnosisResult outToTargetResult = performPingDiagnosis(outNode, targetIp, "出口->目标");
|
||||
results.add(outToTargetResult);
|
||||
}
|
||||
|
||||
// 7. 构建诊断报告
|
||||
Map<String, Object> diagnosisReport = new HashMap<>();
|
||||
diagnosisReport.put("forwardId", id);
|
||||
diagnosisReport.put("forwardName", forward.getName());
|
||||
diagnosisReport.put("tunnelType", tunnel.getType() == TUNNEL_TYPE_PORT_FORWARD ? "端口转发" : "隧道转发");
|
||||
diagnosisReport.put("results", results);
|
||||
diagnosisReport.put("timestamp", System.currentTimeMillis());
|
||||
|
||||
return R.ok(diagnosisReport);
|
||||
}
|
||||
|
||||
/**
|
||||
* 从地址字符串中提取IP地址
|
||||
* 支持格式: ip:port, [ipv6]:port, domain:port
|
||||
*/
|
||||
private String extractIpFromAddress(String address) {
|
||||
if (address == null || address.trim().isEmpty()) {
|
||||
return null;
|
||||
}
|
||||
|
||||
address = address.trim();
|
||||
|
||||
// IPv6格式: [ipv6]:port
|
||||
if (address.startsWith("[")) {
|
||||
int closeBracket = address.indexOf(']');
|
||||
if (closeBracket > 1) {
|
||||
return address.substring(1, closeBracket);
|
||||
}
|
||||
}
|
||||
|
||||
// IPv4或域名格式: ip:port 或 domain:port
|
||||
int lastColon = address.lastIndexOf(':');
|
||||
if (lastColon > 0) {
|
||||
return address.substring(0, lastColon);
|
||||
}
|
||||
|
||||
// 如果没有端口,直接返回地址
|
||||
return address;
|
||||
}
|
||||
|
||||
/**
|
||||
* 执行ping诊断
|
||||
*
|
||||
* @param node 执行ping的节点
|
||||
* @param targetIp 目标IP地址
|
||||
* @param description 诊断描述
|
||||
* @return 诊断结果
|
||||
*/
|
||||
private DiagnosisResult performPingDiagnosis(Node node, String targetIp, String description) {
|
||||
try {
|
||||
// 构建ping请求数据
|
||||
JSONObject pingData = new JSONObject();
|
||||
pingData.put("ip", targetIp);
|
||||
pingData.put("count", 4);
|
||||
|
||||
// 发送ping命令到节点
|
||||
GostDto gostResult = WebSocketServer.send_msg(node.getId(), pingData, "Ping");
|
||||
|
||||
DiagnosisResult result = new DiagnosisResult();
|
||||
result.setNodeId(node.getId());
|
||||
result.setNodeName(node.getName());
|
||||
result.setTargetIp(targetIp);
|
||||
result.setDescription(description);
|
||||
result.setTimestamp(System.currentTimeMillis());
|
||||
|
||||
if (gostResult != null && "OK".equals(gostResult.getMsg())) {
|
||||
// 尝试解析ping响应数据
|
||||
try {
|
||||
if (gostResult.getData() != null) {
|
||||
JSONObject pingResponse = (JSONObject) gostResult.getData();
|
||||
boolean success = pingResponse.getBooleanValue("success");
|
||||
|
||||
result.setSuccess(success);
|
||||
if (success) {
|
||||
result.setMessage("ping成功");
|
||||
result.setAverageTime(pingResponse.getDoubleValue("averageTime"));
|
||||
result.setPacketLoss(pingResponse.getDoubleValue("packetLoss"));
|
||||
} else {
|
||||
result.setMessage(pingResponse.getString("errorMessage"));
|
||||
result.setAverageTime(-1.0);
|
||||
result.setPacketLoss(100.0);
|
||||
}
|
||||
} else {
|
||||
// 没有详细数据,使用默认值
|
||||
result.setSuccess(true);
|
||||
result.setMessage("ping成功");
|
||||
result.setAverageTime(0.0);
|
||||
result.setPacketLoss(0.0);
|
||||
}
|
||||
} catch (Exception e) {
|
||||
// 解析响应数据失败,但ping命令本身成功了
|
||||
result.setSuccess(true);
|
||||
result.setMessage("ping成功,但无法解析详细数据");
|
||||
result.setAverageTime(0.0);
|
||||
result.setPacketLoss(0.0);
|
||||
}
|
||||
} else {
|
||||
result.setSuccess(false);
|
||||
result.setMessage(gostResult != null ? gostResult.getMsg() : "节点无响应");
|
||||
result.setAverageTime(-1.0);
|
||||
result.setPacketLoss(100.0);
|
||||
}
|
||||
|
||||
return result;
|
||||
} catch (Exception e) {
|
||||
DiagnosisResult result = new DiagnosisResult();
|
||||
result.setNodeId(node.getId());
|
||||
result.setNodeName(node.getName());
|
||||
result.setTargetIp(targetIp);
|
||||
result.setDescription(description);
|
||||
result.setSuccess(false);
|
||||
result.setMessage("诊断执行异常: " + e.getMessage());
|
||||
result.setTimestamp(System.currentTimeMillis());
|
||||
result.setAverageTime(-1.0);
|
||||
result.setPacketLoss(100.0);
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取当前用户信息
|
||||
*/
|
||||
@@ -1142,4 +1316,20 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
|
||||
return new NodeInfo(true, errorMessage, null, null);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 诊断结果数据类
|
||||
*/
|
||||
@Data
|
||||
public static class DiagnosisResult {
|
||||
private Long nodeId;
|
||||
private String nodeName;
|
||||
private String targetIp;
|
||||
private String description;
|
||||
private boolean success;
|
||||
private String message;
|
||||
private double averageTime;
|
||||
private double packetLoss;
|
||||
private long timestamp;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -430,4 +430,125 @@ public class NodeServiceImpl extends ServiceImpl<NodeMapper, Node> implements No
|
||||
throw new RuntimeException(ERROR_PORT_ORDER_INVALID);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 检查和修复节点状态
|
||||
* 如果没有指定节点ID,则检查所有节点
|
||||
* 如果指定了节点ID,则只检查该节点
|
||||
*
|
||||
* @param params 包含节点ID的参数(可选)
|
||||
* @return 检查结果
|
||||
*/
|
||||
@Override
|
||||
public R checkAndFixNodeStatus(java.util.Map<String, Object> params) {
|
||||
try {
|
||||
java.util.List<CheckResult> results = new java.util.ArrayList<>();
|
||||
|
||||
if (params != null && params.containsKey("nodeId")) {
|
||||
// 检查指定节点
|
||||
Long nodeId = Long.valueOf(params.get("nodeId").toString());
|
||||
CheckResult result = checkSingleNodeStatus(nodeId);
|
||||
results.add(result);
|
||||
} else {
|
||||
// 检查所有节点
|
||||
List<Node> allNodes = this.list();
|
||||
for (Node node : allNodes) {
|
||||
CheckResult result = checkSingleNodeStatus(node.getId());
|
||||
results.add(result);
|
||||
}
|
||||
}
|
||||
|
||||
// 统计结果
|
||||
long totalNodes = results.size();
|
||||
long inconsistentNodes = results.stream()
|
||||
.mapToLong(r -> r.isFixed() ? 1 : 0)
|
||||
.sum();
|
||||
long connectedNodes = results.stream()
|
||||
.mapToLong(r -> r.isConnected() ? 1 : 0)
|
||||
.sum();
|
||||
|
||||
java.util.Map<String, Object> response = new java.util.HashMap<>();
|
||||
response.put("totalNodes", totalNodes);
|
||||
response.put("connectedNodes", connectedNodes);
|
||||
response.put("inconsistentNodes", inconsistentNodes);
|
||||
response.put("details", results);
|
||||
|
||||
return R.ok(response);
|
||||
|
||||
} catch (Exception e) {
|
||||
return R.err("检查节点状态时发生错误:" + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 检查单个节点的状态
|
||||
*
|
||||
* @param nodeId 节点ID
|
||||
* @return 检查结果
|
||||
*/
|
||||
private CheckResult checkSingleNodeStatus(Long nodeId) {
|
||||
CheckResult result = new CheckResult();
|
||||
result.setNodeId(nodeId);
|
||||
|
||||
try {
|
||||
Node node = this.getById(nodeId);
|
||||
if (node == null) {
|
||||
result.setNodeName("未知");
|
||||
result.setConnected(false);
|
||||
result.setDatabaseStatus(0);
|
||||
result.setActualStatus(false);
|
||||
result.setFixed(false);
|
||||
result.setMessage("节点不存在");
|
||||
return result;
|
||||
}
|
||||
|
||||
result.setNodeName(node.getName());
|
||||
result.setDatabaseStatus(node.getStatus());
|
||||
|
||||
// 调用WebSocketServer的静态方法检查实际连接状态
|
||||
boolean actualConnected = com.admin.common.utils.WebSocketServer.checkAndFixNodeStatus(this, nodeId);
|
||||
result.setConnected(actualConnected);
|
||||
result.setActualStatus(actualConnected);
|
||||
|
||||
// 重新查询节点状态,看是否被修复了
|
||||
Node updatedNode = this.getById(nodeId);
|
||||
boolean wasFixed = (updatedNode.getStatus() != node.getStatus());
|
||||
result.setFixed(wasFixed);
|
||||
result.setFinalStatus(updatedNode.getStatus());
|
||||
|
||||
if (wasFixed) {
|
||||
result.setMessage(String.format("状态已修复:%d -> %d",
|
||||
node.getStatus(), updatedNode.getStatus()));
|
||||
} else if (actualConnected && updatedNode.getStatus() == 1) {
|
||||
result.setMessage("状态正常");
|
||||
} else if (!actualConnected && updatedNode.getStatus() == 0) {
|
||||
result.setMessage("状态正常");
|
||||
} else {
|
||||
result.setMessage("状态可能存在异常");
|
||||
}
|
||||
|
||||
} catch (Exception e) {
|
||||
result.setConnected(false);
|
||||
result.setActualStatus(false);
|
||||
result.setFixed(false);
|
||||
result.setMessage("检查失败:" + e.getMessage());
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
/**
|
||||
* 节点状态检查结果
|
||||
*/
|
||||
@lombok.Data
|
||||
public static class CheckResult {
|
||||
private Long nodeId;
|
||||
private String nodeName;
|
||||
private boolean connected;
|
||||
private int databaseStatus;
|
||||
private boolean actualStatus;
|
||||
private boolean fixed;
|
||||
private int finalStatus;
|
||||
private String message;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user