mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 04:38:48 +00:00
Merge reviewed bounded flapping cleanup repair
Change-source: pulse-maintainer
This commit is contained in:
commit
971ebde8af
5 changed files with 193 additions and 13 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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++
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue