mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-28 07:36:38 +08:00
Compare commits
12 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 02ff215f99 | |||
| 2a4e7777ab | |||
| eac94a5719 | |||
| 96fcd0fc57 | |||
| 06869aedfd | |||
| 6a201131a3 | |||
| 1130a55ef5 | |||
| 265cd0a50e | |||
| e7ffa77b15 | |||
| 51cbd4b9de | |||
| e122e7460d | |||
| 7c898154b3 |
@@ -252,7 +252,14 @@ func (h *forwardHandler) Handle(ctx context.Context, conn net.Conn, opts ...hand
|
||||
}
|
||||
defer cc.Close()
|
||||
|
||||
xnet.Transport(conn, cc)
|
||||
if err := xnet.Transport(conn, cc); err != nil {
|
||||
if marker := target.Marker(); marker != nil {
|
||||
marker.Mark()
|
||||
h.options.Logger.Debugf("[handler.transport] transport failed, marked node=%s count=%d err=%v",
|
||||
target.Addr, marker.Count(), err)
|
||||
}
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -204,17 +204,11 @@ func (p *chainHop) Select(ctx context.Context, opts ...hop.SelectOption) *chain.
|
||||
return nodes[0]
|
||||
}
|
||||
|
||||
// For single-node case: bypass selector/FailFilter to ensure availability.
|
||||
// The marker system still works for metrics, but we don't block requests
|
||||
// based on recent failures - the connection will be attempted regardless.
|
||||
// This matches upstream go-gost/x behavior.
|
||||
if len(nodes) == 1 {
|
||||
return nodes[0]
|
||||
}
|
||||
|
||||
// Multi-node case: use selector with FailFilter for proper failover.
|
||||
// Use selector with FailFilter for proper failover.
|
||||
// FailFilter will exclude recently-failed nodes, allowing traffic to
|
||||
// be routed to healthy alternatives.
|
||||
// Note: FailFilter has a safety guard (len <= 1 returns as-is) to ensure
|
||||
// the last remaining node is never permanently blocked.
|
||||
if s := p.options.selector; s != nil {
|
||||
log.Debugf("[hop.Select] calling selector.Select with %d nodes", len(nodes))
|
||||
if node := s.Select(ctx, nodes...); node != nil {
|
||||
|
||||
@@ -2,11 +2,17 @@ package socket
|
||||
|
||||
import (
|
||||
"os"
|
||||
"sync"
|
||||
|
||||
"github.com/go-gost/x/config"
|
||||
)
|
||||
|
||||
// configMutex 保护配置文件的并发写入
|
||||
var configMutex sync.Mutex
|
||||
|
||||
func saveConfig() {
|
||||
configMutex.Lock()
|
||||
defer configMutex.Unlock()
|
||||
|
||||
file := "gost.json"
|
||||
|
||||
|
||||
@@ -466,7 +466,13 @@ func (w *WebSocketReporter) handleReceivedMessage(messageType int, message []byt
|
||||
}
|
||||
|
||||
if cmdMsg.Type != "call" {
|
||||
w.routeCommand(cmdMsg)
|
||||
// TcpPing 诊断命令异步执行,避免阻塞其他命令
|
||||
// 其他状态变更命令保持同步,确保顺序执行
|
||||
if cmdMsg.Type == "TcpPing" {
|
||||
go w.routeCommand(cmdMsg)
|
||||
} else {
|
||||
w.routeCommand(cmdMsg)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// 处理普通消息
|
||||
@@ -477,7 +483,13 @@ func (w *WebSocketReporter) handleReceivedMessage(messageType int, message []byt
|
||||
return
|
||||
}
|
||||
if cmdMsg.Type != "call" {
|
||||
w.routeCommand(cmdMsg)
|
||||
// TcpPing 诊断命令异步执行,避免阻塞其他命令
|
||||
// 其他状态变更命令保持同步,确保顺序执行
|
||||
if cmdMsg.Type == "TcpPing" {
|
||||
go w.routeCommand(cmdMsg)
|
||||
} else {
|
||||
w.routeCommand(cmdMsg)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -497,6 +509,7 @@ func (w *WebSocketReporter) routeCommand(cmd CommandMessage) {
|
||||
fmt.Println("🔔 收到命令: ", string(jsonBytes))
|
||||
var err error
|
||||
var response CommandResponse
|
||||
var needSaveConfig bool // 标记是否需要保存配置(只有状态变更命令才需要)
|
||||
|
||||
// 传递 requestId
|
||||
response.RequestId = cmd.RequestId
|
||||
@@ -506,65 +519,81 @@ func (w *WebSocketReporter) routeCommand(cmd CommandMessage) {
|
||||
case "AddService":
|
||||
err = w.handleAddService(cmd.Data)
|
||||
response.Type = "AddServiceResponse"
|
||||
needSaveConfig = true
|
||||
case "UpdateService":
|
||||
err = w.handleUpdateService(cmd.Data)
|
||||
response.Type = "UpdateServiceResponse"
|
||||
needSaveConfig = true
|
||||
case "DeleteService":
|
||||
err = w.handleDeleteService(cmd.Data)
|
||||
response.Type = "DeleteServiceResponse"
|
||||
needSaveConfig = true
|
||||
case "PauseService":
|
||||
err = w.handlePauseService(cmd.Data)
|
||||
response.Type = "PauseServiceResponse"
|
||||
needSaveConfig = true
|
||||
case "ResumeService":
|
||||
err = w.handleResumeService(cmd.Data)
|
||||
response.Type = "ResumeServiceResponse"
|
||||
needSaveConfig = true
|
||||
|
||||
// Chain 相关命令
|
||||
case "AddChains":
|
||||
err = w.handleAddChain(cmd.Data)
|
||||
response.Type = "AddChainsResponse"
|
||||
needSaveConfig = true
|
||||
case "UpdateChains":
|
||||
err = w.handleUpdateChain(cmd.Data)
|
||||
response.Type = "UpdateChainsResponse"
|
||||
needSaveConfig = true
|
||||
case "DeleteChains":
|
||||
err = w.handleDeleteChain(cmd.Data)
|
||||
response.Type = "DeleteChainsResponse"
|
||||
needSaveConfig = true
|
||||
|
||||
// Limiter 相关命令
|
||||
case "AddLimiters":
|
||||
err = w.handleAddLimiter(cmd.Data)
|
||||
response.Type = "AddLimitersResponse"
|
||||
needSaveConfig = true
|
||||
case "UpdateLimiters":
|
||||
err = w.handleUpdateLimiter(cmd.Data)
|
||||
response.Type = "UpdateLimitersResponse"
|
||||
needSaveConfig = true
|
||||
case "DeleteLimiters":
|
||||
err = w.handleDeleteLimiter(cmd.Data)
|
||||
response.Type = "DeleteLimitersResponse"
|
||||
needSaveConfig = true
|
||||
|
||||
// TCP Ping 诊断命令
|
||||
// TCP Ping 诊断命令(只读,不需要保存配置)
|
||||
case "TcpPing":
|
||||
var tcpPingResult TcpPingResponse
|
||||
tcpPingResult, err = w.handleTcpPing(cmd.Data)
|
||||
response.Type = "TcpPingResponse"
|
||||
response.Data = tcpPingResult
|
||||
// needSaveConfig = false (默认值)
|
||||
|
||||
// Protocol blocking switches
|
||||
case "SetProtocol":
|
||||
err = w.handleSetProtocol(cmd.Data)
|
||||
response.Type = "SetProtocolResponse"
|
||||
needSaveConfig = true
|
||||
|
||||
default:
|
||||
err = fmt.Errorf("未知命令类型: %s", cmd.Type)
|
||||
response.Type = "UnknownCommandResponse"
|
||||
}
|
||||
|
||||
// 只有状态变更命令才保存配置
|
||||
if needSaveConfig {
|
||||
saveConfig()
|
||||
}
|
||||
|
||||
// 发送响应
|
||||
if err != nil {
|
||||
saveConfig()
|
||||
response.Success = false
|
||||
response.Message = err.Error()
|
||||
} else {
|
||||
saveConfig()
|
||||
response.Success = true
|
||||
response.Message = "OK"
|
||||
}
|
||||
|
||||
@@ -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<ForwardMapper, Forward> impl
|
||||
.filter(ct -> ct.getChainType() == 3)
|
||||
.toList();
|
||||
|
||||
List<DiagnosisResult> results = new ArrayList<>();
|
||||
List<CompletableFuture<DiagnosisResult>> 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<ForwardMapper, Forward> 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<ForwardMapper, Forward> 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<ForwardMapper, Forward> impl
|
||||
// 2. 链路测试
|
||||
for (int i = 0; i < chainNodesList.size(); i++) {
|
||||
List<ChainTunnel> 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<ForwardMapper, Forward> 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<ForwardMapper, Forward> 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<DiagnosisResult> results = futures.stream()
|
||||
.map(CompletableFuture::join)
|
||||
.collect(Collectors.toList());
|
||||
|
||||
// 构建诊断报告
|
||||
Map<String, Object> diagnosisReport = new HashMap<>();
|
||||
diagnosisReport.put("forwardId", id);
|
||||
|
||||
@@ -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<TunnelMapper, Tunnel> impleme
|
||||
.filter(ct -> ct.getChainType() == 3)
|
||||
.toList();
|
||||
|
||||
List<DiagnosisResult> results = new ArrayList<>();
|
||||
List<CompletableFuture<DiagnosisResult>> 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<TunnelMapper, Tunnel> 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<TunnelMapper, Tunnel> impleme
|
||||
|
||||
for (int i = 0; i < chainNodesList.size(); i++) {
|
||||
List<ChainTunnel> 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<TunnelMapper, Tunnel> 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<TunnelMapper, Tunnel> 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<DiagnosisResult> results = futures.stream()
|
||||
.map(CompletableFuture::join)
|
||||
.collect(Collectors.toList());
|
||||
|
||||
Map<String, Object> diagnosisReport = new HashMap<>();
|
||||
diagnosisReport.put("tunnelId", tunnelId);
|
||||
diagnosisReport.put("tunnelName", tunnel.getName());
|
||||
|
||||
Reference in New Issue
Block a user