Compare commits

...

27 Commits

Author SHA1 Message Date
sagit 02ff215f99 Merge pull request #28 from Sagit-chu/opencode/lucky-eagle
fix(gost): process WebSocket commands concurrently to prevent diagnos…
2026-02-05 14:25:24 +08:00
root 2a4e7777ab fix(gost): run TcpPing commands concurrently without config save race
When multiple TcpPing requests are sent in parallel for diagnosing
multiple remote addresses, the Go agent was processing them serially.
This caused later requests to timeout (10s) while waiting for earlier
requests to complete.

Changes:
- Only TcpPing commands run in goroutines for parallel execution
- TcpPing (read-only diagnostic) no longer triggers saveConfig()
- Other state-mutating commands remain synchronous with config save
- Add mutex to saveConfig() to protect concurrent file writes
2026-02-05 06:18:56 +00:00
sagit eac94a5719 Merge pull request #26 from Sagit-chu/opencode/hidden-pixel
fix(diagnose): parallelize TCP ping diagnostics to prevent timeout ca…
2026-02-05 12:57:10 +08:00
root 96fcd0fc57 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.
2026-02-05 04:52:42 +00:00
sagit 06869aedfd Merge pull request #25 from Sagit-chu/opencode/kind-sailor
fix(gost): mark node failed when transport detects relay error
2026-02-05 12:21:32 +08:00
sagit 6a201131a3 Merge branch 'main' into opencode/kind-sailor 2026-02-05 12:17:27 +08:00
root 1130a55ef5 fix(gost): mark node failed when transport detects relay error
When using relay connector with noDelay=false (default), connection
errors to the final target are deferred until first read/write during
Transport(). Previously the Transport() return value was ignored,
causing the marker to never be called for unreachable targets.

Now we capture the Transport() error and mark the node as failed,
enabling failover for subsequent connections.
2026-02-05 04:09:57 +00:00
root 265cd0a50e Revert "fix(backend): enable noDelay for relay connector to fix chain failover"
This reverts commit 51cbd4b9de.
2026-02-05 04:06:46 +00:00
sagit e7ffa77b15 Merge pull request #24 from Sagit-chu/opencode/kind-sailor
fix(backend): enable noDelay for relay connector to fix chain failover
2026-02-05 11:14:46 +08:00
root 51cbd4b9de fix(backend): enable noDelay for relay connector to fix chain failover
When using relay connector with noDelay=false (default), connection
errors are deferred until first read/write. This prevents the forwarder
marker from being called, causing failover to never trigger.

Setting nodelay=true ensures connection errors propagate immediately,
allowing proper failover behavior when chain targets are unreachable.
2026-02-05 03:11:52 +00:00
sagit e122e7460d Merge pull request #23 from Sagit-chu/opencode/crisp-cabin
fix(gost): remove single-node optimization to enable forwarder failover
2026-02-05 10:18:00 +08:00
root 7c898154b3 fix(gost): remove single-node optimization to enable forwarder failover
The single-node bypass in hop.Select() was preventing FailFilter from
being applied when retry excludes reduced available nodes to one.
This caused failed forwarder nodes to keep being selected instead of
failing over to healthy alternatives.

FailFilter's built-in safety guard (len <= 1 returns as-is) ensures
the last remaining node is never permanently blocked.
2026-02-05 02:14:37 +00:00
sagit 1d19d68019 Merge pull request #22 from Sagit-chu/feat/failover-debug-logging
feat(gost): add debug logging for failover mechanism analysis
2026-02-05 09:09:09 +08:00
root 09c58e2298 feat(gost): add debug logging for failover mechanism analysis
Add debug logs to trace failover behavior:
- FailFilter.Filter(): log node name, fail count, maxFails, timeSince, failTimeout
- hop.Select(): log excludeNodes list, node selection results
- handler retry loop: log maxRetries, selected nodes, dial failures

This helps diagnose issues where failover between multiple target nodes
is not working as expected.
2026-02-05 01:06:56 +00:00
sagit 583905b7ed Merge pull request #21 from Sagit-chu/opencode/neon-nebula
fix(gost): use chain.NewNode() to properly initialize marker for fail…
2026-02-05 07:29:15 +08:00
sagit 6e3f045b9b Merge branch 'main' into opencode/neon-nebula 2026-02-05 07:26:49 +08:00
root ec41202b3c fix(gost): use chain.NewNode() to properly initialize marker for failover
When creating temporary Node instances with struct literals like
&chain.Node{Addr: host}, the marker field was not initialized.
Only chain.NewNode() properly initializes marker = selector.NewFailMarker().

