Advisory-lock DB migrations against concurrent replica startup
Migrate's check-then-apply loop against schema_migrations had no locking: two replicas booting at once against a fresh or partially-migrated database could both pass the "not yet applied" check for the same file and race applying it, crashing whichever lost the duplicate-key insert (confirmed: reverting the lock fails the new test 10/10 on a duplicate-key violation, racing as early as the CREATE TABLE IF NOT EXISTS schema_migrations statement itself). Hold a Postgres advisory lock for Migrate's whole run, on a dedicated connection reserved via db.Conn so lock and unlock happen on the same session. Blocking (pg_advisory_lock), unlike the archiver/notifier's pg_try_advisory_lock: on boot there's no later tick to defer to, so a second replica should wait for the first to finish migrating rather than skip ahead. Adds internal/db's first test file, exercising two concurrent Migrate calls against a fresh schema. Still open: the new-incident-insert race on a webhook for a brand-new groupKey, noted in the chart's updated comment. Login rate limiting staying in-process, diluted across replicas, is an accepted tradeoff. Co-authored-by: Claude <noreply@anthropic.com>
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"embed"
|
||||
"fmt"
|
||||
@@ -62,6 +63,16 @@ func Open(dsn string) (*sql.DB, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// migrationLockKey is the Postgres advisory lock Migrate holds for its whole
|
||||
// run. Two replicas starting at once would otherwise race the check-then-apply
|
||||
// loop below against schema_migrations: the loser could crash on a
|
||||
// duplicate-key insert, or contend with the winner's uncommitted DDL. Blocking
|
||||
// (pg_advisory_lock, not pg_try_advisory_lock as the archiver and notifier
|
||||
// use): on boot there is no later tick to defer to, so the right behaviour is
|
||||
// to wait for the other replica to finish migrating, not to skip ahead and
|
||||
// start serving against an unmigrated schema.
|
||||
const migrationLockKey int64 = 7265_0003
|
||||
|
||||
// Migrate applies every embedded migration that has not been applied yet, in
|
||||
// filename order, recording each in schema_migrations.
|
||||
//
|
||||
@@ -69,6 +80,22 @@ func Open(dsn string) (*sql.DB, error) {
|
||||
// migration that failed half way used to leave the schema in whatever state it
|
||||
// had reached. Postgres has transactional DDL, so the rollback is real.
|
||||
func Migrate(db *sql.DB) error {
|
||||
ctx := context.Background()
|
||||
conn, err := db.Conn(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("migrate: acquire connection: %w", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_lock($1)", migrationLockKey); err != nil {
|
||||
return fmt.Errorf("migrate: acquire advisory lock: %w", err)
|
||||
}
|
||||
defer func() {
|
||||
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", migrationLockKey); err != nil {
|
||||
log.Printf("migrate: release advisory lock: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS schema_migrations (
|
||||
version TEXT PRIMARY KEY,
|
||||
applied_at BIGINT NOT NULL DEFAULT FLOOR(EXTRACT(EPOCH FROM now()))::bigint
|
||||
|
||||
@@ -0,0 +1,125 @@
|
||||
package db_test
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"git.ryuvia.com/niklas/terdut-server/internal/db"
|
||||
|
||||
_ "github.com/jackc/pgx/v5/stdlib"
|
||||
)
|
||||
|
||||
// TERDUT_TEST_DSN must point at a database the test role may create schemas
|
||||
// in; see internal/api/testdb_test.go for the fuller rationale this mirrors.
|
||||
// An unset DSN fails rather than skips, deliberately.
|
||||
const testDSNEnv = "TERDUT_TEST_DSN"
|
||||
|
||||
// TestMigrate_ConcurrentCallersDoNotRace reproduces two replicas starting at
|
||||
// once against a brand-new, unmigrated schema: both call db.Migrate at the
|
||||
// same time. Before migrationLockKey, the loser could crash on a
|
||||
// duplicate-key insert into schema_migrations, or contend with the winner's
|
||||
// uncommitted DDL; with the advisory lock, one blocks until the other
|
||||
// finishes and both return cleanly.
|
||||
func TestMigrate_ConcurrentCallersDoNotRace(t *testing.T) {
|
||||
dsn := os.Getenv(testDSNEnv)
|
||||
if dsn == "" {
|
||||
t.Fatalf("%s is not set: these tests need Postgres.\n"+
|
||||
"Run `make test-db` for a local one, then\n"+
|
||||
" export %s=postgres://terdut:terdut@localhost:5432/terdut_test?sslmode=disable",
|
||||
testDSNEnv, testDSNEnv)
|
||||
}
|
||||
|
||||
schema := fmt.Sprintf("migrate_race_%d", os.Getpid())
|
||||
admin, err := sql.Open("pgx", dsn)
|
||||
if err != nil {
|
||||
t.Fatalf("connect to %s: %v", testDSNEnv, err)
|
||||
}
|
||||
defer admin.Close()
|
||||
if _, err := admin.Exec("CREATE SCHEMA " + schema); err != nil {
|
||||
t.Fatalf("create schema %s: %v", schema, err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
cleanup, err := sql.Open("pgx", dsn)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer cleanup.Close()
|
||||
if _, err := cleanup.Exec("DROP SCHEMA " + schema + " CASCADE"); err != nil {
|
||||
t.Logf("drop schema %s: %v", schema, err)
|
||||
}
|
||||
})
|
||||
|
||||
scoped := withSearchPath(dsn, schema)
|
||||
|
||||
const callers = 2
|
||||
errs := make([]error, callers)
|
||||
var wg sync.WaitGroup
|
||||
for i := range callers {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
database, err := db.Open(scoped)
|
||||
if err != nil {
|
||||
errs[i] = fmt.Errorf("open: %w", err)
|
||||
return
|
||||
}
|
||||
defer database.Close()
|
||||
errs[i] = db.Migrate(database)
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
for i, err := range errs {
|
||||
if err != nil {
|
||||
t.Fatalf("Migrate #%d: %v", i, err)
|
||||
}
|
||||
}
|
||||
|
||||
entries, err := os.ReadDir("migrations")
|
||||
if err != nil {
|
||||
t.Fatalf("read migrations dir: %v", err)
|
||||
}
|
||||
var want int
|
||||
for _, e := range entries {
|
||||
if !e.IsDir() && strings.HasSuffix(e.Name(), ".sql") {
|
||||
want++
|
||||
}
|
||||
}
|
||||
|
||||
check, err := sql.Open("pgx", scoped)
|
||||
if err != nil {
|
||||
t.Fatalf("connect for verification: %v", err)
|
||||
}
|
||||
defer check.Close()
|
||||
|
||||
var got int
|
||||
if err := check.QueryRow("SELECT COUNT(*) FROM schema_migrations").Scan(&got); err != nil {
|
||||
t.Fatalf("count schema_migrations: %v", err)
|
||||
}
|
||||
if got != want {
|
||||
t.Fatalf("schema_migrations has %d row(s) after two concurrent Migrate calls, want %d (one per migration file, no duplicates)", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// withSearchPath pins a DSN to one schema. Copied from
|
||||
// internal/api/testdb_test.go rather than shared: that helper lives in the
|
||||
// api_test package, a separate compiled package this one cannot import.
|
||||
func withSearchPath(dsn, schema string) string {
|
||||
opt := "-csearch_path=" + schema
|
||||
|
||||
if strings.HasPrefix(dsn, "postgres://") || strings.HasPrefix(dsn, "postgresql://") {
|
||||
u, err := url.Parse(dsn)
|
||||
if err == nil {
|
||||
q := u.Query()
|
||||
q.Set("options", opt)
|
||||
u.RawQuery = q.Encode()
|
||||
return u.String()
|
||||
}
|
||||
}
|
||||
return dsn + " options='" + opt + "'"
|
||||
}
|
||||
Reference in New Issue
Block a user