diff --git a/charts/terdut-server/templates/deployment.yaml b/charts/terdut-server/templates/deployment.yaml index f7bdfc3..f394ff1 100644 --- a/charts/terdut-server/templates/deployment.yaml +++ b/charts/terdut-server/templates/deployment.yaml @@ -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: diff --git a/internal/api/advisory_lock_test.go b/internal/api/advisory_lock_test.go new file mode 100644 index 0000000..789b802 --- /dev/null +++ b/internal/api/advisory_lock_test.go @@ -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") + } +} diff --git a/internal/api/archiver.go b/internal/api/archiver.go index 606ec2c..d2bb24d 100644 --- a/internal/api/archiver.go +++ b/internal/api/archiver.go @@ -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 } diff --git a/internal/api/helpers.go b/internal/api/helpers.go index 48bfc70..b566a40 100644 --- a/internal/api/helpers.go +++ b/internal/api/helpers.go @@ -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() +} diff --git a/internal/api/notifier.go b/internal/api/notifier.go index 37d6f86..f34cb03 100644 --- a/internal/api/notifier.go +++ b/internal/api/notifier.go @@ -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 }