mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 12:47:49 +00:00
fix(notifications): start queue processing after destination config
A persisted notification that was already due at startup could be terminally cancelled as 'delivery disabled'. NewNotificationManagerWithDataDir woke the queue worker during construction, before the owner had applied saved webhook/email/Apprise configuration, so the still-empty destination list looked like a removed destination. Install the queue processor without waking it. The worker ticker and later enqueues start delivery once configuration is present. This removes the order/state-dependent failure in TestResolvedGroupingOrdinaryRestart and drops a persisted grouped recovery in production. Also wait past the queue poll interval in the grouping contract test, add a regression test for the startup ordering, and update the notifications subsystem contract's processor-attachment semantics. Change-source: pulse-maintainer
This commit is contained in:
parent
50dc7f3e17
commit
f235fc37f2
4 changed files with 110 additions and 6 deletions
|
|
@ -428,9 +428,22 @@ occurrence.
|
|||
That same queue boundary also owns processor attachment semantics. The
|
||||
canonical queue may persist pending notifications before a delivery processor is
|
||||
configured, but it must not mark those entries sending, failed, or sent until a
|
||||
processor exists. When a processor is attached, the queue owner must wake the
|
||||
pending backlog through the same canonical batch path instead of relying on a
|
||||
separate direct-send shortcut or waiting for an unrelated timer tick.
|
||||
processor exists. An explicit processor attachment must wake the pending backlog
|
||||
through the same canonical batch path instead of relying on a separate
|
||||
direct-send shortcut or waiting for an unrelated timer tick.
|
||||
|
||||
Manager construction is the one attachment that must not wake the worker.
|
||||
Construction runs before the owner has applied saved destination configuration
|
||||
(webhooks, email, Apprise), so an already-due persisted job would be evaluated
|
||||
against a still-empty destination list and terminally cancelled as
|
||||
`ErrNotificationDeliverySkipped`. Construction records the processor without
|
||||
waking it; the worker's own ticker and any later enqueue begin delivery once
|
||||
configuration is present. This preserves the disabled-destination cancellation
|
||||
policy at processing time without letting startup order decide it.
|
||||
`TestRestartKeepsDueDeliveryPendingUntilConfigured` in
|
||||
`resolved_grouping_contract_test.go` seeds a due persisted delivery, constructs
|
||||
the manager before restoring its webhook, and proves the row stays `pending`
|
||||
until configuration is applied and then delivers.
|
||||
Alert delivery cooldown is also owned at this boundary. Normal alert delivery
|
||||
must suppress duplicate sends for the same active alert occurrence when
|
||||
cooldown is disabled or still active. A manager-admitted increase above the
|
||||
|
|
|
|||
|
|
@ -755,9 +755,14 @@ func NewNotificationManagerWithDataDir(publicURL string, dataDir string) *Notifi
|
|||
// Create webhook client after NotificationManager is initialized
|
||||
nm.webhookClient = nm.createSecureWebhookClient(WebhookTimeout)
|
||||
|
||||
// Wire up queue processor if queue is available
|
||||
// Wire up the queue processor without waking the worker. Construction runs
|
||||
// before the owner applies saved destination configuration (webhooks,
|
||||
// email, Apprise); waking here would let an already-due persisted job be
|
||||
// processed while the destination list is still empty and be terminally
|
||||
// cancelled as "delivery disabled". The worker's ticker and later enqueues
|
||||
// begin delivery once configuration is present.
|
||||
if queue != nil {
|
||||
queue.SetProcessor(nm.ProcessQueuedNotification)
|
||||
queue.installProcessor(nm.ProcessQueuedNotification)
|
||||
}
|
||||
|
||||
// Start periodic cleanup of old lastNotified entries (every 1 hour)
|
||||
|
|
|
|||
|
|
@ -1682,6 +1682,19 @@ func (nq *NotificationQueue) SetProcessor(processor func(*QueuedNotification) er
|
|||
}
|
||||
}
|
||||
|
||||
// installProcessor records the notification processor without waking the worker.
|
||||
// Manager construction uses this instead of SetProcessor because construction
|
||||
// runs before the owner applies saved destination configuration (webhooks,
|
||||
// email, Apprise). Waking here would let an already-due persisted job be
|
||||
// processed while the destination list is still empty and be terminally
|
||||
// cancelled as "delivery disabled". The worker's ticker and any later enqueue
|
||||
// begin delivery once configuration is present.
|
||||
func (nq *NotificationQueue) installProcessor(processor func(*QueuedNotification) error) {
|
||||
nq.mu.Lock()
|
||||
nq.processor = processor
|
||||
nq.mu.Unlock()
|
||||
}
|
||||
|
||||
// processBatch processes a batch of pending notifications concurrently
|
||||
func (nq *NotificationQueue) processBatch() {
|
||||
const batchLimit = 20
|
||||
|
|
|
|||
|
|
@ -102,7 +102,9 @@ func TestResolvedGroupingOrdinaryRestart(t *testing.T) {
|
|||
var recovery payload
|
||||
select {
|
||||
case recovery = <-received:
|
||||
case <-time.After(5 * time.Second):
|
||||
case <-time.After(12 * time.Second):
|
||||
// The queue processes on its own ticker; the wait must exceed that
|
||||
// interval so a grouped recovery is not raced by the poll period.
|
||||
t.Fatal("no grouped recovery")
|
||||
}
|
||||
if recovery.Event != "resolved" || len(recovery.Alerts) != 15 {
|
||||
|
|
@ -230,6 +232,77 @@ func TestResolvedGroupingDisabledBurstRetainsRateLimit(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
// A restart must not terminally cancel a due persisted delivery before the
|
||||
// owner has applied saved destination configuration. Manager construction
|
||||
// happens before saved webhooks/email/Apprise are loaded in the monitor
|
||||
// startup path; the worker must not treat the still-empty destination list as
|
||||
// a permanent policy decision.
|
||||
func TestRestartKeepsDueDeliveryPendingUntilConfigured(t *testing.T) {
|
||||
received := make(chan struct{}, 4)
|
||||
server := newResolvedContractServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
received <- struct{}{}
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
dir := t.TempDir()
|
||||
hook := WebhookConfig{ID: "ops", Name: "ops", URL: server.URL, Enabled: true}
|
||||
configJSON, err := json.Marshal(hook)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
seed, err := NewNotificationQueue(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
due := time.Now().Add(-time.Second)
|
||||
if err := seed.Enqueue(&QueuedNotification{
|
||||
ID: "restart-due",
|
||||
Type: "webhook_resolved",
|
||||
Status: QueueStatusPending,
|
||||
Config: configJSON,
|
||||
Alerts: []*alerts.Alert{{ID: "restart-due", ResourceName: "restart-due", StartTime: time.Now().Add(-time.Minute)}},
|
||||
MaxAttempts: 3,
|
||||
NextRetryAt: &due,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := seed.Stop(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Production order: construct the manager, then apply saved configuration.
|
||||
m := NewNotificationManagerWithDataDir("", dir)
|
||||
defer m.Stop()
|
||||
m.webhookClient = server.Client()
|
||||
if err := m.UpdateAllowedPrivateCIDRs("127.0.0.1/32"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// A premature startup wake must not consume the due job before the
|
||||
// destination configuration is restored.
|
||||
time.Sleep(250 * time.Millisecond)
|
||||
queue := m.GetQueue()
|
||||
if queue == nil {
|
||||
t.Fatal("queue unavailable")
|
||||
}
|
||||
var status string
|
||||
if err := queue.db.QueryRow(`SELECT status FROM notification_queue WHERE id = 'restart-due'`).Scan(&status); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if status != string(QueueStatusPending) {
|
||||
t.Fatalf("due delivery status before configuration = %s, want pending", status)
|
||||
}
|
||||
|
||||
m.AddWebhook(hook)
|
||||
queue.processBatch()
|
||||
select {
|
||||
case <-received:
|
||||
case <-time.After(5 * time.Second):
|
||||
stats, _ := queue.GetQueueStats()
|
||||
t.Fatalf("configured due delivery not sent: %v", stats)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResolvedGroupingPagerDutyKeepsIndividualKeys(t *testing.T) {
|
||||
hook := WebhookConfig{ID: "pd", Service: "pagerduty", Enabled: true}
|
||||
batch := contractGroupedAlerts()
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue