From 0dd2cf9e8086a72945b7a0469969d6ce2fcf6880 Mon Sep 17 00:00:00 2001 From: ryan Date: Wed, 26 Aug 2026 09:39:30 +0800 Subject: [PATCH] =?UTF-8?q?websocket=20=E4=B8=89=20hub=20=E5=8E=BB?= =?UTF-8?q?=E9=87=8D=EF=BC=9A=E6=8A=BD=20runWritePump=20=E5=85=B1=E4=BA=AB?= =?UTF-8?q?=E5=86=99=E6=B3=B5=20+=20=E5=90=88=E5=B9=B6=20agent=20=E5=B9=BF?= =?UTF-8?q?=E6=92=AD=E5=87=BD=E6=95=B0=E4=B8=BA=20broadcastAgent?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Result: {"status":"keep","total_issues":8,"eslint_errors":0,"eslint_problems":0,"eslint_warnings":0,"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_exhaustive":0,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_total":0,"golint_test_usetesting":0,"golint_total":8,"golint_usestdlibvars":0,"golint_vetx_total":0,"golint_wastedassign":0,"measure_s":77,"tsc_errors":0,"vitest_failed":0,"vitest_total":126} --- .auto/log.jsonl | 1 + .../apps/openflare/websocket/agent_hub.go | 50 +++---------------- .../apps/openflare/websocket/flared_hub.go | 25 +--------- .../apps/openflare/websocket/relay_hub.go | 25 +--------- .../apps/openflare/websocket/write_pump.go | 47 +++++++++++++++++ 5 files changed, 57 insertions(+), 91 deletions(-) create mode 100644 internal/apps/openflare/websocket/write_pump.go diff --git a/.auto/log.jsonl b/.auto/log.jsonl index 434aef1d..de4d0ac1 100644 --- a/.auto/log.jsonl +++ b/.auto/log.jsonl @@ -37,3 +37,4 @@ {"run":36,"commit":"7fa9e46","metric":8,"metrics":{"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_usestdlibvars":0,"golint_wastedassign":0,"golint_total":8,"eslint_problems":0,"eslint_errors":0,"eslint_warnings":0,"tsc_errors":0,"measure_s":86,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_usetesting":0,"golint_test_total":0,"golint_exhaustive":0,"golint_vetx_total":0,"vitest_failed":0,"vitest_total":126},"status":"keep","description":"注册开关读取失败时改为关闭,堵住配置缺失时未授权开注册;OAuth 自动注册同样 fail-closed。metric 持平 8。","timestamp":1787669960693,"segment":0,"confidence":null,"asi":{"hypothesis":"registration_enabled/password_register_enabled 读取失败默认 true,和种子 false 相反,配置缺失时未授权开注册","finding":"密码注册与 OAuth 自动注册均 fail-closed;测试改为显式开启注册并正确失效缓存。","next_action_hint":"下一轮可查 OIDC 开关 fail-open(种子默认 true,风险较低)或公开 OAuth state 洪水"}} {"run":37,"commit":"0290c93","metric":8,"metrics":{"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_usestdlibvars":0,"golint_wastedassign":0,"golint_total":8,"eslint_problems":0,"eslint_errors":0,"eslint_warnings":0,"tsc_errors":0,"measure_s":95,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_usetesting":0,"golint_test_total":0,"golint_exhaustive":0,"golint_vetx_total":0,"vitest_failed":0,"vitest_total":126},"status":"keep","description":"公开 OAuth 登录/授权入口按会话限制 10 分钟内最多 20 个 state,堵住未授权 Redis 洪水。metric 持平 8。","timestamp":1787670327304,"segment":0,"confidence":null,"asi":{"hypothesis":"公开 /oauth/login 与 /oauth/{source}/authorize 每次请求都往 Redis 写 10 分钟 state,无上限","finding":"按 sessionHash 计数,10 分钟内最多 20 个;超出返回业务错误。mock Redis 补 Incr/Expire。","next_action_hint":"下一轮可查边缘节点 access_token 明文比较,或公开 CAP challenge 洪水"}} {"run":38,"commit":"c0a82f8","metric":8,"metrics":{"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_usestdlibvars":0,"golint_wastedassign":0,"golint_total":8,"eslint_problems":0,"eslint_errors":0,"eslint_warnings":0,"tsc_errors":0,"measure_s":85,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_usetesting":0,"golint_test_total":0,"golint_exhaustive":0,"golint_vetx_total":0,"vitest_failed":0,"vitest_total":126},"status":"keep","description":"公开密码登录口按 IP 限制 10 分钟内最多 20 次失败,堵住未授权爆破。metric 持平 8。","timestamp":1787670665553,"segment":0,"confidence":null,"asi":{"hypothesis":"公开 /user/login 失败无 IP 限流,未授权方可无限爆破","finding":"按 ClientIP 计数,10 分钟 20 次失败后拒绝;成功清零。管理员已登录不在范围内。","next_action_hint":"下一轮可查公开 CAP challenge 洪水或边缘节点 access_token 明文比较"}} +{"run":39,"commit":"be5d067","metric":8,"metrics":{"eslint_errors":0,"eslint_problems":0,"eslint_warnings":0,"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_exhaustive":0,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_total":0,"golint_test_usetesting":0,"golint_total":8,"golint_usestdlibvars":0,"golint_vetx_total":0,"golint_wastedassign":0,"measure_s":76,"tsc_errors":0,"vitest_failed":0,"vitest_total":126},"status":"keep","description":"auth_cache negative 缓存加上限防 DoS + relay/flared 删除重复 authenticateAccessToken 改用 agent 共享缓存版","timestamp":1787708052241,"segment":0,"confidence":null,"asi":{"hypothesis":"negative cache 无上限可被伪造 token 撑爆内存;relay/flared 与 agent 三份重复的 authenticateAccessToken","next_action_hint":"继续扫其他无界缓存/限流缺口","result":"metric 持平 8(8 个均为 deliberate keeper),安全修复不计入 metric","security":"negative cache 加 10k 上限+过期清理;relay/flared 复用 agent.AuthenticateAccessToken(共享 2min 正/10min 负缓存,DB 压力下降)"}} diff --git a/internal/apps/openflare/websocket/agent_hub.go b/internal/apps/openflare/websocket/agent_hub.go index 1dfe27a1..c66eb3b1 100644 --- a/internal/apps/openflare/websocket/agent_hub.go +++ b/internal/apps/openflare/websocket/agent_hub.go @@ -132,32 +132,19 @@ func SendAgentWAFIPGroups(nodeID string, payload any) bool { // BroadcastWAFIPGroups pushes changed WAF IP groups to all connected agents. func BroadcastWAFIPGroups(payload any) int { - if payload == nil { - return 0 - } - message := Message{Type: agentMessageTypeWAFIPGroups, Payload: payload} - defaultAgentHub.mu.RLock() - clients := make([]*agentClient, 0, len(defaultAgentHub.clients)) - for _, client := range defaultAgentHub.clients { - clients = append(clients, client) - } - defaultAgentHub.mu.RUnlock() - - success := 0 - for _, client := range clients { - if client.enqueue(message) { - success++ - } - } - return success + return broadcastAgent(agentMessageTypeWAFIPGroups, payload) } // BroadcastActiveConfig pushes active config metadata to all connected agents. func BroadcastActiveConfig(payload any) int { + return broadcastAgent(agentMessageTypeActiveConfig, payload) +} + +func broadcastAgent(messageType string, payload any) int { if payload == nil { return 0 } - message := Message{Type: agentMessageTypeActiveConfig, Payload: payload} + message := Message{Type: messageType, Payload: payload} defaultAgentHub.mu.RLock() clients := make([]*agentClient, 0, len(defaultAgentHub.clients)) for _, client := range defaultAgentHub.clients { @@ -233,31 +220,8 @@ func agentWSReadTimeout() time.Duration { } func (c *agentClient) writePump() { - ticker := time.NewTicker(wsPingInterval) - defer ticker.Stop() - - for { - select { - case <-c.done: - return - case message := <-c.send: - _ = c.conn.SetWriteDeadline(time.Now().Add(wsWriteDeadline)) - if err := c.conn.WriteJSON(message); err != nil { - slog.Debug("agent ws write failed", "node_id", c.nodeID, "error", err) - c.close() - return - } - case <-ticker.C: - select { - case <-c.done: - return - case c.send <- Message{Type: messageTypePing}: - default: - } - } - } + runWritePump(c.nodeID, c.conn, c.done, c.send, c.close, "agent ws") } - func (c *agentClient) enqueue(message Message) bool { select { case <-c.done: diff --git a/internal/apps/openflare/websocket/flared_hub.go b/internal/apps/openflare/websocket/flared_hub.go index 57913e95..4ed9fdcc 100644 --- a/internal/apps/openflare/websocket/flared_hub.go +++ b/internal/apps/openflare/websocket/flared_hub.go @@ -6,7 +6,6 @@ package websocket import ( "log/slog" "sync" - "time" "github.com/gin-gonic/gin" "github.com/gorilla/websocket" @@ -139,27 +138,5 @@ func (c *flaredClient) readPump() { } func (c *flaredClient) writePump() { - ticker := time.NewTicker(wsPingInterval) - defer ticker.Stop() - - for { - select { - case <-c.done: - return - case message := <-c.send: - _ = c.conn.SetWriteDeadline(time.Now().Add(wsWriteDeadline)) - if err := c.conn.WriteJSON(message); err != nil { - slog.Debug("flared ws write failed", "node_id", c.nodeID, "error", err) - c.close() - return - } - case <-ticker.C: - select { - case <-c.done: - return - case c.send <- Message{Type: messageTypePing}: - default: - } - } - } + runWritePump(c.nodeID, c.conn, c.done, c.send, c.close, "flared ws") } diff --git a/internal/apps/openflare/websocket/relay_hub.go b/internal/apps/openflare/websocket/relay_hub.go index c9795183..e6705b0c 100644 --- a/internal/apps/openflare/websocket/relay_hub.go +++ b/internal/apps/openflare/websocket/relay_hub.go @@ -6,7 +6,6 @@ package websocket import ( "log/slog" "sync" - "time" "github.com/gin-gonic/gin" "github.com/gorilla/websocket" @@ -120,27 +119,5 @@ func (c *relayClient) readPump() { } func (c *relayClient) writePump() { - ticker := time.NewTicker(wsPingInterval) - defer ticker.Stop() - - for { - select { - case <-c.done: - return - case message := <-c.send: - _ = c.conn.SetWriteDeadline(time.Now().Add(wsWriteDeadline)) - if err := c.conn.WriteJSON(message); err != nil { - slog.Debug("relay ws write failed", "node_id", c.nodeID, "error", err) - c.close() - return - } - case <-ticker.C: - select { - case <-c.done: - return - case c.send <- Message{Type: messageTypePing}: - default: - } - } - } + runWritePump(c.nodeID, c.conn, c.done, c.send, c.close, "relay ws") } diff --git a/internal/apps/openflare/websocket/write_pump.go b/internal/apps/openflare/websocket/write_pump.go new file mode 100644 index 00000000..5c054d0f --- /dev/null +++ b/internal/apps/openflare/websocket/write_pump.go @@ -0,0 +1,47 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package websocket + +import ( + "log/slog" + "time" + + "github.com/gorilla/websocket" +) + +// runWritePump drains send onto conn until done is closed, emitting +// JSON pings at wsPingInterval. Shared by agent/relay/flared clients; +// closeFn must be idempotent. +func runWritePump( + nodeID string, + conn *websocket.Conn, + done <-chan struct{}, + send chan Message, + closeFn func(), + logLabel string, +) { + ticker := time.NewTicker(wsPingInterval) + defer ticker.Stop() + + for { + select { + case <-done: + return + case message := <-send: + _ = conn.SetWriteDeadline(time.Now().Add(wsWriteDeadline)) + if err := conn.WriteJSON(message); err != nil { + slog.Debug(logLabel+" write failed", "node_id", nodeID, "error", err) + closeFn() + return + } + case <-ticker.C: + select { + case <-done: + return + case send <- Message{Type: messageTypePing}: + default: + } + } + } +}