From d813c936e464eb09bb9ae7e6c7ac33716acde378 Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Sat, 3 Oct 2026 00:33:08 +0100 Subject: [PATCH] Retire alert flapping episodes when cleanup expires their hold Both cleanup owners removed only the deadline, leaving a latched episode that could suppress notifications without a new cooldown. Share dispatch expiry with cleanup, preserve unexpired policies and active occurrence identity, and verify renewed delivery, bounded rearming, diagnosis and one-shot callbacks without bypassing ordinary repeat-notification cooldown. Change-source: pulse-maintainer --- .../v6/internal/subsystems/alerts.md | 12 ++ internal/alerts/active_cleanup.go | 6 +- internal/alerts/flapping_threshold_test.go | 161 ++++++++++++++++++ internal/alerts/notification_policy.go | 22 ++- internal/alerts/tracking_cleanup.go | 5 +- 5 files changed, 193 insertions(+), 13 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/alerts.md b/docs/release-control/v6/internal/subsystems/alerts.md index 3cc8f2321..4b002b9e4 100644 --- a/docs/release-control/v6/internal/subsystems/alerts.md +++ b/docs/release-control/v6/internal/subsystems/alerts.md @@ -1500,6 +1500,18 @@ actually enforced at dispatch. `internal/alerts/flapping_threshold_test.go` pins these boundaries through normal configuration updates and notification callbacks, alongside the existing flapping cooldown and one-shot callback tests. +Ordinary retention cleanup and hourly tracking cleanup retire a served +suppression deadline, its flapping latch and its observations atomically under +the manager lock, using the same expiry rule as dispatch. Neither sweep may +leave an active occurrence in unbounded flapping suppression after deleting its +deadline. The next burst starts a fresh observation window and can arm a new +bounded cooldown and one-shot callback. Unexpired cooldowns, other keys' policy +and pending bursts without a deadline remain unchanged; expiry does not clear, +acknowledge or change the identity of an active alert. The cleanup regression +controls in `internal/alerts/flapping_threshold_test.go` exercise both sweeps, +drained and retained windows, dispatch callbacks and delivery diagnosis. These +modeled-time controls are not installed notification-destination acceptance. + ### Monitor-only delivery is terminal Monitor-only alerts remain visible, but neither a firing nor a recovery diff --git a/internal/alerts/active_cleanup.go b/internal/alerts/active_cleanup.go index ccf3040e4..9d0a19af6 100644 --- a/internal/alerts/active_cleanup.go +++ b/internal/alerts/active_cleanup.go @@ -130,10 +130,8 @@ func (m *Manager) Cleanup(maxAge time.Duration) { } } - for id, suppressUntil := range m.suppressedUntil { - if now.After(suppressUntil) { - delete(m.suppressedUntil, id) - } + for id := range m.suppressedUntil { + m.expireSuppressionNoLock(id, now) } cutoff := now.Add(-1 * time.Hour) diff --git a/internal/alerts/flapping_threshold_test.go b/internal/alerts/flapping_threshold_test.go index 25872dca0..bd75698a5 100644 --- a/internal/alerts/flapping_threshold_test.go +++ b/internal/alerts/flapping_threshold_test.go @@ -2,6 +2,8 @@ package alerts import ( "fmt" + "reflect" + "sync/atomic" "testing" "time" ) @@ -152,3 +154,162 @@ func TestFlappingLoweredThresholdBoundsRetainedHistory(t *testing.T) { t.Fatalf("lowered threshold not enforced with bounded history: suppressed=%v transitioned=%v history=%d", suppressed, transitioned, historySize) } } + +// Model elapsed cooldown time, but exercise the real cleanup, dispatch, +// one-shot callback and diagnosis paths. A sweep must not turn a timed hold +// into an indefinitely latched episode, whether its observation window drained +// or still contains the previous burst. +func TestFlappingCleanupReleasesAndRearmsDelivery(t *testing.T) { + for _, sweep := range []struct { + name string + run func(*Manager) + }{ + {"ordinary", func(m *Manager) { m.Cleanup(time.Hour) }}, + {"hourly", (*Manager).cleanupStaleMaps}, + } { + for _, window := range []int{60, 3600} { + t.Run(fmt.Sprintf("%s/window_%d", sweep.name, window), func(t *testing.T) { + m := NewManagerWithDataDir(t.TempDir()) + t.Cleanup(m.Stop) + cfg := m.GetConfig() + cfg.Enabled = true + cfg.ActivationState = ActivationActive + cfg.AutoAcknowledgeAfterHours = 0 + cfg.Schedule.Cooldown = 0 + cfg.Schedule.QuietHours.Enabled = false + cfg.FlappingEnabled = true + cfg.FlappingThreshold = 3 + cfg.FlappingWindowSeconds = window + cfg.FlappingCooldownMinutes = 15 + m.UpdateConfig(cfg) + + _, alert := testNewCanonicalAlert("vm-cleanup", "vm-cleanup-cpu", "vm", "cpu") + alert.Level = AlertLevelWarning + alert.StartTime = time.Now() + alert.LastSeen = alert.StartTime + m.mu.Lock() + m.setActiveAlertNoLock(alert.ID, alert) + m.mu.Unlock() + key := canonicalTrackingKeyForAlert(alert) + var deliveries atomic.Int32 + m.SetAlertCallback(func(*Alert) { deliveries.Add(1) }) + transitions := make(chan string, 4) + m.SetFlappingDetectedCallback(func(a *Alert, trackingKey string) { + // Re-enter the manager: the one-shot callback must remain outside its lock. + active := m.GetActiveAlerts() + if len(active) != 1 || a.ID != alert.ID || trackingKey != key { + t.Error("flapping callback lost the continuing occurrence identity") + } + transitions <- trackingKey + }) + dispatch := func() bool { + m.mu.Lock() + defer m.mu.Unlock() + return m.dispatchAlert(alert, false) + } + receiveTransition := func() { + t.Helper() + select { + case <-transitions: + case <-time.After(2 * time.Second): + t.Error("new flapping episode did not notify its one-shot callback") + } + } + for attempt := 1; attempt <= cfg.FlappingThreshold; attempt++ { + if got, want := dispatch(), attempt < cfg.FlappingThreshold; got != want { + t.Fatalf("initial dispatch %d = %v, want %v", attempt, got, want) + } + } + receiveTransition() + + m.mu.Lock() + for i := range m.flappingHistory[key] { + m.flappingHistory[key][i] = m.flappingHistory[key][i].Add(-16 * time.Minute) + } + m.suppressedUntil[key] = time.Now().Add(-time.Second) + m.mu.Unlock() + sweep.run(m) + diagnosis, found := m.DiagnoseAlertDelivery(alert.ID) + if !found || diagnosis.FlappingActive || diagnosis.FlappingHistoryInWindow != 0 || + diagnosis.SuppressedUntil != nil || diagnosis.Reason != AlertDeliveryReasonCooldown { + t.Errorf("served episode not retired coherently: %+v", diagnosis) + } + for attempt := 1; attempt <= cfg.FlappingThreshold; attempt++ { + if got, want := dispatch(), attempt < cfg.FlappingThreshold; got != want { + t.Errorf("post-cleanup dispatch %d = %v, want %v", attempt, got, want) + } + } + receiveTransition() + diagnosis, _ = m.DiagnoseAlertDelivery(alert.ID) + if diagnosis.Reason != AlertDeliveryReasonFlapping || !diagnosis.FlappingActive || + diagnosis.SuppressedUntil == nil || !diagnosis.SuppressedUntil.After(time.Now()) { + t.Errorf("new storm has no bounded cooldown: %+v", diagnosis) + } + deadline := diagnosis.SuppressedUntil + for range 20 { + sweep.run(m) + if dispatch() { + t.Error("cleanup released an unexpired cooldown") + } + } + final, _ := m.DiagnoseAlertDelivery(alert.ID) + if !reflect.DeepEqual(final.SuppressedUntil, deadline) || final.FlappingHistoryInWindow != 3 || + deliveries.Load() != 4 || len(transitions) != 0 { + t.Errorf("cooldown changed during cleanup: before=%+v after=%+v deliveries=%d extra callbacks=%d", + diagnosis, final, deliveries.Load(), len(transitions)) + } + active := m.GetActiveAlerts() + if len(active) != 1 || active[0].ID != alert.ID || !active[0].StartTime.Equal(alert.StartTime) || active[0].Acknowledged { + t.Fatal("suppression expiry changed or removed the still-active occurrence") + } + t.Logf("two bounded episodes: %d dispatch callbacks; both cleanup and delivery diagnosis preserve occurrence %s", deliveries.Load(), alert.ID) + }) + } + } +} + +func TestFlappingCleanupKeepsUnexpiredAndOtherKeys(t *testing.T) { + for _, sweep := range []struct { + name string + run func(*Manager) + }{ + {"ordinary", func(m *Manager) { m.Cleanup(time.Hour) }}, + {"hourly", (*Manager).cleanupStaleMaps}, + } { + t.Run(sweep.name, func(t *testing.T) { + m := NewManagerWithDataDir(t.TempDir()) + t.Cleanup(m.Stop) + now := time.Now() + m.mu.Lock() + m.config.AutoAcknowledgeAfterHours = 0 + for _, key := range []string{"expired-episode", "ongoing-episode", "ordinary-suppression", "pending-burst"} { + m.activeAlerts[key] = &Alert{ID: key, StartTime: now, LastSeen: now} + m.flappingHistory[key] = []time.Time{now} + } + m.flappingActive["expired-episode"] = true + m.flappingActive["ongoing-episode"] = true + m.suppressedUntil["expired-episode"] = now.Add(-time.Second) + m.suppressedUntil["ongoing-episode"] = now.Add(time.Hour) + m.suppressedUntil["ordinary-suppression"] = now.Add(time.Hour) + m.mu.Unlock() + sweep.run(m) + m.mu.RLock() + defer m.mu.RUnlock() + if m.flappingActive["expired-episode"] || len(m.flappingHistory["expired-episode"]) != 0 { + t.Error("expired suppression kept its episode") + } + if _, exists := m.suppressedUntil["expired-episode"]; exists { + t.Error("expired deadline retained") + } + if !m.flappingActive["ongoing-episode"] || !m.suppressedUntil["ongoing-episode"].Equal(now.Add(time.Hour)) || + !m.suppressedUntil["ordinary-suppression"].Equal(now.Add(time.Hour)) || m.flappingActive["ordinary-suppression"] { + t.Error("cleanup altered another key's unexpired policy") + } + for _, key := range []string{"ongoing-episode", "ordinary-suppression", "pending-burst"} { + if !reflect.DeepEqual(m.flappingHistory[key], []time.Time{now}) { + t.Errorf("cleanup altered %s observations", key) + } + } + }) + } +} diff --git a/internal/alerts/notification_policy.go b/internal/alerts/notification_policy.go index 4356f95f7..2fc298450 100644 --- a/internal/alerts/notification_policy.go +++ b/internal/alerts/notification_policy.go @@ -103,12 +103,7 @@ func (m *Manager) checkFlappingLocked(trackingKey string) (suppress bool, justTr if now.Before(until) { return true, false } - // The cooldown has been served. Clear the latch so a later episode can - // open a fresh one; flappingActive is otherwise only ever set to true, - // which would leave every subsequent episode with no cooldown at all. - delete(m.suppressedUntil, trackingKey) - delete(m.flappingActive, trackingKey) - delete(m.flappingHistory, trackingKey) + m.expireSuppressionNoLock(trackingKey, now) } // Record this state change @@ -162,6 +157,21 @@ func (m *Manager) checkFlappingLocked(trackingKey string) (suppress bool, justTr return false, false } +// expireSuppressionNoLock retires a served suppression and its episode together. +// Cleanup must use the same reset as dispatch: deleting only the deadline leaves +// a latched flapping episode that can suppress delivery without ever rearming a +// cooldown or notifying the flapping callback. Caller MUST hold m.mu. +func (m *Manager) expireSuppressionNoLock(trackingKey string, now time.Time) bool { + until, exists := m.suppressedUntil[trackingKey] + if !exists || now.Before(until) { + return false + } + delete(m.suppressedUntil, trackingKey) + delete(m.flappingActive, trackingKey) + delete(m.flappingHistory, trackingKey) + return true +} + func (m *Manager) dispatchAlert(alert *Alert, async bool) bool { callbacks := m.getAlertCallbacks() if len(callbacks) == 0 || alert == nil { diff --git a/internal/alerts/tracking_cleanup.go b/internal/alerts/tracking_cleanup.go index 7b7294b6d..a32860006 100644 --- a/internal/alerts/tracking_cleanup.go +++ b/internal/alerts/tracking_cleanup.go @@ -45,9 +45,8 @@ func (m *Manager) cleanupStaleMaps() { } } - for alertID, suppressUntil := range m.suppressedUntil { - if now.After(suppressUntil) { - delete(m.suppressedUntil, alertID) + for alertID := range m.suppressedUntil { + if m.expireSuppressionNoLock(alertID, now) { cleaned++ } }