diff --git a/go-backend/internal/http/handler/control_plane.go b/go-backend/internal/http/handler/control_plane.go index a366b06..86b691d 100644 --- a/go-backend/internal/http/handler/control_plane.go +++ b/go-backend/internal/http/handler/control_plane.go @@ -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 diff --git a/vite-frontend/package.json b/vite-frontend/package.json index b8ac354..cc7fa3d 100644 --- a/vite-frontend/package.json +++ b/vite-frontend/package.json @@ -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" } } diff --git a/vite-frontend/src/api/index.ts b/vite-frontend/src/api/index.ts index 4df128b..3e84015 100644 --- a/vite-frontend/src/api/index.ts +++ b/vite-frontend/src/api/index.ts @@ -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("/tunnel/diagnose", { tunnelId }); + Network.post( + "/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("/forward/diagnose", { forwardId }); + Network.post( + "/forward/diagnose", + { forwardId }, + { timeout: 60 * 1000 }, + ); // 转发排序操作 export const updateForwardOrder = (data: {