mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-06 18:06:36 +08:00
fix(ws): stabilize node connectivity with ping/pong keepalive
This commit is contained in:
@@ -54,6 +54,12 @@ type pendingRequest struct {
|
|||||||
ch chan CommandResult
|
ch chan CommandResult
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
wsPingPeriod = 15 * time.Second
|
||||||
|
wsPongWait = 45 * time.Second
|
||||||
|
wsWriteWait = 5 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
type CommandResult struct {
|
type CommandResult struct {
|
||||||
Type string `json:"type"`
|
Type string `json:"type"`
|
||||||
Success bool `json:"success"`
|
Success bool `json:"success"`
|
||||||
@@ -120,12 +126,19 @@ func (s *Server) handleAdmin(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
cw := &connWrap{conn: conn}
|
cw := &connWrap{conn: conn}
|
||||||
|
_ = conn.SetReadDeadline(time.Now().Add(wsPongWait))
|
||||||
|
conn.SetPongHandler(func(string) error {
|
||||||
|
return conn.SetReadDeadline(time.Now().Add(wsPongWait))
|
||||||
|
})
|
||||||
|
done := make(chan struct{})
|
||||||
|
go startKeepalive(cw, done)
|
||||||
|
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
s.admins[cw] = struct{}{}
|
s.admins[cw] = struct{}{}
|
||||||
s.mu.Unlock()
|
s.mu.Unlock()
|
||||||
|
|
||||||
defer func() {
|
defer func() {
|
||||||
|
close(done)
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
delete(s.admins, cw)
|
delete(s.admins, cw)
|
||||||
s.mu.Unlock()
|
s.mu.Unlock()
|
||||||
@@ -145,6 +158,12 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
cw := &connWrap{conn: conn}
|
cw := &connWrap{conn: conn}
|
||||||
|
_ = conn.SetReadDeadline(time.Now().Add(wsPongWait))
|
||||||
|
conn.SetPongHandler(func(string) error {
|
||||||
|
return conn.SetReadDeadline(time.Now().Add(wsPongWait))
|
||||||
|
})
|
||||||
|
done := make(chan struct{})
|
||||||
|
go startKeepalive(cw, done)
|
||||||
|
|
||||||
version := r.URL.Query().Get("version")
|
version := r.URL.Query().Get("version")
|
||||||
httpVal := parseIntDefault(r.URL.Query().Get("http"), 0)
|
httpVal := parseIntDefault(r.URL.Query().Get("http"), 0)
|
||||||
@@ -165,6 +184,7 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64
|
|||||||
s.broadcastStatus(nodeID, 1)
|
s.broadcastStatus(nodeID, 1)
|
||||||
|
|
||||||
defer func() {
|
defer func() {
|
||||||
|
close(done)
|
||||||
needOfflineBroadcast := false
|
needOfflineBroadcast := false
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
current, ok := s.nodes[nodeID]
|
current, ok := s.nodes[nodeID]
|
||||||
@@ -442,3 +462,27 @@ func parseIntDefault(v string, fallback int) int {
|
|||||||
}
|
}
|
||||||
return x
|
return x
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func startKeepalive(cw *connWrap, done <-chan struct{}) {
|
||||||
|
if cw == nil || cw.conn == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
ticker := time.NewTicker(wsPingPeriod)
|
||||||
|
defer ticker.Stop()
|
||||||
|
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-done:
|
||||||
|
return
|
||||||
|
case <-ticker.C:
|
||||||
|
cw.mu.Lock()
|
||||||
|
_ = cw.conn.SetWriteDeadline(time.Now().Add(wsWriteWait))
|
||||||
|
err := cw.conn.WriteMessage(websocket.PingMessage, nil)
|
||||||
|
cw.mu.Unlock()
|
||||||
|
if err != nil {
|
||||||
|
_ = cw.conn.Close()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -91,6 +91,11 @@ type TcpPingResponse struct {
|
|||||||
RequestId string `json:"requestId,omitempty"`
|
RequestId string `json:"requestId,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
reporterReadWait = 60 * time.Second
|
||||||
|
reporterWriteWait = 5 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
type WebSocketReporter struct {
|
type WebSocketReporter struct {
|
||||||
url string
|
url string
|
||||||
addr string // 保存服务器地址
|
addr string // 保存服务器地址
|
||||||
@@ -243,6 +248,14 @@ func (w *WebSocketReporter) connect() error {
|
|||||||
|
|
||||||
w.conn = conn
|
w.conn = conn
|
||||||
w.connected = true
|
w.connected = true
|
||||||
|
_ = conn.SetReadDeadline(time.Now().Add(reporterReadWait))
|
||||||
|
conn.SetPingHandler(func(appData string) error {
|
||||||
|
_ = conn.SetReadDeadline(time.Now().Add(reporterReadWait))
|
||||||
|
return conn.WriteControl(websocket.PongMessage, []byte(appData), time.Now().Add(reporterWriteWait))
|
||||||
|
})
|
||||||
|
conn.SetPongHandler(func(string) error {
|
||||||
|
return conn.SetReadDeadline(time.Now().Add(reporterReadWait))
|
||||||
|
})
|
||||||
|
|
||||||
// 设置关闭处理器来检测连接状态
|
// 设置关闭处理器来检测连接状态
|
||||||
w.conn.SetCloseHandler(func(code int, text string) error {
|
w.conn.SetCloseHandler(func(code int, text string) error {
|
||||||
@@ -383,7 +396,7 @@ func (w *WebSocketReporter) receiveMessages() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// 设置读取超时
|
// 设置读取超时
|
||||||
conn.SetReadDeadline(time.Now().Add(30 * time.Second))
|
conn.SetReadDeadline(time.Now().Add(reporterReadWait))
|
||||||
|
|
||||||
messageType, message, err := conn.ReadMessage()
|
messageType, message, err := conn.ReadMessage()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
Reference in New Issue
Block a user