From 96fcd0fc579f2959838fac3a7bd1ab4ffe3f9669 Mon Sep 17 00:00:00 2001 From: root Date: Thu, 5 Feb 2026 04:52:42 +0000 Subject: [PATCH] fix(diagnose): parallelize TCP ping diagnostics to prevent timeout cascade Previously, forward/tunnel diagnosis executed TCP pings sequentially, causing total time to accumulate. If the first remote address timed out (5s), subsequent checks could push total time beyond the frontend's 30s timeout, resulting in diagnosis failure even for healthy endpoints. Now all diagnostic tasks run in parallel using CompletableFuture, so total time equals max(individual ping time) instead of sum. --- .../service/impl/ForwardServiceImpl.java | 133 ++++++++++++------ .../admin/service/impl/TunnelServiceImpl.java | 120 ++++++++++------ 2 files changed, 164 insertions(+), 89 deletions(-) diff --git a/springboot-backend/src/main/java/com/admin/service/impl/ForwardServiceImpl.java b/springboot-backend/src/main/java/com/admin/service/impl/ForwardServiceImpl.java index 4246f09..0d8339c 100644 --- a/springboot-backend/src/main/java/com/admin/service/impl/ForwardServiceImpl.java +++ b/springboot-backend/src/main/java/com/admin/service/impl/ForwardServiceImpl.java @@ -22,6 +22,7 @@ import org.springframework.transaction.annotation.Transactional; import javax.annotation.Resource; import java.util.*; +import java.util.concurrent.CompletableFuture; import java.util.stream.Collectors; /** @@ -513,10 +514,10 @@ public class ForwardServiceImpl extends ServiceImpl impl .filter(ct -> ct.getChainType() == 3) .toList(); - List results = new ArrayList<>(); + List> futures = new ArrayList<>(); String[] remoteAddresses = forward.getRemoteAddr().split(","); - // 根据隧道类型执行不同的诊断策略 + // 根据隧道类型执行不同的诊断策略(并行执行所有诊断任务) if (tunnel.getType() == 1) { // 端口转发:入口节点直接TCP ping目标地址 for (ChainTunnel inNode : inNodes) { @@ -526,12 +527,18 @@ public class ForwardServiceImpl extends ServiceImpl impl String targetIp = extractIpFromAddress(remoteAddress); int targetPort = extractPortFromAddress(remoteAddress); if (targetIp != null && targetPort != -1) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - node, targetIp, targetPort, - "入口(" + node.getName() + ")->目标(" + remoteAddress + ")" - ); - result.setFromChainType(1); - results.add(result); + final Node finalNode = node; + final String finalTargetIp = targetIp; + final int finalTargetPort = targetPort; + final String finalRemoteAddress = remoteAddress; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalNode, finalTargetIp, finalTargetPort, + "入口(" + finalNode.getName() + ")->目标(" + finalRemoteAddress + ")" + ); + result.setFromChainType(1); + return result; + })); } } } @@ -547,27 +554,37 @@ public class ForwardServiceImpl extends ServiceImpl impl for (ChainTunnel firstChainNode : chainNodesList.getFirst()) { Node toNode = nodeService.getById(firstChainNode.getNodeId()); if (toNode != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - fromNode, GostUtil.selectDialHost(fromNode, toNode), firstChainNode.getPort(), - "入口(" + fromNode.getName() + ")->第1跳(" + toNode.getName() + ")" - ); - result.setFromChainType(1); - result.setToChainType(2); - result.setToInx(firstChainNode.getInx()); - results.add(result); + final Node finalFromNode = fromNode; + final Node finalToNode = toNode; + final ChainTunnel finalFirstChainNode = firstChainNode; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalFirstChainNode.getPort(), + "入口(" + finalFromNode.getName() + ")->第1跳(" + finalToNode.getName() + ")" + ); + result.setFromChainType(1); + result.setToChainType(2); + result.setToInx(finalFirstChainNode.getInx()); + return result; + })); } } } else if (!outNodes.isEmpty()) { for (ChainTunnel outNode : outNodes) { Node toNode = nodeService.getById(outNode.getNodeId()); if (toNode != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - fromNode, GostUtil.selectDialHost(fromNode, toNode), outNode.getPort(), - "入口(" + fromNode.getName() + ")->出口(" + toNode.getName() + ")" - ); - result.setFromChainType(1); - result.setToChainType(3); - results.add(result); + final Node finalFromNode = fromNode; + final Node finalToNode = toNode; + final ChainTunnel finalOutNode = outNode; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalOutNode.getPort(), + "入口(" + finalFromNode.getName() + ")->出口(" + finalToNode.getName() + ")" + ); + result.setFromChainType(1); + result.setToChainType(3); + return result; + })); } } } @@ -577,6 +594,7 @@ public class ForwardServiceImpl extends ServiceImpl impl // 2. 链路测试 for (int i = 0; i < chainNodesList.size(); i++) { List currentHop = chainNodesList.get(i); + final int hopIndex = i; for (ChainTunnel currentNode : currentHop) { Node fromNode = nodeService.getById(currentNode.getNodeId()); @@ -586,29 +604,41 @@ public class ForwardServiceImpl extends ServiceImpl impl for (ChainTunnel nextNode : chainNodesList.get(i + 1)) { Node toNode = nodeService.getById(nextNode.getNodeId()); if (toNode != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - fromNode, GostUtil.selectDialHost(fromNode, toNode), nextNode.getPort(), - "第" + (i + 1) + "跳(" + fromNode.getName() + ")->第" + (i + 2) + "跳(" + toNode.getName() + ")" - ); - result.setFromChainType(2); - result.setFromInx(currentNode.getInx()); - result.setToChainType(2); - result.setToInx(nextNode.getInx()); - results.add(result); + final Node finalFromNode = fromNode; + final Node finalToNode = toNode; + final ChainTunnel finalCurrentNode = currentNode; + final ChainTunnel finalNextNode = nextNode; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalNextNode.getPort(), + "第" + (hopIndex + 1) + "跳(" + finalFromNode.getName() + ")->第" + (hopIndex + 2) + "跳(" + finalToNode.getName() + ")" + ); + result.setFromChainType(2); + result.setFromInx(finalCurrentNode.getInx()); + result.setToChainType(2); + result.setToInx(finalNextNode.getInx()); + return result; + })); } } } else if (!outNodes.isEmpty()) { for (ChainTunnel outNode : outNodes) { Node toNode = nodeService.getById(outNode.getNodeId()); if (toNode != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - fromNode, GostUtil.selectDialHost(fromNode, toNode), outNode.getPort(), - "第" + (i + 1) + "跳(" + fromNode.getName() + ")->出口(" + toNode.getName() + ")" - ); - result.setFromChainType(2); - result.setFromInx(currentNode.getInx()); - result.setToChainType(3); - results.add(result); + final Node finalFromNode = fromNode; + final Node finalToNode = toNode; + final ChainTunnel finalCurrentNode = currentNode; + final ChainTunnel finalOutNode = outNode; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalOutNode.getPort(), + "第" + (hopIndex + 1) + "跳(" + finalFromNode.getName() + ")->出口(" + finalToNode.getName() + ")" + ); + result.setFromChainType(2); + result.setFromInx(finalCurrentNode.getInx()); + result.setToChainType(3); + return result; + })); } } } @@ -624,18 +654,29 @@ public class ForwardServiceImpl extends ServiceImpl impl String targetIp = extractIpFromAddress(remoteAddress); int targetPort = extractPortFromAddress(remoteAddress); if (targetIp != null && targetPort != -1) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - node, targetIp, targetPort, - "出口(" + node.getName() + ")->目标(" + remoteAddress + ")" - ); - result.setFromChainType(3); - results.add(result); + final Node finalNode = node; + final String finalTargetIp = targetIp; + final int finalTargetPort = targetPort; + final String finalRemoteAddress = remoteAddress; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalNode, finalTargetIp, finalTargetPort, + "出口(" + finalNode.getName() + ")->目标(" + finalRemoteAddress + ")" + ); + result.setFromChainType(3); + return result; + })); } } } } } + // 等待所有诊断任务完成并收集结果 + List results = futures.stream() + .map(CompletableFuture::join) + .collect(Collectors.toList()); + // 构建诊断报告 Map diagnosisReport = new HashMap<>(); diagnosisReport.put("forwardId", id); diff --git a/springboot-backend/src/main/java/com/admin/service/impl/TunnelServiceImpl.java b/springboot-backend/src/main/java/com/admin/service/impl/TunnelServiceImpl.java index 36187e3..314bb6f 100644 --- a/springboot-backend/src/main/java/com/admin/service/impl/TunnelServiceImpl.java +++ b/springboot-backend/src/main/java/com/admin/service/impl/TunnelServiceImpl.java @@ -22,6 +22,7 @@ import org.springframework.transaction.annotation.Transactional; import javax.annotation.Resource; import java.util.*; +import java.util.concurrent.CompletableFuture; import java.util.stream.Collectors; /** @@ -669,17 +670,20 @@ public class TunnelServiceImpl extends ServiceImpl impleme .filter(ct -> ct.getChainType() == 3) .toList(); - List results = new ArrayList<>(); + List> futures = new ArrayList<>(); if (tunnel.getType() == 1) { for (ChainTunnel inNode : inNodes) { Node node = nodeService.getById(inNode.getNodeId()); if (node != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - node, "www.google.com", 443, "入口(" + node.getName() + ")->外网" - ); - result.setFromChainType(1); // 入口 - results.add(result); + final Node finalNode = node; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalNode, "www.google.com", 443, "入口(" + finalNode.getName() + ")->外网" + ); + result.setFromChainType(1); + return result; + })); } } } else if (tunnel.getType() == 2) { @@ -691,27 +695,37 @@ public class TunnelServiceImpl extends ServiceImpl impleme for (ChainTunnel firstChainNode : chainNodesList.getFirst()) { Node toNode = nodeService.getById(firstChainNode.getNodeId()); if (toNode != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - fromNode, GostUtil.selectDialHost(fromNode, toNode), firstChainNode.getPort(), - "入口(" + fromNode.getName() + ")->第1跳(" + toNode.getName() + ")" - ); - result.setFromChainType(1); // 入口 - result.setToChainType(2); // 链 - result.setToInx(firstChainNode.getInx()); - results.add(result); + final Node finalFromNode = fromNode; + final Node finalToNode = toNode; + final ChainTunnel finalFirstChainNode = firstChainNode; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalFirstChainNode.getPort(), + "入口(" + finalFromNode.getName() + ")->第1跳(" + finalToNode.getName() + ")" + ); + result.setFromChainType(1); + result.setToChainType(2); + result.setToInx(finalFirstChainNode.getInx()); + return result; + })); } } } else if (!outNodes.isEmpty()) { for (ChainTunnel outNode : outNodes) { Node toNode = nodeService.getById(outNode.getNodeId()); if (toNode != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - fromNode, GostUtil.selectDialHost(fromNode, toNode), outNode.getPort(), - "入口(" + fromNode.getName() + ")->出口(" + toNode.getName() + ")" - ); - result.setFromChainType(1); - result.setToChainType(3); - results.add(result); + final Node finalFromNode = fromNode; + final Node finalToNode = toNode; + final ChainTunnel finalOutNode = outNode; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalOutNode.getPort(), + "入口(" + finalFromNode.getName() + ")->出口(" + finalToNode.getName() + ")" + ); + result.setFromChainType(1); + result.setToChainType(3); + return result; + })); } } } @@ -720,6 +734,7 @@ public class TunnelServiceImpl extends ServiceImpl impleme for (int i = 0; i < chainNodesList.size(); i++) { List currentHop = chainNodesList.get(i); + final int hopIndex = i; for (ChainTunnel currentNode : currentHop) { Node fromNode = nodeService.getById(currentNode.getNodeId()); @@ -729,29 +744,41 @@ public class TunnelServiceImpl extends ServiceImpl impleme for (ChainTunnel nextNode : chainNodesList.get(i + 1)) { Node toNode = nodeService.getById(nextNode.getNodeId()); if (toNode != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - fromNode, GostUtil.selectDialHost(fromNode, toNode), nextNode.getPort(), - "第" + (i + 1) + "跳(" + fromNode.getName() + ")->第" + (i + 2) + "跳(" + toNode.getName() + ")" - ); - result.setFromChainType(2); - result.setFromInx(currentNode.getInx()); - result.setToChainType(2); - result.setToInx(nextNode.getInx()); - results.add(result); + final Node finalFromNode = fromNode; + final Node finalToNode = toNode; + final ChainTunnel finalCurrentNode = currentNode; + final ChainTunnel finalNextNode = nextNode; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalNextNode.getPort(), + "第" + (hopIndex + 1) + "跳(" + finalFromNode.getName() + ")->第" + (hopIndex + 2) + "跳(" + finalToNode.getName() + ")" + ); + result.setFromChainType(2); + result.setFromInx(finalCurrentNode.getInx()); + result.setToChainType(2); + result.setToInx(finalNextNode.getInx()); + return result; + })); } } } else if (!outNodes.isEmpty()) { for (ChainTunnel outNode : outNodes) { Node toNode = nodeService.getById(outNode.getNodeId()); if (toNode != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - fromNode, GostUtil.selectDialHost(fromNode, toNode), outNode.getPort(), - "第" + (i + 1) + "跳(" + fromNode.getName() + ")->出口(" + toNode.getName() + ")" - ); - result.setFromChainType(2); - result.setFromInx(currentNode.getInx()); - result.setToChainType(3); - results.add(result); + final Node finalFromNode = fromNode; + final Node finalToNode = toNode; + final ChainTunnel finalCurrentNode = currentNode; + final ChainTunnel finalOutNode = outNode; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalOutNode.getPort(), + "第" + (hopIndex + 1) + "跳(" + finalFromNode.getName() + ")->出口(" + finalToNode.getName() + ")" + ); + result.setFromChainType(2); + result.setFromInx(finalCurrentNode.getInx()); + result.setToChainType(3); + return result; + })); } } } @@ -761,15 +788,22 @@ public class TunnelServiceImpl extends ServiceImpl impleme for (ChainTunnel outNode : outNodes) { Node node = nodeService.getById(outNode.getNodeId()); if (node != null) { - DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( - node, "www.google.com", 443, "出口(" + node.getName() + ")->外网" - ); - result.setFromChainType(3); - results.add(result); + final Node finalNode = node; + futures.add(CompletableFuture.supplyAsync(() -> { + DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( + finalNode, "www.google.com", 443, "出口(" + finalNode.getName() + ")->外网" + ); + result.setFromChainType(3); + return result; + })); } } } + List results = futures.stream() + .map(CompletableFuture::join) + .collect(Collectors.toList()); + Map diagnosisReport = new HashMap<>(); diagnosisReport.put("tunnelId", tunnelId); diagnosisReport.put("tunnelName", tunnel.getName());