Without a valid marker:
- Failed nodes cannot be marked (marker.Mark() is no-op on nil)
- Subsequent selections cannot filter out failed nodes
- Failover mechanism completely fails

Fixed locations:
- sniffer.go dial(): &chain.Node{Addr: host} -> chain.NewNode("", host)
- sniffer.go dialTLS(): &chain.Node{Addr: host} -> chain.NewNode("", host)
- local/handler.go: target := &chain.Node{} -> var target *chain.Node
- remote/handler.go: &chain.Node{Addr: host} -> chain.NewNode("", host)
2026-02-04 23:10:56 +00:00
sagit 1ee7dea8b4 Merge pull request #20 from Sagit-chu/Sagit-chu-patch-1
change beta to main
2026-02-04 16:55:44 +08:00
sagit 2d69350bab change beta to main 2026-02-04 16:54:15 +08:00
sagit c984e5b62a docs: remove stable installation instructions
docs: remove stable installation instructions
2026-02-04 16:53:06 +08:00
sagit 5b79b11101 Merge branch 'beta' into opencode/calm-sailor 2026-02-04 16:50:28 +08:00
root 2c22e600f7 docs: remove stable installation instructions 2026-02-04 08:46:41 +00:00
sagit aef284c474 Merge pull request #18 from Sagit-chu/opencode/sunny-wizard
fix(gost): sync agent version with release tag
2026-02-04 16:29:24 +08:00
root 0443cd9ceb fix(gost): sync agent version with release tag
- Change version.go default to 'dev' for local development
- Use version variable in WebSocket reporter instead of hardcoded '2.0.2'
- Inject version via -ldflags in CI build from tag name
2026-02-04 08:22:37 +00:00
sagit 3337422775 Merge pull request #17 from Sagit-chu/opencode/curious-nebula
fix(gost): add fallback when FailFilter excludes all nodes
2026-02-04 15:52:19 +08:00
root 3e046fc80e fix(gost): restore single-node bypass and preserve FailFilter backoff
Address reviewer feedback from PR #14 fix:

1. Single-node case: Bypass selector/FailFilter to ensure availability.
   This matches upstream go-gost/x behavior - single nodes should always
   be attempted regardless of recent failures.

2. Multi-node case: Preserve FailFilter's backoff contract. When all nodes
   are marked as failed, return nil to signal 'no healthy nodes' rather
   than falling back to a known-bad node. This prevents hammering unhealthy
   nodes and respects the failTimeout window.

The handler's retry loop with ExcludeNodes context handles the multi-node
failover properly - this change ensures hop.Select() provides correct
information about node health status.

