fix: agent ws connect panic
This commit is contained in:
@@ -77,6 +77,7 @@ type Agent struct {
|
|||||||
wsConn *websocket.Conn
|
wsConn *websocket.Conn
|
||||||
wsMu sync.Mutex
|
wsMu sync.Mutex
|
||||||
stopCh chan struct{}
|
stopCh chan struct{}
|
||||||
|
wsStopCh chan struct{} // 用于停止当前 WebSocket 相关的 goroutine
|
||||||
}
|
}
|
||||||
|
|
||||||
// generateMachineID 生成机器识别码
|
// generateMachineID 生成机器识别码
|
||||||
@@ -190,6 +191,7 @@ func (a *Agent) connectWS() error {
|
|||||||
|
|
||||||
a.wsMu.Lock()
|
a.wsMu.Lock()
|
||||||
a.wsConn = conn
|
a.wsConn = conn
|
||||||
|
a.wsStopCh = make(chan struct{})
|
||||||
a.wsMu.Unlock()
|
a.wsMu.Unlock()
|
||||||
|
|
||||||
log.Info("WebSocket 已连接")
|
log.Info("WebSocket 已连接")
|
||||||
@@ -202,6 +204,10 @@ func (a *Agent) connectWS() error {
|
|||||||
func (a *Agent) closeWS() {
|
func (a *Agent) closeWS() {
|
||||||
a.wsMu.Lock()
|
a.wsMu.Lock()
|
||||||
defer a.wsMu.Unlock()
|
defer a.wsMu.Unlock()
|
||||||
|
if a.wsStopCh != nil {
|
||||||
|
close(a.wsStopCh)
|
||||||
|
a.wsStopCh = nil
|
||||||
|
}
|
||||||
if a.wsConn != nil {
|
if a.wsConn != nil {
|
||||||
a.wsConn.Close()
|
a.wsConn.Close()
|
||||||
a.wsConn = nil
|
a.wsConn = nil
|
||||||
@@ -209,6 +215,8 @@ func (a *Agent) closeWS() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (a *Agent) readWS() {
|
func (a *Agent) readWS() {
|
||||||
|
defer a.closeWS() // 读取结束时关闭连接,停止 heartbeatLoop
|
||||||
|
|
||||||
for {
|
for {
|
||||||
a.wsMu.Lock()
|
a.wsMu.Lock()
|
||||||
conn := a.wsConn
|
conn := a.wsConn
|
||||||
@@ -327,10 +335,20 @@ func (a *Agent) heartbeatLoop() {
|
|||||||
ticker := time.NewTicker(time.Duration(a.config.Interval) * time.Second)
|
ticker := time.NewTicker(time.Duration(a.config.Interval) * time.Second)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
|
|
||||||
|
a.wsMu.Lock()
|
||||||
|
wsStopCh := a.wsStopCh
|
||||||
|
a.wsMu.Unlock()
|
||||||
|
|
||||||
|
if wsStopCh == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-a.stopCh:
|
case <-a.stopCh:
|
||||||
return
|
return
|
||||||
|
case <-wsStopCh:
|
||||||
|
return
|
||||||
case <-ticker.C:
|
case <-ticker.C:
|
||||||
a.wsMu.Lock()
|
a.wsMu.Lock()
|
||||||
conn := a.wsConn
|
conn := a.wsConn
|
||||||
|
|||||||
@@ -413,9 +413,8 @@ func (c *AgentController) wsReadPump(ac *services.AgentConnection, agent *models
|
|||||||
c.wsManager.Unregister(agent.ID)
|
c.wsManager.Unregister(agent.ID)
|
||||||
}()
|
}()
|
||||||
|
|
||||||
// 检查连接是否有效
|
// 检查连接是否有效(可能是旧连接被新连接替换)
|
||||||
if ac == nil || ac.IsClosed() {
|
if ac == nil || ac.IsClosed() {
|
||||||
logger.Warnf("[AgentWS] Agent #%d 连接无效", agent.ID)
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user