test(db): 恢复链在真实 PostgreSQL 上跑往返与三种失败回滚,-yes 闸门测试真正走到闸门,e2e 按文档做一次同机恢复
This commit is contained in:
7 files changed
+599
-16
No files matched your search
@@ -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
|
||||
|
||||
+9
-2
@@ -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 <container>' FELIS_TEST_PG_URL=... go test -tags pgint ./internal/pgint/
|
||||
```
|
||||
|
||||
Build the CLI:
|
||||
|
||||
|
||||
+49
-11
@@ -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:[email protected]: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:[email protected]: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")
|
||||
|
||||
+57
-2
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 <service container>`, 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
|
||||
}
|
||||
@@ -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
|
||||
|
||||
Reference in new issue
Block a user