Revalidate queued alert delivery against current quiet hours

Keep continuous and late replay held without spending a provider attempt. Split mixed batches atomically so eligible alerts retain admitted destinations, occurrence links and retry budgets, and bind saved monitor policy before activating persisted work. Add local receipt, cancellation, rollback, reconstruction and race regression coverage with the owning contracts and verification routes.

Change-source: pulse-maintainer
This commit is contained in:
pulse-triage[bot] 2026-10-01 23:05:07 +01:00
parent 680ac756e6
commit 5213b794b2
13 changed files with 1009 additions and 14 deletions

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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(&notifications.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.

View file

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

View file

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

View file

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

View file

@ -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(&copy)
}
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(&copy)
}()
}
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")
}
}