42180948d1
Both background loops run unconditionally on every instance with no coordination between them, which the chart's replicas: 1 + strategy: Recreate exists specifically to paper over: with more than one replica, every one of them would sweep and deliver notifications independently, and two overlapping during a rollout would both page for the same incident. Add withAdvisoryLock, which takes a Postgres advisory lock on a dedicated connection and runs a pass only if it gets the lock, otherwise skipping until the next tick. Wire StartArchiver and StartNotifier through it with their own lock keys, so Sweep and NotifySweep themselves are untouched and every existing test calling them directly keeps working unchanged. This also closes the notifier's double-delivery race in passing: two replicas can no longer both be inside deliverPending at once, since only one can hold notifierLockKey at a time. Deliberately not addressed here, and still blocking a replica count above 1: the in-memory login rate limiter, the unlocked migration runner, and the new-incident-insert race on a webhook for a brand-new groupKey. Noted in the updated chart comment. Co-authored-by: Claude <noreply@anthropic.com>
264 lines
8.6 KiB
Go
264 lines
8.6 KiB
Go
package api
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"log"
|
|
|
|
"time"
|
|
)
|
|
|
|
const (
|
|
// sweepInterval is how often the background sweeper runs.
|
|
sweepInterval = 15 * time.Minute
|
|
|
|
// expiryGrace absorbs clock skew and notification latency before an alert
|
|
// whose ends_at watermark has passed is treated as stale.
|
|
expiryGrace = 5 * time.Minute
|
|
|
|
// archiverLockKey is the Postgres advisory lock the sweeper takes for the
|
|
// duration of each pass, so that running more than one replica does not run
|
|
// the sweep concurrently on all of them. Its value has no meaning beyond
|
|
// being distinct from notifierLockKey.
|
|
archiverLockKey int64 = 7265_0001
|
|
)
|
|
|
|
// StartArchiver runs the alert sweeper until ctx is cancelled, starting with an
|
|
// immediate pass so a restart reconciles state right away.
|
|
// archiveAfter and staleAfter are the values the server started with. They are
|
|
// the fallback, not the setting: each pass reads the current value from the
|
|
// settings table, so an administrator's change takes effect on the next tick
|
|
// instead of at the next restart.
|
|
//
|
|
// Each pass runs under archiverLockKey (see withAdvisoryLock), so that on more
|
|
// than one replica only whichever instance's tick takes the lock first actually
|
|
// sweeps; the rest skip that tick rather than racing the same pass.
|
|
func StartArchiver(ctx context.Context, db *sql.DB, archiveAfter, staleAfter time.Duration, notify NotifyConfig) {
|
|
ticker := time.NewTicker(sweepInterval)
|
|
defer ticker.Stop()
|
|
|
|
sweep := func() {
|
|
withAdvisoryLock(ctx, db, archiverLockKey, "sweeper", func() {
|
|
Sweep(ctx, db, archiveAfter, staleAfter, notify)
|
|
})
|
|
}
|
|
|
|
sweep()
|
|
for {
|
|
select {
|
|
case <-ticker.C:
|
|
sweep()
|
|
case <-ctx.Done():
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
// Sweep runs a single pass, in dependency order: reconcile the dead man's
|
|
// switches, expire stale firing alerts, close the incidents that leaves with
|
|
// nothing firing, then archive whatever has been settled long enough. Running
|
|
// them in one pass means an alert can go stale and its incident can close and
|
|
// archive without waiting three ticks.
|
|
//
|
|
// The switches go first because they hand expireStale the alerts it must not
|
|
// touch: a heartbeat answers to its own, much tighter, timeout, and the generic
|
|
// staleness rules would otherwise resolve it as 'expiry' long before that.
|
|
// Exported so tests can drive a pass without waiting on the ticker.
|
|
func Sweep(ctx context.Context, db *sql.DB, archiveAfter, staleAfter time.Duration, notify NotifyConfig) {
|
|
settings := NewSettings(db)
|
|
staleAfter = settings.Duration(ctx, SettingStaleAfter, staleAfter)
|
|
archiveAfter = settings.Duration(ctx, SettingArchiveAfter, archiveAfter)
|
|
|
|
heartbeats := sweepDeadman(ctx, db, notify)
|
|
expireStale(ctx, db, staleAfter, heartbeats)
|
|
resolveSettledIncidents(ctx, db)
|
|
archiveResolved(ctx, db, archiveAfter)
|
|
archiveResolvedIncidents(ctx, db, archiveAfter)
|
|
purgeAckTokens(ctx, db)
|
|
purgeSessions(ctx, db)
|
|
}
|
|
|
|
// expireStale resolves firing alerts that Alertmanager has stopped refreshing.
|
|
//
|
|
// A resolved webhook is otherwise the only way out of the firing state, so a
|
|
// notification that is dropped, silenced, or lost to a restart would pin the
|
|
// alert as firing forever. Two independent signals mark an alert stale:
|
|
//
|
|
// - ends_at, the "valid until" watermark Alertmanager sets on outgoing firing
|
|
// notifications, has passed (plus expiryGrace for clock skew). Absent on
|
|
// rows whose payload carried no ends_at, hence the second signal.
|
|
// - received_at is older than staleAfter. Alertmanager re-sends firing
|
|
// notifications every repeat_interval, making received_at a liveness
|
|
// heartbeat — provided staleAfter exceeds that interval.
|
|
//
|
|
// Alerts in skip are left alone: they are dead man's switch heartbeats, whose
|
|
// liveness sweepDeadman has already judged against a timeout of its own.
|
|
//
|
|
// The matching rows are collected before the update rather than updated in bulk,
|
|
// because each one owes its incident a timeline entry.
|
|
func expireStale(ctx context.Context, db *sql.DB, staleAfter time.Duration, skip map[int64]bool) {
|
|
now := time.Now()
|
|
|
|
found, err := staleAlertIDs(ctx, db, now, staleAfter)
|
|
if err != nil {
|
|
log.Printf("sweeper: find stale: %v", err)
|
|
return
|
|
}
|
|
|
|
ids := make([]int64, 0, len(found))
|
|
for _, id := range found {
|
|
if !skip[id] {
|
|
ids = append(ids, id)
|
|
}
|
|
}
|
|
if len(ids) == 0 {
|
|
return
|
|
}
|
|
|
|
args := &sqlArgs{}
|
|
source := args.add(resolutionExpiry)
|
|
idList := make([]any, len(ids))
|
|
for i, id := range ids {
|
|
idList[i] = id
|
|
}
|
|
if _, err := db.ExecContext(ctx, `
|
|
UPDATE alerts
|
|
SET status = 'resolved',
|
|
resolution_source = `+source+`,
|
|
ends_at = COALESCE(ends_at, `+nowEpoch+`)
|
|
WHERE id IN (`+args.addList(idList)+`)`, args.all()...); err != nil {
|
|
log.Printf("sweeper: expire stale: %v", err)
|
|
return
|
|
}
|
|
log.Printf("sweeper: expired %d stale firing alert(s)", len(ids))
|
|
|
|
for _, id := range ids {
|
|
incidentID, err := openIncidentForAlert(ctx, db, id)
|
|
if err != nil {
|
|
log.Printf("sweeper: incident for alert %d: %v", id, err)
|
|
continue
|
|
}
|
|
if incidentID == 0 {
|
|
continue
|
|
}
|
|
alertID := id
|
|
if err := logEvent(ctx, db, incidentID, evAlertResolved, nil, &alertID, nil); err != nil {
|
|
log.Printf("sweeper: log expiry event: %v", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// staleAlertIDs reads the ids in one go and closes the cursor before the caller
|
|
// writes. Under SQLite's single connection an open read would have blocked the
|
|
// update outright; with a pool it is no longer a deadlock, but reading the set
|
|
// first still keeps the write off a cursor the same transaction is walking.
|
|
func staleAlertIDs(ctx context.Context, db *sql.DB, now time.Time, staleAfter time.Duration) ([]int64, error) {
|
|
rows, err := db.QueryContext(ctx, `
|
|
SELECT id FROM alerts
|
|
WHERE status = 'firing'
|
|
AND archived_at IS NULL
|
|
AND ((ends_at IS NOT NULL AND ends_at < $1) OR received_at < $2)`,
|
|
now.Add(-expiryGrace).Unix(), now.Add(-staleAfter).Unix())
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var ids []int64
|
|
for rows.Next() {
|
|
var id int64
|
|
if err := rows.Scan(&id); err != nil {
|
|
return nil, err
|
|
}
|
|
ids = append(ids, id)
|
|
}
|
|
return ids, rows.Err()
|
|
}
|
|
|
|
// resolveSettledIncidents closes incidents whose alerts have all stopped firing.
|
|
// This is the cascade from alerts up to the work item, and it is what turns an
|
|
// expiry into a closed incident rather than one that sits open forever.
|
|
func resolveSettledIncidents(ctx context.Context, db *sql.DB) {
|
|
ids, err := settledIncidentIDs(ctx, db)
|
|
if err != nil {
|
|
log.Printf("sweeper: find settled incidents: %v", err)
|
|
return
|
|
}
|
|
|
|
resolved := 0
|
|
for _, id := range ids {
|
|
ok, err := resolveIfSettled(ctx, db, id)
|
|
if err != nil {
|
|
log.Printf("sweeper: resolve incident %d: %v", id, err)
|
|
continue
|
|
}
|
|
if ok {
|
|
resolved++
|
|
}
|
|
}
|
|
if resolved > 0 {
|
|
log.Printf("sweeper: resolved %d settled incident(s)", resolved)
|
|
}
|
|
}
|
|
|
|
func settledIncidentIDs(ctx context.Context, db *sql.DB) ([]int64, error) {
|
|
rows, err := db.QueryContext(ctx, `
|
|
SELECT i.id
|
|
FROM incidents i
|
|
WHERE i.resolved_at IS NULL
|
|
AND EXISTS (SELECT 1 FROM incident_alerts ia WHERE ia.incident_id = i.id)
|
|
AND NOT EXISTS (SELECT 1
|
|
FROM incident_alerts ia
|
|
JOIN alerts a ON a.id = ia.alert_id
|
|
WHERE ia.incident_id = i.id
|
|
AND a.status = 'firing')`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var ids []int64
|
|
for rows.Next() {
|
|
var id int64
|
|
if err := rows.Scan(&id); err != nil {
|
|
return nil, err
|
|
}
|
|
ids = append(ids, id)
|
|
}
|
|
return ids, rows.Err()
|
|
}
|
|
|
|
// archiveResolved hides resolved alerts that have been settled for archiveAfter.
|
|
func archiveResolved(ctx context.Context, db *sql.DB, archiveAfter time.Duration) {
|
|
cutoff := time.Now().Add(-archiveAfter).Unix()
|
|
res, err := db.ExecContext(ctx,
|
|
`UPDATE alerts SET archived_at = `+nowEpoch+`
|
|
WHERE status = 'resolved'
|
|
AND archived_at IS NULL
|
|
AND COALESCE(ends_at, received_at) < $1`, cutoff)
|
|
if err != nil {
|
|
log.Printf("archiver: %v", err)
|
|
return
|
|
}
|
|
if n, _ := res.RowsAffected(); n > 0 {
|
|
log.Printf("archiver: archived %d resolved alert(s)", n)
|
|
}
|
|
}
|
|
|
|
// archiveResolvedIncidents does the same for the work items, on the same clock.
|
|
func archiveResolvedIncidents(ctx context.Context, db *sql.DB, archiveAfter time.Duration) {
|
|
cutoff := time.Now().Add(-archiveAfter).Unix()
|
|
res, err := db.ExecContext(ctx,
|
|
`UPDATE incidents SET archived_at = `+nowEpoch+`
|
|
WHERE resolved_at IS NOT NULL
|
|
AND archived_at IS NULL
|
|
AND resolved_at < $1`, cutoff)
|
|
if err != nil {
|
|
log.Printf("archiver: incidents: %v", err)
|
|
return
|
|
}
|
|
if n, _ := res.RowsAffected(); n > 0 {
|
|
log.Printf("archiver: archived %d resolved incident(s)", n)
|
|
}
|
|
}
|