fix(hub): synchronise system status access (#2452)

This commit is contained in:
user01010111
2026-09-29 11:58:46 +13:00
committed by GitHub
parent 083b28a1fc
commit 09a277f074
4 changed files with 26 additions and 10 deletions

View File

@@ -47,7 +47,7 @@ func bindNetworkMonitorsEvents(hub *Hub) {
// If connected, run the monitor immediately. Paused systems may be absent // If connected, run the monitor immediately. Paused systems may be absent
// from the manager; their monitors will sync when they reconnect. // from the manager; their monitors will sync when they reconnect.
system, err := hub.sm.GetSystem(e.Record.GetString("system")) 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) go hub.upsertNetworkMonitor(e.Record, true)
} }
return nil return nil

View File

@@ -41,7 +41,8 @@ type System struct {
Id string `db:"id"` Id string `db:"id"`
Host string `db:"host"` Host string `db:"host"`
Port string `db:"port"` 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 manager *SystemManager // Manager that this system belongs to
client atomic.Pointer[ssh.Client] // SSH client for fetching data client atomic.Pointer[ssh.Client] // SSH client for fetching data
sshTransport *transport.SSHTransport // SSH transport for requests sshTransport *transport.SSHTransport // SSH transport for requests
@@ -65,6 +66,22 @@ type System struct {
lastSavedMonitorProbe map[string]int64 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 { func (sm *SystemManager) NewSystem(systemId string) *System {
system := &System{ system := &System{
Id: systemId, Id: systemId,
@@ -101,7 +118,7 @@ func (sys *System) StartUpdater() {
// update immediately if system is not paused (only for ws connections) // update immediately if system is not paused (only for ws connections)
// we'll wait a minute before connecting via SSH to prioritize 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 { if err := sys.update(); err != nil {
_ = sys.setDown(err) _ = sys.setDown(err)
} }
@@ -134,7 +151,7 @@ func (sys *System) StartUpdater() {
// update updates the system data and records. // update updates the system data and records.
func (sys *System) update() error { func (sys *System) update() error {
if sys.Status == paused { if sys.GetStatus() == paused {
sys.handlePaused() sys.handlePaused()
return nil 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. // encountered during the process of updating the system status.
// It is a no-op if the system's context has been cancelled. // It is a no-op if the system's context has been cancelled.
func (sys *System) setDown(originalError error) error { func (sys *System) setDown(originalError error) error {
if sys.Status == down || sys.Status == paused { if status := sys.GetStatus(); status == down || status == paused {
return nil return nil
} }
// the updater can race shutdown, and the app may already be disposed by the // 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. // 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 { func (sys *System) runSSHOperation(timeout time.Duration, retries int, operation func(*ssh.Session) (bool, error)) error {
for attempt := 0; attempt <= retries; attempt++ { 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 { if err := sys.createSSHClient(); err != nil {
return err return err
} }

View File

@@ -209,8 +209,7 @@ func (sm *SystemManager) onRecordAfterUpdateSuccess(e *core.RecordEvent) error {
prevStatus := pending prevStatus := pending
system, ok := sm.systems.GetOk(e.Record.Id) system, ok := sm.systems.GetOk(e.Record.Id)
if ok { if ok {
prevStatus = system.Status prevStatus = system.swapStatus(newStatus)
system.Status = newStatus
} }
switch newStatus { switch newStatus {
@@ -335,7 +334,7 @@ func (sm *SystemManager) AddRecord(record *core.Record, system *System) (err err
} }
// Populate system from record // Populate system from record
system.Status = record.GetString("status") system.swapStatus(record.GetString("status"))
system.Host = record.GetString("host") system.Host = record.GetString("host")
system.Port = record.GetString("port") system.Port = record.GetString("port")

View File

@@ -38,7 +38,7 @@ func (sm *SystemManager) GetSystemStatusFromStore(systemID string) string {
if !ok { if !ok {
return "" return ""
} }
return sys.Status return sys.GetStatus()
} }
// TESTING ONLY: GetSystemContextFromStore returns the context and cancel function for a system // TESTING ONLY: GetSystemContextFromStore returns the context and cancel function for a system