Fixes intermittent forwarding failures introduced by #14.
2026-02-04 07:46:40 +00:00
root 0273bc6921 docs: add AGENTS.md for go-gost/x/registry 2026-02-04 06:36:03 +00:00
14 changed files with 298 additions and 132 deletions
+2 -2
View File
@@ -100,11 +100,11 @@ jobs:
- name: Build GOST binary (AMD64)
working-directory: ./go-gost
run: CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -ldflags="-s -w" -o gost-amd64
run: CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -ldflags="-s -w -X main.version=${{ needs.check-version.outputs.version }}" -o gost-amd64
- name: Build GOST binary (ARM64)
working-directory: ./go-gost
run: CGO_ENABLED=0 GOOS=linux GOARCH=arm64 go build -ldflags="-s -w" -o gost-arm64
run: CGO_ENABLED=0 GOOS=linux GOARCH=arm64 go build -ldflags="-s -w -X main.version=${{ needs.check-version.outputs.version }}" -o gost-arm64
- name: Compress with UPX
working-directory: ./go-gost
+2 -13
View File
@@ -16,24 +16,13 @@
---
### Docker Compose部署
#### 快速部署
面板端(稳定版):
面板端:
```bash
curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/panel_install.sh -o panel_install.sh && chmod +x panel_install.sh && ./panel_install.sh
```
节点端(稳定版):
节点端:
```bash
curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/install.sh -o install.sh && chmod +x install.sh && ./install.sh
```
面板端(开发版):
```bash
curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/beta/panel_install.sh -o panel_install.sh && chmod +x panel_install.sh && ./panel_install.sh
```
节点端(开发版):
```bash
curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/beta/install.sh -o install.sh && chmod +x install.sh && ./install.sh
```
#### 默认管理员账号
+1 -1
View File
@@ -119,7 +119,7 @@ func main() {
log := xlogger.NewLogger()
logger.SetDefault(log)
wsReporter := socket.StartWebSocketReporterWithConfig(config.Addr, config.Secret, config.Http, config.Tls, config.Socks, "2.0.2")
wsReporter := socket.StartWebSocketReporterWithConfig(config.Addr, config.Secret, config.Http, config.Tls, config.Socks, version)
defer wsReporter.Stop()
service.SetHTTPReportURL(config.Addr, config.Secret)
+1 -1
View File
@@ -1,5 +1,5 @@
package main
var (
version = "3.1.0"
version = "dev"
)
+16 -2
View File
@@ -192,22 +192,27 @@ func (h *forwardHandler) Handle(ctx context.Context, conn net.Conn, opts ...hand
var lastErr error
var cc net.Conn
h.options.Logger.Debugf("[handler.retry] starting retry loop: maxRetries=%d", maxRetries)
for attempt := 0; attempt < maxRetries; attempt++ {
// Select a target node, excluding previously tried nodes
selectCtx := ctxvalue.ContextWithExcludeNodes(ctx, triedNodes)
target := &chain.Node{}
var target *chain.Node
if h.hop != nil {
target = h.hop.Select(selectCtx,
hop.ProtocolSelectOption(proto),
)
}
if target == nil {
h.options.Logger.Debugf("[handler.retry] attempt=%d target=nil, triedNodes=%v", attempt, triedNodes)
if lastErr != nil {
return lastErr
}
return errors.New("node not available")
}
h.options.Logger.Debugf("[handler.retry] attempt=%d selected node=%s addr=%s", attempt, target.Name, target.Addr)
// Track this node as tried
triedNodes = append(triedNodes, target.Addr)
@@ -233,6 +238,8 @@ func (h *forwardHandler) Handle(ctx context.Context, conn net.Conn, opts ...hand
// Mark node as failed for future selections
if marker := target.Marker(); marker != nil {
marker.Mark()
h.options.Logger.Debugf("[handler.retry] attempt=%d dial failed, marked node=%s count=%d err=%v",
attempt, target.Addr, marker.Count(), err)
}
lastErr = err
// Try next node
@@ -245,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
}
+1 -3
View File
@@ -225,9 +225,7 @@ func (h *forwardHandler) Handle(ctx context.Context, conn net.Conn, opts ...hand
selectCtx := ctxvalue.ContextWithExcludeNodes(ctx, triedNodes)
var target *chain.Node
if host != "" {
target = &chain.Node{
Addr: host,
}
target = chain.NewNode("", host)
}
if h.hop != nil {
target = h.hop.Select(selectCtx,
+18 -4
View File
@@ -149,6 +149,9 @@ func (p *chainHop) Select(ctx context.Context, opts ...hop.SelectOption) *chain.
excludeSet[addr] = true
}
// Debug logging for failover analysis
log.Debugf("[hop.Select] excludeNodes=%v, totalNodes=%d", excludeNodes, len(p.Nodes()))
var nodes []*chain.Node
for _, node := range p.Nodes() {
if node == nil {
@@ -201,11 +204,22 @@ func (p *chainHop) Select(ctx context.Context, opts ...hop.SelectOption) *chain.
return nodes[0]
}
// Always go through selector for proper FailFilter evaluation,
// even when there's only one node. This ensures failed nodes
// can be filtered out properly.
// 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 {
return s.Select(ctx, nodes...)
log.Debugf("[hop.Select] calling selector.Select with %d nodes", len(nodes))
if node := s.Select(ctx, nodes...); node != nil {
log.Debugf("[hop.Select] selected node=%s addr=%s", node.Name, node.Addr)
return node
}
// All nodes filtered out by FailFilter - all are marked as failed.
// Return nil to signal "no healthy nodes available" to the caller.
// The handler's retry loop will handle this appropriately.
log.Debugf("all %d nodes filtered out by FailFilter, no healthy nodes available", len(nodes))
return nil
}
// Fallback: return first node if no selector configured
+2 -6
View File
@@ -263,9 +263,7 @@ func (h *Sniffer) dial(ctx context.Context, conn net.Conn, req *http.Request, ho
// Select a node, excluding previously tried nodes
selectCtx := ctxvalue.ContextWithExcludeNodes(ctx, triedNodes)
node = &chain.Node{
Addr: host,
}
node = chain.NewNode("", host)
if ho.Hop != nil {
node = ho.Hop.Select(selectCtx,
hop.ClientIPSelectOption(net.ParseIP(ro.ClientIP)),
@@ -903,9 +901,7 @@ func (h *Sniffer) dialTLS(ctx context.Context, host string, ho *HandleOptions) (
node = nil
if host != "" {
node = &chain.Node{
Addr: host,
}
node = chain.NewNode("", host)
}
if ho.Hop != nil {
node = ho.Hop.Select(selectCtx,
+29
View File
@@ -0,0 +1,29 @@
# GO-GOST REGISTRY KNOWLEDGE BASE
**Generated:** Wed Feb 04 2026
## OVERVIEW
Central registration point for all pluggable GOST components (handlers, listeners, dialers, etc.).
Allows the configuration system to resolve string types (e.g., "socks5") to actual Go implementations.
## STRUCTURE
One file per component type, exporting a standard Registry interface.
```
go-gost/x/registry/
├── handler.go # RegisterHandler(name, newFunc)
├── listener.go # RegisterListener(name, newFunc)
├── dialer.go # RegisterDialer(name, newFunc)
└── ... # Same pattern for auth, bypass, admission
```
## WHERE TO LOOK
| Task | Location | Notes |
|------|----------|-------|
| Register a new component | `go-gost/x/registry/{type}.go` | Use `Register{Type}(name, creator)` |
| Component lookup | `go-gost/x/registry/{type}.go` | `Get{Type}(name)` returns the creator function |
| Default registrations | `go-gost/x/` (init functions) | Most components register themselves in their package `init()` |
## CONVENTIONS
- Thread-safe maps used for storage.
- Names are case-sensitive (usually lowercase).
- Components must be registered *before* the configuration parser runs (usually done via `import _ "..."` in `main.go`).
+22 -6
View File
@@ -2,8 +2,10 @@ package selector
import (
"context"
"fmt"
"time"
"github.com/go-gost/core/chain"
"github.com/go-gost/core/metadata"
"github.com/go-gost/core/selector"
mdutil "github.com/go-gost/x/metadata/util"
@@ -24,11 +26,12 @@ func FailFilter[T any](maxFails int, timeout time.Duration) selector.Filter[T] {
}
// Filter filters dead objects.
// Note: We intentionally do NOT skip filtering when len(vs) <= 1.
// This ensures that even a single dead node gets filtered out,
// allowing the caller to know that no healthy nodes are available
// and potentially trigger failover behavior.
// For single-node case, skip filtering to ensure availability (matches upstream).
// For multi-node case, filter out failed nodes to enable failover.
func (f *failFilter[T]) Filter(ctx context.Context, vs ...T) []T {
if len(vs) <= 1 {
return vs
}
var l []T
for _, v := range vs {
maxFails := f.maxFails
@@ -52,8 +55,21 @@ func (f *failFilter[T]) Filter(ctx context.Context, vs ...T) []T {
if mi, _ := any(v).(selector.Markable); mi != nil {
if marker := mi.Marker(); marker != nil {
if marker.Count() < int64(maxFails) ||
time.Since(marker.Time()) >= failTimeout {
count := marker.Count()
timeSince := time.Since(marker.Time())
passed := count < int64(maxFails) || timeSince >= failTimeout
// Debug logging for failover analysis
nodeName := "unknown"
nodeAddr := "unknown"
if node, ok := any(v).(*chain.Node); ok {
nodeName = node.Name
nodeAddr = node.Addr
}
fmt.Printf("[FailFilter] node=%s addr=%s count=%d maxFails=%d timeSince=%v failTimeout=%v passed=%v\n",
nodeName, nodeAddr, count, maxFails, timeSince, failTimeout, passed)
if passed {
l = append(l, v)
}
continue
+6
View File
@@ -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"
+34 -5
View File
@@ -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());