Guard the archiver and notifier passes with a Postgres advisory lock
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>
This commit is contained in:
@@ -11,8 +11,11 @@ spec:
|
||||
matchLabels:
|
||||
{{- include "terdut-server.selectorLabels" . | nindent 6 }}
|
||||
# Recreate, not RollingUpdate, even though the PVC that forced it is gone: the
|
||||
# sweeper and the notifier are unsynchronised singletons, and two replicas
|
||||
# overlapping during a rollout would both page for the same incident.
|
||||
# sweeper and notifier now take a Postgres advisory lock for each pass, so two
|
||||
# replicas overlapping during a rollout no longer both page for the same
|
||||
# incident, but DB migrations and new-incident creation on first webhook are
|
||||
# still unguarded — a second replica starting concurrently with the first can
|
||||
# still race either of those.
|
||||
strategy:
|
||||
type: Recreate
|
||||
template:
|
||||
|
||||
@@ -0,0 +1,122 @@
|
||||
package api
|
||||
|
||||
// This file is internal (package api, not api_test) because withAdvisoryLock is
|
||||
// unexported and these tests exercise its locking semantics directly rather than
|
||||
// through the full StartArchiver/StartNotifier loop, which would make the "does
|
||||
// not run while held" case timing-dependent instead of deterministic. It opens a
|
||||
// plain connection to TERDUT_TEST_DSN rather than reusing testdb_test.go's
|
||||
// newTestDB, since that helper lives in the separate, already-compiled
|
||||
// api_test package and a Postgres advisory lock needs no schema or migration
|
||||
// to exercise.
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
_ "github.com/jackc/pgx/v5/stdlib"
|
||||
)
|
||||
|
||||
// advisoryTestDB opens a plain, unmigrated connection to the test database. An
|
||||
// unset DSN fails rather than skips, matching testdb_test.go's rationale: a
|
||||
// suite that quietly tests nothing is worse than one that does not run.
|
||||
func advisoryTestDB(t *testing.T) *sql.DB {
|
||||
t.Helper()
|
||||
|
||||
dsn := os.Getenv("TERDUT_TEST_DSN")
|
||||
if dsn == "" {
|
||||
t.Fatalf("TERDUT_TEST_DSN is not set: these tests need Postgres.\n" +
|
||||
"Run `make test-db` for a local one, then\n" +
|
||||
" export TERDUT_TEST_DSN=postgres://terdut:terdut@localhost:5432/terdut_test?sslmode=disable")
|
||||
}
|
||||
|
||||
db, err := sql.Open("pgx", dsn)
|
||||
if err != nil {
|
||||
t.Fatalf("connect to TERDUT_TEST_DSN: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { db.Close() })
|
||||
return db
|
||||
}
|
||||
|
||||
func TestWithAdvisoryLock_RunsWhenFree(t *testing.T) {
|
||||
db := advisoryTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
ran := false
|
||||
withAdvisoryLock(ctx, db, archiverLockKey, "test", func() { ran = true })
|
||||
|
||||
if !ran {
|
||||
t.Fatal("fn did not run although the lock was free")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithAdvisoryLock_SkipsWhileHeldElsewhere(t *testing.T) {
|
||||
db := advisoryTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
// Hold the lock on a connection of our own, standing in for another
|
||||
// replica mid-pass.
|
||||
holder, err := db.Conn(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("acquire holder connection: %v", err)
|
||||
}
|
||||
defer holder.Close()
|
||||
if _, err := holder.ExecContext(ctx, "SELECT pg_advisory_lock($1)", archiverLockKey); err != nil {
|
||||
t.Fatalf("pre-acquire lock: %v", err)
|
||||
}
|
||||
|
||||
ran := false
|
||||
withAdvisoryLock(ctx, db, archiverLockKey, "test", func() { ran = true })
|
||||
if ran {
|
||||
t.Fatal("fn ran although another connection already held the lock")
|
||||
}
|
||||
|
||||
if _, err := holder.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", archiverLockKey); err != nil {
|
||||
t.Fatalf("release held lock: %v", err)
|
||||
}
|
||||
|
||||
// Now that the holder released it, the next caller should get it.
|
||||
ran = false
|
||||
withAdvisoryLock(ctx, db, archiverLockKey, "test", func() { ran = true })
|
||||
if !ran {
|
||||
t.Fatal("fn did not run after the other connection released the lock")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithAdvisoryLock_ReleasesAfterFnReturns(t *testing.T) {
|
||||
db := advisoryTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
withAdvisoryLock(ctx, db, notifierLockKey, "test", func() {})
|
||||
|
||||
// If the first call had leaked the lock, this one would see it held and
|
||||
// skip, leaving ran false.
|
||||
ran := false
|
||||
withAdvisoryLock(ctx, db, notifierLockKey, "test", func() { ran = true })
|
||||
if !ran {
|
||||
t.Fatal("fn did not run on a later call: the earlier call leaked its lock")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithAdvisoryLock_KeysAreIndependent(t *testing.T) {
|
||||
db := advisoryTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
holder, err := db.Conn(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("acquire holder connection: %v", err)
|
||||
}
|
||||
defer holder.Close()
|
||||
if _, err := holder.ExecContext(ctx, "SELECT pg_advisory_lock($1)", archiverLockKey); err != nil {
|
||||
t.Fatalf("pre-acquire archiver lock: %v", err)
|
||||
}
|
||||
defer holder.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", archiverLockKey)
|
||||
|
||||
// Holding archiverLockKey must not block notifierLockKey.
|
||||
ran := false
|
||||
withAdvisoryLock(ctx, db, notifierLockKey, "test", func() { ran = true })
|
||||
if !ran {
|
||||
t.Fatal("fn did not run under a different key although only archiverLockKey was held")
|
||||
}
|
||||
}
|
||||
@@ -15,6 +15,12 @@ const (
|
||||
// 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
|
||||
@@ -23,15 +29,25 @@ const (
|
||||
// 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(ctx, db, archiveAfter, staleAfter, notify)
|
||||
sweep := func() {
|
||||
withAdvisoryLock(ctx, db, archiverLockKey, "sweeper", func() {
|
||||
Sweep(ctx, db, archiveAfter, staleAfter, notify)
|
||||
})
|
||||
}
|
||||
|
||||
sweep()
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
Sweep(ctx, db, archiveAfter, staleAfter, notify)
|
||||
sweep()
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1,8 +1,11 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"log"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
@@ -77,3 +80,39 @@ func decodeJSON(r *http.Request, v any) error {
|
||||
func errResp(msg string) map[string]string {
|
||||
return map[string]string{"error": msg}
|
||||
}
|
||||
|
||||
// withAdvisoryLock runs fn only if it can take the named Postgres advisory lock on a
|
||||
// dedicated connection, and skips fn otherwise. This is what keeps the archiver and
|
||||
// notifier safe to run on more than one replica: whichever instance's tick gets there
|
||||
// first does the work; the rest see the lock held and simply wait for their next tick
|
||||
// instead of running the same pass concurrently.
|
||||
//
|
||||
// pg_try_advisory_lock is session-scoped, so taking and releasing it must happen on the
|
||||
// same connection, reserved via db.Conn rather than borrowed from the pool's shared
|
||||
// connections fn itself may use — and released (unlocked, then closed) before returning,
|
||||
// since a session lock otherwise outlives this call and leaks onto whatever reuses the
|
||||
// pooled connection next.
|
||||
func withAdvisoryLock(ctx context.Context, db *sql.DB, key int64, name string, fn func()) {
|
||||
conn, err := db.Conn(ctx)
|
||||
if err != nil {
|
||||
log.Printf("%s: advisory lock: acquire connection: %v", name, err)
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
var locked bool
|
||||
if err := conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", key).Scan(&locked); err != nil {
|
||||
log.Printf("%s: advisory lock: %v", name, err)
|
||||
return
|
||||
}
|
||||
if !locked {
|
||||
return // another replica is already running this pass
|
||||
}
|
||||
defer func() {
|
||||
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", key); err != nil {
|
||||
log.Printf("%s: advisory unlock: %v", name, err)
|
||||
}
|
||||
}()
|
||||
|
||||
fn()
|
||||
}
|
||||
|
||||
@@ -38,6 +38,13 @@ const (
|
||||
ackTokenTTL = 24 * time.Hour
|
||||
)
|
||||
|
||||
// notifierLockKey is the Postgres advisory lock the notifier takes for the
|
||||
// duration of each pass, so that running more than one replica does not
|
||||
// deliver (or double-deliver) the same notification from more than one of
|
||||
// them at once. Its value has no meaning beyond being distinct from
|
||||
// archiverLockKey.
|
||||
const notifierLockKey int64 = 7265_0002
|
||||
|
||||
// Notification kinds, recording why a push was sent.
|
||||
const (
|
||||
notifyTriggered = "triggered"
|
||||
@@ -97,6 +104,10 @@ var notifyClient = &http.Client{Timeout: 10 * time.Second}
|
||||
|
||||
// StartNotifier delivers queued notifications until ctx is cancelled, starting
|
||||
// with an immediate pass so a restart flushes whatever the last one left behind.
|
||||
//
|
||||
// Each pass runs under notifierLockKey (see withAdvisoryLock), so that on more
|
||||
// than one replica only whichever instance's tick takes the lock first actually
|
||||
// delivers; the rest skip that tick rather than racing the same pass.
|
||||
func StartNotifier(ctx context.Context, db *sql.DB, cfg NotifyConfig) {
|
||||
if !cfg.enabled() {
|
||||
log.Print("notifier: disabled (no ntfy URL configured)")
|
||||
@@ -107,11 +118,17 @@ func StartNotifier(ctx context.Context, db *sql.DB, cfg NotifyConfig) {
|
||||
ticker := time.NewTicker(notifyInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
NotifySweep(ctx, db, cfg)
|
||||
sweep := func() {
|
||||
withAdvisoryLock(ctx, db, notifierLockKey, "notifier", func() {
|
||||
NotifySweep(ctx, db, cfg)
|
||||
})
|
||||
}
|
||||
|
||||
sweep()
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
NotifySweep(ctx, db, cfg)
|
||||
sweep()
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user