refactor(internal/komari/ws): streamline status request logic and enhance read loop handling

refactor(internal/node/tracker): improve initial refresh gating and snapshot handling
This commit is contained in:
2026-08-07 18:33:34 +08:00
parent 58b600baf6
commit fcb0751267
2 changed files with 76 additions and 51 deletions
+32 -36
View File
@@ -101,12 +101,6 @@ func (w *WSClient) Connect() error {
w.reconnectCount.Store(0) w.reconnectCount.Store(0)
w.running.Store(true) w.running.Store(true)
// Request the initial node status snapshot.
if err := w.requestStatus(conn); err != nil {
postLog.Error("Failed to send initial node status request: " + err.Error())
return err
}
postLog.Info("Connected to Komari WebSocket") postLog.Info("Connected to Komari WebSocket")
go w.readLoop() go w.readLoop()
@@ -180,9 +174,15 @@ func (w *WSClient) requestStatus(conn *websocket.Conn) error {
return conn.WriteMessage(websocket.TextMessage, payload) return conn.WriteMessage(websocket.TextMessage, payload)
} }
// readLoop reads JSON-RPC responses and feeds the node tracker. // readLoop repeatedly pulls a fresh status snapshot and feeds the node tracker.
// Komari never pushes unsolicited messages on /api/rpc2, so the loop re-requests // Komari never pushes unsolicited messages on /api/rpc2, so the loop issues one
// a fresh snapshot whenever the connection stays quiet for a refresh interval. // request per refresh interval and reads exactly one response.
//
// It must never continue reading on the same connection after a read deadline
// timeout: gorilla/websocket marks a connection as corrupt once a read times
// out and returns the stored error on every subsequent read, which would make
// the loop spin until gorilla panics. A timed-out read therefore triggers a
// reconnect instead.
func (w *WSClient) readLoop() { func (w *WSClient) readLoop() {
defer func() { defer func() {
if r := recover(); r != nil { if r := recover(); r != nil {
@@ -193,7 +193,7 @@ func (w *WSClient) readLoop() {
const ( const (
refreshInterval = 5 * time.Second refreshInterval = 5 * time.Second
readTimeout = refreshInterval + 5*time.Second readTimeout = refreshInterval + 10*time.Second
) )
for { for {
@@ -211,6 +211,13 @@ func (w *WSClient) readLoop() {
return return
} }
// Pull a fresh snapshot, then wait for its single response.
if err := w.requestStatus(conn); err != nil {
postLog.Warning("Failed to request node status refresh: " + err.Error())
w.tryReconnect()
return
}
if err := conn.SetReadDeadline(time.Now().Add(readTimeout)); err != nil { if err := conn.SetReadDeadline(time.Now().Add(readTimeout)); err != nil {
postLog.Warning("Komari WebSocket set read deadline error: " + err.Error()) postLog.Warning("Komari WebSocket set read deadline error: " + err.Error())
w.tryReconnect() w.tryReconnect()
@@ -220,13 +227,10 @@ func (w *WSClient) readLoop() {
_, data, err := conn.ReadMessage() _, data, err := conn.ReadMessage()
if err != nil { if err != nil {
if netErr, ok := err.(net.Error); ok && netErr.Timeout() { if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
// No push within the window — pull a fresh snapshot. postLog.Warning("Komari WebSocket read timed out waiting for node status; reconnecting")
if rerr := w.requestStatus(conn); rerr != nil { } else {
postLog.Warning("Failed to request node status refresh: " + rerr.Error()) postLog.Warning("Komari WebSocket read error: " + err.Error())
}
continue
} }
postLog.Warning("Komari WebSocket read error: " + err.Error())
w.tryReconnect() w.tryReconnect()
return return
} }
@@ -234,19 +238,20 @@ func (w *WSClient) readLoop() {
var rpcResp jsonRpcResponse var rpcResp jsonRpcResponse
if err := json.Unmarshal(data, &rpcResp); err != nil { if err := json.Unmarshal(data, &rpcResp); err != nil {
postLog.Warning("Failed to parse Komari WebSocket message: " + err.Error()) postLog.Warning("Failed to parse Komari WebSocket message: " + err.Error())
continue } else if rpcResp.Error != nil {
}
if rpcResp.Error != nil {
postLog.Warning("Komari WebSocket RPC error: " + rpcResp.Error.Message) postLog.Warning("Komari WebSocket RPC error: " + rpcResp.Error.Message)
continue } else if len(rpcResp.Result) > 0 {
} online, reports := parseLatestStatus(rpcResp.Result)
if len(rpcResp.Result) == 0 { if tracker := node.GetTracker(); tracker != nil {
continue tracker.UpdateStatus(online, reports)
}
} }
online, reports := parseLatestStatus(rpcResp.Result) // Wait for the next refresh cycle before pulling again.
if tracker := node.GetTracker(); tracker != nil { select {
tracker.UpdateStatus(online, reports) case <-w.stopCh:
return
case <-time.After(refreshInterval):
} }
} }
} }
@@ -313,10 +318,7 @@ func (w *WSClient) tryReconnect() {
} }
// Exponential backoff. // Exponential backoff.
delay := time.Duration(1<<(count-1)) * time.Second delay := min(time.Duration(1<<(count-1))*time.Second, 30*time.Second)
if delay > 30*time.Second {
delay = 30 * time.Second
}
postLog.Info(fmt.Sprintf("Attempting Komari WebSocket reconnect %d/%d in %v...", postLog.Info(fmt.Sprintf("Attempting Komari WebSocket reconnect %d/%d in %v...",
count, w.maxRetries, delay)) count, w.maxRetries, delay))
@@ -339,12 +341,6 @@ func (w *WSClient) tryReconnect() {
w.reconnectCount.Store(0) w.reconnectCount.Store(0)
// Request a fresh status snapshot after reconnect.
if err := w.requestStatus(conn); err != nil {
postLog.Error("Failed to send status request after reconnect: " + err.Error())
return
}
// Re-fetch node list after reconnect. // Re-fetch node list after reconnect.
client := GetClient() client := GetClient()
if client != nil { if client != nil {
+44 -15
View File
@@ -78,11 +78,12 @@ type Tracker struct {
callbacks []StatusChangeCallback callbacks []StatusChangeCallback
// Initial-refresh gating: status-change notifications are suppressed until // Initial-refresh gating: status-change notifications are suppressed until
// every node from the Komari node list has delivered a status update, so the // the initial status snapshot has been received, so the startup state is
// startup state is treated as a baseline rather than a flood of changes. // treated as a baseline rather than a flood of changes.
knownUUIDs map[string]bool // uuids from the Komari node list knownUUIDs map[string]bool // uuids from the Komari node list
reportedUUIDs map[string]bool // node-list uuids that have reported status reportedUUIDs map[string]bool // node-list uuids that have reported status
bootstrapDone bool // true once the initial refresh has finished receivedSnapshot bool // true once at least one snapshot has arrived
bootstrapDone bool // true once the initial refresh has finished
} }
var globalTracker *Tracker var globalTracker *Tracker
@@ -228,9 +229,7 @@ func (t *Tracker) UpdateStatus(onlineUUIDs []string, reports map[string]Report)
t.onlineSet = newOnlineSet t.onlineSet = newOnlineSet
// Track which node-list nodes have delivered a status update. The initial // Track which node-list nodes have delivered a status update.
// refresh is considered complete once every known node has reported at
// least once (either online or offline).
for _, uuid := range onlineUUIDs { for _, uuid := range onlineUUIDs {
if t.knownUUIDs[uuid] { if t.knownUUIDs[uuid] {
t.reportedUUIDs[uuid] = true t.reportedUUIDs[uuid] = true
@@ -242,12 +241,19 @@ func (t *Tracker) UpdateStatus(onlineUUIDs []string, reports map[string]Report)
} }
} }
// Changes detected while the initial refresh is still in progress represent // The first received snapshot establishes the startup baseline, so its
// the nodes' startup state rather than real transitions, so suppress them. // changes are always suppressed. Afterward, changes are suppressed only
suppressChanges := !t.bootstrapDone // until the initial refresh is judged complete.
if !t.bootstrapDone && t.allKnownNodesReported() { firstSnapshot := !t.receivedSnapshot
if len(onlineUUIDs) > 0 || len(reports) > 0 {
t.receivedSnapshot = true
}
suppressChanges := !t.bootstrapDone || firstSnapshot
if !t.bootstrapDone && t.refreshComplete(reports) {
t.bootstrapDone = true t.bootstrapDone = true
postLog.Info(fmt.Sprintf("All %d nodes refreshed their status; status notifications enabled", len(t.knownUUIDs))) postLog.Info(fmt.Sprintf("Initial node status refresh complete (%d/%d nodes reported); status notifications enabled",
len(t.reportedUUIDs), len(t.knownUUIDs)))
} }
t.mu.Unlock() t.mu.Unlock()
@@ -274,9 +280,32 @@ func (t *Tracker) allKnownNodesReported() bool {
return true return true
} }
// refreshComplete reports whether the initial status refresh is finished.
// Komari returns the full latest-status batch in a single snapshot, so the
// refresh is complete once a snapshot covering every node we have ever seen
// has been processed. Node-list nodes that never report a status are treated
// as permanently offline rather than blocking notifications forever.
// Must be called with the lock held.
func (t *Tracker) refreshComplete(currentReports map[string]Report) bool {
if t.allKnownNodesReported() {
return true
}
if len(currentReports) == 0 {
return false
}
for uuid := range t.reportedUUIDs {
if _, ok := currentReports[uuid]; !ok {
return false
}
}
return true
}
// CompleteBootstrap force-completes the initial status refresh gate. It is a // CompleteBootstrap force-completes the initial status refresh gate. It is a
// safety net so that notifications cannot be blocked forever by a node that // safety net so that notifications cannot be blocked forever when no status
// never reports a status. Calling it again is a no-op. // snapshot is ever received. It does not mark a snapshot as received, so the
// first snapshot that does arrive is still treated as the baseline. Calling it
// again is a no-op.
func (t *Tracker) CompleteBootstrap() { func (t *Tracker) CompleteBootstrap() {
t.mu.Lock() t.mu.Lock()
defer t.mu.Unlock() defer t.mu.Unlock()