Files
terdut-server/internal/api/archiver.go
T
Niklas Ye 74359c72ab
CI / chart (pull_request) Successful in 1s
CI / security (pull_request) Successful in 14s
CI / test (pull_request) Successful in 1m57s
Give each team its own dead man's switches, and the UI a team to show
The rest of #4. Two halves that belong together because they are the
same sentence from opposite ends: a team decides which of its alerts are
heartbeats, and the UI has to be able to say which team it is talking
about.

Switches were three environment variables, which made them one setting
for the whole install. That was the last piece of the alerting path a
team could not control: it could take its own alerts on its own key and
still not say which of them were heartbeats, or how long a silence had
to last. They are a row per team now, edited by an owner through
PUT /api/teams/{teamID}/deadman, and the sweeper runs each team against
its own matchers, timeout and severity.

The environment variables become the starting point rather than the
setting. Every team without a configuration is seeded from them at
startup, so an upgrade keeps watching exactly what it was watching, and
SeedDeadmanConfigs never overwrites -- a redeploy must not put the
environment's value back over an owner's edit. A team created later
watches nothing until somebody says otherwise: inheriting an
install-wide heartbeat would page a new team about a source it has never
heard of, and a switch nobody chose is the kind that gets muted rather
than fixed.

A matcher string with no alertname in it is refused at the door instead
of stored. Storing it would produce a switch that watches nothing
silently, which is the exact failure the feature exists to prevent.

NewRouter and Sweep lose their DeadmanConfig parameter -- there is no
longer one answer to hand them. The type stays, because parsing a
matcher string is still parsing a matcher string.

The UI side: rows in the queue carry a team badge, the filter row gains
a team chip per team, and "on call now" shows one card per team. All
three appear only when the viewer is in more than one team -- otherwise
they are the same word repeated down a list, which is noise rather than
information, and the single-team install reads exactly as it did before
teams existed.

Verified against a live two-team server as well as in tests: the
combined queue labelled by team, the team_id filter, a heartbeat that is
a heartbeat in one team and an ordinary alert in another, and a new
team's switches starting empty while the upgraded team keeps the
environment's.

Claude-Session: https://claude.ai/code/session_01RHPj4ggeFdEjKKfm4SHbD7
2026-09-20 15:18:22 +02:00

240 lines
7.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
)
// StartArchiver runs the alert sweeper until ctx is cancelled, starting with an
// immediate pass so a restart reconciles state right away.
func StartArchiver(ctx context.Context, db *sql.DB, archiveAfter, staleAfter time.Duration, notify NotifyConfig) {
ticker := time.NewTicker(sweepInterval)
defer ticker.Stop()
Sweep(ctx, db, archiveAfter, staleAfter, notify)
for {
select {
case <-ticker.C:
Sweep(ctx, db, archiveAfter, staleAfter, notify)
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) {
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)
}
}