diff --git a/go-backend/internal/http/handler/handler.go b/go-backend/internal/http/handler/handler.go index c90cd56..7c87742 100644 --- a/go-backend/internal/http/handler/handler.go +++ b/go-backend/internal/http/handler/handler.go @@ -1234,6 +1234,21 @@ func nullableNullInt64(v sql.NullInt64) interface{} { return nil } +// flowCryptoCache caches AES crypto instances by secret to avoid per-request SHA256+GCM init. +var flowCryptoCache sync.Map + +func getOrCreateFlowCrypto(secret string) *security.AESCrypto { + if v, ok := flowCryptoCache.Load(secret); ok { + return v.(*security.AESCrypto) + } + c, err := security.NewAESCrypto(secret) + if err != nil { + return nil + } + flowCryptoCache.Store(secret, c) + return c +} + func readAndDecryptFlowBody(body io.ReadCloser, secret string) (string, error) { defer body.Close() raw, err := io.ReadAll(body) @@ -1254,8 +1269,8 @@ func readAndDecryptFlowBody(body io.ReadCloser, secret string) (string, error) { return text, nil } - crypto, err := security.NewAESCrypto(secret) - if err != nil { + crypto := getOrCreateFlowCrypto(secret) + if crypto == nil { return text, nil } plain, err := crypto.Decrypt(wrap.Data) diff --git a/go-backend/internal/ws/server.go b/go-backend/internal/ws/server.go index 1e24116..52ff3bf 100644 --- a/go-backend/internal/ws/server.go +++ b/go-backend/internal/ws/server.go @@ -39,6 +39,7 @@ type nodeSession struct { nodeID int64 secret string conn *connWrap + crypto *security.AESCrypto // 缓存的 AES 加密器,避免每条消息重建 } type commandResponse struct { @@ -211,7 +212,12 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64 _ = old.conn.conn.Close() delete(s.byConn, old.conn.conn) } - ns := &nodeSession{nodeID: nodeID, secret: secret, conn: cw} + // 初始化 AES 加密器并缓存(仅创建一次) + var nodeCrypto *security.AESCrypto + if strings.TrimSpace(secret) != "" { + nodeCrypto, _ = security.NewAESCrypto(secret) + } + ns := &nodeSession{nodeID: nodeID, secret: secret, conn: cw, crypto: nodeCrypto} s.nodes[nodeID] = ns s.byConn[conn] = ns s.mu.Unlock() @@ -251,7 +257,7 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64 return } - msg := decryptIfNeeded(payload, secret) + msg := decryptIfNeeded(payload, ns.crypto, secret) s.tryResolvePending(nodeID, msg) var parsed struct { @@ -259,6 +265,26 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64 } if json.Unmarshal([]byte(msg), &parsed) == nil && parsed.Type != "" { switch parsed.Type { + case "metric": + // Agent 新版指标消息:{type:"metric", data:{...}} + var envelope struct { + Data json.RawMessage `json:"data"` + } + if err := json.Unmarshal([]byte(msg), &envelope); err == nil && len(envelope.Data) > 0 { + // 解析 SystemInfo 并调用 hook + var sysInfo SystemInfo + if json.Unmarshal(envelope.Data, &sysInfo) == nil { + s.mu.RLock() + onMetric := s.onNodeMetric + s.mu.RUnlock() + if onMetric != nil { + go onMetric(nodeID, sysInfo) + } + } + // 广播内层 data 给前端(保持平坦结构兼容性) + s.broadcastTyped(nodeID, "metric", string(envelope.Data)) + } + continue case "UpgradeProgress": s.broadcastTyped(nodeID, "upgrade_progress", msg) continue @@ -270,6 +296,7 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64 } } + // 兼容旧版 Agent:无 type 字段的系统信息消息 if looksLikeSystemInfoMessage(msg) { var sysInfo SystemInfo if err := json.Unmarshal([]byte(msg), &sysInfo); err == nil { @@ -372,13 +399,8 @@ func (s *Server) SendCommand(nodeID int64, cmdType string, data interface{}, tim } messageData := rawCmd - if strings.TrimSpace(ns.secret) != "" { - crypto, err := security.NewAESCrypto(ns.secret) - if err != nil { - cleanup() - return CommandResult{}, err - } - encrypted, err := crypto.Encrypt(rawCmd) + if ns.crypto != nil { + encrypted, err := ns.crypto.Encrypt(rawCmd) if err != nil { cleanup() return CommandResult{}, err @@ -428,6 +450,11 @@ func (s *Server) tryResolvePending(nodeID int64, message string) { return } + // 快速短路:指标消息永远不含 requestId,跳过完整 JSON 解析 + if !strings.Contains(message, "\"requestId\"") { + return + } + var resp commandResponse if err := json.Unmarshal([]byte(message), &resp); err != nil { return @@ -545,18 +572,22 @@ func (s *Server) broadcastToAdmins(message string) { } } -func decryptIfNeeded(payload []byte, secret string) string { +func decryptIfNeeded(payload []byte, crypto *security.AESCrypto, secret string) string { text := string(payload) var wrap encryptedMessage if err := json.Unmarshal(payload, &wrap); err != nil || !wrap.Encrypted || strings.TrimSpace(wrap.Data) == "" { return text } - crypto, err := security.NewAESCrypto(secret) - if err != nil { + // 优先使用缓存的 crypto 实例 + c := crypto + if c == nil && strings.TrimSpace(secret) != "" { + c, _ = security.NewAESCrypto(secret) + } + if c == nil { return text } - plain, err := crypto.Decrypt(wrap.Data) + plain, err := c.Decrypt(wrap.Data) if err != nil { return text } diff --git a/go-gost/x/socket/websocket_reporter.go b/go-gost/x/socket/websocket_reporter.go index b2758f0..e95d0b9 100644 --- a/go-gost/x/socket/websocket_reporter.go +++ b/go-gost/x/socket/websocket_reporter.go @@ -9,6 +9,7 @@ import ( "encoding/json" "fmt" "io" + "math/rand" "net" "net/http" "net/url" @@ -146,6 +147,9 @@ type ServiceMonitorCheckResult struct { const ( reporterReadWait = 60 * time.Second reporterWriteWait = 5 * time.Second + wsPingInterval = 20 * time.Second // 独立 WebSocket ping 间隔 + initialBackoff = 2 * time.Second // 重连初始退避 + maxBackoff = 2 * time.Minute // 重连最大退避 ) type WebSocketReporter struct { @@ -155,15 +159,15 @@ type WebSocketReporter struct { version string // 保存版本号 preferredWSScheme string conn *websocket.Conn - reconnectTime time.Duration + curBackoff time.Duration // 当前重连退避间隔 pingInterval time.Duration configInterval time.Duration ctx context.Context cancel context.CancelFunc connected bool - connecting bool // 新增:正在连接状态 - connMutex sync.Mutex // 新增:连接状态锁 - aesCrypto *crypto.AESCrypto // 新增:AES加密器 + connecting bool // 正在连接状态 + connMutex sync.Mutex // 连接状态锁 + aesCrypto *crypto.AESCrypto // AES加密器 } var wsDial = func(dialer *websocket.Dialer, rawURL string) (*websocket.Conn, *http.Response, error) { @@ -185,7 +189,7 @@ func NewWebSocketReporter(serverURL string, secret string) *WebSocketReporter { return &WebSocketReporter{ url: serverURL, - reconnectTime: 5 * time.Second, // 重连间隔 + curBackoff: initialBackoff, // 当前退避间隔 pingInterval: 5 * time.Second, // 指标上报间隔 configInterval: 10 * time.Minute, // 配置上报间隔 ctx: ctx, @@ -204,10 +208,17 @@ func (w *WebSocketReporter) Start() { // Stop 停止WebSocket报告器 func (w *WebSocketReporter) Stop() { w.cancel() + w.connMutex.Lock() if w.conn != nil { w.conn.Close() } + w.connMutex.Unlock() +} +// backoffWithJitter 返回带随机抖动的退避时间(±25%) +func backoffWithJitter(base time.Duration) time.Duration { + jitter := time.Duration(float64(base) * (0.75 + rand.Float64()*0.5)) + return jitter } // run 主运行循环 @@ -224,23 +235,32 @@ func (w *WebSocketReporter) run() { if needConnect { if err := w.connect(); err != nil { - fmt.Printf("❌ WebSocket连接失败: %v,%v后重试\n", err, w.reconnectTime) + wait := backoffWithJitter(w.curBackoff) + fmt.Printf("❌ WebSocket连接失败: %v,%v后重试\n", err, wait) + // 指数退避:翻倍当前退避间隔,上限 maxBackoff + w.curBackoff *= 2 + if w.curBackoff > maxBackoff { + w.curBackoff = maxBackoff + } select { - case <-time.After(w.reconnectTime): + case <-time.After(wait): continue case <-w.ctx.Done(): return } } + // 连接成功:重置退避 + w.curBackoff = initialBackoff } // 连接成功,开始发送消息 if w.connected { w.handleConnection() } else { + wait := backoffWithJitter(w.curBackoff) // 如果连接失败,等待重试 select { - case <-time.After(w.reconnectTime): + case <-time.After(wait): continue case <-w.ctx.Done(): return @@ -473,15 +493,34 @@ func (w *WebSocketReporter) handleConnection() { // 启动消息接收goroutine go w.receiveMessages() - // 主发送循环 - ticker := time.NewTicker(w.pingInterval) - defer ticker.Stop() + // 指标上报 ticker + metricTicker := time.NewTicker(w.pingInterval) + defer metricTicker.Stop() + + // 独立 WebSocket keepalive ping ticker + pingTicker := time.NewTicker(wsPingInterval) + defer pingTicker.Stop() for { select { case <-w.ctx.Done(): return - case <-ticker.C: + + case <-pingTicker.C: + // 发送 WebSocket ping 保活,独立于指标上报 + w.connMutex.Lock() + conn := w.conn + isConnected := w.connected + w.connMutex.Unlock() + if !isConnected || conn == nil { + return + } + if err := conn.WriteControl(websocket.PingMessage, nil, time.Now().Add(reporterWriteWait)); err != nil { + fmt.Printf("❌ 发送WebSocket ping失败: %v,准备重连\n", err) + return + } + + case <-metricTicker.C: // 检查连接状态 w.connMutex.Lock() isConnected := w.connected @@ -548,6 +587,37 @@ func (w *WebSocketReporter) collectSystemInfo() SystemInfo { } } +// encryptPayload 加密 JSON 数据,返回加密后的消息字节(若加密失败则回退到原始数据) +func (w *WebSocketReporter) encryptPayload(jsonData []byte) []byte { + if w.aesCrypto == nil { + return jsonData + } + + encryptedData, err := w.aesCrypto.Encrypt(jsonData) + if err != nil { + fmt.Printf("⚠️ 加密失败,发送原始数据: %v\n", err) + return jsonData + } + + encryptedMessage := map[string]interface{}{ + "encrypted": true, + "data": encryptedData, + "timestamp": time.Now().Unix(), + } + messageData, err := json.Marshal(encryptedMessage) + if err != nil { + fmt.Printf("⚠️ 序列化加密消息失败,发送原始数据: %v\n", err) + return jsonData + } + return messageData +} + +// metricEnvelope wraps SystemInfo with a type field for fast identification on the panel side. +type metricEnvelope struct { + Type string `json:"type"` + Data SystemInfo `json:"data"` +} + // sendSystemInfo 发送系统信息 func (w *WebSocketReporter) sendSystemInfo(sysInfo SystemInfo) error { w.connMutex.Lock() @@ -557,42 +627,19 @@ func (w *WebSocketReporter) sendSystemInfo(sysInfo SystemInfo) error { return fmt.Errorf("连接未建立") } - // 转换为JSON - jsonData, err := json.Marshal(sysInfo) + // 使用 type:"metric" 信封包装,Panel 可通过 type 字段直接识别指标消息 + envelope := metricEnvelope{Type: "metric", Data: sysInfo} + jsonData, err := json.Marshal(envelope) if err != nil { return fmt.Errorf("序列化系统信息失败: %v", err) } - var messageData []byte + messageData := w.encryptPayload(jsonData) - // 如果有加密器,则加密数据 - if w.aesCrypto != nil { - encryptedData, err := w.aesCrypto.Encrypt(jsonData) - if err != nil { - fmt.Printf("⚠️ 加密失败,发送原始数据: %v\n", err) - messageData = jsonData - } else { - // 创建加密消息包装器 - encryptedMessage := map[string]interface{}{ - "encrypted": true, - "data": encryptedData, - "timestamp": time.Now().Unix(), - } - messageData, err = json.Marshal(encryptedMessage) - if err != nil { - fmt.Printf("⚠️ 序列化加密消息失败,发送原始数据: %v\n", err) - messageData = jsonData - } - } - } else { - messageData = jsonData - } - - // 设置写入超时 w.conn.SetWriteDeadline(time.Now().Add(5 * time.Second)) if err := w.conn.WriteMessage(websocket.TextMessage, messageData); err != nil { - w.connected = false // 标记连接已断开 + w.connected = false return fmt.Errorf("写入消息失败: %v", err) } @@ -601,23 +648,19 @@ func (w *WebSocketReporter) sendSystemInfo(sysInfo SystemInfo) error { // receiveMessages 接收服务端发送的消息 func (w *WebSocketReporter) receiveMessages() { + // 获取连接引用一次即可,连接生命周期由 handleConnection 管理 + w.connMutex.Lock() + conn := w.conn + w.connMutex.Unlock() + if conn == nil { + return + } + for { select { case <-w.ctx.Done(): return default: - w.connMutex.Lock() - conn := w.conn - connected := w.connected - w.connMutex.Unlock() - - if conn == nil || !connected { - return - } - - // 设置读取超时 - conn.SetReadDeadline(time.Now().Add(reporterReadWait)) - messageType, message, err := conn.ReadMessage() if err != nil { if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) { @@ -705,12 +748,8 @@ func (w *WebSocketReporter) handleReceivedMessage(messageType int, message []byt } if cmdMsg.Type != "call" { - // 其他状态变更命令保持同步,确保顺序执行 - if cmdMsg.Type == "TcpPing" || cmdMsg.Type == "ServiceMonitorCheck" || cmdMsg.Type == "UpgradeAgent" || cmdMsg.Type == "RollbackAgent" { - go w.routeCommand(cmdMsg) - } else { - w.routeCommand(cmdMsg) - } + // 所有命令统一异步执行,避免阻塞消息接收循环 + go w.routeCommand(cmdMsg) } } else { // 处理普通消息 @@ -721,12 +760,8 @@ func (w *WebSocketReporter) handleReceivedMessage(messageType int, message []byt return } if cmdMsg.Type != "call" { - // 其他状态变更命令保持同步,确保顺序执行 - if cmdMsg.Type == "TcpPing" || cmdMsg.Type == "ServiceMonitorCheck" || cmdMsg.Type == "UpgradeAgent" || cmdMsg.Type == "RollbackAgent" { - go w.routeCommand(cmdMsg) - } else { - w.routeCommand(cmdMsg) - } + // 所有命令统一异步执行,避免阻塞消息接收循环 + go w.routeCommand(cmdMsg) } } @@ -1400,30 +1435,7 @@ func (w *WebSocketReporter) sendResponse(response CommandResponse) { return } - var messageData []byte - - // 如果有加密器,则加密数据 - if w.aesCrypto != nil { - encryptedData, err := w.aesCrypto.Encrypt(jsonData) - if err != nil { - fmt.Printf("⚠️ 加密响应失败,发送原始数据: %v\n", err) - messageData = jsonData - } else { - // 创建加密消息包装器 - encryptedMessage := map[string]interface{}{ - "encrypted": true, - "data": encryptedData, - "timestamp": time.Now().Unix(), - } - messageData, err = json.Marshal(encryptedMessage) - if err != nil { - fmt.Printf("⚠️ 序列化加密响应失败,发送原始数据: %v\n", err) - messageData = jsonData - } - } - } else { - messageData = jsonData - } + messageData := w.encryptPayload(jsonData) // 检查消息大小,如果超过10MB则记录警告 if len(messageData) > 10*1024*1024 { diff --git a/plans/054-agent-panel-comm-optimization.md b/plans/054-agent-panel-comm-optimization.md new file mode 100644 index 0000000..c4addd8 --- /dev/null +++ b/plans/054-agent-panel-comm-optimization.md @@ -0,0 +1,147 @@ +# Agent-Panel 通信优化:提升稳定性与效率 + +## 背景 + +Agent(`go-gost/x/socket/websocket_reporter.go`)与 Panel(`go-backend/internal/ws/server.go`)之间通过 WebSocket 进行实时通信,包括指标上报(每 5s)、命令下发/响应、和流量上报(HTTP)。经过代码审查,以下是发现的问题和优化建议。 + +--- + +## 发现的问题 + +### 1. Keepalive 时序不匹配 —— 导致误断连 + +| 参数 | Agent 侧 | Panel 侧 | +|------|----------|----------| +| Read deadline | `reporterReadWait` = 60s | `wsPongWait` = 45s | +| Ping 发送间隔 | 无主动 ping(靠指标数据 5s 续命) | `wsPingPeriod` = 15s | +| Write timeout | `reporterWriteWait` = 5s | `wsWriteWait` = 5s | + +**问题**:Panel 每 15s 发 ping,Agent read deadline 60s,但 Panel pong deadline 只有 45s。如果 Agent 的指标消息被延迟(网络抖动),Panel 可能因 pong 超时而关闭连接。两侧的超时参数缺乏协调设计。 + +### 2. 固定重连间隔 —— 无退避策略 + +Agent 断线后以固定 5s 间隔重试(`reconnectTime = 5 * time.Second`),在 Panel 长时间不可用(升级、网络故障)的情况下,会产生大量无用连接尝试。 + +### 3. Panel 侧每次解密都重建 AES 加密器 + +`ws/server.go` 的 `decryptIfNeeded()` 和 `SendCommand()` 每次调用都 `security.NewAESCrypto(secret)` 重新创建 cipher(SHA256 + AES-GCM 初始化),对于高频指标消息(5s/次 × N 节点),有不必要的 CPU 开销。 + +### 4. 指标消息使用 JSON Text 格式传输 + +每 5s 发送一次包含 13 个字段的 SystemInfo JSON,加密后还需 base64 编码,一条消息约 300-500 bytes(加密后约 700 bytes)。对于大量节点场景,存在优化空间。 + +### 5. `receiveMessages` 紧循环中有频繁锁竞争 + +`receiveMessages()` 在每次 `ReadMessage()` 前都要 `Lock/Unlock connMutex` 检查连接状态,但 `ReadMessage` 本身是阻塞的,实际不需要在循环外检查。 + +### 6. 状态变更命令阻塞读消息循环 + +`routeCommand` 中的 Service/Chain/Limiter CRUD 命令是同步执行的,包括 `saveConfig()` 文件写入。执行期间会阻塞 `receiveMessages` 的读取循环。 + +--- + +## 推荐的优化方案(按优先级排列) + +### P0 — 高收益、低风险 + +#### 优化 1:协调 Keepalive 参数 + +**文件**:`websocket_reporter.go` + +- Agent 增加独立的 WebSocket ping 发送(每 20s),不依赖指标数据来维持连接 +- 统一 read deadline 设置,确保两侧 read timeout > 2×ping interval + +#### 优化 2:指数退避重连 + +**文件**:`websocket_reporter.go` + +- 初始间隔 2s,按指数退避增长至最大 2 分钟 +- 连接成功后立即重置退避 +- 增加随机抖动(jitter)避免大量 Agent 同时重连 + +#### 优化 3:Panel 侧缓存 AES 加密器 + +**文件**:`ws/server.go` + +- 将 `AESCrypto` 实例缓存在 `nodeSession` 中,避免每条消息重建 +- `SendCommand` 复用缓存实例 + +### P1 — 中等收益 + +#### 优化 4:减少 `receiveMessages` 锁竞争 + +**文件**:`websocket_reporter.go` + +- 将连接状态检查移到循环外,只在出错/关闭时通过 channel 通知退出 +- 用 `context.WithCancel` 代替锁检查 `connected` flag 来控制生命周期 + +#### 优化 5:异步化状态变更命令处理 + +**文件**:`websocket_reporter.go` + +- 所有命令统一异步执行(通过 goroutine + response channel),避免阻塞 readLoop +- 当前只有 TcpPing/ServiceMonitorCheck/UpgradeAgent/RollbackAgent 是异步的 + +--- + +## 具体代码变更 + +### Agent 侧 (`go-gost/x/socket`) + +--- + +#### [MODIFY] [websocket_reporter.go](file:///Users/sagit/Documents/github/flvx/go-gost/x/socket/websocket_reporter.go) + +1. **指数退避重连**:将 `reconnectTime` 从固定 `5s` 改为动态退避字段,增加 `curBackoff/maxBackoff` 字段 +2. **独立 Ping 发送**:在 `handleConnection()` 中增加 WebSocket ping ticker(20s),独立于指标上报 +3. **减少锁竞争**:`receiveMessages` 中只在循环入口检查一次连接,此后靠 `ReadMessage` 的 error 退出 +4. **统一命令异步化**:所有 `routeCommand` 调用统一使用 goroutine + +--- + +### Panel 侧 (`go-backend/internal/ws`) + +--- + +#### [MODIFY] [server.go](file:///Users/sagit/Documents/github/flvx/go-backend/internal/ws/server.go) + +1. **缓存 AES 加密器**:在 `nodeSession` 中增加 `crypto *security.AESCrypto` 字段,节点连接时初始化 +2. **`decryptIfNeeded` 接收 crypto 参数**而非 secret 字符串 +3. **`SendCommand` 使用缓存 crypto** 实例 + +--- + +## Verification Plan + +### Automated Tests + +```bash +# 运行现有 agent 侧单元测试(验证不回归) +(cd go-gost/x && go test ./socket/... -v -count=1) + +# 运行现有流量上报测试 +(cd go-gost/x && go test ./service/... -v -count=1) + +# 运行 panel 侧全部测试 +(cd go-backend && go test ./... -count=1) +``` + +### Manual Verification + +> [!IMPORTANT] +> 本次改动涉及实时通信核心路径,建议在 staging 环境部署后观察至少 30 分钟: +> 1. 检查节点在面板中状态是否正常显示为在线 +> 2. 手动停止面板后观察 Agent 日志,确认重连间隔呈指数增长 +> 3. 恢复面板后确认 Agent 能自动恢复连接并恢复指标上报 +> 4. 通过面板下发命令(如添加/删除 service),确认命令执行成功 + +--- + +## 任务清单 + +- [x] 优化 1:Agent 增加独立 WebSocket ping 发送 +- [x] 优化 2:Agent 指数退避重连 +- [x] 优化 3:Panel 缓存 AES 加密器 +- [x] 优化 4:Agent 减少 receiveMessages 锁竞争 +- [x] 优化 5:Agent 命令处理统一异步化 +- [x] 运行现有测试验证不回归