Merge pull request #26 from Sagit-chu/opencode/hidden-pixel

fix(diagnose): parallelize TCP ping diagnostics to prevent timeout ca…
This commit is contained in:
sagit
2026-02-05 12:57:10 +08:00
committed by GitHub
2 changed files with 164 additions and 89 deletions
@@ -22,6 +22,7 @@ import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource; import javax.annotation.Resource;
import java.util.*; import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors; import java.util.stream.Collectors;
/** /**
@@ -513,10 +514,10 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
.filter(ct -> ct.getChainType() == 3) .filter(ct -> ct.getChainType() == 3)
.toList(); .toList();
List<DiagnosisResult> results = new ArrayList<>(); List<CompletableFuture<DiagnosisResult>> futures = new ArrayList<>();
String[] remoteAddresses = forward.getRemoteAddr().split(","); String[] remoteAddresses = forward.getRemoteAddr().split(",");
// 根据隧道类型执行不同的诊断策略 // 根据隧道类型执行不同的诊断策略(并行执行所有诊断任务)
if (tunnel.getType() == 1) { if (tunnel.getType() == 1) {
// 端口转发:入口节点直接TCP ping目标地址 // 端口转发:入口节点直接TCP ping目标地址
for (ChainTunnel inNode : inNodes) { for (ChainTunnel inNode : inNodes) {
@@ -526,12 +527,18 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
String targetIp = extractIpFromAddress(remoteAddress); String targetIp = extractIpFromAddress(remoteAddress);
int targetPort = extractPortFromAddress(remoteAddress); int targetPort = extractPortFromAddress(remoteAddress);
if (targetIp != null && targetPort != -1) { if (targetIp != null && targetPort != -1) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalNode = node;
node, targetIp, targetPort, final String finalTargetIp = targetIp;
"入口(" + node.getName() + ")->目标(" + remoteAddress + ")" final int finalTargetPort = targetPort;
); final String finalRemoteAddress = remoteAddress;
result.setFromChainType(1); futures.add(CompletableFuture.supplyAsync(() -> {
results.add(result); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
finalNode, finalTargetIp, finalTargetPort,
"入口(" + finalNode.getName() + ")->目标(" + finalRemoteAddress + ")"
);
result.setFromChainType(1);
return result;
}));
} }
} }
} }
@@ -547,27 +554,37 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
for (ChainTunnel firstChainNode : chainNodesList.getFirst()) { for (ChainTunnel firstChainNode : chainNodesList.getFirst()) {
Node toNode = nodeService.getById(firstChainNode.getNodeId()); Node toNode = nodeService.getById(firstChainNode.getNodeId());
if (toNode != null) { if (toNode != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalFromNode = fromNode;
fromNode, GostUtil.selectDialHost(fromNode, toNode), firstChainNode.getPort(), final Node finalToNode = toNode;
"入口(" + fromNode.getName() + ")->第1跳(" + toNode.getName() + ")" final ChainTunnel finalFirstChainNode = firstChainNode;
); futures.add(CompletableFuture.supplyAsync(() -> {
result.setFromChainType(1); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setToChainType(2); finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalFirstChainNode.getPort(),
result.setToInx(firstChainNode.getInx()); "入口(" + finalFromNode.getName() + ")->第1跳(" + finalToNode.getName() + ")"
results.add(result); );
result.setFromChainType(1);
result.setToChainType(2);
result.setToInx(finalFirstChainNode.getInx());
return result;
}));
} }
} }
} else if (!outNodes.isEmpty()) { } else if (!outNodes.isEmpty()) {
for (ChainTunnel outNode : outNodes) { for (ChainTunnel outNode : outNodes) {
Node toNode = nodeService.getById(outNode.getNodeId()); Node toNode = nodeService.getById(outNode.getNodeId());
if (toNode != null) { if (toNode != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalFromNode = fromNode;
fromNode, GostUtil.selectDialHost(fromNode, toNode), outNode.getPort(), final Node finalToNode = toNode;
"入口(" + fromNode.getName() + ")->出口(" + toNode.getName() + ")" final ChainTunnel finalOutNode = outNode;
); futures.add(CompletableFuture.supplyAsync(() -> {
result.setFromChainType(1); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setToChainType(3); finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalOutNode.getPort(),
results.add(result); "入口(" + finalFromNode.getName() + ")->出口(" + finalToNode.getName() + ")"
);
result.setFromChainType(1);
result.setToChainType(3);
return result;
}));
} }
} }
} }
@@ -577,6 +594,7 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
// 2. 链路测试 // 2. 链路测试
for (int i = 0; i < chainNodesList.size(); i++) { for (int i = 0; i < chainNodesList.size(); i++) {
List<ChainTunnel> currentHop = chainNodesList.get(i); List<ChainTunnel> currentHop = chainNodesList.get(i);
final int hopIndex = i;
for (ChainTunnel currentNode : currentHop) { for (ChainTunnel currentNode : currentHop) {
Node fromNode = nodeService.getById(currentNode.getNodeId()); Node fromNode = nodeService.getById(currentNode.getNodeId());
@@ -586,29 +604,41 @@ public class ForwardServiceImpl extends ServiceImpl<ForwardMapper, Forward> impl
for (ChainTunnel nextNode : chainNodesList.get(i + 1)) { for (ChainTunnel nextNode : chainNodesList.get(i + 1)) {
Node toNode = nodeService.getById(nextNode.getNodeId()); Node toNode = nodeService.getById(nextNode.getNodeId());
if (toNode != null) { if (toNode != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalFromNode = fromNode;
fromNode, GostUtil.selectDialHost(fromNode, toNode), nextNode.getPort(), final Node finalToNode = toNode;
"第" + (i + 1) + "跳(" + fromNode.getName() + ")->第" + (i + 2) + "跳(" + toNode.getName() + ")" final ChainTunnel finalCurrentNode = currentNode;
); final ChainTunnel finalNextNode = nextNode;
result.setFromChainType(2); futures.add(CompletableFuture.supplyAsync(() -> {
result.setFromInx(currentNode.getInx()); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setToChainType(2); finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalNextNode.getPort(),
result.setToInx(nextNode.getInx()); "第" + (hopIndex + 1) + "跳(" + finalFromNode.getName() + ")->第" + (hopIndex + 2) + "跳(" + finalToNode.getName() + ")"
results.add(result); );
result.setFromChainType(2);
result.setFromInx(finalCurrentNode.getInx());
result.setToChainType(2);
result.setToInx(finalNextNode.getInx());
return result;
}));
} }
} }
} else if (!outNodes.isEmpty()) { } else if (!outNodes.isEmpty()) {
for (ChainTunnel outNode : outNodes) { for (ChainTunnel outNode : outNodes) {
Node toNode = nodeService.getById(outNode.getNodeId()); Node toNode = nodeService.getById(outNode.getNodeId());
if (toNode != null) { if (toNode != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalFromNode = fromNode;
fromNode, GostUtil.selectDialHost(fromNode, toNode), outNode.getPort(), final Node finalToNode = toNode;
"第" + (i + 1) + "跳(" + fromNode.getName() + ")->出口(" + toNode.getName() + ")" final ChainTunnel finalCurrentNode = currentNode;
); final ChainTunnel finalOutNode = outNode;
result.setFromChainType(2); futures.add(CompletableFuture.supplyAsync(() -> {
result.setFromInx(currentNode.getInx()); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setToChainType(3); finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalOutNode.getPort(),
results.add(result); "第" + (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<ForwardMapper, Forward> impl
String targetIp = extractIpFromAddress(remoteAddress); String targetIp = extractIpFromAddress(remoteAddress);
int targetPort = extractPortFromAddress(remoteAddress); int targetPort = extractPortFromAddress(remoteAddress);
if (targetIp != null && targetPort != -1) { if (targetIp != null && targetPort != -1) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalNode = node;
node, targetIp, targetPort, final String finalTargetIp = targetIp;
"出口(" + node.getName() + ")->目标(" + remoteAddress + ")" final int finalTargetPort = targetPort;
); final String finalRemoteAddress = remoteAddress;
result.setFromChainType(3); futures.add(CompletableFuture.supplyAsync(() -> {
results.add(result); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
finalNode, finalTargetIp, finalTargetPort,
"出口(" + finalNode.getName() + ")->目标(" + finalRemoteAddress + ")"
);
result.setFromChainType(3);
return result;
}));
} }
} }
} }
} }
} }
// 等待所有诊断任务完成并收集结果
List<DiagnosisResult> results = futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
// 构建诊断报告 // 构建诊断报告
Map<String, Object> diagnosisReport = new HashMap<>(); Map<String, Object> diagnosisReport = new HashMap<>();
diagnosisReport.put("forwardId", id); diagnosisReport.put("forwardId", id);
@@ -22,6 +22,7 @@ import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource; import javax.annotation.Resource;
import java.util.*; import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors; import java.util.stream.Collectors;
/** /**
@@ -669,17 +670,20 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
.filter(ct -> ct.getChainType() == 3) .filter(ct -> ct.getChainType() == 3)
.toList(); .toList();
List<DiagnosisResult> results = new ArrayList<>(); List<CompletableFuture<DiagnosisResult>> futures = new ArrayList<>();
if (tunnel.getType() == 1) { if (tunnel.getType() == 1) {
for (ChainTunnel inNode : inNodes) { for (ChainTunnel inNode : inNodes) {
Node node = nodeService.getById(inNode.getNodeId()); Node node = nodeService.getById(inNode.getNodeId());
if (node != null) { if (node != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalNode = node;
node, "www.google.com", 443, "入口(" + node.getName() + ")->外网" futures.add(CompletableFuture.supplyAsync(() -> {
); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setFromChainType(1); // 入口 finalNode, "www.google.com", 443, "入口(" + finalNode.getName() + ")->外网"
results.add(result); );
result.setFromChainType(1);
return result;
}));
} }
} }
} else if (tunnel.getType() == 2) { } else if (tunnel.getType() == 2) {
@@ -691,27 +695,37 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
for (ChainTunnel firstChainNode : chainNodesList.getFirst()) { for (ChainTunnel firstChainNode : chainNodesList.getFirst()) {
Node toNode = nodeService.getById(firstChainNode.getNodeId()); Node toNode = nodeService.getById(firstChainNode.getNodeId());
if (toNode != null) { if (toNode != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalFromNode = fromNode;
fromNode, GostUtil.selectDialHost(fromNode, toNode), firstChainNode.getPort(), final Node finalToNode = toNode;
"入口(" + fromNode.getName() + ")->第1跳(" + toNode.getName() + ")" final ChainTunnel finalFirstChainNode = firstChainNode;
); futures.add(CompletableFuture.supplyAsync(() -> {
result.setFromChainType(1); // 入口 DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setToChainType(2); // 链 finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalFirstChainNode.getPort(),
result.setToInx(firstChainNode.getInx()); "入口(" + finalFromNode.getName() + ")->第1跳(" + finalToNode.getName() + ")"
results.add(result); );
result.setFromChainType(1);
result.setToChainType(2);
result.setToInx(finalFirstChainNode.getInx());
return result;
}));
} }
} }
} else if (!outNodes.isEmpty()) { } else if (!outNodes.isEmpty()) {
for (ChainTunnel outNode : outNodes) { for (ChainTunnel outNode : outNodes) {
Node toNode = nodeService.getById(outNode.getNodeId()); Node toNode = nodeService.getById(outNode.getNodeId());
if (toNode != null) { if (toNode != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalFromNode = fromNode;
fromNode, GostUtil.selectDialHost(fromNode, toNode), outNode.getPort(), final Node finalToNode = toNode;
"入口(" + fromNode.getName() + ")->出口(" + toNode.getName() + ")" final ChainTunnel finalOutNode = outNode;
); futures.add(CompletableFuture.supplyAsync(() -> {
result.setFromChainType(1); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setToChainType(3); finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalOutNode.getPort(),
results.add(result); "入口(" + finalFromNode.getName() + ")->出口(" + finalToNode.getName() + ")"
);
result.setFromChainType(1);
result.setToChainType(3);
return result;
}));
} }
} }
} }
@@ -720,6 +734,7 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
for (int i = 0; i < chainNodesList.size(); i++) { for (int i = 0; i < chainNodesList.size(); i++) {
List<ChainTunnel> currentHop = chainNodesList.get(i); List<ChainTunnel> currentHop = chainNodesList.get(i);
final int hopIndex = i;
for (ChainTunnel currentNode : currentHop) { for (ChainTunnel currentNode : currentHop) {
Node fromNode = nodeService.getById(currentNode.getNodeId()); Node fromNode = nodeService.getById(currentNode.getNodeId());
@@ -729,29 +744,41 @@ public class TunnelServiceImpl extends ServiceImpl<TunnelMapper, Tunnel> impleme
for (ChainTunnel nextNode : chainNodesList.get(i + 1)) { for (ChainTunnel nextNode : chainNodesList.get(i + 1)) {
Node toNode = nodeService.getById(nextNode.getNodeId()); Node toNode = nodeService.getById(nextNode.getNodeId());
if (toNode != null) { if (toNode != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalFromNode = fromNode;
fromNode, GostUtil.selectDialHost(fromNode, toNode), nextNode.getPort(), final Node finalToNode = toNode;
"第" + (i + 1) + "跳(" + fromNode.getName() + ")->第" + (i + 2) + "跳(" + toNode.getName() + ")" final ChainTunnel finalCurrentNode = currentNode;
); final ChainTunnel finalNextNode = nextNode;
result.setFromChainType(2); futures.add(CompletableFuture.supplyAsync(() -> {
result.setFromInx(currentNode.getInx()); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setToChainType(2); finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalNextNode.getPort(),
result.setToInx(nextNode.getInx()); "第" + (hopIndex + 1) + "跳(" + finalFromNode.getName() + ")->第" + (hopIndex + 2) + "跳(" + finalToNode.getName() + ")"
results.add(result); );
result.setFromChainType(2);
result.setFromInx(finalCurrentNode.getInx());
result.setToChainType(2);
result.setToInx(finalNextNode.getInx());
return result;
}));
} }
} }
} else if (!outNodes.isEmpty()) { } else if (!outNodes.isEmpty()) {
for (ChainTunnel outNode : outNodes) { for (ChainTunnel outNode : outNodes) {
Node toNode = nodeService.getById(outNode.getNodeId()); Node toNode = nodeService.getById(outNode.getNodeId());
if (toNode != null) { if (toNode != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalFromNode = fromNode;
fromNode, GostUtil.selectDialHost(fromNode, toNode), outNode.getPort(), final Node finalToNode = toNode;
"第" + (i + 1) + "跳(" + fromNode.getName() + ")->出口(" + toNode.getName() + ")" final ChainTunnel finalCurrentNode = currentNode;
); final ChainTunnel finalOutNode = outNode;
result.setFromChainType(2); futures.add(CompletableFuture.supplyAsync(() -> {
result.setFromInx(currentNode.getInx()); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setToChainType(3); finalFromNode, GostUtil.selectDialHost(finalFromNode, finalToNode), finalOutNode.getPort(),
results.add(result); "第" + (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<TunnelMapper, Tunnel> impleme
for (ChainTunnel outNode : outNodes) { for (ChainTunnel outNode : outNodes) {
Node node = nodeService.getById(outNode.getNodeId()); Node node = nodeService.getById(outNode.getNodeId());
if (node != null) { if (node != null) {
DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck( final Node finalNode = node;
node, "www.google.com", 443, "出口(" + node.getName() + ")->外网" futures.add(CompletableFuture.supplyAsync(() -> {
); DiagnosisResult result = performTcpPingDiagnosisWithConnectionCheck(
result.setFromChainType(3); finalNode, "www.google.com", 443, "出口(" + finalNode.getName() + ")->外网"
results.add(result); );
result.setFromChainType(3);
return result;
}));
} }
} }
} }
List<DiagnosisResult> results = futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
Map<String, Object> diagnosisReport = new HashMap<>(); Map<String, Object> diagnosisReport = new HashMap<>();
diagnosisReport.put("tunnelId", tunnelId); diagnosisReport.put("tunnelId", tunnelId);
diagnosisReport.put("tunnelName", tunnel.getName()); diagnosisReport.put("tunnelName", tunnel.getName());