fix(diagnosis): parallelize runtime checks and raise API timeouts

This commit is contained in:
sagitchu
2026-02-28 18:46:36 +08:00
parent b01dbdb6e5
commit de21a55f37
3 changed files with 222 additions and 59 deletions
+209 -57
View File
@@ -8,6 +8,7 @@ import (
"sort"
"strconv"
"strings"
"sync"
"time"
"go-backend/internal/http/client"
@@ -30,6 +31,19 @@ type diagnosisTarget struct {
Port int
}
type diagnosisWorkItem struct {
fromNodeID int64
targetIP string
targetPort int
description string
metadata map[string]interface{}
toNode chainNodeRecord
hasChainHop bool
ipPreference string
}
const diagnosisMaxConcurrency = 8
func (h *Handler) resolveForwardAccess(r *http.Request, forwardID int64) (*forwardRecord, int64, int, error) {
userID, roleID, err := userRoleFromRequest(r)
if err != nil {
@@ -304,7 +318,7 @@ func (h *Handler) sendNodeCommand(nodeID int64, commandType string, data interfa
if nodeErr == nil && node != nil && node.IsRemote == 1 {
result, err = h.sendRemoteNodeCommand(node, commandType, data)
} else {
result, err = h.wsServer.SendCommand(nodeID, commandType, data, 12*time.Second)
result, err = h.wsServer.SendCommand(nodeID, commandType, data, 6*time.Second)
}
if err == nil {
return result, nil
@@ -386,16 +400,21 @@ func (h *Handler) diagnoseForwardRuntime(forward *forwardRecord) (map[string]int
ipPreference := h.repo.GetTunnelIPPreference(forward.TunnelID)
inNodes, chainHops, outNodes := splitChainNodeGroups(chainRows)
results := make([]map[string]interface{}, 0, len(chainRows)*2+len(targets))
nodeCache := map[int64]*nodeRecord{}
workItems := make([]diagnosisWorkItem, 0, len(chainRows)*2+len(targets))
switch tunnel.Type {
case 1:
for _, inNode := range inNodes {
for _, target := range targets {
description := fmt.Sprintf("入口(%s)->目标(%s)", inNode.NodeName, target.Address)
h.appendPathDiagnosis(&results, nodeCache, inNode.NodeID, target.IP, target.Port, description, map[string]interface{}{
"fromChainType": 1,
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: inNode.NodeID,
targetIP: target.IP,
targetPort: target.Port,
description: description,
metadata: map[string]interface{}{
"fromChainType": 1,
},
})
}
}
@@ -404,19 +423,33 @@ func (h *Handler) diagnoseForwardRuntime(forward *forwardRecord) (map[string]int
if len(chainHops) > 0 {
for _, firstNode := range chainHops[0] {
description := fmt.Sprintf("入口(%s)->第1跳(%s)", inNode.NodeName, firstNode.NodeName)
h.appendChainHopDiagnosis(&results, nodeCache, inNode.NodeID, firstNode, description, map[string]interface{}{
"fromChainType": 1,
"toChainType": 2,
"toInx": firstNode.Inx,
}, ipPreference)
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: inNode.NodeID,
toNode: firstNode,
hasChainHop: true,
ipPreference: ipPreference,
description: description,
metadata: map[string]interface{}{
"fromChainType": 1,
"toChainType": 2,
"toInx": firstNode.Inx,
},
})
}
} else {
for _, outNode := range outNodes {
description := fmt.Sprintf("入口(%s)->出口(%s)", inNode.NodeName, outNode.NodeName)
h.appendChainHopDiagnosis(&results, nodeCache, inNode.NodeID, outNode, description, map[string]interface{}{
"fromChainType": 1,
"toChainType": 3,
}, ipPreference)
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: inNode.NodeID,
toNode: outNode,
hasChainHop: true,
ipPreference: ipPreference,
description: description,
metadata: map[string]interface{}{
"fromChainType": 1,
"toChainType": 3,
},
})
}
}
}
@@ -426,21 +459,35 @@ func (h *Handler) diagnoseForwardRuntime(forward *forwardRecord) (map[string]int
if i+1 < len(chainHops) {
for _, nextNode := range chainHops[i+1] {
description := fmt.Sprintf("第%d跳(%s)->第%d跳(%s)", i+1, currentNode.NodeName, i+2, nextNode.NodeName)
h.appendChainHopDiagnosis(&results, nodeCache, currentNode.NodeID, nextNode, description, map[string]interface{}{
"fromChainType": 2,
"fromInx": currentNode.Inx,
"toChainType": 2,
"toInx": nextNode.Inx,
}, ipPreference)
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: currentNode.NodeID,
toNode: nextNode,
hasChainHop: true,
ipPreference: ipPreference,
description: description,
metadata: map[string]interface{}{
"fromChainType": 2,
"fromInx": currentNode.Inx,
"toChainType": 2,
"toInx": nextNode.Inx,
},
})
}
} else {
for _, outNode := range outNodes {
description := fmt.Sprintf("第%d跳(%s)->出口(%s)", i+1, currentNode.NodeName, outNode.NodeName)
h.appendChainHopDiagnosis(&results, nodeCache, currentNode.NodeID, outNode, description, map[string]interface{}{
"fromChainType": 2,
"fromInx": currentNode.Inx,
"toChainType": 3,
}, ipPreference)
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: currentNode.NodeID,
toNode: outNode,
hasChainHop: true,
ipPreference: ipPreference,
description: description,
metadata: map[string]interface{}{
"fromChainType": 2,
"fromInx": currentNode.Inx,
"toChainType": 3,
},
})
}
}
}
@@ -449,8 +496,14 @@ func (h *Handler) diagnoseForwardRuntime(forward *forwardRecord) (map[string]int
for _, outNode := range outNodes {
for _, target := range targets {
description := fmt.Sprintf("出口(%s)->目标(%s)", outNode.NodeName, target.Address)
h.appendPathDiagnosis(&results, nodeCache, outNode.NodeID, target.IP, target.Port, description, map[string]interface{}{
"fromChainType": 3,
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: outNode.NodeID,
targetIP: target.IP,
targetPort: target.Port,
description: description,
metadata: map[string]interface{}{
"fromChainType": 3,
},
})
}
}
@@ -458,13 +511,21 @@ func (h *Handler) diagnoseForwardRuntime(forward *forwardRecord) (map[string]int
for _, inNode := range inNodes {
for _, target := range targets {
description := fmt.Sprintf("入口(%s)->目标(%s)", inNode.NodeName, target.Address)
h.appendPathDiagnosis(&results, nodeCache, inNode.NodeID, target.IP, target.Port, description, map[string]interface{}{
"fromChainType": 1,
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: inNode.NodeID,
targetIP: target.IP,
targetPort: target.Port,
description: description,
metadata: map[string]interface{}{
"fromChainType": 1,
},
})
}
}
}
results := h.runDiagnosisWorkItems(workItems)
payload := map[string]interface{}{
"forwardName": forward.Name,
"timestamp": time.Now().UnixMilli(),
@@ -497,15 +558,20 @@ func (h *Handler) diagnoseTunnelRuntime(tunnelID int64) (map[string]interface{},
ipPreference := h.repo.GetTunnelIPPreference(tunnelID)
inNodes, chainHops, outNodes := splitChainNodeGroups(chainRows)
results := make([]map[string]interface{}, 0, len(chainRows)*2)
nodeCache := map[int64]*nodeRecord{}
workItems := make([]diagnosisWorkItem, 0, len(chainRows)*2)
switch tunnel.Type {
case 1:
for _, inNode := range inNodes {
description := fmt.Sprintf("入口(%s)->外网", inNode.NodeName)
h.appendPathDiagnosis(&results, nodeCache, inNode.NodeID, "www.bing.com", 443, description, map[string]interface{}{
"fromChainType": 1,
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: inNode.NodeID,
targetIP: "www.bing.com",
targetPort: 443,
description: description,
metadata: map[string]interface{}{
"fromChainType": 1,
},
})
}
case 2:
@@ -513,19 +579,33 @@ func (h *Handler) diagnoseTunnelRuntime(tunnelID int64) (map[string]interface{},
if len(chainHops) > 0 {
for _, firstNode := range chainHops[0] {
description := fmt.Sprintf("入口(%s)->第1跳(%s)", inNode.NodeName, firstNode.NodeName)
h.appendChainHopDiagnosis(&results, nodeCache, inNode.NodeID, firstNode, description, map[string]interface{}{
"fromChainType": 1,
"toChainType": 2,
"toInx": firstNode.Inx,
}, ipPreference)
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: inNode.NodeID,
toNode: firstNode,
hasChainHop: true,
ipPreference: ipPreference,
description: description,
metadata: map[string]interface{}{
"fromChainType": 1,
"toChainType": 2,
"toInx": firstNode.Inx,
},
})
}
} else {
for _, outNode := range outNodes {
description := fmt.Sprintf("入口(%s)->出口(%s)", inNode.NodeName, outNode.NodeName)
h.appendChainHopDiagnosis(&results, nodeCache, inNode.NodeID, outNode, description, map[string]interface{}{
"fromChainType": 1,
"toChainType": 3,
}, ipPreference)
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: inNode.NodeID,
toNode: outNode,
hasChainHop: true,
ipPreference: ipPreference,
description: description,
metadata: map[string]interface{}{
"fromChainType": 1,
"toChainType": 3,
},
})
}
}
}
@@ -535,21 +615,35 @@ func (h *Handler) diagnoseTunnelRuntime(tunnelID int64) (map[string]interface{},
if i+1 < len(chainHops) {
for _, nextNode := range chainHops[i+1] {
description := fmt.Sprintf("第%d跳(%s)->第%d跳(%s)", i+1, currentNode.NodeName, i+2, nextNode.NodeName)
h.appendChainHopDiagnosis(&results, nodeCache, currentNode.NodeID, nextNode, description, map[string]interface{}{
"fromChainType": 2,
"fromInx": currentNode.Inx,
"toChainType": 2,
"toInx": nextNode.Inx,
}, ipPreference)
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: currentNode.NodeID,
toNode: nextNode,
hasChainHop: true,
ipPreference: ipPreference,
description: description,
metadata: map[string]interface{}{
"fromChainType": 2,
"fromInx": currentNode.Inx,
"toChainType": 2,
"toInx": nextNode.Inx,
},
})
}
} else {
for _, outNode := range outNodes {
description := fmt.Sprintf("第%d跳(%s)->出口(%s)", i+1, currentNode.NodeName, outNode.NodeName)
h.appendChainHopDiagnosis(&results, nodeCache, currentNode.NodeID, outNode, description, map[string]interface{}{
"fromChainType": 2,
"fromInx": currentNode.Inx,
"toChainType": 3,
}, ipPreference)
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: currentNode.NodeID,
toNode: outNode,
hasChainHop: true,
ipPreference: ipPreference,
description: description,
metadata: map[string]interface{}{
"fromChainType": 2,
"fromInx": currentNode.Inx,
"toChainType": 3,
},
})
}
}
}
@@ -557,19 +651,33 @@ func (h *Handler) diagnoseTunnelRuntime(tunnelID int64) (map[string]interface{},
for _, outNode := range outNodes {
description := fmt.Sprintf("出口(%s)->外网", outNode.NodeName)
h.appendPathDiagnosis(&results, nodeCache, outNode.NodeID, "www.bing.com", 443, description, map[string]interface{}{
"fromChainType": 3,
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: outNode.NodeID,
targetIP: "www.bing.com",
targetPort: 443,
description: description,
metadata: map[string]interface{}{
"fromChainType": 3,
},
})
}
default:
for _, inNode := range inNodes {
description := fmt.Sprintf("入口(%s)->外网", inNode.NodeName)
h.appendPathDiagnosis(&results, nodeCache, inNode.NodeID, "www.bing.com", 443, description, map[string]interface{}{
"fromChainType": 1,
workItems = append(workItems, diagnosisWorkItem{
fromNodeID: inNode.NodeID,
targetIP: "www.bing.com",
targetPort: 443,
description: description,
metadata: map[string]interface{}{
"fromChainType": 1,
},
})
}
}
results := h.runDiagnosisWorkItems(workItems)
payload := map[string]interface{}{
"tunnelName": tunnelName,
"tunnelType": map[bool]string{true: "端口转发", false: "隧道转发"}[tunnel.Type == 1],
@@ -628,6 +736,50 @@ func resolveDiagnosisTargets(remoteAddr string) ([]diagnosisTarget, error) {
return targets, nil
}
func (h *Handler) runDiagnosisWorkItems(workItems []diagnosisWorkItem) []map[string]interface{} {
results := make([]map[string]interface{}, len(workItems))
if len(workItems) == 0 {
return results
}
workerLimit := diagnosisMaxConcurrency
if workerLimit < 1 {
workerLimit = 1
}
if workerLimit > len(workItems) {
workerLimit = len(workItems)
}
semaphore := make(chan struct{}, workerLimit)
var wg sync.WaitGroup
for i := range workItems {
wg.Add(1)
semaphore <- struct{}{}
go func(index int) {
defer wg.Done()
defer func() { <-semaphore }()
workItem := workItems[index]
single := make([]map[string]interface{}, 0, 1)
nodeCache := map[int64]*nodeRecord{}
if workItem.hasChainHop {
h.appendChainHopDiagnosis(&single, nodeCache, workItem.fromNodeID, workItem.toNode, workItem.description, workItem.metadata, workItem.ipPreference)
} else {
h.appendPathDiagnosis(&single, nodeCache, workItem.fromNodeID, workItem.targetIP, workItem.targetPort, workItem.description, workItem.metadata)
}
if len(single) == 0 {
results[index] = newDiagnosisResultItem(workItem.fromNodeID, workItem.targetIP, workItem.targetPort, workItem.description, workItem.metadata)
results[index]["success"] = false
results[index]["message"] = "诊断任务未返回结果"
return
}
results[index] = single[0]
}(i)
}
wg.Wait()
return results
}
func (h *Handler) cachedNode(nodeCache map[int64]*nodeRecord, nodeID int64) (*nodeRecord, error) {
if node, ok := nodeCache[nodeID]; ok {
return node, nil
+3
View File
@@ -80,5 +80,8 @@
"vite": "npm:rolldown-vite@^7.3.1",
"vite-plugin-pwa": "^1.1.0",
"vite-tsconfig-paths": "^6.0.5"
},
"overrides": {
"serialize-javascript": "7.0.3"
}
}
+10 -2
View File
@@ -118,7 +118,11 @@ export const updateTunnel = (data: TunnelMutationPayload) =>
export const deleteTunnel = (id: number) =>
Network.post("/tunnel/delete", { id });
export const diagnoseTunnel = (tunnelId: number) =>
Network.post<TunnelDiagnosisApiData>("/tunnel/diagnose", { tunnelId });
Network.post<TunnelDiagnosisApiData>(
"/tunnel/diagnose",
{ tunnelId },
{ timeout: 60 * 1000 },
);
export const updateTunnelOrder = (data: {
tunnels: Array<{ id: number; inx: number }>;
}) => Network.post("/tunnel/update-order", data);
@@ -159,7 +163,11 @@ export const resumeForwardService = (forwardId: number) =>
// 转发诊断操作
export const diagnoseForward = (forwardId: number) =>
Network.post<ForwardDiagnosisApiData>("/forward/diagnose", { forwardId });
Network.post<ForwardDiagnosisApiData>(
"/forward/diagnose",
{ forwardId },
{ timeout: 60 * 1000 },
);
// 转发排序操作
export const updateForwardOrder = (data: {