From 056eb639a69d40e1237a222d8dcd57f810972786 Mon Sep 17 00:00:00 2001 From: Pulse Autonomous Maintainer Date: Fri, 21 Aug 2026 11:15:06 +0100 Subject: [PATCH] fix(alerts): prevent disabled PBS offline notifications --- internal/alerts/pbs.go | 14 +- .../alerts/pbs_offline_notification_test.go | 213 ++++++++++++++++++ 2 files changed, 225 insertions(+), 2 deletions(-) diff --git a/internal/alerts/pbs.go b/internal/alerts/pbs.go index 6e4017756..d29e6511b 100644 --- a/internal/alerts/pbs.go +++ b/internal/alerts/pbs.go @@ -115,7 +115,7 @@ func (m *Manager) CheckPBS(pbs models.PBSInstance) { } else { // Check if PBS is offline first (similar to nodes) if pbsOffline { - m.checkPBSOffline(pbs) + m.checkPBSOfflineWithThresholds(pbs, thresholds) } else { // Clear any existing offline alert if PBS is back online m.clearPBSOfflineAlert(pbs) @@ -141,11 +141,21 @@ func (m *Manager) CheckPBS(pbs models.PBSInstance) { // checkPBSOffline creates an alert for offline PBS instances func (m *Manager) checkPBSOffline(pbs models.PBSInstance) { + m.mu.RLock() + thresholds := m.resolveResourceThresholds("pbs", pbs.ID) + m.mu.RUnlock() + m.checkPBSOfflineWithThresholds(pbs, thresholds) +} + +// checkPBSOfflineWithThresholds evaluates one observation against the policy +// snapshot CheckPBS resolved when that observation began. A concurrent config +// update applies to the next observation instead of changing eligibility +// between this observation's pre-dispatch gate and lifecycle evaluation. +func (m *Manager) checkPBSOfflineWithThresholds(pbs models.PBSInstance, thresholds ThresholdConfig) { m.mu.Lock() delete(m.offlineRecoveryConfirmations, canonicalConnectivityStateID(pbs.ID)) m.mu.Unlock() - thresholds := m.resolveResourceThresholds("pbs", pbs.ID) spec, err := buildCanonicalConnectivitySpec(pbs.ID, pbs.Name, unifiedresources.ResourceTypePBS, AlertLevelCritical, 3, thresholds.Disabled || thresholds.DisableConnectivity) if err != nil { log.Warn(). diff --git a/internal/alerts/pbs_offline_notification_test.go b/internal/alerts/pbs_offline_notification_test.go index f133d845e..36632cda6 100644 --- a/internal/alerts/pbs_offline_notification_test.go +++ b/internal/alerts/pbs_offline_notification_test.go @@ -1,13 +1,226 @@ package alerts import ( + "sync" "testing" "time" alertspecs "github.com/rcourtman/pulse-go-rewrite/internal/alerts/specs" "github.com/rcourtman/pulse-go-rewrite/internal/models" + "github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources" ) +func newPBSOfflinePolicyFixture(t *testing.T) (models.PBSInstance, *unifiedresources.MonitorAdapter, string) { + t.Helper() + + pbs := models.PBSInstance{ + ID: "pbs-monitor-id", + Name: "pbs-main", + Host: "https://pbs.example.invalid:8007", + Status: "offline", + ConnectionHealth: "error", + LastSeen: time.Now(), + } + adapter := unifiedresources.NewMonitorAdapter(unifiedresources.NewRegistry(nil)) + adapter.PopulateFromSnapshot(models.StateSnapshot{PBSInstances: []models.PBSInstance{pbs}}) + canonicalID, ok := adapter.ResolveCanonicalResourceID(pbs.ID) + if !ok { + t.Fatalf("registry did not resolve PBS monitor ID %q", pbs.ID) + } + if canonicalID == "" || canonicalID == pbs.ID { + t.Fatalf("canonical PBS ID %q did not differ from monitor ID %q", canonicalID, pbs.ID) + } + registryResourceFound := false + for _, resource := range adapter.GetAll() { + if resource.Type == unifiedresources.ResourceTypePBS && resource.ID == canonicalID { + registryResourceFound = true + break + } + } + if !registryResourceFound { + t.Fatalf("canonical PBS resource %q was not present in the registry", canonicalID) + } + return pbs, adapter, canonicalID +} + +func configurePBSOfflinePolicyTestManager(t *testing.T, canonicalID string, disabled bool) (*Manager, chan *Alert) { + t.Helper() + + m := newTestManager(t) + cfg := m.GetConfig() + cfg.ActivationState = ActivationActive + cfg.TimeThresholds["pbs"] = 0 + cfg.Overrides = map[string]ThresholdConfig{} + if disabled { + cfg.Overrides[canonicalID] = ThresholdConfig{DisableConnectivity: true} + } + m.UpdateConfig(cfg) + + dispatched := make(chan *Alert, 4) + m.SetAlertCallback(func(alert *Alert) { + dispatched <- alert + }) + return m, dispatched +} + +func checkPBSRepeatedly(m *Manager, pbs models.PBSInstance, count int) { + for range count { + m.CheckPBS(pbs) + } +} + +func requirePBSDispatch(t *testing.T, dispatched <-chan *Alert) *Alert { + t.Helper() + select { + case alert := <-dispatched: + return alert + case <-time.After(2 * time.Second): + t.Fatal("timed out waiting for PBS alert callback") + return nil + } +} + +func requireNoPBSDispatch(t *testing.T, dispatched <-chan *Alert) { + t.Helper() + select { + case alert := <-dispatched: + t.Fatalf("unexpected alert callback for %s", alert.ID) + case <-time.After(100 * time.Millisecond): + } +} + +func requireSinglePBSAlert(t *testing.T, m *Manager, alertType string, level AlertLevel) Alert { + t.Helper() + active := m.GetActiveAlerts() + if len(active) != 1 { + t.Fatalf("active alerts = %d, want 1: %+v", len(active), active) + } + if active[0].Type != alertType || active[0].Level != level { + t.Fatalf("active alert = type %q level %q, want type %q level %q", active[0].Type, active[0].Level, alertType, level) + } + return active[0] +} + +func TestCheckPBSOfflineCanonicalOverrideBlocksAlertAndNotification(t *testing.T) { + pbs, adapter, canonicalID := newPBSOfflinePolicyFixture(t) + + t.Run("legacy ID lookup reproduces disabled-policy notification leak", func(t *testing.T) { + m, dispatched := configurePBSOfflinePolicyTestManager(t, canonicalID, true) + + // This control deliberately omits the registry resolver. It models the + // pre-fix lookup, which checked only pbs.ID even though the UI persisted + // the override under the canonical registry ID. + checkPBSRepeatedly(m, pbs, 3) + + alert := requireSinglePBSAlert(t, m, "offline", AlertLevelCritical) + if alert.ResourceID != pbs.ID { + t.Fatalf("offline alert resource ID = %q, want monitor ID %q", alert.ResourceID, pbs.ID) + } + dispatch := requirePBSDispatch(t, dispatched) + if dispatch.Type != "offline" || dispatch.Level != AlertLevelCritical { + t.Fatalf("dispatch = type %q level %q, want critical offline", dispatch.Type, dispatch.Level) + } + }) + + t.Run("canonical override blocks offline flaps but preserves PBS metrics", func(t *testing.T) { + m, dispatched := configurePBSOfflinePolicyTestManager(t, canonicalID, true) + m.SetResourceIntentIdentityResolver(adapter.ResolveCanonicalResourceID) + + checkPBSRepeatedly(m, pbs, 3) + if active := m.GetActiveAlerts(); len(active) != 0 { + t.Fatalf("disabled PBS offline policy created active alerts: %+v", active) + } + m.mu.RLock() + _, tracked := m.offlineConfirmations[pbs.ID] + m.mu.RUnlock() + if tracked { + t.Fatalf("disabled PBS offline policy retained confirmation tracking for %q", pbs.ID) + } + requireNoPBSDispatch(t, dispatched) + + // Exercise an offline -> healthy -> offline flap. The healthy sample + // also proves DisableConnectivity does not suppress unrelated metrics. + healthy := pbs + healthy.Status = "online" + healthy.ConnectionHealth = "healthy" + healthy.CPU = 99 + m.CheckPBS(healthy) + metricAlert := requireSinglePBSAlert(t, m, "cpu", AlertLevelCritical) + metricDispatch := requirePBSDispatch(t, dispatched) + if metricDispatch.ID != metricAlert.ID || metricDispatch.Type != "cpu" || metricDispatch.Level != AlertLevelCritical { + t.Fatalf("metric dispatch = %+v, want critical CPU alert %q", metricDispatch, metricAlert.ID) + } + + healthy.CPU = 0 + m.CheckPBS(healthy) + if active := m.GetActiveAlerts(); len(active) != 0 { + t.Fatalf("healthy CPU recovery left active alerts: %+v", active) + } + checkPBSRepeatedly(m, pbs, 3) + if active := m.GetActiveAlerts(); len(active) != 0 { + t.Fatalf("disabled PBS policy created an alert after a reconnect flap: %+v", active) + } + requireNoPBSDispatch(t, dispatched) + }) + + t.Run("enabled offline policy still alerts and recovers", func(t *testing.T) { + m, dispatched := configurePBSOfflinePolicyTestManager(t, canonicalID, false) + m.SetResourceIntentIdentityResolver(adapter.ResolveCanonicalResourceID) + + checkPBSRepeatedly(m, pbs, 3) + alert := requireSinglePBSAlert(t, m, "offline", AlertLevelCritical) + dispatch := requirePBSDispatch(t, dispatched) + if dispatch.ID != alert.ID || dispatch.Type != "offline" || dispatch.Level != AlertLevelCritical { + t.Fatalf("enabled dispatch = %+v, want critical offline alert %q", dispatch, alert.ID) + } + + healthy := pbs + healthy.Status = "online" + healthy.ConnectionHealth = "healthy" + checkPBSRepeatedly(m, healthy, offlineRecoveryConfirmationsDefault) + if active := m.GetActiveAlerts(); len(active) != 0 { + t.Fatalf("recovered PBS left active alerts: %+v", active) + } + }) + + t.Run("concurrent policy updates keep the pre-dispatch gate race-free", func(t *testing.T) { + m, _ := configurePBSOfflinePolicyTestManager(t, canonicalID, false) + m.SetResourceIntentIdentityResolver(adapter.ResolveCanonicalResourceID) + m.SetAlertCallback(nil) + + var wg sync.WaitGroup + wg.Add(2) + go func() { + defer wg.Done() + for i := 0; i < 200; i++ { + cfg := m.GetConfig() + cfg.Overrides = map[string]ThresholdConfig{} + if i%2 == 0 { + cfg.Overrides[canonicalID] = ThresholdConfig{DisableConnectivity: true} + } + m.UpdateConfig(cfg) + } + }() + go func() { + defer wg.Done() + for i := 0; i < 200; i++ { + m.CheckPBS(pbs) + } + }() + wg.Wait() + + cfg := m.GetConfig() + cfg.Overrides = map[string]ThresholdConfig{ + canonicalID: {DisableConnectivity: true}, + } + m.UpdateConfig(cfg) + checkPBSRepeatedly(m, pbs, 3) + if active := m.GetActiveAlerts(); len(active) != 0 { + t.Fatalf("final disabled policy left active alerts after concurrent updates: %+v", active) + } + }) +} + func TestCheckPBSOfflineDoesNotRenotifyExistingAlert(t *testing.T) { m := newTestManager(t) m.config.ActivationState = ActivationActive