mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
websocket 三 client 结构体去重:嵌入共享 wsClientCore(close/enqueue 单份实现)
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}
This commit is contained in:
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user