diff --git a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md index 7c7ba657d..6ee1881f3 100644 --- a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md +++ b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md @@ -794,6 +794,14 @@ installer download and the agent's subsequent Pulse TLS connection. ## Shared Boundaries +The shared monitor constructor also composes notification bootstrap: it loads +saved alert policy and destinations, binds the alert owner's immutable +quiet-hours policy provider, then activates persistent queue delivery. +This ordering is notifications/alerts-owned and must not change agent enrollment +or transport. `internal/monitoring/monitor_notification_startup_test.go` checks +the actual constructor and reconstruction, including holding due persisted +notifications under a saved continuous schedule without provider attempts. + - Local administrator setup synchronises the router authorizer before returning a browser session, so API Access remains available to manage agent credentials. This changes neither agent token scopes nor command-policy intent, and must not grant an unrelated identity access to credential management. ### Notification recovery reload ownership diff --git a/docs/release-control/v6/internal/subsystems/alerts.md b/docs/release-control/v6/internal/subsystems/alerts.md index eb93915d9..3cc8f2321 100644 --- a/docs/release-control/v6/internal/subsystems/alerts.md +++ b/docs/release-control/v6/internal/subsystems/alerts.md @@ -1988,14 +1988,25 @@ quiet-hours policy; nonexistent clock minutes do not shift the configured start or end. Initial queued replay uses the first real minute outside the current non-full-day window, including an unselected day reached at midnight. This does not reinterpret the day toggles as the previous evening's ownership. -Full-day windows retain the existing daily replay boundary; this clock repair -is not current-policy revalidation of already queued work or external receipt. +Full-day windows retain the existing daily revalidation boundary; it is not an +invented open minute or permission to send during a continuous schedule. `TestQuietHoursDaylightSavingClockAndReplay` in `internal/alerts/quiet_hours_test.go` checks London, New York and Lord Howe (one-hour and half-hour changes), inclusive boundaries and exact UTC replay. `TestQuietHoursOvernightSelectedCalendarDays` checks overnight day selection. Quiet-hours location fallback is read-only; configuration updates own the location cache, including when public suppression helpers run concurrently. +`Manager.QuietHoursNotificationPolicy` exposes an immutable read-only snapshot +to the notification owner before it takes queue or per-alert delivery locks. +Initial admission and queued replay share civil-clock/category semantics, +including repeated/skipped minutes and inclusive ends. The queue revalidates +every due attempt against current policy rather than trusting old replay +metadata. Policy snapshots neither mutate alert lifecycle/acknowledgement nor +select destinations; cancellation and actual delivery remain queue-owned. +`TestQuietHoursDeliverySnapshotIsCurrentIndependentAndReadOnly` and +`TestQuietHoursDeliverySnapshotConcurrentConfigUpdates` in +`internal/alerts/quiet_hours_test.go` cover snapshot isolation and concurrent +configuration reads. This establishes source scheduling, not external receipt. Quiet-hours suppression also applies to alert delivery lifecycle, not only the initial raised notification. Resolved notifications must not fan out when the alert was never notified or was already acknowledged, and monitoring-driven diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index f8151068f..d0764f50b 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -1799,6 +1799,19 @@ service-history reads plus denial/recovery without fabricated samples. ## Current State +### Saved quiet-hours policy before queue activation + +`Monitor.New` binds `alerts.Manager.QuietHoursNotificationPolicy` to the +notification manager after saved policy/destinations load and before +`StartQueueProcessing`. The notification owner revalidates due stored work and +owns all partitioning, delivery and cancellation; monitoring neither duplicates +the schedule nor changes destination selection or operational lifecycle. +`TestNewRevalidatesPersistedQuietHoursBeforeDelivery` in +`internal/monitoring/monitor_notification_startup_test.go` reconstructs the +actual monitor twice with due persisted work under a full-day saved schedule, +checking pending state and no local delivery/audit. This is source bootstrap +proof, not installed restart or provider acceptance. + TrueNAS REST alert arguments may be object, scalar, array, null or absent. Non-object arguments must not abort snapshot collection or discard alerts. Only typed object fields supply disk identity and SMART counters; neither scalar diff --git a/docs/release-control/v6/internal/subsystems/notifications.md b/docs/release-control/v6/internal/subsystems/notifications.md index f542cf02b..55c36f3cf 100644 --- a/docs/release-control/v6/internal/subsystems/notifications.md +++ b/docs/release-control/v6/internal/subsystems/notifications.md @@ -73,6 +73,7 @@ or displayed a notification. 10. `internal/notifications/delivery_health.go` 11. `internal/notifications/deadman_config.go` 12. `internal/notifications/failure_class.go` +13. `internal/notifications/quiet_hours_queue.go` ## Shared Boundaries @@ -141,6 +142,44 @@ stable opaque routing identities and must not expose credentials. ## Current State +### Current-policy quiet-hours replay + +The persistent queue checks every due provider attempt against an immutable +snapshot from `alerts.Manager.QuietHoursNotificationPolicy`, installed by +`NotificationManager.SetQuietHoursPolicyProvider` before saved work is activated +in the monitor constructor. The snapshot is obtained before queue/database or +per-alert delivery locks; its evaluator does not call back into the alert owner. +Stored replay timestamps are wake conditions, not evidence that the current +schedule permits delivery. Full-day schedules recheck at the retained daily +boundary, and late replay into another quiet interval remains held. + +`quiet_hours_queue.go` reloads the pending row after acquiring delivery gates, +preserving cancellations that rewrote a grouped row after worker discovery. +An all-held batch moves only its wake time. A mixed batch atomically retains the +held original and creates one ready child with the admitted destination bytes, +occurrence/transition links, creation time and remaining attempt budget. Child +IDs bind the parent and exact ready payload; held work does not gain an attempt, +audit entry, firing receipt or false success. Prior attempt audits remain on the +original batch, while actual child attempts carry their own linked receipts. +Already-due recovery groups cannot reopen their grouping window while held. +Partition persistence or ambiguous linkage failure leaves the original intact +and makes no provider attempt. Due-time refresh prevents stale worker snapshots +from repeatedly postponing or claiming a now-held row. Actual retries retain +the existing claim/cancellation gates and provider enabled checks. + +This does not accelerate a future stored deadline after a schedule edit; the +current rule is applied at the next admitted wake. It does not change initial +target selection, mutable recovery tag admission, acknowledgement semantics, +operator retry budgets or at-least-once delivery after an interrupted send. + +`quiet_hours_queue_test.go` proves continuous/late/DST replay holds, mixed +immediate/held delivery, exact admitted destination and operational linkage, +budget retention, transaction rollback, durable reopen, cancellation and stale +concurrent snapshots, grouped recoveries and all six provider/event families. +`internal/monitoring/monitor_notification_startup_test.go` exercises the actual +constructor with due persisted work across two startups. These are local source +and HTTP-receipt proofs, not installed provider or natural-day acceptance. + ### Persistent resolved-alert grouping `SendResolvedAlert` applies the configured grouping window to recoveries as diff --git a/docs/release-control/v6/internal/subsystems/registry.json b/docs/release-control/v6/internal/subsystems/registry.json index ea6d02710..36dfea172 100644 --- a/docs/release-control/v6/internal/subsystems/registry.json +++ b/docs/release-control/v6/internal/subsystems/registry.json @@ -1623,6 +1623,7 @@ "internal/monitoring/docker_metric_presence_test.go", "internal/monitoring/monitor_host_agent_removal_lifecycle_test.go", "internal/monitoring/monitor_host_agents_test.go", + "internal/monitoring/monitor_notification_startup_test.go", "internal/monitoring/physical_disk_roundtrip_test.go", "scripts/installtests/agent_state_dir_lifecycle_test.go", "scripts/installtests/install_ps1_test.go", @@ -6565,6 +6566,7 @@ "internal/monitoring/monitor_docker_test.go", "internal/monitoring/monitor_host_agent_removal_lifecycle_test.go", "internal/monitoring/monitor_host_agents_test.go", + "internal/monitoring/monitor_notification_startup_test.go", "internal/monitoring/monitor_mock_alerts_test.go", "internal/monitoring/monitor_pve_cluster_refresh_test.go", "internal/monitoring/monitor_pve_guest_lxc_test.go", diff --git a/internal/alerts/notification_policy.go b/internal/alerts/notification_policy.go index cfa5d1950..4356f95f7 100644 --- a/internal/alerts/notification_policy.go +++ b/internal/alerts/notification_policy.go @@ -481,14 +481,18 @@ func quietHoursCategoryForAlert(alert *Alert) string { } func (m *Manager) shouldSuppressNotification(alert *Alert) (bool, string) { - if alert == nil { - return false, "" - } // Only quiet-hours matches are replayable. Operator policy is checked by // each caller before this helper, because it must drop rather than defer. if !m.isInQuietHours() { return false, "" } + return quietHoursSuppressionForAlert(alert, m.config.Schedule.QuietHours.Suppress) +} + +func quietHoursSuppressionForAlert(alert *Alert, suppress QuietHoursSuppression) (bool, string) { + if alert == nil { + return false, "" + } if alert.Level != AlertLevelCritical { return true, "non-critical" @@ -497,15 +501,15 @@ func (m *Manager) shouldSuppressNotification(alert *Alert) (bool, string) { category := quietHoursCategoryForAlert(alert) switch category { case "performance": - if m.config.Schedule.QuietHours.Suppress.Performance { + if suppress.Performance { return true, category } case "storage": - if m.config.Schedule.QuietHours.Suppress.Storage { + if suppress.Storage { return true, category } case "offline": - if m.config.Schedule.QuietHours.Suppress.Offline { + if suppress.Offline { return true, category } } @@ -516,13 +520,42 @@ func (m *Manager) shouldSuppressNotification(alert *Alert) (bool, string) { func (m *Manager) quietHoursReplayAt() time.Time { now := m.policyNow() clock, valid := m.quietHoursClock() - if !valid || !clock.contains(now) { + if !valid { + return now.Add(time.Minute).UTC() + } + return clock.replayAt(now) +} + +// QuietHoursNotificationPolicy returns an immutable snapshot of the current +// quiet-hours rule. The delivery owner obtains it before taking queue or +// per-alert delivery locks, so evaluation never calls back into this manager +// while cancellation is waiting for those locks. It does not mutate alert +// lifecycle, acknowledgement, metadata or destination selection. +func (m *Manager) QuietHoursNotificationPolicy() func(*Alert, time.Time) *time.Time { + m.mu.RLock() + clock, valid := m.quietHoursClock() + suppress := m.config.Schedule.QuietHours.Suppress + m.mu.RUnlock() + return func(alert *Alert, now time.Time) *time.Time { + if !valid || !clock.contains(now) { + return nil + } + if suppressed, _ := quietHoursSuppressionForAlert(alert, suppress); !suppressed { + return nil + } + replayAt := clock.replayAt(now) + return &replayAt + } +} + +func (clock quietHoursClock) replayAt(now time.Time) time.Time { + if !clock.contains(now) { return now.Add(time.Minute).UTC() } // A full-day window has no non-quiet clock minute. Preserve its existing - // daily replay boundary; queued-policy revalidation is a separate delivery - // concern, not permission to invent an end to a continuous schedule. + // daily revalidation boundary. The queue must check policy again there, + // not invent an end to a continuous schedule or attempt a provider send. if (clock.end-clock.start+24*60)%(24*60) == 24*60-1 { local := now.In(clock.location) endExclusive := time.Date(local.Year(), local.Month(), local.Day(), clock.end/60, clock.end%60, 0, 0, clock.location).Add(time.Minute) diff --git a/internal/alerts/quiet_hours_test.go b/internal/alerts/quiet_hours_test.go index cd737bff9..b6ab51181 100644 --- a/internal/alerts/quiet_hours_test.go +++ b/internal/alerts/quiet_hours_test.go @@ -1,6 +1,7 @@ package alerts import ( + "reflect" "sync" "testing" "time" @@ -8,6 +9,74 @@ import ( alertspecs "github.com/rcourtman/pulse-go-rewrite/internal/alerts/specs" ) +func TestQuietHoursDeliverySnapshotIsCurrentIndependentAndReadOnly(t *testing.T) { + now := time.Date(2026, 10, 25, 1, 45, 0, 0, time.UTC) + quiet := QuietHours{Enabled: true, Start: "01:30", End: "02:00", Timezone: "Europe/London", + Days: map[string]bool{"sunday": true}, Suppress: QuietHoursSuppression{Storage: true}} + m := fixedQuietHoursTestManager(now, quiet) + defer m.Stop() + policy := m.QuietHoursNotificationPolicy() + for _, tc := range []struct { + kind string + level AlertLevel + held bool + }{ + {"cpu", AlertLevelWarning, true}, + {"cpu", AlertLevelCritical, false}, + {"connectivity", AlertLevelCritical, false}, + {"disk-health", AlertLevelCritical, true}, + } { + alert := &Alert{ID: tc.kind, Type: tc.kind, Level: tc.level, Acknowledged: true, + Metadata: map[string]interface{}{MetadataQuietHoursReplayAt: "old-admission", "unchanged": "value"}} + before := alert.Clone() + at := policy(alert, now) + if (at != nil) != tc.held || (at != nil && !at.Equal(time.Date(2026, 10, 25, 2, 1, 0, 0, time.UTC))) { + t.Fatalf("%s level=%s held=%t replay=%v", tc.kind, tc.level, tc.held, at) + } + if !reflect.DeepEqual(before, alert) { + t.Fatal("delivery policy mutated acknowledgement or original occurrence metadata") + } + } + if at := policy(nil, now); at != nil { + t.Fatalf("nil alert acquired a hold: %v", at) + } + config := m.GetConfig() + config.Schedule.QuietHours = QuietHours{} + m.UpdateConfig(config) + alert := &Alert{Type: "cpu", Level: AlertLevelWarning} + if at := policy(alert, now); at == nil { + t.Fatal("immutable policy snapshot changed after configuration update") + } + if at := m.QuietHoursNotificationPolicy()(alert, now); at != nil { + t.Fatal("fresh delivery snapshot did not pick up disabled quiet hours") + } +} + +func TestQuietHoursDeliverySnapshotConcurrentConfigUpdates(t *testing.T) { + now := time.Date(2026, 10, 2, 0, 0, 0, 0, time.UTC) + m := newManagerWithQuietHoursSuppress(QuietHoursSuppression{}) + defer m.Stop() + config := m.GetConfig() + var wg sync.WaitGroup + for i := 0; i < 8; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + for range 20 { + if i == 0 { + m.UpdateConfig(config) + } else { + at := m.QuietHoursNotificationPolicy()(&Alert{Type: "cpu", Level: AlertLevelWarning}, now) + if at == nil || !at.After(now) { + t.Errorf("continuous schedule acquired an invented open minute: %v", at) + } + } + } + }(i) + } + wg.Wait() +} + func fixedQuietHoursTestManager(now time.Time, quietHours QuietHours) *Manager { m := NewManager() m.now = func() time.Time { return now } diff --git a/internal/monitoring/monitor.go b/internal/monitoring/monitor.go index ce0c51b77..b0f3e2e5f 100644 --- a/internal/monitoring/monitor.go +++ b/internal/monitoring/monitor.go @@ -1931,6 +1931,7 @@ func New(cfg *config.Config) (*Monitor, error) { // The queue worker starts with the manager but has no processor until all // saved destinations are installed. Activate it only after the last load; // a slow migration or config read must not cancel due persisted work. + m.notificationMgr.SetQuietHoursPolicyProvider(m.alertManager.QuietHoursNotificationPolicy) m.notificationMgr.StartQueueProcessing() // In mock mode the canonical sampler owns demo chart history by default. diff --git a/internal/monitoring/monitor_notification_startup_test.go b/internal/monitoring/monitor_notification_startup_test.go index 129abb1fe..34897b605 100644 --- a/internal/monitoring/monitor_notification_startup_test.go +++ b/internal/monitoring/monitor_notification_startup_test.go @@ -1,14 +1,106 @@ package monitoring import ( + "encoding/json" + "net/http" + "net/http/httptest" "reflect" + "sync/atomic" "testing" + "time" "github.com/rcourtman/pulse-go-rewrite/internal/alerts" "github.com/rcourtman/pulse-go-rewrite/internal/config" "github.com/rcourtman/pulse-go-rewrite/internal/notifications" ) +// A due persisted row must receive current saved quiet-hours policy BEFORE the +// autonomous queue is activated, on initial construction and reconstruction. +// This observes local HTTP acceptance only, not an installed destination. +func TestNewRevalidatesPersistedQuietHoursBeforeDelivery(t *testing.T) { + dir := t.TempDir() + t.Setenv("PULSE_DATA_DIR", dir) + var deliveries atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + deliveries.Add(1) + w.WriteHeader(http.StatusOK) + })) + defer server.Close() + quiet := alerts.QuietHours{Enabled: true, Start: "00:00", End: "23:59", Timezone: "UTC", + Days: map[string]bool{"sunday": true, "monday": true, "tuesday": true, + "wednesday": true, "thursday": true, "friday": true, "saturday": true}} + persistence := config.NewConfigPersistence(dir) + saved := alerts.AlertConfig{Enabled: true, ActivationState: alerts.ActivationActive, + Schedule: alerts.ScheduleConfig{QuietHours: quiet}} + if err := persistence.SaveAlertConfig(saved); err != nil { + t.Fatal(err) + } + hook := notifications.WebhookConfig{ID: "saved-local-ops", Name: "saved-local-ops", Enabled: true, URL: server.URL, Service: "generic"} + if err := persistence.SaveWebhooks([]notifications.WebhookConfig{hook}); err != nil { + t.Fatal(err) + } + configJSON, err := json.Marshal(hook) + if err != nil { + t.Fatal(err) + } + for startup := 0; startup < 2; startup++ { + seed, err := notifications.NewNotificationQueue(dir) + if err != nil { + t.Fatal(err) + } + due := time.Now().Add(-time.Hour) + id := "saved-quiet-first" + if startup == 1 { + id = "saved-quiet-second" + } + if err := seed.Enqueue(¬ifications.QueuedNotification{ID: id, Type: "webhook", Config: configJSON, + Alerts: []*alerts.Alert{{ID: id, Type: "cpu", Level: alerts.AlertLevelWarning, StartTime: due}}, + NextRetryAt: &due}); err != nil { + t.Fatal(err) + } + if err := seed.Stop(); err != nil { + t.Fatal(err) + } + m, err := New(&config.Config{DataPath: dir}) + if err != nil { + t.Fatal(err) + } + n := m.GetNotificationManager() + if err := n.UpdateAllowedPrivateCIDRs("127.0.0.1/32"); err != nil { + m.Stop() + t.Fatal(err) + } + q := n.GetQueue() + deadline := time.Now().Add(5 * time.Second) + for { + pending, err := q.GetPending(10) + if err != nil { + m.Stop() + t.Fatal(err) + } + if len(pending) == 0 { + break + } + if time.Now().After(deadline) { + m.Stop() + t.Fatal("persisted quiet notification was not revalidated") + } + time.Sleep(10 * time.Millisecond) + } + stats, err := q.GetQueueStats() + if err != nil || stats["pending"] != startup+1 || stats["sent"] != 0 || stats["dlq"] != 0 || deliveries.Load() != 0 { + m.Stop() + t.Fatalf("QUIET_BOOTSTRAP: saved current schedule did not hold every persisted row: stats=%v deliveries=%d err=%v", stats, deliveries.Load(), err) + } + logs, err := q.GetDeliveryLog(time.Now().Add(-time.Hour), 10) + if err != nil || len(logs) != 0 { + m.Stop() + t.Fatalf("quiet bootstrap manufactured provider attempts: %+v, %v", logs, err) + } + m.Stop() + } +} + // Cover the real monitor constructor, not a fixture which reapplies manager // setters after restart. Saving here uses persistence directly: this is not // browser/API save acceptance, a process restart, or destination receipt proof. diff --git a/internal/notifications/notifications.go b/internal/notifications/notifications.go index e9f85bc28..744ebbd5f 100644 --- a/internal/notifications/notifications.go +++ b/internal/notifications/notifications.go @@ -792,6 +792,20 @@ func (n *NotificationManager) StartQueueProcessing() { } } +// SetQuietHoursPolicyProvider binds persistent replay to the alert owner's +// current schedule. Install it before StartQueueProcessing during bootstrap. +// Each provider call must return a lock-independent, read-only evaluator. +func (n *NotificationManager) SetQuietHoursPolicyProvider(provider func() func(*alerts.Alert, time.Time) *time.Time) { + n.mu.RLock() + queue := n.queue + n.mu.RUnlock() + if queue != nil { + queue.mu.Lock() + queue.quietHoursPolicy = provider + queue.mu.Unlock() + } +} + // SetTenantIdentityResolver installs an org-backed resolver for the tenant // identity stamped into webhook payloads. It overrides the // PULSE_TENANT_ID / PULSE_TENANT_NAME environment defaults; multi-tenant diff --git a/internal/notifications/queue.go b/internal/notifications/queue.go index e365ec2b0..47c28fda1 100644 --- a/internal/notifications/queue.go +++ b/internal/notifications/queue.go @@ -187,8 +187,10 @@ type NotificationQueue struct { cleanupTicker *time.Ticker notifyChan chan struct{} // Signal when new notifications are added processor func(*QueuedNotification) error // Notification processor function - deliveryHealthChanged func() // Reconcile the monitoring-owned delivery alert after health-changing transitions - workerSem chan struct{} // Semaphore for limiting concurrent workers + quietHoursPolicy func() func(*alerts.Alert, time.Time) *time.Time + now func() time.Time // Queue scheduling clock; nil uses wall time. + deliveryHealthChanged func() // Reconcile the monitoring-owned delivery alert after health-changing transitions + workerSem chan struct{} // Semaphore for limiting concurrent workers } // notificationDeliveryGate orders delivery and cancellation for a single @@ -1156,7 +1158,7 @@ func (nq *NotificationQueue) GetPending(limit int) ([]*QueuedNotification, error } // scanNotification scans a database row into a QueuedNotification -func (nq *NotificationQueue) scanNotification(rows *sql.Rows) (*QueuedNotification, error) { +func (nq *NotificationQueue) scanNotification(rows interface{ Scan(...any) error }) (*QueuedNotification, error) { var notif QueuedNotification var alertsJSON, configJSON, linksJSON string var lastAttempt, nextRetryAt, completedAt *int64 @@ -1757,6 +1759,16 @@ func (nq *NotificationQueue) processNotification(notif *QueuedNotification) { Msg("Skipping cancelled notification") return } + // Obtain policy before the delivery gates: the policy owner may have + // lifecycle callbacks that cancel queued work. The returned evaluator is + // immutable and must not acquire any owner locks. + nq.mu.RLock() + policyProvider := nq.quietHoursPolicy + nq.mu.RUnlock() + var quietHoursPolicy func(*alerts.Alert, time.Time) *time.Time + if policyProvider != nil { + quietHoursPolicy = policyProvider() + } releaseDeliveryGates := nq.acquireAlertDeliveryGates(alertIdentifiersFromAlerts(notif.Alerts), false) healthChanged := false defer func() { @@ -1765,6 +1777,15 @@ func (nq *NotificationQueue) processNotification(notif *QueuedNotification) { nq.notifyDeliveryHealthChanged() } }() + if quietHoursPolicy != nil { + ready, err := nq.prepareQuietHoursDelivery(notif, quietHoursPolicy) + if err != nil { + log.Error().Err(err).Str("id", notif.ID).Msg("Failed to revalidate queued quiet-hours delivery") + } + if err != nil || !ready { + return + } + } // Atomically claim the pending row. A concurrent resolution may have // cancelled it while it was waiting for its per-alert delivery gate. diff --git a/internal/notifications/quiet_hours_queue.go b/internal/notifications/quiet_hours_queue.go new file mode 100644 index 000000000..f6eaee573 --- /dev/null +++ b/internal/notifications/quiet_hours_queue.go @@ -0,0 +1,180 @@ +package notifications + +import ( + "crypto/sha256" + "database/sql" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "strings" + "time" + + "github.com/rcourtman/pulse-go-rewrite/internal/alerts" + "github.com/rcourtman/pulse-go-rewrite/internal/operationaltrust" +) + +// prepareQuietHoursDelivery runs under the original batch's delivery gates, +// before claiming a provider attempt. Reloading the row preserves cancellation +// rewrites made after GetPending. Held work keeps the original row; a mixed +// batch gets one ready child, committed atomically with the held remainder. +// A hold is scheduling, not an attempt, failure, receipt or lifecycle change. +func (nq *NotificationQueue) prepareQuietHoursDelivery( + notif *QueuedNotification, + policy func(*alerts.Alert, time.Time) *time.Time, +) (bool, error) { + nq.mu.Lock() + defer nq.mu.Unlock() + now := time.Now() + if nq.now != nil { + now = nq.now() + } + current, err := nq.scanNotification(nq.db.QueryRow(` + SELECT id, type, method, status, alerts, config, attempts, max_attempts, + last_attempt, last_error, created_at, next_retry_at, completed_at, + payload_bytes, operational_links + FROM notification_queue WHERE id = ? AND status = 'pending'`, notif.ID)) + if errors.Is(err, sql.ErrNoRows) { + return false, nil + } + if err != nil { + return false, fmt.Errorf("reload quiet-hours notification: %w", err) + } + // Another worker may already have postponed this row since GetPending. + if current.NextRetryAt != nil && current.NextRetryAt.After(now) { + return false, nil + } + var ready, held []*alerts.Alert + var replayAt *time.Time + for _, alert := range current.Alerts { + at := policy(alert, now) + if at == nil || !at.After(now) { + ready = append(ready, alert) + continue + } + held = append(held, alert) + if replayAt == nil || at.Before(*replayAt) { + value := at.UTC() + replayAt = &value + } + } + if len(held) == 0 { + *notif = *current + return true, nil + } + // The queue stores second precision. Never round a policy deadline down + // into a provider send (or a busy retry loop) before it becomes eligible. + deadline := replayAt.Truncate(time.Second) + if deadline.Before(*replayAt) { + deadline = deadline.Add(time.Second) + } + method := current.Method + if strings.HasPrefix(method, "resolved-group") { + // This batch has reached its admitted grouping deadline. It must not + // reopen that window and absorb later recoveries while quiet-held. + method = "quiet-hours-replay" + } + if len(ready) == 0 { + _, err := nq.db.Exec(`UPDATE notification_queue + SET next_retry_at = ?, method = ? WHERE id = ? AND status = 'pending'`, deadline.Unix(), method, current.ID) + if err != nil { + return false, fmt.Errorf("postpone quiet-hours notification: %w", err) + } + current.NextRetryAt = &deadline + current.Method = method + *notif = *current + return false, nil + } + + readyJSON, err := json.Marshal(ready) + if err != nil { + return false, fmt.Errorf("encode ready quiet-hours alerts: %w", err) + } + heldJSON, err := json.Marshal(held) + if err != nil { + return false, fmt.Errorf("encode held quiet-hours alerts: %w", err) + } + // Include the parent identity and exact ready payload. Restarts and duplicate + // worker snapshots cannot generate a second child for the same partition. + sum := sha256.Sum256(append([]byte(current.ID+"\x00"), readyJSON...)) + childID := current.ID + ":quiet:" + hex.EncodeToString(sum[:8]) + readyLinks, heldLinks, err := partitionQuietHoursLinks(current.Links, ready, held, childID) + if err != nil { + return false, err + } + readyLinksJSON, err := json.Marshal(readyLinks) + if err != nil { + return false, fmt.Errorf("encode ready quiet-hours links: %w", err) + } + heldLinksJSON, err := json.Marshal(heldLinks) + if err != nil { + return false, fmt.Errorf("encode held quiet-hours links: %w", err) + } + + tx, err := nq.db.Begin() + if err != nil { + return false, fmt.Errorf("begin quiet-hours partition: %w", err) + } + defer tx.Rollback() + // Copy the admitted destination bytes and attempt budget, not the current + // destination configuration. Previous attempts remain in the original audit + // trail; splitting never resets the retry budget or fabricates a new audit. + _, err = tx.Exec(`INSERT INTO notification_queue + (id, type, method, status, alerts, operational_links, config, attempts, + max_attempts, last_attempt, last_error, created_at, next_retry_at, + completed_at, payload_bytes) + SELECT ?, type, method, status, ?, ?, config, attempts, max_attempts, + last_attempt, last_error, created_at, NULL, completed_at, payload_bytes + FROM notification_queue WHERE id = ? AND status = 'pending'`, childID, string(readyJSON), string(readyLinksJSON), current.ID) + if err != nil { + return false, fmt.Errorf("persist ready quiet-hours partition: %w", err) + } + _, err = tx.Exec(`UPDATE notification_queue + SET alerts = ?, operational_links = ?, next_retry_at = ?, method = ? + WHERE id = ? AND status = 'pending'`, string(heldJSON), string(heldLinksJSON), deadline.Unix(), method, current.ID) + if err != nil { + return false, fmt.Errorf("persist held quiet-hours partition: %w", err) + } + if err := tx.Commit(); err != nil { + return false, fmt.Errorf("commit quiet-hours partition: %w", err) + } + current.ID = childID + current.Alerts = ready + current.Links = readyLinks + current.NextRetryAt = nil + *notif = *current + return true, nil +} + +func partitionQuietHoursLinks( + links []operationaltrust.NotificationLink, + ready, held []*alerts.Alert, + childID string, +) ([]operationaltrust.NotificationLink, []operationaltrust.NotificationLink, error) { + keys := func(batch []*alerts.Alert) map[string]bool { + result := make(map[string]bool, len(batch)) + for _, alert := range batch { + if alert != nil && alert.OperationalRecord != nil && alert.LatestTransition != nil { + result[alert.OperationalRecord.ID+"\x00"+alert.LatestTransition.ID] = true + } + } + return result + } + readyKeys, heldKeys := keys(ready), keys(held) + var readyLinks, heldLinks []operationaltrust.NotificationLink + for _, link := range links { + key := link.OperationalRecordID + "\x00" + link.TransitionID + switch { + case readyKeys[key] && !heldKeys[key]: + link = link.Clone() + link.NotificationID = childID + readyLinks = append(readyLinks, link) + case heldKeys[key] && !readyKeys[key]: + heldLinks = append(heldLinks, link.Clone()) + default: + // Do not silently discard or misattribute a historical linkage. + return nil, nil, fmt.Errorf("quiet-hours partition has an ambiguous operational link") + } + } + return readyLinks, heldLinks, nil +} diff --git a/internal/notifications/quiet_hours_queue_test.go b/internal/notifications/quiet_hours_queue_test.go new file mode 100644 index 000000000..cc08e7607 --- /dev/null +++ b/internal/notifications/quiet_hours_queue_test.go @@ -0,0 +1,512 @@ +package notifications + +import ( + "database/sql" + "encoding/json" + "net/http" + "net/http/httptest" + "path/filepath" + "reflect" + "strings" + "sync" + "testing" + "time" + + "github.com/rcourtman/pulse-go-rewrite/internal/alerts" + "github.com/rcourtman/pulse-go-rewrite/internal/operationaltrust" +) + +// Use the real persistent schema and delivery processor, without a background +// ticker racing these controlled replay instants. The monitor constructor test +// covers installing this policy on the autonomous worker after saved-config load. +func newQuietReplayQueue(t *testing.T, dir string, now *time.Time) *NotificationQueue { + t.Helper() + db, err := sql.Open("sqlite", queueSQLiteDSN(filepath.Join(dir, "replay.db"))) + if err != nil { + t.Fatal(err) + } + db.SetMaxOpenConns(1) + q := &NotificationQueue{db: db, deliveryGates: make(map[string]*notificationDeliveryGate), + notifyChan: make(chan struct{}, 100), now: func() time.Time { return *now }} + if err := q.initSchema(); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = db.Close() }) + return q +} + +func quietReplaySchedule(start, end, zone string) alerts.QuietHours { + return alerts.QuietHours{Enabled: true, Start: start, End: end, Timezone: zone, + Days: map[string]bool{"sunday": true, "monday": true, "tuesday": true, + "wednesday": true, "thursday": true, "friday": true, "saturday": true}} +} + +func quietReplayAlert(id string, level alerts.AlertLevel, start time.Time) *alerts.Alert { + return &alerts.Alert{ID: id, Type: "cpu", Level: level, StartTime: start, + ResourceID: id, ResourceName: id, + OperationalRecord: &operationaltrust.OperationalRecord{ID: "record-" + id}, + LatestTransition: &operationaltrust.LifecycleTransition{ID: "transition-" + id, + To: operationaltrust.OperationalOpen, CauseKey: "cause-" + id}} +} + +type quietReplayReceipt struct { + Event string `json:"event"` + Alerts []*alerts.Alert `json:"alerts"` +} + +type quietReplaySink struct { + mu sync.Mutex + receipts []quietReplayReceipt +} + +func (s *quietReplaySink) all() []quietReplayReceipt { + s.mu.Lock() + defer s.mu.Unlock() + return append([]quietReplayReceipt(nil), s.receipts...) +} + +func newQuietReplayNotifier(t *testing.T, q *NotificationQueue, policyOwner *alerts.Manager) (*NotificationManager, WebhookConfig, *quietReplaySink) { + t.Helper() + sink := &quietReplaySink{} + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + defer r.Body.Close() + var receipt quietReplayReceipt + if err := json.NewDecoder(r.Body).Decode(&receipt); err != nil { + t.Errorf("decode local delivery: %v", err) + w.WriteHeader(http.StatusBadRequest) + return + } + sink.mu.Lock() + sink.receipts = append(sink.receipts, receipt) + sink.mu.Unlock() + w.WriteHeader(http.StatusOK) + })) + t.Cleanup(server.Close) + hook := WebhookConfig{ID: "admitted-ops", Enabled: true, Service: "generic", URL: server.URL} + n := &NotificationManager{enabled: true, queue: q, webhooks: []WebhookConfig{hook}, + webhookClient: server.Client(), webhookRateLimits: make(map[string]*webhookRateLimit), + lastNotified: make(map[string]notificationRecord), deliveryReceipts: make(map[string]struct{})} + if err := n.UpdateAllowedPrivateCIDRs("127.0.0.1/32"); err != nil { + t.Fatal(err) + } + n.SetQuietHoursPolicyProvider(policyOwner.QuietHoursNotificationPolicy) + q.processor = n.ProcessQueuedNotification + return n, hook, sink +} + +func loadQuietReplay(t *testing.T, q *NotificationQueue, id string) *QueuedNotification { + t.Helper() + n, err := q.scanNotification(q.db.QueryRow(`SELECT id, type, method, status, alerts, config, + attempts, max_attempts, last_attempt, last_error, created_at, next_retry_at, + completed_at, payload_bytes, operational_links FROM notification_queue WHERE id = ?`, id)) + if err != nil { + t.Fatal(err) + } + return n +} + +func quietReplayRows(t *testing.T, q *NotificationQueue) []*QueuedNotification { + t.Helper() + rows, err := q.db.Query(`SELECT id, type, method, status, alerts, config, + attempts, max_attempts, last_attempt, last_error, created_at, next_retry_at, + completed_at, payload_bytes, operational_links FROM notification_queue ORDER BY id`) + if err != nil { + t.Fatal(err) + } + defer rows.Close() + var result []*QueuedNotification + for rows.Next() { + n, err := q.scanNotification(rows) + if err != nil { + t.Fatal(err) + } + result = append(result, n) + } + if err := rows.Err(); err != nil { + t.Fatal(err) + } + return result +} + +func seedQuietReplay(t *testing.T, q *NotificationQueue, hook WebhookConfig, now time.Time, batch ...*alerts.Alert) *QueuedNotification { + t.Helper() + config, err := json.Marshal(hook) + if err != nil { + t.Fatal(err) + } + due := now.Add(-time.Hour) + n := &QueuedNotification{ID: "admitted-batch", Type: "webhook", DestinationID: "webhook:" + hook.ID, + Alerts: batch, Config: config, MaxAttempts: 3, CreatedAt: now.Add(-24 * time.Hour), NextRetryAt: &due} + if err := q.Enqueue(n); err != nil { + t.Fatal(err) + } + return n +} + +func setQuietReplaySchedule(m *alerts.Manager, quiet alerts.QuietHours) { + config := m.GetConfig() + config.Schedule.QuietHours = quiet + m.UpdateConfig(config) +} + +func TestQueuedQuietHoursRevalidatesCurrentSchedule(t *testing.T) { + for _, tc := range []struct { + name, now, start, end, zone string + }{ + {"continuous_midnight", "2026-10-02T00:00:00Z", "00:00", "23:59", "UTC"}, + {"continuous_offset_boundary", "2026-10-02T12:00:00Z", "12:00", "11:59", "UTC"}, + {"late_replay_next_night", "2026-10-02T22:30:00Z", "22:00", "06:00", "UTC"}, + {"edited_longer_window", "2026-10-02T07:00:00Z", "22:00", "08:00", "UTC"}, + {"repeated_end_minute", "2026-10-25T01:45:00Z", "01:30", "02:00", "Europe/London"}, + } { + t.Run(tc.name, func(t *testing.T) { + now, err := time.Parse(time.RFC3339, tc.now) + if err != nil { + t.Fatal(err) + } + owner := alerts.NewManagerWithDataDir(t.TempDir()) + defer owner.Stop() + setQuietReplaySchedule(owner, quietReplaySchedule(tc.start, tc.end, tc.zone)) + q := newQuietReplayQueue(t, t.TempDir(), &now) + _, hook, sink := newQuietReplayNotifier(t, q, owner) + alert := quietReplayAlert("warning", alerts.AlertLevelWarning, now.Add(-2*time.Hour)) + alert.Metadata = map[string]interface{}{alerts.MetadataQuietHoursReplayAt: now.Add(-time.Hour).Format(time.RFC3339)} + n := seedQuietReplay(t, q, hook, now, alert) + before := loadQuietReplay(t, q, n.ID) + stale := *before + q.processNotification(n) + after := loadQuietReplay(t, q, before.ID) + if got := sink.all(); len(got) != 0 { + t.Fatalf("QUIET_REPLAY_LEAK: held notification reached provider: %+v", got) + } + if after.Status != QueueStatusPending || after.Attempts != 0 || after.LastAttempt != nil || + after.NextRetryAt == nil || !after.NextRetryAt.After(now) { + t.Fatalf("hold counted as delivery or did not move replay: %+v", after) + } + if !reflect.DeepEqual(before.Alerts, after.Alerts) || string(before.Config) != string(after.Config) || + !reflect.DeepEqual(before.Links, after.Links) { + t.Fatal("quiet hold changed occurrence, destination or operational identity") + } + for range 3 { + copy := stale + q.processNotification(©) + } + if repeated := loadQuietReplay(t, q, before.ID); !reflect.DeepEqual(after, repeated) { + t.Fatal("stale worker snapshot consumed an attempt or moved a held row again") + } + logs, err := q.GetDeliveryLog(before.CreatedAt, 10) + if err != nil || len(logs) != 0 { + t.Fatalf("quiet hold manufactured delivery audit: %+v, %v", logs, err) + } + // A continuous schedule never gets an invented open minute. Once the + // operator disables it, the next admitted wake reuses the same row. + setQuietReplaySchedule(owner, alerts.QuietHours{}) + now = after.NextRetryAt.Add(time.Second) + q.processNotification(&stale) + if got := sink.all(); len(got) != 1 || len(got[0].Alerts) != 1 || got[0].Alerts[0].ID != alert.ID || + !got[0].Alerts[0].StartTime.Equal(alert.StartTime) { + t.Fatalf("eligible replay did not preserve occurrence: %+v", got) + } + if sent := loadQuietReplay(t, q, before.ID); sent.Status != QueueStatusSent || sent.Attempts != 1 { + t.Fatalf("eligible replay attempt: %+v", sent) + } + }) + } +} + +func TestQueuedQuietHoursMixedBatchPreservesDeliveryAndBudget(t *testing.T) { + now := time.Date(2026, 10, 2, 1, 0, 0, 0, time.UTC) + owner := alerts.NewManagerWithDataDir(t.TempDir()) + defer owner.Stop() + setQuietReplaySchedule(owner, quietReplaySchedule("22:00", "06:00", "UTC")) + q := newQuietReplayQueue(t, t.TempDir(), &now) + notifier, hook, sink := newQuietReplayNotifier(t, q, owner) + warning := quietReplayAlert("warning", alerts.AlertLevelWarning, now.Add(-time.Hour)) + critical := quietReplayAlert("critical", alerts.AlertLevelCritical, now.Add(-time.Minute)) + n := seedQuietReplay(t, q, hook, now, warning, critical) + for range 2 { + if err := q.IncrementAttemptAndSetStatus(n.ID, QueueStatusSending); err != nil { + t.Fatal(err) + } + } + if err := q.ScheduleRetry(n.ID, 2); err != nil { + t.Fatal(err) + } + if _, err := q.db.Exec(`UPDATE notification_queue SET next_retry_at = ?, last_error = 'destination unavailable' WHERE id = ?`, now.Add(-time.Second).Unix(), n.ID); err != nil { + t.Fatal(err) + } + before := loadQuietReplay(t, q, n.ID) + if err := q.RecordAudit(before, false, "destination unavailable"); err != nil { + t.Fatal(err) + } + originalID := n.ID + q.processNotification(n) + got := sink.all() + if len(got) != 1 || len(got[0].Alerts) != 1 || got[0].Alerts[0].ID != critical.ID { + t.Fatalf("QUIET_MIXED_BATCH: immediately eligible subset was delayed or held member leaked: %+v", got) + } + rows := quietReplayRows(t, q) + if len(rows) != 2 { + t.Fatalf("mixed batch rows=%d, want held original plus ready child", len(rows)) + } + held := loadQuietReplay(t, q, originalID) + if len(held.Alerts) != 1 || held.Alerts[0].ID != warning.ID || held.Status != QueueStatusPending || + held.Attempts != 2 || held.MaxAttempts != 3 || !reflect.DeepEqual(held.LastAttempt, before.LastAttempt) || + !reflect.DeepEqual(held.LastError, before.LastError) { + t.Fatalf("held partition lost its admission/attempt history: %+v", held) + } + if held.NextRetryAt == nil || !held.NextRetryAt.Equal(time.Date(2026, 10, 2, 6, 1, 0, 0, time.UTC)) { + t.Fatalf("held replay deadline=%v, want inclusive end at 06:01", held.NextRetryAt) + } + child := loadQuietReplay(t, q, n.ID) + if child.ID == originalID || child.Attempts != 3 || child.Status != QueueStatusSent || child.MaxAttempts != 3 { + t.Fatalf("ready partition reset budget or did not complete: %+v", child) + } + for _, row := range []*QueuedNotification{held, child} { + if string(row.Config) != string(before.Config) || row.DestinationID != before.DestinationID || !row.CreatedAt.Equal(before.CreatedAt) || len(row.Links) != 1 { + t.Fatalf("partition lost admitted destination, time or links: %+v", row) + } + link := row.Links[0] + if err := link.Validate(); err != nil || link.NotificationID != row.ID || + link.OperationalRecordID != row.Alerts[0].OperationalRecord.ID || + link.TransitionID != row.Alerts[0].LatestTransition.ID || link.DestinationID != before.DestinationID { + t.Fatalf("partition misattributed operational receipt: %+v, %v", link, err) + } + } + if held.Links[0].DeliveryState != operationaltrust.NotificationRetrying || child.Links[0].DeliveryState != operationaltrust.NotificationDelivered { + t.Fatal("held link gained a delivery or lost its existing retry evidence") + } + jobs := buildNotificationDeliveryJobs(EmailConfig{}, []WebhookConfig{hook}, AppriseConfig{}, []*alerts.Alert{warning, critical}, eventResolved, now) + filtered := notifier.filterResolvedJobsByDeliveryReceipt(jobs) + if len(filtered) != 1 || len(filtered[0].Alerts) != 1 || filtered[0].Alerts[0].ID != critical.ID { + t.Fatalf("held member gained a firing receipt: %+v", filtered) + } + // Release at its real schedule boundary: only the held occurrence gets its + // third provider attempt, with no repetition of the already sent child. + now = *held.NextRetryAt + q.processNotification(held) + if got := sink.all(); len(got) != 2 || got[1].Alerts[0].ID != warning.ID || len(got[1].Alerts) != 1 { + t.Fatalf("released partition delivery: %+v", got) + } + if done := loadQuietReplay(t, q, originalID); done.Attempts != 3 || done.Status != QueueStatusSent { + t.Fatalf("held retry budget after release: %+v", done) + } +} + +func TestQueuedQuietHoursPartitionRollbackAndRestart(t *testing.T) { + for _, failure := range []string{"ready_insert", "held_rewrite"} { + t.Run(failure, func(t *testing.T) { + now := time.Date(2026, 10, 2, 1, 0, 0, 0, time.UTC) + owner := alerts.NewManagerWithDataDir(t.TempDir()) + defer owner.Stop() + setQuietReplaySchedule(owner, quietReplaySchedule("22:00", "06:00", "UTC")) + dir := t.TempDir() + q := newQuietReplayQueue(t, dir, &now) + _, hook, sink := newQuietReplayNotifier(t, q, owner) + n := seedQuietReplay(t, q, hook, now, + quietReplayAlert("warning", alerts.AlertLevelWarning, now.Add(-time.Hour)), + quietReplayAlert("critical", alerts.AlertLevelCritical, now.Add(-time.Minute))) + before := loadQuietReplay(t, q, n.ID) + trigger := `CREATE TRIGGER reject_partition BEFORE INSERT ON notification_queue WHEN NEW.id != 'admitted-batch' BEGIN SELECT RAISE(ABORT, 'fixture partition persistence failed'); END` + if failure == "held_rewrite" { + trigger = `CREATE TRIGGER reject_partition BEFORE UPDATE OF alerts ON notification_queue BEGIN SELECT RAISE(ABORT, 'fixture partition persistence failed'); END` + } + if _, err := q.db.Exec(trigger); err != nil { + t.Fatal(err) + } + q.processNotification(n) + if got := sink.all(); len(got) != 0 { + t.Fatalf("QUIET_PARTITION_FAILURE: persistence failure still attempted delivery: %+v", got) + } + if rows := quietReplayRows(t, q); len(rows) != 1 || !reflect.DeepEqual(rows[0], before) { + t.Fatalf("failed transaction did not retain original complete batch: %+v", rows) + } + if _, err := q.db.Exec("DROP TRIGGER reject_partition"); err != nil { + t.Fatal(err) + } + q.processNotification(n) + if got := sink.all(); len(got) != 1 || len(got[0].Alerts) != 1 || got[0].Alerts[0].ID != "critical" { + t.Fatalf("successful partition: %+v", got) + } + if err := q.db.Close(); err != nil { + t.Fatal(err) + } + q = newQuietReplayQueue(t, dir, &now) + _, _, restartedSink := newQuietReplayNotifier(t, q, owner) + rows := quietReplayRows(t, q) + if len(rows) != 2 { + t.Fatalf("reopen lost partition: %+v", rows) + } + held := loadQuietReplay(t, q, before.ID) + if count, err := q.CancelByAlertIdentifiers([]string{"warning"}); err != nil || count != 1 { + t.Fatalf("cancellation of held member: %d, %v", count, err) + } + now = held.NextRetryAt.Add(time.Second) + q.processNotification(held) + if got := restartedSink.all(); len(got) != 0 { + t.Fatalf("cancelled held member resurrected after reopen: %+v", got) + } + if cancelled := loadQuietReplay(t, q, before.ID); cancelled.Status != QueueStatusCancelled || cancelled.Attempts != 0 { + t.Fatalf("cancelled member acquired attempt: %+v", cancelled) + } + }) + } +} + +func TestQueuedQuietHoursConcurrentSnapshotsAndCancellation(t *testing.T) { + now := time.Date(2026, 10, 2, 1, 0, 0, 0, time.UTC) + owner := alerts.NewManagerWithDataDir(t.TempDir()) + defer owner.Stop() + setQuietReplaySchedule(owner, quietReplaySchedule("22:00", "06:00", "UTC")) + q := newQuietReplayQueue(t, t.TempDir(), &now) + _, hook, sink := newQuietReplayNotifier(t, q, owner) + n := seedQuietReplay(t, q, hook, now, + quietReplayAlert("cancelled", alerts.AlertLevelWarning, now.Add(-time.Hour)), + quietReplayAlert("held", alerts.AlertLevelWarning, now.Add(-time.Hour)), + quietReplayAlert("ready", alerts.AlertLevelCritical, now.Add(-time.Minute))) + snapshotTaken := make(chan struct{}) + continueDelivery := make(chan struct{}) + q.quietHoursPolicy = func() func(*alerts.Alert, time.Time) *time.Time { + policy := owner.QuietHoursNotificationPolicy() + close(snapshotTaken) + <-continueDelivery + return policy + } + done := make(chan struct{}) + go func() { q.processNotification(n); close(done) }() + <-snapshotTaken + if count, err := q.CancelByAlertIdentifiers([]string{"cancelled"}); err != nil || count != 1 { + t.Fatalf("cancellation between snapshot and claim: %d, %v", count, err) + } + close(continueDelivery) + <-done + q.quietHoursPolicy = owner.QuietHoursNotificationPolicy + held := loadQuietReplay(t, q, "admitted-batch") + var wg sync.WaitGroup + for range 12 { + wg.Add(1) + go func() { + defer wg.Done() + copy := *held + q.processNotification(©) + }() + } + wg.Wait() + if got := sink.all(); len(got) != 1 || len(got[0].Alerts) != 1 || got[0].Alerts[0].ID != "ready" { + t.Fatalf("QUIET_CANCELLATION: stale snapshots lost cancellation or duplicated delivery: %+v", got) + } + if len(held.Alerts) != 1 || held.Alerts[0].ID != "held" || held.Attempts != 0 || len(held.Links) != 1 || held.Links[0].TransitionID != "transition-held" { + t.Fatalf("held remainder resurrected removed member or consumed budget: %+v", held) + } + if count, err := q.CancelByAlertIdentifiers([]string{"held"}); err != nil || count != 1 { + t.Fatalf("post-partition cancellation: %d, %v", count, err) + } + now = held.NextRetryAt.Add(time.Second) + q.processNotification(held) + if got := sink.all(); len(got) != 1 { + t.Fatalf("post-partition cancellation leaked delivery: %+v", got) + } +} + +func TestQueuedQuietHoursResolvedKeepsUndeliveredReceipt(t *testing.T) { + now := time.Date(2026, 10, 2, 1, 0, 0, 0, time.UTC) + owner := alerts.NewManagerWithDataDir(t.TempDir()) + defer owner.Stop() + setQuietReplaySchedule(owner, quietReplaySchedule("22:00", "06:00", "UTC")) + q := newQuietReplayQueue(t, t.TempDir(), &now) + notifier, hook, sink := newQuietReplayNotifier(t, q, owner) + warning := quietReplayAlert("warning", alerts.AlertLevelWarning, now.Add(-2*time.Hour)) + critical := quietReplayAlert("critical", alerts.AlertLevelCritical, now.Add(-time.Hour)) + batch := []*alerts.Alert{warning, critical} + firing := notificationDeliveryJob{Type: "webhook", Event: eventAlert, Alerts: batch, WebhookConfig: &hook} + notifier.recordSuccessfulDelivery(firing, now.Add(-time.Hour)) + for _, alert := range batch { + alert.LatestTransition.To = operationaltrust.OperationalResolved + annotateResolvedMetadata(alert, now.Add(-time.Minute)) + } + n := seedQuietReplay(t, q, hook, now, batch...) + if _, err := q.db.Exec(`UPDATE notification_queue SET type = 'webhook_resolved', method = 'resolved-group' WHERE id = ?`, n.ID); err != nil { + t.Fatal(err) + } + q.processNotification(n) + if got := sink.all(); len(got) != 1 || got[0].Event != "resolved" || len(got[0].Alerts) != 1 || got[0].Alerts[0].ID != critical.ID { + t.Fatalf("QUIET_RECOVERY: grouped recovery leaked held member or lost ready member: %+v", got) + } + held := loadQuietReplay(t, q, "admitted-batch") + if held.Method != "quiet-hours-replay" || held.Attempts != 0 { + t.Fatalf("postponed recovery reopened grouping window or consumed attempt: %+v", held) + } + resolved := notificationDeliveryJob{Type: "webhook", Event: eventResolved, Alerts: batch, WebhookConfig: &hook} + if eligible := notifier.filterResolvedJobsByDeliveryReceipt([]notificationDeliveryJob{resolved}); len(eligible) != 1 || len(eligible[0].Alerts) != 1 || eligible[0].Alerts[0].ID != warning.ID { + t.Fatalf("held recovery lost the original occurrence receipt: %+v", eligible) + } + if count, err := q.CancelByAlertIdentifiers([]string{warning.ID}); err != nil || count != 0 { + t.Fatalf("firing cancellation consumed admitted recovery: %d, %v", count, err) + } + now = *held.NextRetryAt + q.processNotification(held) + if got := sink.all(); len(got) != 2 || got[1].Event != "resolved" || len(got[1].Alerts) != 1 || got[1].Alerts[0].ID != warning.ID { + t.Fatalf("held recovery did not complete once: %+v", got) + } + if eligible := notifier.filterResolvedJobsByDeliveryReceipt([]notificationDeliveryJob{resolved}); len(eligible) != 0 { + t.Fatalf("delivered recovery did not consume original receipts: %+v", eligible) + } +} + +func TestQueuedQuietHoursProviderFamiliesAndLegacyRows(t *testing.T) { + for _, kind := range []string{"email", "webhook", "apprise", "email_resolved", "webhook_resolved", "apprise_resolved"} { + t.Run(kind, func(t *testing.T) { + now := time.Date(2026, 10, 2, 0, 0, 0, 0, time.UTC) + owner := alerts.NewManagerWithDataDir(t.TempDir()) + defer owner.Stop() + setQuietReplaySchedule(owner, quietReplaySchedule("00:00", "23:59", "UTC")) + q := newQuietReplayQueue(t, t.TempDir(), &now) + notifier, hook, sink := newQuietReplayNotifier(t, q, owner) + notifier.emailConfig.Enabled = true + notifier.appriseConfig.Enabled = true + n := seedQuietReplay(t, q, hook, now, quietReplayAlert("legacy-warning", alerts.AlertLevelWarning, now.Add(-time.Hour))) + // Legacy payloads have no operational links or quiet replay metadata; + // they are still governed by the current policy, not a schema marker. + if _, err := q.db.Exec(`UPDATE notification_queue SET type = ?, config = '{}', operational_links = '[]' WHERE id = ?`, kind, n.ID); err != nil { + t.Fatal(err) + } + q.processNotification(n) + after := loadQuietReplay(t, q, n.ID) + if after.Status != QueueStatusPending || after.Attempts != 0 || after.NextRetryAt == nil || !after.NextRetryAt.After(now) { + t.Fatalf("QUIET_PROVIDER_FAMILY: %s bypassed quiet revalidation: %+v", kind, after) + } + if got := sink.all(); len(got) != 0 { + t.Fatalf("legacy quiet row leaked: %+v", got) + } + }) + } +} + +func TestQueuedQuietHoursRejectsAmbiguousLinks(t *testing.T) { + now := time.Date(2026, 10, 2, 1, 0, 0, 0, time.UTC) + owner := alerts.NewManagerWithDataDir(t.TempDir()) + defer owner.Stop() + setQuietReplaySchedule(owner, quietReplaySchedule("22:00", "06:00", "UTC")) + q := newQuietReplayQueue(t, t.TempDir(), &now) + _, hook, sink := newQuietReplayNotifier(t, q, owner) + n := seedQuietReplay(t, q, hook, now, + quietReplayAlert("held", alerts.AlertLevelWarning, now.Add(-time.Hour)), + quietReplayAlert("ready", alerts.AlertLevelCritical, now.Add(-time.Minute))) + links := append([]operationaltrust.NotificationLink(nil), n.Links...) + links[0].OperationalRecordID = "orphaned-record" + encoded, err := json.Marshal(links) + if err != nil { + t.Fatal(err) + } + if _, err := q.db.Exec(`UPDATE notification_queue SET operational_links = ? WHERE id = ?`, string(encoded), n.ID); err != nil { + t.Fatal(err) + } + before := loadQuietReplay(t, q, n.ID) + ready, err := q.prepareQuietHoursDelivery(n, owner.QuietHoursNotificationPolicy()) + if ready || err == nil || !strings.Contains(err.Error(), "ambiguous operational link") { + t.Fatalf("orphaned operational link was silently discarded: ready=%t err=%v", ready, err) + } + if got := sink.all(); len(got) != 0 || !reflect.DeepEqual(before, loadQuietReplay(t, q, before.ID)) { + t.Fatal("rejected linkage still changed batch or attempted delivery") + } +}