From e0a089baabf4cb8774fc6bdee551db8eef11268d Mon Sep 17 00:00:00 2001 From: Pulse Test Date: Thu, 27 Aug 2026 23:19:03 +0100 Subject: [PATCH] Coalesce repeated alert delivery holds --- .../v6/internal/subsystems/alerts.md | 15 +- internal/alerts/eventlog/eventlog.go | 176 ++++++++++++++- .../alerts/eventlog/eventlog_snapshot_test.go | 207 ++++++++++++++++++ 3 files changed, 393 insertions(+), 5 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/alerts.md b/docs/release-control/v6/internal/subsystems/alerts.md index aa8ec4ede..57798f811 100644 --- a/docs/release-control/v6/internal/subsystems/alerts.md +++ b/docs/release-control/v6/internal/subsystems/alerts.md @@ -1086,7 +1086,20 @@ SQLite-backed store under the alerts data directory; ephemeral managers record nothing unless a store is installed explicitly. Resolution, acknowledgement, unacknowledgement, escalation, flapping detection, dispatch, quiet-hours deferral, and suppression append immutable events without changing lifecycle -or delivery behavior. Snapshot-bearing lifecycle transitions commit +or delivery behavior. Suppression and deferral are recorded as outcome +episodes: the first decision is immutable, identical reevaluations of that +unchanged outcome append no duplicate row, and a changed reason, details, +resource presentation, intervening lifecycle/dispatch event, or later return +to that outcome opens a new episode. Lifecycle, escalation, and dispatch +events are never coalesced. This keeps delivery activity explanatory instead +of poll-frequency-shaped, prevents unchanged reevaluations from crowding a +real outcome change out of the non-blocking diagnostic buffer, and bounds +diagnostic growth without erasing a meaningful decision transition. A failed +write forgets its admission key so a recovered store can accept the outcome +again. Reads apply the same episode projection to +redundant diagnostic rows written by older versions, so upgraded installations +become readable immediately without rewriting their immutable event records. +Snapshot-bearing lifecycle transitions commit synchronously before downstream lifecycle projections run; they never share the droppable diagnostic buffer because alert history is reconstructed from them. High-volume notification decisions remain non-blocking and fail-open for diff --git a/internal/alerts/eventlog/eventlog.go b/internal/alerts/eventlog/eventlog.go index 471633dab..8f93e918d 100644 --- a/internal/alerts/eventlog/eventlog.go +++ b/internal/alerts/eventlog/eventlog.go @@ -1,6 +1,8 @@ // Package eventlog is the append-only alert event log: every alert lifecycle -// transition and notification decision — including suppressions, with the -// mechanism that held them — is recorded as one immutable event. It exists so +// transition and distinct notification-decision episode — including +// suppressions, with the mechanism that held them — is recorded as one +// immutable event. Re-evaluating an unchanged suppression or deferral does not +// create another episode. It exists so // "why didn't I get notified?" is answerable from durable data instead of // being reconstructed from logs (docs/ALERT_ENGINE_EVOLUTION.md, Phase 0). // @@ -112,6 +114,8 @@ type Store struct { failureMu sync.RWMutex lastError error retention time.Duration + episodeMu sync.Mutex + episodes map[string]string } func sqliteDSN(path string) string { @@ -171,6 +175,7 @@ func openDSN(dbPath, dsn string) (*Store, error) { events: make(chan Event, appendBufferSize), stop: make(chan struct{}), retention: defaultRetention, + episodes: make(map[string]string), } if err := s.initSchema(); err != nil { db.Close() @@ -258,6 +263,24 @@ func (s *Store) Append(event Event) { if event.OccurredAt.IsZero() { event.OccurredAt = time.Now() } + if event.AlertID != "" && coalescesUnchangedDeliveryEpisode(event.Type) { + signature := deliveryEpisodeSignature(event) + s.episodeMu.Lock() + if s.episodes[event.AlertID] == signature { + s.episodeMu.Unlock() + return + } + select { + case s.events <- event: + s.episodes[event.AlertID] = signature + s.appended.Add(1) + default: + s.dropped.Add(1) + } + s.episodeMu.Unlock() + return + } + s.clearDeliveryEpisode(event.AlertID) select { case s.events <- event: s.appended.Add(1) @@ -276,6 +299,7 @@ func (s *Store) AppendDurable(event Event) error { if event.OccurredAt.IsZero() { event.OccurredAt = time.Now() } + s.clearDeliveryEpisode(event.AlertID) if err := s.insertBatch([]Event{event}); err != nil { s.recordWriteFailure(err) return err @@ -327,6 +351,7 @@ func (s *Store) writeLoop() { } } if err := s.insertBatch(batch); err != nil { + s.forgetFailedDeliveryEpisodes(batch) s.recordWriteFailure(err) log.Error().Err(err).Int("events", len(batch)).Msg("alert event log write failed") continue @@ -340,6 +365,7 @@ func (s *Store) writeLoop() { select { case event := <-s.events: if err := s.insertBatch([]Event{event}); err != nil { + s.forgetFailedDeliveryEpisodes([]Event{event}) s.recordWriteFailure(err) log.Error().Err(err).Msg("alert event log final drain write failed") continue @@ -391,6 +417,34 @@ func (s *Store) insertBatch(batch []Event) error { return err } defer stmt.Close() + decisionStmt, err := tx.Prepare(` + INSERT INTO alert_events + (occurred_at, event_type, alert_id, resource_id, resource_name, alert_type, level, reason, message, details, snapshot) + SELECT ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ? + WHERE NOT EXISTS ( + SELECT 1 FROM alert_events AS latest + WHERE latest.id = ( + SELECT candidate.id FROM alert_events AS candidate + WHERE candidate.alert_id = ? + ORDER BY candidate.occurred_at DESC, candidate.id DESC + LIMIT 1 + ) + AND latest.event_type = ? + AND latest.resource_id = ? + AND latest.resource_name = ? + AND latest.alert_type = ? + AND latest.level = ? + AND latest.reason = ? + AND latest.message = ? + AND latest.details = ? + AND latest.snapshot = ? + ) + `) + if err != nil { + _ = tx.Rollback() + return err + } + defer decisionStmt.Close() var importStmt *sql.Stmt for _, event := range batch { if event.Type != TypeHistoryImported { @@ -443,6 +497,19 @@ func (s *Store) insertBatch(batch []Event) error { occurredAt, snapshot, )...) + } else if event.AlertID != "" && coalescesUnchangedDeliveryEpisode(event.Type) { + _, insertErr = decisionStmt.Exec(append(args, + event.AlertID, + event.Type, + event.ResourceID, + event.ResourceName, + event.AlertType, + event.Level, + event.Reason, + event.Message, + details, + snapshot, + )...) } else { _, insertErr = stmt.Exec(args...) } @@ -458,6 +525,67 @@ func (s *Store) insertBatch(batch []Event) error { return tx.Commit() } +func coalescesUnchangedDeliveryEpisode(eventType string) bool { + return eventType == TypeNotificationSuppressed || eventType == TypeNotificationDeferred +} + +type deliveryEpisodeIdentity struct { + Type string `json:"type"` + ResourceID string `json:"resourceId"` + ResourceName string `json:"resourceName"` + AlertType string `json:"alertType"` + Level string `json:"level"` + Reason string `json:"reason"` + Message string `json:"message"` + Details map[string]string `json:"details"` + Snapshot string `json:"snapshot"` +} + +func deliveryEpisodeSignature(event Event) string { + details := event.Details + if len(details) == 0 { + details = nil + } + encoded, _ := json.Marshal(deliveryEpisodeIdentity{ + Type: event.Type, + ResourceID: event.ResourceID, + ResourceName: event.ResourceName, + AlertType: event.AlertType, + Level: event.Level, + Reason: event.Reason, + Message: event.Message, + Details: details, + Snapshot: string(event.Snapshot), + }) + return string(encoded) +} + +func (s *Store) clearDeliveryEpisode(alertID string) { + if s == nil || alertID == "" { + return + } + s.episodeMu.Lock() + delete(s.episodes, alertID) + s.episodeMu.Unlock() +} + +func (s *Store) forgetFailedDeliveryEpisodes(events []Event) { + if s == nil { + return + } + s.episodeMu.Lock() + defer s.episodeMu.Unlock() + for _, event := range events { + if !coalescesUnchangedDeliveryEpisode(event.Type) { + continue + } + signature := deliveryEpisodeSignature(event) + if s.episodes[event.AlertID] == signature { + delete(s.episodes, event.AlertID) + } + } +} + func (s *Store) pruneOld() { if s.retention <= 0 { return @@ -492,7 +620,47 @@ func (s *Store) Flush() error { return nil } -// Query returns matching events, newest first. +// deliveryEpisodeQuerySource projects legacy poll-frequency duplicates as the +// single decision episode they represent. It compares each diagnostic hold to +// the immediately preceding event for that alert, including events outside the +// caller's filter, so a lifecycle or dispatch transition always opens a new +// episode. New writes are already coalesced in insertBatch; this read projection +// keeps upgraded stores immediately useful without rewriting immutable rows. +const deliveryEpisodeQuerySource = `( + SELECT current.* + FROM alert_events AS current + WHERE NOT ( + current.alert_id <> '' + AND current.event_type IN ('notification_suppressed', 'notification_deferred') + AND EXISTS ( + SELECT 1 + FROM alert_events AS previous + WHERE previous.id = ( + SELECT candidate.id + FROM alert_events AS candidate + WHERE candidate.alert_id = current.alert_id + AND ( + candidate.occurred_at < current.occurred_at + OR (candidate.occurred_at = current.occurred_at AND candidate.id < current.id) + ) + ORDER BY candidate.occurred_at DESC, candidate.id DESC + LIMIT 1 + ) + AND previous.event_type = current.event_type + AND previous.resource_id = current.resource_id + AND previous.resource_name = current.resource_name + AND previous.alert_type = current.alert_type + AND previous.level = current.level + AND previous.reason = current.reason + AND previous.message = current.message + AND previous.details = current.details + AND previous.snapshot = current.snapshot + ) + ) +) AS alert_events` + +// Query returns matching events, newest first. Repeated unchanged suppression +// and deferral rows written by earlier versions are projected as one episode. func (s *Store) Query(filter Filter) ([]Event, error) { if s == nil { return nil, nil @@ -508,7 +676,7 @@ func (s *Store) Query(filter Filter) ([]Event, error) { limit = maxQueryLimit } - query := "SELECT id, occurred_at, event_type, alert_id, resource_id, resource_name, alert_type, level, reason, message, details, snapshot FROM alert_events" + query := "SELECT id, occurred_at, event_type, alert_id, resource_id, resource_name, alert_type, level, reason, message, details, snapshot FROM " + deliveryEpisodeQuerySource if len(where) > 0 { query += " WHERE " + strings.Join(where, " AND ") } diff --git a/internal/alerts/eventlog/eventlog_snapshot_test.go b/internal/alerts/eventlog/eventlog_snapshot_test.go index ce22ea934..eef8773b2 100644 --- a/internal/alerts/eventlog/eventlog_snapshot_test.go +++ b/internal/alerts/eventlog/eventlog_snapshot_test.go @@ -4,6 +4,7 @@ import ( "database/sql" "encoding/json" "path/filepath" + "sync" "testing" "time" ) @@ -84,3 +85,209 @@ func TestSnapshotRoundTrip(t *testing.T) { t.Fatalf("snapshot round trip: got %s want %s", events[0].Snapshot, snapshot) } } + +func TestUnchangedDeliveryDecisionIsOneEpisode(t *testing.T) { + store, err := OpenInMemory() + if err != nil { + t.Fatal(err) + } + defer store.Close() + + started := time.Date(2026, 8, 27, 20, 0, 0, 0, time.UTC) + base := Event{ + OccurredAt: started, + Type: TypeNotificationSuppressed, + AlertID: "alert-1", + ResourceID: "node/pve-1", + ResourceName: "pve-1", + AlertType: "cpu", + Level: "warning", + Reason: "notifications_inactive", + Message: "Notification suppressed: alert delivery is not turned on.", + Details: map[string]string{"activationState": "pending_review"}, + } + store.Append(base) + store.Flush() + repeated := base + repeated.OccurredAt = started.Add(10 * time.Second) + store.Append(repeated) + store.Flush() + if got := store.appended.Load(); got != 1 { + t.Fatalf("unchanged suppression admitted %d diagnostic events, want one", got) + } + + events, err := store.Query(Filter{AlertID: base.AlertID}) + if err != nil { + t.Fatal(err) + } + if len(events) != 1 { + t.Fatalf("unchanged suppression produced %d events, want one episode", len(events)) + } + if !events[0].OccurredAt.Equal(started) { + t.Fatalf("episode timestamp = %s, want first decision at %s", events[0].OccurredAt, started) + } + + changed := repeated + changed.OccurredAt = started.Add(20 * time.Second) + changed.Reason = "quiet_hours:critical" + changed.Message = "Notification deferred by quiet hours." + changed.Type = TypeNotificationDeferred + changed.Details = map[string]string{"replayAt": started.Add(time.Hour).Format(time.RFC3339)} + store.Append(changed) + store.Flush() + repeatedDeferred := changed + repeatedDeferred.OccurredAt = started.Add(25 * time.Second) + store.Append(repeatedDeferred) + store.Append(Event{ + OccurredAt: started.Add(30 * time.Second), + Type: TypeNotificationDispatched, + AlertID: base.AlertID, + Reason: "ready", + }) + resuppressed := base + resuppressed.OccurredAt = started.Add(40 * time.Second) + store.Append(resuppressed) + store.Flush() + if got := store.appended.Load(); got != 4 { + t.Fatalf("changed delivery lifecycle admitted %d diagnostic events, want four", got) + } + + events, err = store.Query(Filter{AlertID: base.AlertID}) + if err != nil { + t.Fatal(err) + } + if len(events) != 4 { + t.Fatalf("changed delivery lifecycle produced %d events, want four distinct episodes", len(events)) + } + if events[0].Type != TypeNotificationSuppressed || + events[1].Type != TypeNotificationDispatched || + events[2].Type != TypeNotificationDeferred || + events[3].Type != TypeNotificationSuppressed { + t.Fatalf("delivery episode order = %+v", events) + } +} + +func TestQueryProjectsLegacyRepeatedDeliveryRowsAsEpisodes(t *testing.T) { + store, err := OpenInMemory() + if err != nil { + t.Fatal(err) + } + defer store.Close() + + started := time.Date(2026, 8, 27, 20, 0, 0, 0, time.UTC) + insert := func(alertID string, occurredAt time.Time, eventType, reason string) { + t.Helper() + if _, err := store.db.Exec(` + INSERT INTO alert_events + (occurred_at, event_type, alert_id, resource_id, resource_name, alert_type, level, reason, message, details, snapshot) + VALUES (?, ?, ?, 'node/pve-1', 'pve-1', 'cpu', 'warning', ?, 'held', '{}', '') + `, occurredAt.Format(time.RFC3339Nano), eventType, alertID, reason); err != nil { + t.Fatal(err) + } + } + + insert("older-alert", started.Add(-time.Minute), TypeNotificationDeferred, "quiet_hours") + insert("legacy-alert", started, TypeNotificationSuppressed, "notifications_inactive") + insert("legacy-alert", started.Add(10*time.Second), TypeNotificationSuppressed, "notifications_inactive") + insert("legacy-alert", started.Add(20*time.Second), TypeNotificationDispatched, "ready") + insert("legacy-alert", started.Add(30*time.Second), TypeNotificationSuppressed, "notifications_inactive") + insert("legacy-alert", started.Add(40*time.Second), TypeNotificationSuppressed, "notifications_inactive") + + events, err := store.Query(Filter{Limit: 4}) + if err != nil { + t.Fatal(err) + } + if len(events) != 4 { + t.Fatalf("legacy delivery rows projected as %d events, want four distinct episodes", len(events)) + } + if !events[0].OccurredAt.Equal(started.Add(30*time.Second)) || + events[1].Type != TypeNotificationDispatched || + !events[2].OccurredAt.Equal(started) || + events[3].AlertID != "older-alert" { + t.Fatalf("legacy delivery episode order = %+v", events) + } +} + +func TestDeliveryEpisodeCoalescingSurvivesRestart(t *testing.T) { + dir := t.TempDir() + store, err := Open(dir) + if err != nil { + t.Fatal(err) + } + started := time.Date(2026, 8, 27, 20, 0, 0, 0, time.UTC) + event := Event{ + OccurredAt: started, + Type: TypeNotificationSuppressed, + AlertID: "restart-alert", + Reason: "notifications_inactive", + Message: "held", + Details: map[string]string{"activationState": "pending_review"}, + } + store.Append(event) + if err := store.Flush(); err != nil { + t.Fatal(err) + } + store.Close() + + reopened, err := Open(dir) + if err != nil { + t.Fatal(err) + } + defer reopened.Close() + repeated := event + repeated.OccurredAt = started.Add(10 * time.Second) + reopened.Append(repeated) + if err := reopened.Flush(); err != nil { + t.Fatal(err) + } + + var rows int + if err := reopened.db.QueryRow( + `SELECT COUNT(*) FROM alert_events WHERE alert_id = ?`, + event.AlertID, + ).Scan(&rows); err != nil { + t.Fatal(err) + } + if rows != 1 { + t.Fatalf("restart appended %d unchanged suppression rows, want one episode", rows) + } +} + +func TestConcurrentUnchangedDeliveryDecisionAdmitsOneEpisode(t *testing.T) { + store, err := OpenInMemory() + if err != nil { + t.Fatal(err) + } + defer store.Close() + + event := Event{ + OccurredAt: time.Date(2026, 8, 27, 20, 0, 0, 0, time.UTC), + Type: TypeNotificationSuppressed, + AlertID: "concurrent-alert", + Reason: "acknowledged", + Message: "held", + } + var writers sync.WaitGroup + for range 100 { + writers.Add(1) + go func() { + defer writers.Done() + store.Append(event) + }() + } + writers.Wait() + if err := store.Flush(); err != nil { + t.Fatal(err) + } + + if got := store.appended.Load(); got != 1 { + t.Fatalf("concurrent reevaluations admitted %d diagnostic events, want one", got) + } + events, err := store.Query(Filter{AlertID: event.AlertID}) + if err != nil { + t.Fatal(err) + } + if len(events) != 1 { + t.Fatalf("concurrent reevaluations wrote %d episodes, want one", len(events)) + } +}