From 9a892541460954c845b5b6fb59c76c7bdb7357bd Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Sun, 27 Sep 2026 01:58:35 +0800 Subject: [PATCH] =?UTF-8?q?test(db):=20=E6=81=A2=E5=A4=8D=E9=93=BE?= =?UTF-8?q?=E5=9C=A8=E7=9C=9F=E5=AE=9E=20PostgreSQL=20=E4=B8=8A=E8=B7=91?= =?UTF-8?q?=E5=BE=80=E8=BF=94=E4=B8=8E=E4=B8=89=E7=A7=8D=E5=A4=B1=E8=B4=A5?= =?UTF-8?q?=E5=9B=9E=E6=BB=9A=EF=BC=8C-yes=20=E9=97=B8=E9=97=A8=E6=B5=8B?= =?UTF-8?q?=E8=AF=95=E7=9C=9F=E6=AD=A3=E8=B5=B0=E5=88=B0=E9=97=B8=E9=97=A8?= =?UTF-8?q?=EF=BC=8Ce2e=20=E6=8C=89=E6=96=87=E6=A1=A3=E5=81=9A=E4=B8=80?= =?UTF-8?q?=E6=AC=A1=E5=90=8C=E6=9C=BA=E6=81=A2=E5=A4=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .github/workflows/ci.yml | 4 + CONTRIBUTING.md | 11 +- cmd/felis/db_test.go | 60 +++- deploy/e2e_check.sh | 59 +++- docs/troubleshooting.md | 2 +- internal/pgint/dbrestore_test.go | 477 +++++++++++++++++++++++++++++++ internal/pgint/pgint_test.go | 2 + 7 files changed, 599 insertions(+), 16 deletions(-) create mode 100644 internal/pgint/dbrestore_test.go diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 6b1cee1..f85df98 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -72,6 +72,9 @@ jobs: # The business stores' SQL against a real PostgreSQL (internal/pgint): the unit suites run # on fakes, and PGRepo drifted from them three times while those stayed green. 13 is the # oldest server a supported distribution installs (EL9), 18 the newest (Arch). + # `felis db backup` and `restore` run there too, with the tools inside the service + # container, as production runs them inside felis-postgres: the runner's own client is one + # major version, and pg_dump refuses a newer server. pgint: runs-on: ubuntu-latest strategy: @@ -102,6 +105,7 @@ jobs: - run: go test -race -tags pgint -count=1 ./internal/pgint/ env: FELIS_TEST_PG_URL: postgres://felis:pgint@localhost:5432/felis_pgint?sslmode=disable + FELIS_TEST_PG_EXEC: docker exec -i ${{ job.services.postgres.id }} shell: runs-on: ubuntu-latest diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 6417bf4..61d1eb2 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -78,8 +78,15 @@ FELIS_TEST_PG_URL='postgres://felis:***@127.0.0.1:5432/felis_pgint?sslmode=disab ``` Run it after touching anything under `internal/api/pgrepo.go`, `internal/submit`, -or `internal/build` that speaks SQL: the fakes encode the contract, and this -suite exists to catch the drift between the fakes and the real queries. +`internal/build` or `internal/dbbackup` that speaks SQL: the fakes encode the +contract, and this suite exists to catch the drift between the fakes and the real +queries. The `felis db backup` and `restore` tests also run `pg_dump`, `pg_restore` +and `psql`, which must be the server's major version. For a server in a container, +run them in it, as production does in felis-postgres: + +```bash +FELIS_TEST_PG_EXEC='docker exec -i ' FELIS_TEST_PG_URL=... go test -tags pgint ./internal/pgint/ +``` Build the CLI: diff --git a/cmd/felis/db_test.go b/cmd/felis/db_test.go index ae75302..a381847 100644 --- a/cmd/felis/db_test.go +++ b/cmd/felis/db_test.go @@ -29,13 +29,43 @@ func TestDBUsage(t *testing.T) { } } +// TestDBRestoreNeedsYes: without -yes a restore describes the bundle and stops +// before anything reaches the database, even with -force and +// -no-safety-backup, which would otherwise let the replay run at once. func TestDBRestoreNeedsYes(t *testing.T) { - // A bundle that does not exist fails verification (1) before -yes matters; - // the -yes gate itself is exercised against a real bundle in internal/dbbackup - // and on the VM. Here: the refusal path never reaches the config or database. + dir := newPodRig(t) + cfg := podConfig(t, dir) + bundles := filepath.Join(dir, "bundles") var out, errBuf bytes.Buffer - if code := run([]string{"db", "restore", "-dir", t.TempDir(), "missing.tar"}, &out, &errBuf); code != 1 { - t.Fatalf("exit %d, stderr %q", code, errBuf.String()) + if code := run([]string{"db", "backup", "-config", cfg, "-dir", bundles, "-state-dir", "", "-no-servers"}, &out, &errBuf); code != 0 { + t.Fatalf("backup: exit %d: %s", code, errBuf.String()) + } + bundle := strings.TrimSpace(strings.TrimPrefix(out.String(), "felis db backup: wrote ")) + podRuns(t, dir) + + out.Reset() + errBuf.Reset() + code := run([]string{"db", "restore", "-config", cfg, "-dir", bundles, "-force", "-no-safety-backup", filepath.Base(bundle)}, &out, &errBuf) + if code != 2 { + t.Fatalf("exit %d, want 2; stderr %q", code, errBuf.String()) + } + if want := filepath.Base(bundle) + " (manual, taken "; !strings.Contains(errBuf.String(), want) || !strings.Contains(errBuf.String(), "schema 3).") { + t.Errorf("stderr %q does not describe the bundle", errBuf.String()) + } + if !strings.Contains(errBuf.String(), "re-run with -yes") { + t.Errorf("stderr %q does not say how to go on", errBuf.String()) + } + if argv, err := os.ReadFile(filepath.Join(dir, "k3s.args")); err == nil { + t.Errorf("a restore without -yes ran in the database pod:\n%s", argv) + } + + // A bundle that does not verify is refused before -yes is weighed. + errBuf.Reset() + if code := run([]string{"db", "restore", "-config", cfg, "-dir", bundles, "-yes", "missing.tar"}, &out, &errBuf); code != 1 { + t.Errorf("missing bundle: exit %d, want 1; stderr %q", code, errBuf.String()) + } + if _, err := os.Stat(filepath.Join(dir, "k3s.args")); err == nil { + t.Error("a missing bundle reached the database pod") } } @@ -268,6 +298,19 @@ func podRuns(t *testing.T, dir string) []string { return runs } +// podConfig writes an installed host's felis.toml, [database] pointing at the +// pod, into dir. +func podConfig(t *testing.T, dir string) string { + t.Helper() + toml := strings.Replace(installerTOML("example.com", "127.0.0.1"), + `url = "postgres://felis:pw@127.0.0.1:5432/felis?sslmode=disable"`, + `url = "`+podDB.URL+`" +deployment = "`+podDB.Deployment+`"`, 1) + cfg := filepath.Join(dir, "felis.toml") + writeTestFile(t, cfg, toml, 0o600) + return cfg +} + func ranIn(runs []string, prefix string) bool { for _, r := range runs { if strings.HasPrefix(r, prefix) { @@ -284,12 +327,7 @@ func ranIn(runs []string, prefix string) bool { // line. func TestDBBackupAndRestoreRunTheToolsInTheDatabasePod(t *testing.T) { dir := newPodRig(t) - toml := strings.Replace(installerTOML("example.com", "127.0.0.1"), - `url = "postgres://felis:pw@127.0.0.1:5432/felis?sslmode=disable"`, - `url = "`+podDB.URL+`" -deployment = "`+podDB.Deployment+`"`, 1) - cfg := filepath.Join(dir, "felis.toml") - writeTestFile(t, cfg, toml, 0o600) + cfg := podConfig(t, dir) var out, errBuf bytes.Buffer bundles := filepath.Join(dir, "bundles") diff --git a/deploy/e2e_check.sh b/deploy/e2e_check.sh index fa5dd33..f930163 100644 --- a/deploy/e2e_check.sh +++ b/deploy/e2e_check.sh @@ -8,8 +8,9 @@ # sudo bash deploy/e2e_check.sh upgrade # after this commit ran over a release # # It asks what an operator's first minutes ask: the binary runs, the control plane and its -# database are rolled out and ready, a database backup can be taken, the panel answers on -# its NodePort, the proxy answers a Minecraft status ping, and the host timers are there. +# database are rolled out and ready, a database backup can be taken and restored, the panel +# answers on its NodePort, the proxy answers a Minecraft status ping, and the host timers +# are there. # A rerun must also leave the proxy running (it restarts only when what it runs changed) # and keep every earlier answer. set -euo pipefail @@ -50,6 +51,58 @@ check "felis-api is ready (database and cluster reachable)" \ for unit in k3s felis-velocity; do check "${unit} is active" systemctl is-active --quiet "$unit" done +# pod_psql runs one statement in the database's pod, over its socket, as the felis role on +# the felis database. +pod_psql() { + "${KUBECTL[@]}" -n felis exec -i deploy/felis-postgres -c postgres -- \ + psql -X -q -At -v ON_ERROR_STOP=1 -U felis -d felis -c "$1" +} + +# restore_drill walks troubleshooting.md's "Restore on the same host": refused while the +# control plane is connected; with it scaled to 0 the bundle comes back (a row written after +# it is gone) and the database it replaced is kept; migrate up runs; the control plane serves +# again. +restore_drill() { # dir bundle + local dir="$1" bundle="$2" out rc + local sel="app.kubernetes.io/part-of=felis-control-plane,app.kubernetes.io/component in (api,operator)" + if ! pod_psql "INSERT INTO platform_settings (key, value) VALUES ('e2e_restore_drill', '1')" >/dev/null; then + fail "write a row after the bundle" + return + fi + rc=0 + out="$(/usr/local/bin/felis db restore -dir "$dir" -yes "$bundle" 2>&1)" || rc=$? + if [ "$rc" -eq 1 ] && grep -q "other clients are connected to the database" <<<"$out"; then + pass "felis db restore refuses while the control plane is connected" + else + fail "felis db restore refuses while the control plane is connected (exit ${rc}): ${out}" + fi + + "${KUBECTL[@]}" -n felis scale deployment felis-api felis-operator --replicas=0 >/dev/null + for _ in $(seq 60); do + [ -z "$("${KUBECTL[@]}" -n felis get pods -l "$sel" -o name)" ] && break + sleep 2 + done + rc=0 + out="$(/usr/local/bin/felis db restore -dir "$dir" -yes "$bundle" 2>&1)" || rc=$? + if [ "$rc" -eq 0 ]; then + pass "felis db restore replays the bundle with the control plane scaled to 0" + else + fail "felis db restore replays the bundle with the control plane scaled to 0 (exit ${rc}): ${out}" + fi + check "the restore dropped the row written after the bundle" \ + test "$(pod_psql "SELECT count(*) FROM platform_settings WHERE key = 'e2e_restore_drill'")" = 0 + check "the restore kept the database it replaced in a pre-restore bundle" \ + sh -c "ls '${dir}' | grep -q -- '-pre-restore\.tar\$'" + check "felis migrate up runs on the restored database" \ + /usr/local/bin/felis migrate up -config /etc/felis/felis.host.toml + "${KUBECTL[@]}" -n felis scale deployment felis-api felis-operator --replicas=1 >/dev/null + for d in felis-api felis-operator; do + check "deployment ${d} is rolled out again after the restore" "${KUBECTL[@]}" -n felis rollout status "deploy/${d}" --timeout=180s + done + check "felis-api is ready on the restored database" \ + curl -sf --retry 10 --retry-delay 3 --retry-all-errors -o /dev/null "http://${internal}/readyz" +} + # The database runs in k3s; a release may still run it on the host, and the upgrade moved # it. The host has no PostgreSQL client: a bundle that verifies proves felis reaches the # database's pod through kubectl exec, and that pg_dump there reads every table. @@ -59,6 +112,8 @@ if [ "$phase" != release ]; then if out="$(/usr/local/bin/felis db backup -dir "$bundle_dir" -state-dir "" -no-servers 2>&1)"; then bundle="$(printf '%s\n' "$out" | sed -n 's/^felis db backup: wrote //p' | tail -n 1)" check "felis db backup writes a bundle that verifies" /usr/local/bin/felis db verify "$bundle" + # A rerun keeps what the install left; the install and the upgrade restore it. + [ "$phase" = rerun ] || restore_drill "$bundle_dir" "$bundle" else fail "felis db backup writes a bundle: ${out}" fi diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index 1d89c23..812afad 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -2011,7 +2011,7 @@ kubectl -n felis scale deployment felis-api felis-operator --replicas=1 - The replay is one transaction: it drops everything the `felis` role owns and loads the dump. **Any failure rolls back and leaves the database exactly as it was** (`rolled back, the database is unchanged`, with the psql and pg_restore - errors). [GO-TESTED] + errors). [PG-TESTED] - The database before the restore is in the `pre-restore` bundle it names; restoring that one undoes the restore. - `migrate up` brings an older bundle's schema up to the running release. diff --git a/internal/pgint/dbrestore_test.go b/internal/pgint/dbrestore_test.go new file mode 100644 index 0000000..7e25b35 --- /dev/null +++ b/internal/pgint/dbrestore_test.go @@ -0,0 +1,477 @@ +//go:build pgint + +package pgint + +import ( + "archive/tar" + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "maps" + "net/url" + "os" + "path/filepath" + "slices" + "strings" + "testing" + "time" + + "felis.lolicon.best/internal/api" + "felis.lolicon.best/internal/dbbackup" + "felis.lolicon.best/internal/store" +) + +// ---- felis db backup / restore against the real tools and server ---------------- +// +// internal/dbbackup's suite fakes pg_dump, pg_restore and psql with scripts that +// encode what the server is assumed to do: DROP OWNED clears the role's objects, +// and a transaction psql leaves open at EOF is rolled back. These tests run the +// real tools against a real server, as the role production uses (a plain LOGIN +// role owning its database; the superuser is someone else), on a database of +// their own so the rest of the suite is untouched. +// +// FELIS_TEST_PG_EXEC, when set, is the argv prefix the tools run under, the way +// production runs them under `k3s kubectl exec ... --`: CI sets +// `docker exec -i `, which also keeps the tools at the +// server's major version. Unset, the tools come from PATH and connect over the +// URL, so they must match the server. + +const ( + restoreRole = "felis_restore_pgint" + restoreDB = "felis_pgint_restore" +) + +type restoreRig struct { + url string // the role on its database, as felis.toml carries it + suURL string // the suite's superuser on the same database + tools dbbackup.Tools + dir string // bundle directory +} + +func newRestoreRig(t *testing.T) *restoreRig { + t.Helper() + ctx := context.Background() + base, err := url.Parse(os.Getenv("FELIS_TEST_PG_URL")) + if err != nil { + t.Fatal(err) + } + dropAll := func() { + for _, stmt := range []string{ + "DROP DATABASE IF EXISTS " + restoreDB + " WITH (FORCE)", + "DROP ROLE IF EXISTS " + restoreRole, + } { + if _, err := db.ExecContext(ctx, stmt); err != nil { + t.Fatalf("%s: %v", stmt, err) + } + } + } + dropAll() + mustExec(t, "CREATE ROLE "+restoreRole+" LOGIN PASSWORD 'pgint'") + mustExec(t, "CREATE DATABASE "+restoreDB+" OWNER "+restoreRole) + t.Cleanup(dropAll) + + role, su := *base, *base + role.User = url.UserPassword(restoreRole, "pgint") + role.Path, su.Path = "/"+restoreDB, "/"+restoreDB + r := &restoreRig{url: role.String(), suURL: su.String(), dir: t.TempDir()} + if pre := strings.Fields(os.Getenv("FELIS_TEST_PG_EXEC")); len(pre) > 0 { + r.tools = dbbackup.Tools{Exec: pre, Conn: "host=/var/run/postgresql dbname=" + restoreDB + " user=" + restoreRole} + } + + // The schema `felis migrate up` produces, owned by the role. + r.with(t, r.url, func(rdb *store.PostgresDriver) { + ms, err := store.LoadMigrations() + if err != nil { + t.Fatal(err) + } + if _, err := store.Up(ctx, rdb, ms); err != nil { + t.Fatalf("migrate: %v", err) + } + }) + return r +} + +// with runs fn on a pool of its own and closes it before returning: a restore +// refuses to run while any other client is connected. +func (r *restoreRig) with(t *testing.T, dsn string, fn func(*store.PostgresDriver)) { + t.Helper() + drv, err := store.Open(context.Background(), dsn) + if err != nil { + t.Fatal(err) + } + fn(drv) + drv.Close() + r.waitIdle(t) +} + +// waitIdle waits for the database's last session to end: the server ends a +// backend just after its client leaves. +func (r *restoreRig) waitIdle(t *testing.T) { + t.Helper() + for deadline := time.Now().Add(10 * time.Second); ; { + var n int + if err := db.QueryRow(`SELECT count(*) FROM pg_stat_activity WHERE datname = $1 AND backend_type = 'client backend'`, restoreDB).Scan(&n); err != nil { + t.Fatal(err) + } + if n == 0 { + return + } + if time.Now().After(deadline) { + t.Fatalf("%d sessions still on %s", n, restoreDB) + } + time.Sleep(20 * time.Millisecond) + } +} + +func (r *restoreRig) exec(t *testing.T, dsn string, stmts ...string) { + t.Helper() + r.with(t, dsn, func(d *store.PostgresDriver) { + for _, s := range stmts { + if _, err := d.DB().Exec(s); err != nil { + t.Fatalf("%s: %v", s, err) + } + } + }) +} + +// seed fills the database the way a running install does, plus a table of +// incompressible rows sorted last, so the dump is mostly data. +func (r *restoreRig) seed(t *testing.T) { + t.Helper() + r.with(t, r.url, func(d *store.PostgresDriver) { + rp := api.NewPGRepo(d.DB()) + for _, name := range []string{"alice", "bob"} { + u, err := rp.CreateUser(context.Background(), api.CreateUserInput{Username: name, Role: "user"}, "pgint") + if err != nil { + t.Fatal(err) + } + if _, err := d.DB().Exec(`INSERT INTO servers (name, cached_cpu_milli, cached_memory_mb, cached_storage_mb, owner_id) + VALUES ($1, 500, 1024, 2048, $2)`, name+"-smp", u.ID); err != nil { + t.Fatal(err) + } + } + for _, s := range []string{ + `CREATE TABLE zz_pgint_bulk (id int PRIMARY KEY, pad text NOT NULL)`, + `INSERT INTO zz_pgint_bulk SELECT g, md5(random()::text) || md5(random()::text) FROM generate_series(1, 20000) g`, + } { + if _, err := d.DB().Exec(s); err != nil { + t.Fatalf("%s: %v", s, err) + } + } + }) +} + +// fingerprint is every table's rows and every sequence's position, read as the +// superuser so objects of any owner count. The freshness record a restore +// rewrites is left out. +func (r *restoreRig) fingerprint(t *testing.T) map[string]string { + t.Helper() + fp := map[string]string{} + r.with(t, r.suURL, func(d *store.PostgresDriver) { + rows, err := d.DB().Query(`SELECT tablename, tableowner FROM pg_tables WHERE schemaname = 'public' ORDER BY 1`) + if err != nil { + t.Fatal(err) + } + var tables, owners []string + for rows.Next() { + var name, owner string + if err := rows.Scan(&name, &owner); err != nil { + t.Fatal(err) + } + tables, owners = append(tables, name), append(owners, owner) + } + rows.Close() + for i, tbl := range tables { + where := "" + if tbl == "platform_settings" { + where = " WHERE key <> '" + dbbackup.StatusKey + "'" + } + var v string + q := fmt.Sprintf(`SELECT count(*)::text || ' ' || coalesce(md5(string_agg(t::text, E'\n' ORDER BY t::text)), '') FROM public.%q t%s`, tbl, where) + if err := d.DB().QueryRow(q).Scan(&v); err != nil { + t.Fatalf("%s: %v", q, err) + } + fp["table "+tbl] = owners[i] + " " + v + } + srows, err := d.DB().Query(`SELECT sequencename, coalesce(last_value::text, '-') FROM pg_sequences WHERE schemaname = 'public'`) + if err != nil { + t.Fatal(err) + } + for srows.Next() { + var name, v string + if err := srows.Scan(&name, &v); err != nil { + t.Fatal(err) + } + fp["sequence "+name] = v + } + srows.Close() + }) + return fp +} + +// diff names what differs between two fingerprints. +func diff(a, b map[string]string) []string { + var out []string + for _, k := range slices.Sorted(maps.Keys(a)) { + if a[k] != b[k] { + out = append(out, fmt.Sprintf("%s: %q vs %q", k, a[k], b[k])) + } + } + for _, k := range slices.Sorted(maps.Keys(b)) { + if _, ok := a[k]; !ok { + out = append(out, fmt.Sprintf("%s: only in the second (%q)", k, b[k])) + } + } + return out +} + +// change makes every kind of difference a restore must undo: rows removed, +// added and edited, a sequence advanced, a table dropped, one created. +func (r *restoreRig) change(t *testing.T) { + t.Helper() + r.exec(t, r.url, + `DELETE FROM servers WHERE name = 'bob-smp'`, + `UPDATE servers SET cached_memory_mb = 4096 WHERE name = 'alice-smp'`, + `INSERT INTO platform_settings (key, value) VALUES ('pgint_after', '{"x": 1}')`, + `DROP TABLE zz_pgint_bulk`, + `CREATE TABLE pgint_after (id serial PRIMARY KEY)`, + `INSERT INTO pgint_after DEFAULT VALUES`, + ) +} + +func (r *restoreRig) backup(t *testing.T) string { + t.Helper() + bundle, err := dbbackup.Backup(context.Background(), dbbackup.BackupOptions{ + DatabaseURL: r.url, Dir: r.dir, Label: dbbackup.LabelManual, Tools: r.tools}) + if err != nil { + t.Fatalf("backup: %v", err) + } + return bundle +} + +func (r *restoreRig) restore(bundle string, mut func(*dbbackup.RestoreOptions)) (string, error) { + o := dbbackup.RestoreOptions{DatabaseURL: r.url, Bundle: bundle, Dir: r.dir, Tools: r.tools} + if mut != nil { + mut(&o) + } + _, safety, err := dbbackup.Restore(context.Background(), o) + return safety, err +} + +// A bundle taken of a live schema restores it exactly, over rows and objects +// changed since, refuses to while a client is connected, and leaves a safety +// bundle that undoes the restore. +func TestDBRestoreRoundTrip(t *testing.T) { + r := newRestoreRig(t) + r.seed(t) + want := r.fingerprint(t) + bundle := r.backup(t) + m, err := dbbackup.Verify(bundle) + if err != nil { + t.Fatal(err) + } + ms, _ := store.LoadMigrations() + if m.SchemaVersion != ms[len(ms)-1].Version || !strings.HasPrefix(m.PGDumpVersion, "pg_dump (PostgreSQL) ") { + t.Errorf("manifest schema %d pg_dump %q, want schema %d and pg_dump's version", m.SchemaVersion, m.PGDumpVersion, ms[len(ms)-1].Version) + } + + r.change(t) + changed := r.fingerprint(t) + if len(diff(want, changed)) == 0 { + t.Fatal("the change changed nothing") + } + + // A client on the database: refused before anything is taken or replayed. + drv, err := store.Open(context.Background(), r.url) + if err != nil { + t.Fatal(err) + } + _, err = r.restore(bundle, nil) + drv.Close() + if !errors.Is(err, dbbackup.ErrClientsConnected) { + t.Fatalf("restore under a connected client: %v, want ErrClientsConnected", err) + } + if d := diff(changed, r.fingerprint(t)); len(d) > 0 { + t.Fatalf("a refused restore changed the database: %v", d) + } + if all, _ := dbbackup.List(r.dir); len(all) != 1 { + t.Fatalf("bundles after the refusal = %v, want only the one taken", all) + } + r.waitIdle(t) + + safety, err := r.restore(bundle, nil) + if err != nil { + t.Fatalf("restore: %v", err) + } + if d := diff(want, r.fingerprint(t)); len(d) > 0 { + t.Fatalf("restored database differs from the bundle's: %v", d) + } + // The dump's own freshness record predates its bundle; the restore points + // the panel at the newest bundle on disk, the safety one. + r.with(t, r.suURL, func(d *store.PostgresDriver) { + var raw []byte + if err := d.DB().QueryRow(`SELECT value FROM platform_settings WHERE key = $1`, dbbackup.StatusKey).Scan(&raw); err != nil { + t.Fatalf("freshness record: %v", err) + } + var st dbbackup.Status + if err := json.Unmarshal(raw, &st); err != nil || st.Name != filepath.Base(safety) || st.Label != dbbackup.LabelPreRestore { + t.Errorf("freshness record = %s (%v), want %s", raw, err, filepath.Base(safety)) + } + }) + + if _, err := r.restore(safety, func(o *dbbackup.RestoreOptions) { o.SkipSafetyBackup = true }); err != nil { + t.Fatalf("restore the safety bundle: %v", err) + } + if d := diff(changed, r.fingerprint(t)); len(d) > 0 { + t.Fatalf("the safety bundle did not bring back what the restore replaced: %v", d) + } +} + +// A replay that fails part way, in the server or in pg_restore, leaves the +// database as it was: DROP OWNED and every statement after it roll back. +func TestDBRestoreFailureRollsBack(t *testing.T) { + for _, tc := range []struct { + name string + // prepare returns the bundle to restore, after the database changed. + prepare func(t *testing.T, r *restoreRig, bundle string) string + want string + }{ + {"the server refuses a statement", func(t *testing.T, r *restoreRig, bundle string) string { + // Another role's table where the dump creates the role's own: DROP + // OWNED leaves it, and the dump's CREATE TABLE fails after every + // type, function and earlier table was already made again. + r.exec(t, r.suURL, `CREATE TABLE zz_pgint_bulk (id int)`) + return bundle + }, `ERROR: relation "zz_pgint_bulk" already exists`}, + {"the dump ends early", func(t *testing.T, r *restoreRig, bundle string) string { + return truncatedDump(t, bundle) + }, "pg_restore: "}, + {"pg_restore dies between statements", func(t *testing.T, r *restoreRig, bundle string) string { + // Cut in the data, psql is inside a COPY and fails the stream itself. + // Killed (a timeout, the OOM killer) while psql works through the + // indexes, pg_restore leaves whole statements behind, and only the + // COMMIT it never earned keeps them out. + r.tools = killedBeforeIndexes(t, r) + return bundle + }, "exit status 1"}, + } { + t.Run(tc.name, func(t *testing.T) { + r := newRestoreRig(t) + r.seed(t) + bundle := r.backup(t) + r.change(t) + bundle = tc.prepare(t, r, bundle) + before := r.fingerprint(t) + + _, err := r.restore(bundle, func(o *dbbackup.RestoreOptions) { o.SkipSafetyBackup = true }) + if err == nil || !strings.Contains(err.Error(), "rolled back, the database is unchanged") || !strings.Contains(err.Error(), tc.want) { + t.Fatalf("restore: %v, want a rolled back replay and %q", err, tc.want) + } + if d := diff(before, r.fingerprint(t)); len(d) > 0 { + t.Fatalf("the failed replay changed the database: %v", d) + } + }) + } +} + +// killedBeforeIndexes wraps the tools so the replay's pg_restore streams its +// script up to the first CREATE INDEX and then exits 1. Everything else is the +// real tools and server. +func killedBeforeIndexes(t *testing.T, r *restoreRig) dbbackup.Tools { + t.Helper() + wrapper := filepath.Join(t.TempDir(), "killed") + if err := os.WriteFile(wrapper, []byte(`#!/bin/sh +case " $* " in +*" pg_restore "*" --file=- "*) "$@" | sed '/^CREATE INDEX/,$d'; exit 1 ;; +esac +exec "$@" +`), 0o755); err != nil { + t.Fatal(err) + } + tools := r.tools + if len(tools.Exec) == 0 { + // Under an exec prefix the tools connect as Conn. + tools.Conn = r.url + } + tools.Exec = append([]string{wrapper}, tools.Exec...) + return tools +} + +// truncatedDump copies bundle with db.dump cut to three quarters, inside the +// rows of the last table, and a manifest that vouches for the cut dump: the +// bundle verifies and pg_restore reads its table of contents, then dies part +// way through the data. +func truncatedDump(t *testing.T, bundle string) string { + t.Helper() + f, err := os.Open(bundle) + if err != nil { + t.Fatal(err) + } + defer f.Close() + type entry struct { + hdr *tar.Header + data []byte + } + var entries []entry + tr := tar.NewReader(f) + for { + h, err := tr.Next() + if errors.Is(err, io.EOF) { + break + } + if err != nil { + t.Fatal(err) + } + data, err := io.ReadAll(tr) + if err != nil { + t.Fatal(err) + } + entries = append(entries, entry{h, data}) + } + var m dbbackup.Manifest + if entries[0].hdr.Name != "MANIFEST.json" || json.Unmarshal(entries[0].data, &m) != nil { + t.Fatalf("%s does not start with its manifest", bundle) + } + for i := range entries { + if entries[i].hdr.Name != "db.dump" { + continue + } + entries[i].data = entries[i].data[:len(entries[i].data)*3/4] + sum := sha256.Sum256(entries[i].data) + for j := range m.Files { + if m.Files[j].Name == "db.dump" { + m.Files[j].Size, m.Files[j].SHA256 = int64(len(entries[i].data)), hex.EncodeToString(sum[:]) + } + } + } + if entries[0].data, err = json.Marshal(m); err != nil { + t.Fatal(err) + } + var buf bytes.Buffer + tw := tar.NewWriter(&buf) + for _, e := range entries { + e.hdr.Size = int64(len(e.data)) + if err := tw.WriteHeader(e.hdr); err != nil { + t.Fatal(err) + } + if _, err := tw.Write(e.data); err != nil { + t.Fatal(err) + } + } + if err := tw.Close(); err != nil { + t.Fatal(err) + } + out := filepath.Join(t.TempDir(), filepath.Base(bundle)) + if err := os.WriteFile(out, buf.Bytes(), 0o600); err != nil { + t.Fatal(err) + } + if _, err := dbbackup.Verify(out); err != nil { + t.Fatalf("the cut bundle does not verify: %v", err) + } + return out +} diff --git a/internal/pgint/pgint_test.go b/internal/pgint/pgint_test.go index 10c6b18..072c116 100644 --- a/internal/pgint/pgint_test.go +++ b/internal/pgint/pgint_test.go @@ -5,6 +5,8 @@ // run against fakes that encode the CONTRACT — and PGRepo drifted behind that // contract three times (attempt accounting, a missing JOIN, a missing FOR UPDATE) // while every unit test stayed green. +// dbrestore_test.go does the same for `felis db backup` and `restore`, whose +// fakes stand in for pg_dump, pg_restore, psql and the server. // // They run ONLY against a throwaway database whose name contains "pgint": the // harness drops and recreates the public schema and replays the real embedded