Coalesce repeated alert delivery holds

This commit is contained in:
Pulse Test
2026-08-27 23:19:03 +01:00
parent 496801e3d2
commit e0a089baab
3 changed files with 393 additions and 5 deletions
@@ -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
+172 -4
View File
@@ -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 ")
}
@@ -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))
}
}