From 09a277f07487c1e78a82eb331b0e8dcc4906e753 Mon Sep 17 00:00:00 2001 From: user01010111 <12504630+user01010111@users.noreply.github.com> Date: Tue, 29 Sep 2026 11:58:46 +1300 Subject: [PATCH] fix(hub): synchronise system status access (#2452) --- internal/hub/network_monitors.go | 2 +- internal/hub/systems/system.go | 27 ++++++++++++++++---- internal/hub/systems/system_manager.go | 5 ++-- internal/hub/systems/systems_test_helpers.go | 2 +- 4 files changed, 26 insertions(+), 10 deletions(-) diff --git a/internal/hub/network_monitors.go b/internal/hub/network_monitors.go index 98e9e0e88..9fcc86407 100644 --- a/internal/hub/network_monitors.go +++ b/internal/hub/network_monitors.go @@ -47,7 +47,7 @@ func bindNetworkMonitorsEvents(hub *Hub) { // If connected, run the monitor immediately. Paused systems may be absent // from the manager; their monitors will sync when they reconnect. system, err := hub.sm.GetSystem(e.Record.GetString("system")) - if err == nil && system.Status == "up" { + if err == nil && system.GetStatus() == "up" { go hub.upsertNetworkMonitor(e.Record, true) } return nil diff --git a/internal/hub/systems/system.go b/internal/hub/systems/system.go index 349304e21..49d1f5d63 100644 --- a/internal/hub/systems/system.go +++ b/internal/hub/systems/system.go @@ -41,7 +41,8 @@ type System struct { Id string `db:"id"` Host string `db:"host"` Port string `db:"port"` - Status string `db:"status"` + Status string `db:"status"` // Use GetStatus/swapStatus after publishing the system. + statusMu sync.RWMutex // Protects Status and exchanges used by alert transitions. manager *SystemManager // Manager that this system belongs to client atomic.Pointer[ssh.Client] // SSH client for fetching data sshTransport *transport.SSHTransport // SSH transport for requests @@ -65,6 +66,22 @@ type System struct { lastSavedMonitorProbe map[string]int64 } +// GetStatus returns the current monitoring status. +func (sys *System) GetStatus() string { + sys.statusMu.RLock() + defer sys.statusMu.RUnlock() + return sys.Status +} + +// swapStatus updates the status and returns the previous value as one operation. +func (sys *System) swapStatus(status string) string { + sys.statusMu.Lock() + defer sys.statusMu.Unlock() + previous := sys.Status + sys.Status = status + return previous +} + func (sm *SystemManager) NewSystem(systemId string) *System { system := &System{ Id: systemId, @@ -101,7 +118,7 @@ func (sys *System) StartUpdater() { // update immediately if system is not paused (only for ws connections) // we'll wait a minute before connecting via SSH to prioritize ws connections - if sys.Status != paused && sys.ctx.Err() == nil { + if sys.GetStatus() != paused && sys.ctx.Err() == nil { if err := sys.update(); err != nil { _ = sys.setDown(err) } @@ -134,7 +151,7 @@ func (sys *System) StartUpdater() { // update updates the system data and records. func (sys *System) update() error { - if sys.Status == paused { + if sys.GetStatus() == paused { sys.handlePaused() return nil } @@ -603,7 +620,7 @@ func (sys *System) HasUser(app core.App, user *core.Record) bool { // encountered during the process of updating the system status. // It is a no-op if the system's context has been cancelled. func (sys *System) setDown(originalError error) error { - if sys.Status == down || sys.Status == paused { + if status := sys.GetStatus(); status == down || status == paused { return nil } // the updater can race shutdown, and the app may already be disposed by the @@ -885,7 +902,7 @@ func (sys *System) fetchDataViaSSH(options common.DataRequestOptions) (*system.C // The operation can request a retry by returning true as the first return value. func (sys *System) runSSHOperation(timeout time.Duration, retries int, operation func(*ssh.Session) (bool, error)) error { for attempt := 0; attempt <= retries; attempt++ { - if sys.client.Load() == nil || sys.Status == down { + if sys.client.Load() == nil || sys.GetStatus() == down { if err := sys.createSSHClient(); err != nil { return err } diff --git a/internal/hub/systems/system_manager.go b/internal/hub/systems/system_manager.go index 4e65ed55b..6c14b8d9a 100644 --- a/internal/hub/systems/system_manager.go +++ b/internal/hub/systems/system_manager.go @@ -209,8 +209,7 @@ func (sm *SystemManager) onRecordAfterUpdateSuccess(e *core.RecordEvent) error { prevStatus := pending system, ok := sm.systems.GetOk(e.Record.Id) if ok { - prevStatus = system.Status - system.Status = newStatus + prevStatus = system.swapStatus(newStatus) } switch newStatus { @@ -335,7 +334,7 @@ func (sm *SystemManager) AddRecord(record *core.Record, system *System) (err err } // Populate system from record - system.Status = record.GetString("status") + system.swapStatus(record.GetString("status")) system.Host = record.GetString("host") system.Port = record.GetString("port") diff --git a/internal/hub/systems/systems_test_helpers.go b/internal/hub/systems/systems_test_helpers.go index ca61895f2..9e6c48601 100644 --- a/internal/hub/systems/systems_test_helpers.go +++ b/internal/hub/systems/systems_test_helpers.go @@ -38,7 +38,7 @@ func (sm *SystemManager) GetSystemStatusFromStore(systemID string) string { if !ok { return "" } - return sys.Status + return sys.GetStatus() } // TESTING ONLY: GetSystemContextFromStore returns the context and cancel function for a system