diff --git a/cmd/felis/api.go b/cmd/felis/api.go index 2744a00..c2ab081 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -29,7 +29,6 @@ import ( "felis.lolicon.best/internal/reaper" "felis.lolicon.best/internal/registryprune" "felis.lolicon.best/internal/restore" - "felis.lolicon.best/internal/store" "felis.lolicon.best/internal/submit" "k8s.io/apimachinery/pkg/runtime" utilruntime "k8s.io/apimachinery/pkg/util/runtime" @@ -88,7 +87,9 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int { ctx := ctrl.SetupSignalHandler() - drv, err := store.Open(ctx, cfg.Database.URL) + // Before anything serves: an api on a schema it was not built for answers with + // errors, or writes rows the other version cannot read. + drv, err := openStore(ctx, cfg.Database.URL, false) if err != nil { fmt.Fprintf(stderr, "felis api: open database: %v\n", err) return 1 diff --git a/cmd/felis/migrate.go b/cmd/felis/migrate.go index 26b23b6..b67d966 100644 --- a/cmd/felis/migrate.go +++ b/cmd/felis/migrate.go @@ -111,3 +111,26 @@ func hasPending(done map[int]struct{}, migrations []store.Migration) bool { } return false } + +// openStore opens the business database for a command that reads and writes its +// tables, and refuses one whose schema this build was not written against: a newer +// Felis migrated it (a rolled-back binary), or, unless allowPending, migrations this +// build embeds have not run yet (a binary swapped in ahead of `felis migrate up`). +func openStore(ctx context.Context, url string, allowPending bool) (*store.PostgresDriver, error) { + drv, err := store.Open(ctx, url) + if err != nil { + return nil, err + } + s, err := store.ReadSchema(ctx, drv) + if err == nil { + err = s.Err() + if allowPending { + err = s.Newer() + } + } + if err != nil { + drv.Close() + return nil, err + } + return drv, nil +} diff --git a/cmd/felis/offsite.go b/cmd/felis/offsite.go index 56bfece..b852c1a 100644 --- a/cmd/felis/offsite.go +++ b/cmd/felis/offsite.go @@ -16,7 +16,6 @@ import ( "felis.lolicon.best/internal/dbbackup" "felis.lolicon.best/internal/offsite" "felis.lolicon.best/internal/platform" - "felis.lolicon.best/internal/store" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -242,7 +241,7 @@ func runOffsiteSync(cfg *config.Config, env *offsiteEnv, archiveDir, backupPVC, } archiveDir = dir } - drv, err := store.Open(ctx, cfg.Database.URL) + drv, err := openStore(ctx, cfg.Database.URL, false) if err != nil { return offsite.Result{}, fmt.Errorf("open database: %w", err) } @@ -536,7 +535,7 @@ func offsiteFetchWorlds(fs *flag.FlagSet, args []string, stdout, stderr io.Write return 1 } } - drv, err := store.Open(ctx, cfg.Database.URL) + drv, err := openStore(ctx, cfg.Database.URL, false) if err != nil { fmt.Fprintf(stderr, "felis offsite fetch-worlds: open database: %v\n", err) return 1 diff --git a/cmd/felis/reaper.go b/cmd/felis/reaper.go index 031ffef..a57b4d8 100644 --- a/cmd/felis/reaper.go +++ b/cmd/felis/reaper.go @@ -19,7 +19,6 @@ import ( "felis.lolicon.best/internal/mail" "felis.lolicon.best/internal/platform" "felis.lolicon.best/internal/reaper" - "felis.lolicon.best/internal/store" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/runtime" utilruntime "k8s.io/apimachinery/pkg/util/runtime" @@ -71,7 +70,7 @@ func cmdReaper(args []string, stdout, stderr io.Writer) int { return 1 } - drv, err := store.Open(ctx, cfg.Database.URL) + drv, err := openStore(ctx, cfg.Database.URL, false) if err != nil { fmt.Fprintf(stderr, "felis reaper: open database: %v\n", err) return 1 diff --git a/cmd/felis/setup.go b/cmd/felis/setup.go index 085f0fe..edb66ec 100644 --- a/cmd/felis/setup.go +++ b/cmd/felis/setup.go @@ -345,7 +345,9 @@ func openConfiguredSetup(ctx context.Context, cfgPath string) (*configuredSetup, if err != nil { return nil, &setupOpenError{stage: "load config", err: err} } - drv, err := store.Open(ctx, cfg.Database.URL) + // Pending migrations are the preflight's to apply; a newer schema is a rolled-back + // binary, and nothing this console writes would match it. + drv, err := openStore(ctx, cfg.Database.URL, true) if err != nil { return nil, &setupOpenError{stage: "open database", err: err} } diff --git a/cmd/felis/tui_migration.go b/cmd/felis/tui_migration.go index df7480e..f21e2a5 100644 --- a/cmd/felis/tui_migration.go +++ b/cmd/felis/tui_migration.go @@ -33,17 +33,9 @@ func applyMigrations(dbURL string) (int, error) { if _, err := store.Up(ctx, drv, migrations); err != nil { return 0, err } - applied, err := drv.AppliedVersions(ctx) + s, err := store.ReadSchema(ctx, drv) if err != nil { return 0, err } - return len(applied), nil -} - -func totalMigrations() (int, error) { - migrations, err := store.LoadMigrations() - if err != nil { - return 0, err - } - return len(migrations), nil + return s.Applied, nil } diff --git a/cmd/felis/tui_postgres.go b/cmd/felis/tui_postgres.go index f2df545..f92d963 100644 --- a/cmd/felis/tui_postgres.go +++ b/cmd/felis/tui_postgres.go @@ -72,20 +72,14 @@ func parseDBURL(url string) (dbCfg, error) { }, nil } -func countMigrations(dbURL string) (int, error) { +// readSchema lines the database's recorded migrations up with this build's. +func readSchema(dbURL string) (store.Schema, error) { ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() drv, err := store.Open(ctx, dbURL) if err != nil { - return 0, err + return store.Schema{}, err } defer drv.Close() - if err := drv.EnsureVersionTable(ctx); err != nil { - return 0, err - } - applied, err := drv.AppliedVersions(ctx) - if err != nil { - return 0, err - } - return len(applied), nil + return store.ReadSchema(ctx, drv) } diff --git a/cmd/felis/tui_preflight.go b/cmd/felis/tui_preflight.go index 3392e0b..e892ecd 100644 --- a/cmd/felis/tui_preflight.go +++ b/cmd/felis/tui_preflight.go @@ -46,6 +46,7 @@ type pfDBMsg struct{ err error } type pfMigCheckMsg struct { applied int total int + pending bool err error } @@ -84,7 +85,7 @@ func (m *preflightModel) Update(msg tea.Msg) (tea.Model, tea.Cmd) { return m, nil } m.applied, m.total = msg.applied, msg.total - if msg.applied < msg.total { + if msg.pending { m.state = pfApplyMig return m, m.applyMigrations() } @@ -204,12 +205,14 @@ func (m *preflightModel) checkDB() tea.Cmd { func (m *preflightModel) checkMigrations() tea.Cmd { return func() tea.Msg { - applied, err := countMigrations(m.dbURL) - if err != nil { - return pfMigCheckMsg{err: err} + // The sets, not their sizes: a database a newer release migrated can hold as + // many rows as this build has migrations, and must stop here rather than be + // "healed" by an older binary. + s, err := readSchema(m.dbURL) + if err == nil { + err = s.Newer() } - total, err := totalMigrations() - return pfMigCheckMsg{applied: applied, total: total, err: err} + return pfMigCheckMsg{applied: s.Applied, total: s.Total, pending: len(s.Pending) > 0, err: err} } } diff --git a/deploy/bootstrap.sh b/deploy/bootstrap.sh index 9bfe772..3fbb066 100644 --- a/deploy/bootstrap.sh +++ b/deploy/bootstrap.sh @@ -331,6 +331,12 @@ PANEL_TLS_CERT="${STATE_DIR}/panel-tls.crt" PANEL_TLS_KEY="${STATE_DIR}/panel-tls.key" SRC_DIR="/opt/felis/src" HOST_BIN="/usr/local/bin/felis" +# The binary this run replaced (keep_previous_host_binary), and whether the new one has +# been put to use: once migrations start, or the nano service restarts onto it, the old +# one no longer matches what is running and a failed run keeps the new one. +HOST_BIN_PREV="" +HOST_BIN_KEPT=0 +HOST_BIN_IN_USE=0 # Felis's own build toolchain, not /usr/local/go: install_go_toolchain replaces whatever # version sits here, and an operator's Go at the conventional path is not ours to swap. GOROOT_DIR="/opt/felis/go" @@ -437,7 +443,8 @@ on_error() { } cleanup() { - local id path unit + local status=$? id path unit + restore_previous_host_binary "$status" for unit in "${PKG_TIMERS_TO_RESTORE[@]-}"; do [ -n "$unit" ] || continue systemctl start "$unit" >/dev/null 2>&1 || true @@ -455,6 +462,36 @@ cleanup() { rm -f -- "$WATCHDOG_QUIET_FILE" 2>/dev/null || true } +# keep_previous_host_binary copies the felis binary this run is about to replace, once. A +# run that fails before the new binary is in use puts it back (restore_previous_host_binary): +# until then the old binary, the old cluster and the unmigrated database still agree, while +# a new binary left behind runs the host timers against a schema it was not built for, and +# `felis setup` with it would migrate the database under the old control plane. +keep_previous_host_binary() { + [ "$HOST_BIN_KEPT" = 0 ] || return 0 + HOST_BIN_KEPT=1 + [ -x "$HOST_BIN" ] || return 0 + HOST_BIN_PREV="${HOST_BIN}.prev" + rm -f "$HOST_BIN_PREV" + cp "$HOST_BIN" "$HOST_BIN_PREV" +} + +restore_previous_host_binary() { # exit-status + [ -n "$HOST_BIN_PREV" ] && [ -f "$HOST_BIN_PREV" ] || return 0 + if [ "$1" -ne 0 ] && [ "$HOST_BIN_IN_USE" != 1 ]; then + # install(1) onto a fresh file, as everywhere else HOST_BIN is written (SELinux label). + rm -f "$HOST_BIN" + if install -m 0755 "$HOST_BIN_PREV" "$HOST_BIN"; then + command -v restorecon >/dev/null 2>&1 && restorecon "$HOST_BIN" >/dev/null 2>&1 || true + warn "restored the previous felis binary at ${HOST_BIN}; the database was not migrated, so rerunning the installer picks up where this run stopped" + else + warn "could not restore the previous felis binary; it is at ${HOST_BIN_PREV}" + return 0 + fi + fi + rm -f -- "$HOST_BIN_PREV" +} + remember_temp() { TEMP_PATHS+=("$1"); } remember_container() { DOCKER_CONTAINERS+=("$1"); } @@ -1518,6 +1555,7 @@ download_release_binary() { # would carry the source SELinux label instead of type-transitioning to bin_t — see # build_nano_binary for the 203/EXEC this shape avoids. Same-directory staging does not # change that: install(1) still creates the destination and copies. + keep_previous_host_binary rm -f "$HOST_BIN" install -m 0755 "$tmp" "$HOST_BIN" rm -f "$tmp" @@ -1652,6 +1690,7 @@ install_embedded_binary() { mkdir -p "$(dirname "$HOST_BIN")" if [ "$(readlink -f "$src")" != "$(readlink -f "$HOST_BIN" 2>/dev/null || true)" ]; then log "installing current felis binary onto the host (${HOST_BIN})" + keep_previous_host_binary install -m 0755 "$src" "$HOST_BIN" else ok "host binary already installed at ${HOST_BIN}" @@ -1702,6 +1741,7 @@ build_image_from_source() { local cid cid="$(docker create "$FELIS_IMAGE")" remember_container "$cid" + keep_previous_host_binary docker cp "${cid}:/usr/local/bin/felis" "$HOST_BIN" docker rm "$cid" >/dev/null chmod 0755 "$HOST_BIN" @@ -3095,6 +3135,8 @@ run_migrations() { # binary bundles the database into FELIS_DB_BACKUP_DIR first and refuses to migrate # if that fails; a fresh database has nothing to protect and is migrated directly. log "running database migrations (host binary -> 127.0.0.1)" + # From here the database may move forward, and the binary that moved it stays. + HOST_BIN_IN_USE=1 "$HOST_BIN" migrate up -config "${STATE_DIR}/felis.host.toml" "${backup_flags[@]}" ok "migrations applied" } @@ -3814,6 +3856,7 @@ build_nano_binary() { # cannot exec it and felis-nano dies with 203/EXEC. Creating the file fresh at the # destination lets the policy's type transition label it bin_t; restorecon is the belt. mkdir -p "$(dirname "$HOST_BIN")" + keep_previous_host_binary rm -f "$HOST_BIN" install -m 0755 "$staged" "$HOST_BIN" rm -f "$staged" @@ -3956,6 +3999,7 @@ EOF systemctl enable felis-nano # restart, not `enable --now`: on a re-run the service is already active and --now would # leave the OLD binary running against the NEW unit. Converge means converge. + HOST_BIN_IN_USE=1 systemctl restart felis-nano # restart returns as soon as the process is forked. A config the new binary rejects, or a # file it cannot open, only shows once it has exited and the unit sits in auto-restart. diff --git a/deploy/bootstrap_test.sh b/deploy/bootstrap_test.sh index b731790..bbdaba3 100644 --- a/deploy/bootstrap_test.sh +++ b/deploy/bootstrap_test.sh @@ -2005,6 +2005,37 @@ else echo "FAIL BUILDX_NO_DEFAULT_ATTESTATIONS=1 must be exported before the first docker build"; fails=$((fails + 1)) fi +# --- a failed run puts the previous host binary back until the new one is in use ---------- +hbdir="$(mktemp -d)" +run_host_bin() { # exit-status in-use [no-previous] + rm -f "$hbdir"/felis* + [ -n "${3:-}" ] || { printf 'old\n' > "$hbdir/felis"; chmod 0755 "$hbdir/felis"; } + HOST_BIN="$hbdir/felis" bash -c ' + set -e + warn() { printf "WARN: %s\n" "$*"; } + HOST_BIN_PREV=""; HOST_BIN_KEPT=0; HOST_BIN_IN_USE='"$2"' + '"$(awk '/^keep_previous_host_binary\(\) \{/,/^}/' "$BS")"' + '"$(awk '/^restore_previous_host_binary\(\) \{/,/^}/' "$BS")"' + keep_previous_host_binary + rm -f "$HOST_BIN"; printf "new\n" > "$HOST_BIN"; chmod 0755 "$HOST_BIN" + keep_previous_host_binary # a second replacement keeps the first original + restore_previous_host_binary '"$1"' + printf "BIN: %s\n" "$(cat "$HOST_BIN")" + [ -e "$HOST_BIN.prev" ] && echo "PREV LEFT" || true' +} +out="$(run_host_bin 1 0)" +expect "a run that fails before the new binary is used restores the old one" "BIN: old" "$out" +expect "the restore says so" "WARN: restored the previous felis binary" "$out" +out="$(run_host_bin 1 1)" +expect "a run that fails after migrations keeps the new binary" "BIN: new" "$out" +out="$(run_host_bin 0 0)" +expect "a successful run keeps the new binary" "BIN: new" "$out" +case "$out" in *"PREV LEFT"*) echo "FAIL: the previous binary copy must be removed"; fails=$((fails + 1)) ;; esac +out="$(run_host_bin 1 0 fresh)" +expect "a first install has nothing to restore" "BIN: new" "$out" +case "$out" in *WARN:*) echo "FAIL: a first install must not restore the binary it just installed"; fails=$((fails + 1)) ;; esac +rm -rf "$hbdir" + # --------------------------------------------------------------------------------------- if [ "$fails" -eq 0 ]; then diff --git a/deploy/uninstall.sh b/deploy/uninstall.sh index 80c4bf1..a3cf9a3 100644 --- a/deploy/uninstall.sh +++ b/deploy/uninstall.sh @@ -312,7 +312,7 @@ remove_host_files() { userdel "$VELOCITY_USER" >/dev/null 2>&1 || warn "could not remove the ${VELOCITY_USER} user" fi rm -rf "$OPT_DIR" - rm -f "$HOST_BIN" "${HOST_BIN}.new" + rm -f "$HOST_BIN" "${HOST_BIN}.new" "${HOST_BIN}.prev" remove_cloudflared_binary if [ "$PURGE" = 0 ]; then # What describes the removed install goes; what a reinstall reuses stays. Without diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index 6eedb3c..1264096 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -1385,6 +1385,24 @@ exists, so the database is upgraded only when you run `pg_upgrade` yourself. applied; when they are the problem, restore the `pre-migrate` bundle the upgrade took (§16, "Roll back an upgrade that broke the database"). +**Schema guard.** felis-api, the reaper, the off-site copy and `felis migrate up` +compare the migrations the database records with the ones their build embeds, as +sets. A database a newer release migrated stops them with `database schema is +newer than this felis build: it records migration 0026, and this build knows +migrations up to 0025`, so an undo across an upgrade that migrated shows up as +felis-api in CrashLoopBackOff until the `pre-migrate` bundle is restored. A +database still missing migrations stops felis-api, the reaper and the off-site +copy with `database schema is behind this felis build` until `felis migrate up` +runs; `felis setup`'s preflight applies them itself. [PG-TESTED] + +**A failed rerun and the host binary.** The installer replaces +`/usr/local/bin/felis` early (the steps after it run the new binary) and keeps +the old one as `felis.prev` until the new one is in use. A run that fails before +the database migrations start puts the old binary back, so the host timers and +`felis setup` keep matching the database and the control plane that are still +running; a rerun continues from there. Once migrations have started, the new +binary stays. [VM-VERIFIED] + ## 15b. Game images, pinned builds, and moving a world to a newer Minecraft Each release pins the upstream builds it installs in `deploy/game-stack.lock`: diff --git a/internal/pgint/schema_test.go b/internal/pgint/schema_test.go new file mode 100644 index 0000000..49ba250 --- /dev/null +++ b/internal/pgint/schema_test.go @@ -0,0 +1,75 @@ +//go:build pgint + +package pgint + +import ( + "context" + "errors" + "os" + "strings" + "testing" + + "felis.lolicon.best/internal/store" +) + +// The schema guard the api, reaper and offsite copy open the database through: the +// migrated test database passes, a never-migrated one reads as behind (the missing +// schema_migrations table is an empty set, not an error), and one that records a +// version this build does not embed reads as newer. +func TestSchemaGuard(t *testing.T) { + ctx := context.Background() + main, err := store.Open(ctx, os.Getenv("FELIS_TEST_PG_URL")) + if err != nil { + t.Fatal(err) + } + defer main.Close() + if err := store.CheckSchema(ctx, main); err != nil { + t.Fatalf("CheckSchema on the migrated database: %v", err) + } + + const schema = "pgint_schema_guard" + if _, err := db.ExecContext(ctx, "DROP SCHEMA IF EXISTS "+schema+" CASCADE; CREATE SCHEMA "+schema); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _, _ = db.ExecContext(context.Background(), "DROP SCHEMA IF EXISTS "+schema+" CASCADE") }) + dsn := os.Getenv("FELIS_TEST_PG_URL") + sep := "?" + if strings.Contains(dsn, "?") { + sep = "&" + } + other, err := store.Open(ctx, dsn+sep+"search_path="+schema) + if err != nil { + t.Fatal(err) + } + defer other.Close() + + if err := store.CheckSchema(ctx, other); !errors.Is(err, store.ErrSchemaBehind) { + t.Fatalf("CheckSchema on an unmigrated schema = %v, want ErrSchemaBehind", err) + } + + ms, err := store.LoadMigrations() + if err != nil { + t.Fatal(err) + } + if err := other.EnsureVersionTable(ctx); err != nil { + t.Fatal(err) + } + for _, m := range ms { + if _, err := other.DB().ExecContext(ctx, "INSERT INTO schema_migrations (version, name) VALUES ($1, $2)", m.Version, m.Name); err != nil { + t.Fatal(err) + } + } + if err := store.CheckSchema(ctx, other); err != nil { + t.Fatalf("CheckSchema with every version recorded: %v", err) + } + if _, err := other.DB().ExecContext(ctx, "INSERT INTO schema_migrations (version, name) VALUES (9999, 'from_a_newer_release')"); err != nil { + t.Fatal(err) + } + err = store.CheckSchema(ctx, other) + if !errors.Is(err, store.ErrSchemaNewer) || !strings.Contains(err.Error(), "9999") { + t.Fatalf("CheckSchema with a newer version recorded = %v, want ErrSchemaNewer naming 9999", err) + } + if _, err := store.Up(ctx, other, ms); !errors.Is(err, store.ErrSchemaNewer) { + t.Fatalf("Up on a newer schema = %v, want ErrSchemaNewer", err) + } +} diff --git a/internal/store/migrate.go b/internal/store/migrate.go index 4d43f7c..3597b2b 100644 --- a/internal/store/migrate.go +++ b/internal/store/migrate.go @@ -84,7 +84,8 @@ func LoadMigrations() ([]Migration, error) { // Up applies every pending migration in ascending order, exactly once, under // the advisory lock. It is safe to run concurrently from multiple replicas: the -// lock serializes them and AppliedVersions makes the work idempotent. +// lock serializes them and AppliedVersions makes the work idempotent. It refuses +// (ErrSchemaNewer) a database that records a version this build does not embed. func Up(ctx context.Context, d Driver, migrations []Migration) (applied []int, err error) { if err := d.Lock(ctx); err != nil { return nil, fmt.Errorf("acquire migration lock: %w", err) @@ -102,6 +103,12 @@ func Up(ctx context.Context, d Driver, migrations []Migration) (applied []int, e if err != nil { return nil, fmt.Errorf("read applied versions: %w", err) } + // A version this build does not embed means a newer release migrated the database. + // Applying the older build's remaining steps on top would be a guess about a schema + // it never saw, so nothing runs. + if err := CompareSchema(done, migrations).Newer(); err != nil { + return nil, err + } ordered := append([]Migration(nil), migrations...) sort.Slice(ordered, func(i, j int) bool { return ordered[i].Version < ordered[j].Version }) diff --git a/internal/store/migrate_test.go b/internal/store/migrate_test.go index 5e8035e..13834dc 100644 --- a/internal/store/migrate_test.go +++ b/internal/store/migrate_test.go @@ -3,6 +3,7 @@ package store_test import ( "context" "errors" + "fmt" "strings" "testing" @@ -155,3 +156,50 @@ func itoa(v int) string { } return string(buf[i:]) } + +func TestUpRefusesADatabaseANewerBuildMigrated(t *testing.T) { + d := &recordingDriver{already: map[int]struct{}{1: {}, 2: {}, 3: {}}} + ms := []store.Migration{ + {Version: 1, Name: "a", SQL: "y"}, + {Version: 2, Name: "b", SQL: "z"}, + } + _, err := store.Up(context.Background(), d, ms) + if !errors.Is(err, store.ErrSchemaNewer) { + t.Fatalf("Up err = %v, want ErrSchemaNewer", err) + } + if !strings.Contains(err.Error(), "migration 0003") || !strings.Contains(err.Error(), "up to 0002") { + t.Errorf("error does not name the versions: %v", err) + } + if !d.unlocked { + t.Error("the advisory lock was not released") + } +} + +func TestCompareSchemaSeparatesPendingFromUnknown(t *testing.T) { + ms := []store.Migration{{Version: 1}, {Version: 2}, {Version: 4}} + cases := []struct { + name string + applied map[int]struct{} + pending []int + unknown []int + err error + }{ + {"current", map[int]struct{}{1: {}, 2: {}, 4: {}}, nil, nil, nil}, + {"fresh", map[int]struct{}{}, []int{1, 2, 4}, nil, store.ErrSchemaBehind}, + {"behind", map[int]struct{}{1: {}}, []int{2, 4}, nil, store.ErrSchemaBehind}, + // Same row count as "current": a count comparison calls this up to date. + {"newer", map[int]struct{}{1: {}, 2: {}, 5: {}}, []int{4}, []int{5}, store.ErrSchemaNewer}, + } + for _, c := range cases { + s := store.CompareSchema(c.applied, ms) + if fmt.Sprint(s.Pending) != fmt.Sprint(c.pending) || fmt.Sprint(s.Unknown) != fmt.Sprint(c.unknown) { + t.Errorf("%s: pending=%v unknown=%v, want %v %v", c.name, s.Pending, s.Unknown, c.pending, c.unknown) + } + if err := s.Err(); !errors.Is(err, c.err) || (c.err == nil && err != nil) { + t.Errorf("%s: Err() = %v, want %v", c.name, err, c.err) + } + if s.Total != 3 || s.Latest != 4 { + t.Errorf("%s: Total=%d Latest=%d, want 3 4", c.name, s.Total, s.Latest) + } + } +} diff --git a/internal/store/schema.go b/internal/store/schema.go new file mode 100644 index 0000000..95a507a --- /dev/null +++ b/internal/store/schema.go @@ -0,0 +1,109 @@ +package store + +import ( + "context" + "errors" + "fmt" + "sort" + "strings" +) + +// ErrSchemaNewer marks a database that a newer Felis has migrated: it records versions +// this build does not embed. Up only rolls forward, so running this build against it +// would read and write tables whose shape it was never written for. +var ErrSchemaNewer = errors.New("database schema is newer than this felis build") + +// ErrSchemaBehind marks a database that still lacks migrations this build embeds. +var ErrSchemaBehind = errors.New("database schema is behind this felis build") + +// Schema is how a database's recorded migrations line up with the ones embedded in +// this build. Comparing the two sets, rather than counting, is what tells "behind" +// from "migrated by a newer release": both can have the same number of rows. +type Schema struct { + Applied int // embedded migrations the database has + Total int // embedded migrations + Latest int // highest embedded version + Pending []int // embedded, not yet applied, ascending + Unknown []int // applied, but not embedded here, ascending +} + +// CompareSchema lines applied up with migrations. +func CompareSchema(applied map[int]struct{}, migrations []Migration) Schema { + s := Schema{Total: len(migrations)} + known := make(map[int]struct{}, len(migrations)) + for _, m := range migrations { + known[m.Version] = struct{}{} + s.Latest = max(s.Latest, m.Version) + if _, ok := applied[m.Version]; ok { + s.Applied++ + } else { + s.Pending = append(s.Pending, m.Version) + } + } + for v := range applied { + if _, ok := known[v]; !ok { + s.Unknown = append(s.Unknown, v) + } + } + sort.Ints(s.Pending) + sort.Ints(s.Unknown) + return s +} + +// Newer returns ErrSchemaNewer, with the versions and what to do, when a newer Felis +// migrated this database; nil otherwise. +func (s Schema) Newer() error { + if len(s.Unknown) == 0 { + return nil + } + return fmt.Errorf("%w: it records migration %s, and this build knows %s. Run the Felis release that migrated it, or restore the pre-migrate snapshot it took (felis db restore)", + ErrSchemaNewer, versionList(s.Unknown), knownRange(s.Latest)) +} + +// Err is Newer, else ErrSchemaBehind when migrations are pending: the check a server +// makes before it serves anything from the database. +func (s Schema) Err() error { + if err := s.Newer(); err != nil { + return err + } + if len(s.Pending) > 0 { + return fmt.Errorf("%w: migration %s not applied yet. Run `felis migrate up` (the installer does), then start this again", + ErrSchemaBehind, versionList(s.Pending)) + } + return nil +} + +// ReadSchema compares what the database records with the embedded migrations. +func ReadSchema(ctx context.Context, d Driver) (Schema, error) { + migrations, err := LoadMigrations() + if err != nil { + return Schema{}, err + } + applied, err := d.AppliedVersions(ctx) + if err != nil { + return Schema{}, fmt.Errorf("read applied migrations: %w", err) + } + return CompareSchema(applied, migrations), nil +} + +// CheckSchema fails unless the database carries exactly the migrations this build +// embeds. +func CheckSchema(ctx context.Context, d Driver) error { + s, err := ReadSchema(ctx, d) + if err != nil { + return err + } + return s.Err() +} + +func versionList(vs []int) string { + parts := make([]string, len(vs)) + for i, v := range vs { + parts[i] = fmt.Sprintf("%04d", v) + } + return strings.Join(parts, ", ") +} + +func knownRange(latest int) string { + return fmt.Sprintf("migrations up to %04d", latest) +} diff --git a/internal/store/sqldriver.go b/internal/store/sqldriver.go index 2f8fd05..0e6ae98 100644 --- a/internal/store/sqldriver.go +++ b/internal/store/sqldriver.go @@ -3,8 +3,10 @@ package store import ( "context" "database/sql" + "errors" "fmt" + "github.com/jackc/pgx/v5/pgconn" _ "github.com/jackc/pgx/v5/stdlib" // register the "pgx" database/sql driver ) @@ -56,10 +58,15 @@ func (d *PostgresDriver) EnsureVersionTable(ctx context.Context) error { return err } -// AppliedVersions reads the set of recorded versions. +// AppliedVersions reads the set of recorded versions. A database that has never been +// migrated has no schema_migrations table yet, which is an empty set. func (d *PostgresDriver) AppliedVersions(ctx context.Context) (map[int]struct{}, error) { rows, err := d.db.QueryContext(ctx, "SELECT version FROM schema_migrations") if err != nil { + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && pgErr.Code == "42P01" { // undefined_table + return map[int]struct{}{}, nil + } return nil, err } defer rows.Close()