Merge reviewed bounded flapping cleanup repair

Change-source: pulse-maintainer
This commit is contained in:
pulse-triage[bot] 2026-10-03 00:58:10 +01:00
commit 971ebde8af
5 changed files with 193 additions and 13 deletions

View file

@ -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

View file

@ -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)

View file

@ -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)
}
}
})
}
}

View file

@ -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 {

View file

@ -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++
}
}