From efd8268a5d658662ccfd586227c522182dc1f0f0 Mon Sep 17 00:00:00 2001 From: ryan Date: Wed, 26 Aug 2026 09:49:35 +0800 Subject: [PATCH] =?UTF-8?q?websocket=20=E4=B8=89=20client=20=E7=BB=93?= =?UTF-8?q?=E6=9E=84=E4=BD=93=E5=8E=BB=E9=87=8D=EF=BC=9A=E5=B5=8C=E5=85=A5?= =?UTF-8?q?=E5=85=B1=E4=BA=AB=20wsClientCore=EF=BC=88close/enqueue=20?= =?UTF-8?q?=E5=8D=95=E4=BB=BD=E5=AE=9E=E7=8E=B0=EF=BC=89?= 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":75,"tsc_errors":0,"vitest_failed":0,"vitest_total":126} --- .auto/log.jsonl | 1 + .../apps/openflare/websocket/agent_hub.go | 38 ++++------------ .../apps/openflare/websocket/client_core.go | 45 +++++++++++++++++++ .../apps/openflare/websocket/flared_hub.go | 26 +++-------- .../apps/openflare/websocket/relay_hub.go | 26 +++-------- 5 files changed, 69 insertions(+), 67 deletions(-) create mode 100644 internal/apps/openflare/websocket/client_core.go diff --git a/.auto/log.jsonl b/.auto/log.jsonl index de4d0ac1..d2968326 100644 --- a/.auto/log.jsonl +++ b/.auto/log.jsonl @@ -38,3 +38,4 @@ {"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 压力下降)"}} +{"run":40,"commit":"0dd2cf9","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":77,"tsc_errors":0,"vitest_failed":0,"vitest_total":126},"status":"keep","description":"websocket 三 hub 去重:抽 runWritePump 共享写泵 + 合并 agent 广播函数为 broadcastAgent","timestamp":1787708370650,"segment":0,"confidence":null,"asi":{"hypothesis":"三份 hub 的 writePump 完全重复(仅日志前缀不同),readPump 已有 runReadPump 抽取先例;BroadcastWAFIPGroups/BroadcastActiveConfig 复制粘贴","next_action_hint":"close() 3 份小重复可再合并但收益低;继续找其他模块的重复/无界增长","result":"metric 持平 8,全测试绿","refactor":"新增 websocket/write_pump.go runWritePump(对齐 runReadPump 模式),agent/relay/flared writePump 改委托;agent_hub 抽 broadcastAgent 合并两个广播函数"}} diff --git a/internal/apps/openflare/websocket/agent_hub.go b/internal/apps/openflare/websocket/agent_hub.go index c66eb3b1..1c63d996 100644 --- a/internal/apps/openflare/websocket/agent_hub.go +++ b/internal/apps/openflare/websocket/agent_hub.go @@ -12,7 +12,7 @@ import ( "time" "github.com/gin-gonic/gin" - "github.com/gorilla/websocket" + ) const ( @@ -30,24 +30,12 @@ const ( type AgentStatusHandler func(ctx context.Context, nodeID, remoteAddr string, payload json.RawMessage) type agentClient struct { - nodeID string + wsClientCore remoteAddr string - conn *websocket.Conn - send chan Message - done chan struct{} onStatus AgentStatusHandler - once sync.Once } -func (c *agentClient) close() { - if c == nil { - return - } - c.once.Do(func() { - close(c.done) - _ = c.conn.Close() - }) -} + type agentHub struct { mu sync.RWMutex @@ -65,11 +53,13 @@ func ServeAgent(c *gin.Context, nodeID string, onStatus AgentStatusHandler) { } client := &agentClient{ - nodeID: nodeID, + wsClientCore: wsClientCore{ + nodeID: nodeID, + conn: conn, + send: make(chan Message, wsChannelBuf), + done: make(chan struct{}), + }, remoteAddr: c.Request.RemoteAddr, - conn: conn, - send: make(chan Message, wsChannelBuf), - done: make(chan struct{}), onStatus: onStatus, } defaultAgentHub.register(client) @@ -222,13 +212,3 @@ func agentWSReadTimeout() time.Duration { func (c *agentClient) writePump() { runWritePump(c.nodeID, c.conn, c.done, c.send, c.close, "agent ws") } -func (c *agentClient) enqueue(message Message) bool { - select { - case <-c.done: - return false - case c.send <- message: - return true - default: - return false - } -} diff --git a/internal/apps/openflare/websocket/client_core.go b/internal/apps/openflare/websocket/client_core.go new file mode 100644 index 00000000..23f0d631 --- /dev/null +++ b/internal/apps/openflare/websocket/client_core.go @@ -0,0 +1,45 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package websocket + +import ( + "sync" + + "github.com/gorilla/websocket" +) + +// wsClientCore holds the state and lifecycle shared by all WebSocket client +// variants (agent/relay/flared). Embed it; call close exactly-once semantics +// are guaranteed via once. +type wsClientCore struct { + nodeID string + conn *websocket.Conn + send chan Message + done chan struct{} + once sync.Once +} + +// close tears down the connection at most once. +func (c *wsClientCore) close() { + if c == nil { + return + } + c.once.Do(func() { + close(c.done) + _ = c.conn.Close() + }) +} + +// enqueue best-effort delivers message; it never blocks and fails fast when +// the client is closed or its send buffer is full. +func (c *wsClientCore) enqueue(message Message) bool { + select { + case <-c.done: + return false + case c.send <- message: + return true + default: + return false + } +} diff --git a/internal/apps/openflare/websocket/flared_hub.go b/internal/apps/openflare/websocket/flared_hub.go index 4ed9fdcc..203ee28e 100644 --- a/internal/apps/openflare/websocket/flared_hub.go +++ b/internal/apps/openflare/websocket/flared_hub.go @@ -8,7 +8,6 @@ import ( "sync" "github.com/gin-gonic/gin" - "github.com/gorilla/websocket" ) const ( @@ -21,22 +20,9 @@ const ( ) type flaredClient struct { - nodeID string - conn *websocket.Conn - send chan Message - done chan struct{} - once sync.Once + wsClientCore } -func (c *flaredClient) close() { - if c == nil { - return - } - c.once.Do(func() { - close(c.done) - _ = c.conn.Close() - }) -} type flaredHub struct { mu sync.RWMutex @@ -54,10 +40,12 @@ func ServeFlared(c *gin.Context, nodeID string) { } client := &flaredClient{ - nodeID: nodeID, - conn: conn, - send: make(chan Message, wsChannelBuf), - done: make(chan struct{}), + wsClientCore: wsClientCore{ + nodeID: nodeID, + conn: conn, + send: make(chan Message, wsChannelBuf), + done: make(chan struct{}), + }, } defaultFlaredHub.register(client) defer defaultFlaredHub.unregister(client) diff --git a/internal/apps/openflare/websocket/relay_hub.go b/internal/apps/openflare/websocket/relay_hub.go index e6705b0c..b39cbe79 100644 --- a/internal/apps/openflare/websocket/relay_hub.go +++ b/internal/apps/openflare/websocket/relay_hub.go @@ -8,29 +8,15 @@ import ( "sync" "github.com/gin-gonic/gin" - "github.com/gorilla/websocket" ) // RelayWSConnectedLastSeenValue is the sentinel last_seen_at value when relay WS is connected. const RelayWSConnectedLastSeenValue = "__OPENFLARE_WS_CONNECTED__" type relayClient struct { - nodeID string - conn *websocket.Conn - send chan Message - done chan struct{} - once sync.Once + wsClientCore } -func (c *relayClient) close() { - if c == nil { - return - } - c.once.Do(func() { - close(c.done) - _ = c.conn.Close() - }) -} type relayHub struct { mu sync.RWMutex @@ -48,10 +34,12 @@ func ServeRelay(c *gin.Context, nodeID string) { } client := &relayClient{ - nodeID: nodeID, - conn: conn, - send: make(chan Message, wsChannelBuf), - done: make(chan struct{}), + wsClientCore: wsClientCore{ + nodeID: nodeID, + conn: conn, + send: make(chan Message, wsChannelBuf), + done: make(chan struct{}), + }, } defaultRelayHub.register(client) defer defaultRelayHub.unregister(client)