diff --git a/internal/komari/ws.go b/internal/komari/ws.go index c830a69..1a16a6d 100644 --- a/internal/komari/ws.go +++ b/internal/komari/ws.go @@ -101,12 +101,6 @@ func (w *WSClient) Connect() error { w.reconnectCount.Store(0) 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") go w.readLoop() @@ -180,9 +174,15 @@ func (w *WSClient) requestStatus(conn *websocket.Conn) error { return conn.WriteMessage(websocket.TextMessage, payload) } -// readLoop reads JSON-RPC responses and feeds the node tracker. -// Komari never pushes unsolicited messages on /api/rpc2, so the loop re-requests -// a fresh snapshot whenever the connection stays quiet for a refresh interval. +// readLoop repeatedly pulls a fresh status snapshot and feeds the node tracker. +// Komari never pushes unsolicited messages on /api/rpc2, so the loop issues one +// 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() { defer func() { if r := recover(); r != nil { @@ -193,7 +193,7 @@ func (w *WSClient) readLoop() { const ( refreshInterval = 5 * time.Second - readTimeout = refreshInterval + 5*time.Second + readTimeout = refreshInterval + 10*time.Second ) for { @@ -211,6 +211,13 @@ func (w *WSClient) readLoop() { 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 { postLog.Warning("Komari WebSocket set read deadline error: " + err.Error()) w.tryReconnect() @@ -220,13 +227,10 @@ func (w *WSClient) readLoop() { _, data, err := conn.ReadMessage() if err != nil { if netErr, ok := err.(net.Error); ok && netErr.Timeout() { - // No push within the window — pull a fresh snapshot. - if rerr := w.requestStatus(conn); rerr != nil { - postLog.Warning("Failed to request node status refresh: " + rerr.Error()) - } - continue + postLog.Warning("Komari WebSocket read timed out waiting for node status; reconnecting") + } else { + postLog.Warning("Komari WebSocket read error: " + err.Error()) } - postLog.Warning("Komari WebSocket read error: " + err.Error()) w.tryReconnect() return } @@ -234,19 +238,20 @@ func (w *WSClient) readLoop() { var rpcResp jsonRpcResponse if err := json.Unmarshal(data, &rpcResp); err != nil { postLog.Warning("Failed to parse Komari WebSocket message: " + err.Error()) - continue - } - if rpcResp.Error != nil { + } else if rpcResp.Error != nil { postLog.Warning("Komari WebSocket RPC error: " + rpcResp.Error.Message) - continue - } - if len(rpcResp.Result) == 0 { - continue + } else if len(rpcResp.Result) > 0 { + online, reports := parseLatestStatus(rpcResp.Result) + if tracker := node.GetTracker(); tracker != nil { + tracker.UpdateStatus(online, reports) + } } - online, reports := parseLatestStatus(rpcResp.Result) - if tracker := node.GetTracker(); tracker != nil { - tracker.UpdateStatus(online, reports) + // Wait for the next refresh cycle before pulling again. + select { + case <-w.stopCh: + return + case <-time.After(refreshInterval): } } } @@ -313,10 +318,7 @@ func (w *WSClient) tryReconnect() { } // Exponential backoff. - delay := time.Duration(1<<(count-1)) * time.Second - if delay > 30*time.Second { - delay = 30 * time.Second - } + delay := min(time.Duration(1<<(count-1))*time.Second, 30*time.Second) postLog.Info(fmt.Sprintf("Attempting Komari WebSocket reconnect %d/%d in %v...", count, w.maxRetries, delay)) @@ -339,12 +341,6 @@ func (w *WSClient) tryReconnect() { 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. client := GetClient() if client != nil { diff --git a/internal/node/tracker.go b/internal/node/tracker.go index d6980f5..0a2f871 100644 --- a/internal/node/tracker.go +++ b/internal/node/tracker.go @@ -78,11 +78,12 @@ type Tracker struct { callbacks []StatusChangeCallback // Initial-refresh gating: status-change notifications are suppressed until - // every node from the Komari node list has delivered a status update, so the - // startup state is treated as a baseline rather than a flood of changes. - knownUUIDs map[string]bool // uuids from the Komari node list - reportedUUIDs map[string]bool // node-list uuids that have reported status - bootstrapDone bool // true once the initial refresh has finished + // the initial status snapshot has been received, so the startup state is + // treated as a baseline rather than a flood of changes. + knownUUIDs map[string]bool // uuids from the Komari node list + reportedUUIDs map[string]bool // node-list uuids that have reported status + receivedSnapshot bool // true once at least one snapshot has arrived + bootstrapDone bool // true once the initial refresh has finished } var globalTracker *Tracker @@ -228,9 +229,7 @@ func (t *Tracker) UpdateStatus(onlineUUIDs []string, reports map[string]Report) t.onlineSet = newOnlineSet - // Track which node-list nodes have delivered a status update. The initial - // refresh is considered complete once every known node has reported at - // least once (either online or offline). + // Track which node-list nodes have delivered a status update. for _, uuid := range onlineUUIDs { if t.knownUUIDs[uuid] { 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 nodes' startup state rather than real transitions, so suppress them. - suppressChanges := !t.bootstrapDone - if !t.bootstrapDone && t.allKnownNodesReported() { + // The first received snapshot establishes the startup baseline, so its + // changes are always suppressed. Afterward, changes are suppressed only + // until the initial refresh is judged complete. + 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 - 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() @@ -274,9 +280,32 @@ func (t *Tracker) allKnownNodesReported() bool { 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 -// safety net so that notifications cannot be blocked forever by a node that -// never reports a status. Calling it again is a no-op. +// safety net so that notifications cannot be blocked forever when no status +// 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() { t.mu.Lock() defer t.mu.Unlock()