The control plane, as far as identity
Tier 2 exists now. It holds one context of seven, inventory, and does one thing with it: brings its schema up to date. That is step 3 of the substrate bootstrap -- the step the first node cannot get past. Verified against a real PostgreSQL, with the built binary: applied 0001-nodes, reported 'already up to date' on the second run, and the node table is there with the index and the unique constraint the migration asks for. Written in Go, and the image is FROM scratch holding one file. Confirmed by unpacking it. That is the whole argument of ADR 0024: the bundle pins this image by digest and runs it where nothing can check it, so everything in it is something a person has to audit before trusting a first node. Exclusive store ownership is built as a rule about credentials rather than about intentions. There is no mesh-wide connection setting and no way to ask for one -- a context reads MESH_STORE_<ITS OWN NAME> and holds nothing else, so reaching another context's store needs a new variable, which is visible in the declaration that runs it. The migration runner is mostly refusals: an edited migration that already ran, a migration numbered below one that has run, duplicate numbers, misnamed files, empty files. All stop rather than warn, because at the moment any of them is true nobody knows what the database holds. It stops before identity, deliberately. What a node presents to prove who it is has not been decided anywhere, and a migration is the most expensive place in this system to guess. Two tests did not defend what they claimed, and both are fixed rather than removed. One asked only whether Open returned an error, which it did either way -- a bad context name and a missing credential both fail, so deleting the name check changed nothing. The other claimed to prove the migration runs in a transaction, but PostgreSQL already wraps a multi-statement query in one of its own, so it passed with the transaction taken out. What the transaction actually buys is that the schema change and the row recording it commit together, and there is now a test for that which fails when they are split.
This commit is contained in:
@@ -0,0 +1,291 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"testing/fstest"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
)
|
||||
|
||||
// These run against a real PostgreSQL. Not a fake, and for the reason novox/hq ADR 0017 gives
|
||||
// where the host tests the real filesystem: what is being tested here *is* the database's
|
||||
// behaviour — that DDL is transactional, that an advisory lock serialises, that a checksum
|
||||
// mismatch is caught against a record the database actually kept. A fake would assert that the
|
||||
// fake behaves as expected.
|
||||
//
|
||||
// `make check` raises one. Without it these skip, and say so rather than passing.
|
||||
|
||||
func admin(t *testing.T) string {
|
||||
t.Helper()
|
||||
dsn := os.Getenv("MESH_TEST_POSTGRES")
|
||||
if dsn == "" {
|
||||
t.Skip("no MESH_TEST_POSTGRES; run `make check` to raise one")
|
||||
}
|
||||
return dsn
|
||||
}
|
||||
|
||||
// freshStore gives a test its own empty database.
|
||||
//
|
||||
// Its own, rather than a shared one cleaned between tests: these tests are about what a migration
|
||||
// runner does to a schema, and a leftover table from a previous test is indistinguishable from
|
||||
// the bug this whole package exists to catch.
|
||||
func freshStore(t *testing.T) *Store {
|
||||
t.Helper()
|
||||
dsn := admin(t)
|
||||
name := fmt.Sprintf("test_%s_%d", strings.ToLower(strings.NewReplacer(
|
||||
"/", "_", "-", "_").Replace(t.Name())), time.Now().UnixNano()%1_000_000)
|
||||
if len(name) > 60 {
|
||||
name = name[:60]
|
||||
}
|
||||
|
||||
conn, err := pgx.Connect(t.Context(), dsn)
|
||||
if err != nil {
|
||||
t.Fatalf("cannot reach the test PostgreSQL: %v", err)
|
||||
}
|
||||
if _, err := conn.Exec(t.Context(), "create database "+name); err != nil {
|
||||
t.Fatalf("cannot create %s: %v", name, err)
|
||||
}
|
||||
conn.Close(t.Context())
|
||||
|
||||
t.Setenv(Variable("testing"), replaceDatabase(dsn, name))
|
||||
s, err := Open(t.Context(), "testing")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
s.Close()
|
||||
c, err := pgx.Connect(context.Background(), dsn)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer c.Close(context.Background())
|
||||
_, _ = c.Exec(context.Background(), "drop database if exists "+name+" with (force)")
|
||||
})
|
||||
if err := s.Ready(t.Context(), 20*time.Second); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
func replaceDatabase(dsn, name string) string {
|
||||
cut := strings.LastIndex(dsn, "/")
|
||||
rest := ""
|
||||
if q := strings.Index(dsn[cut:], "?"); q >= 0 {
|
||||
rest = dsn[cut+q:]
|
||||
}
|
||||
return dsn[:cut] + "/" + name + rest
|
||||
}
|
||||
|
||||
func TestTheSchemaIsAppliedAndRecorded(t *testing.T) {
|
||||
s := freshStore(t)
|
||||
migrations := []Migration{{Number: 1, Name: "people", SQL: "create table person (id int)", Checksum: "a"}}
|
||||
|
||||
done, err := s.Migrate(t.Context(), migrations)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(done) != 1 {
|
||||
t.Fatalf("applied %d migrations, expected 1", len(done))
|
||||
}
|
||||
|
||||
// Read back from the system rather than trusting the return value (novox/hq ADR 0018).
|
||||
var exists bool
|
||||
if err := s.Pool().QueryRow(t.Context(),
|
||||
`select exists (select 1 from information_schema.tables where table_name = 'person')`,
|
||||
).Scan(&exists); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !exists {
|
||||
t.Error("Migrate reported success and the table is not there")
|
||||
}
|
||||
|
||||
applied, err := s.AppliedMigrations(t.Context())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(applied) != 1 || applied[0].Checksum != "a" {
|
||||
t.Errorf("the record says %+v", applied)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunningTwiceChangesNothing(t *testing.T) {
|
||||
// The bootstrap runs this, and so does every restart of the control plane. A second run that
|
||||
// re-applied the schema would fail on the first `create table`, so a control plane would come
|
||||
// up exactly once.
|
||||
s := freshStore(t)
|
||||
migrations := []Migration{{Number: 1, Name: "people", SQL: "create table person (id int)", Checksum: "a"}}
|
||||
|
||||
if _, err := s.Migrate(t.Context(), migrations); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
done, err := s.Migrate(t.Context(), migrations)
|
||||
if err != nil {
|
||||
t.Fatalf("the second run failed: %v", err)
|
||||
}
|
||||
if len(done) != 0 {
|
||||
t.Errorf("the second run applied %d migrations", len(done))
|
||||
}
|
||||
}
|
||||
|
||||
func TestAFailedMigrationLeavesNothingBehind(t *testing.T) {
|
||||
// Note what this does and does not defend. PostgreSQL wraps a multi-statement simple query in
|
||||
// an implicit transaction of its own, so this passes with this package's transaction removed
|
||||
// — it was checked, and it did. What it defends is the database and driver behaviour relied
|
||||
// on: a driver sending each statement separately would break it, and nothing else would say
|
||||
// so. The property that belongs to this code is the next test.
|
||||
s := freshStore(t)
|
||||
migrations := []Migration{{
|
||||
Number: 1, Name: "half", Checksum: "a",
|
||||
SQL: `create table kept (id int);
|
||||
create table broken (id int) this is not sql;`,
|
||||
}}
|
||||
|
||||
if _, err := s.Migrate(t.Context(), migrations); err == nil {
|
||||
t.Fatal("a migration with a syntax error reported success")
|
||||
}
|
||||
|
||||
var tables int
|
||||
if err := s.Pool().QueryRow(t.Context(),
|
||||
`select count(*) from information_schema.tables where table_name in ('kept','broken')`,
|
||||
).Scan(&tables); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if tables != 0 {
|
||||
t.Errorf("%d table(s) survived a failed migration; it must be all or nothing", tables)
|
||||
}
|
||||
}
|
||||
|
||||
func TestASchemaChangeAndItsRecordCommitTogether(t *testing.T) {
|
||||
// This is what the explicit transaction is for, and all it is for.
|
||||
//
|
||||
// Split the change from the row saying it happened, and a schema moves with nothing recording
|
||||
// it — so the next run finds the migration outstanding and applies it to a database that
|
||||
// already has it. The failure surfaces as a broken migration rather than as a lost record.
|
||||
//
|
||||
// The real case is the process dying between the two, which a test cannot arrange. Standing
|
||||
// in for it: a migration that makes its own record impossible to write.
|
||||
s := freshStore(t)
|
||||
if _, err := s.AppliedMigrations(t.Context()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
migrations := []Migration{{
|
||||
Number: 1, Name: "hostile", Checksum: "a",
|
||||
SQL: "create table kept (id int); drop table migration;",
|
||||
}}
|
||||
if _, err := s.Migrate(t.Context(), migrations); err == nil {
|
||||
t.Fatal("a migration whose record could not be written reported success")
|
||||
}
|
||||
|
||||
var kept, ledger bool
|
||||
if err := s.Pool().QueryRow(t.Context(),
|
||||
`select exists (select 1 from information_schema.tables where table_name = 'kept'),
|
||||
exists (select 1 from information_schema.tables where table_name = 'migration')`,
|
||||
).Scan(&kept, &ledger); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if kept {
|
||||
t.Error("the schema change survived although nothing recorded it; the next run would " +
|
||||
"apply it again, to a database that already has it")
|
||||
}
|
||||
if !ledger {
|
||||
t.Error("the migration record did not come back with the rollback")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAChangedMigrationIsRefusedAgainstARealRecord(t *testing.T) {
|
||||
// The same refusal as the unit test, against a record PostgreSQL actually kept — which is
|
||||
// what the guard protects, and the unit test can only model.
|
||||
s := freshStore(t)
|
||||
first := []Migration{{Number: 1, Name: "people", SQL: "create table person (id int)", Checksum: "a"}}
|
||||
if _, err := s.Migrate(t.Context(), first); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
edited := []Migration{{Number: 1, Name: "people", SQL: "create table person (id bigint)", Checksum: "b"}}
|
||||
if _, err := s.Migrate(t.Context(), edited); err == nil {
|
||||
t.Fatal("a migration edited after it ran was accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestTwoRunnersDoNotRaceEachOther(t *testing.T) {
|
||||
// A restart during a slow migration produces exactly this: two copies of the control plane
|
||||
// migrating one database. Without the advisory lock both read an empty record, both decide
|
||||
// everything is outstanding, and the second fails on `create table` — which looks like a
|
||||
// broken migration rather than a race.
|
||||
s := freshStore(t)
|
||||
migrations := []Migration{{
|
||||
Number: 1, Name: "slow", Checksum: "a",
|
||||
SQL: "create table slow (id int); select pg_sleep(0.4);",
|
||||
}}
|
||||
|
||||
var wg sync.WaitGroup
|
||||
results := make([]error, 2)
|
||||
counts := make([]int, 2)
|
||||
for i := range results {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
done, err := s.Migrate(context.Background(), migrations)
|
||||
results[i], counts[i] = err, len(done)
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
for i, err := range results {
|
||||
if err != nil {
|
||||
t.Errorf("runner %d failed: %v", i, err)
|
||||
}
|
||||
}
|
||||
if counts[0]+counts[1] != 1 {
|
||||
t.Errorf("the migration was applied %d times between two runners; exactly one should have "+
|
||||
"done the work and the other should have found nothing to do", counts[0]+counts[1])
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadedMigrationsApplyToARealDatabase(t *testing.T) {
|
||||
// A migration that parses and does not run is the failure this catches. The unit tests read
|
||||
// files and check names; nothing there executes SQL.
|
||||
s := freshStore(t)
|
||||
migrations, err := LoadMigrations(fstest.MapFS{
|
||||
"migrations/0001-first.sql": {Data: []byte("create table a (id int);")},
|
||||
"migrations/0002-second.sql": {Data: []byte("alter table a add column b text;")},
|
||||
}, "migrations")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
done, err := s.Migrate(t.Context(), migrations)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(done) != 2 || done[0].Number != 1 || done[1].Number != 2 {
|
||||
t.Fatalf("applied %+v", done)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadyRefusesADatabaseThatWillNotAnswer(t *testing.T) {
|
||||
// Ready is what stands between the bootstrap and a control plane that starts against a
|
||||
// database still coming up. It has to give up rather than block for ever, and it has to fail
|
||||
// when nothing is there.
|
||||
admin(t)
|
||||
t.Setenv(Variable("testing"), "postgres://nobody@127.0.0.1:1/nothing")
|
||||
s, err := Open(t.Context(), "testing")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
start := time.Now()
|
||||
if err := s.Ready(t.Context(), 1*time.Second); err == nil {
|
||||
t.Fatal("Ready returned success against a port with nothing on it")
|
||||
}
|
||||
if elapsed := time.Since(start); elapsed > 10*time.Second {
|
||||
t.Errorf("Ready took %s to give up on a 1s budget", elapsed)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,257 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io/fs"
|
||||
"path"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// novox/hq ADR 0013: every schema change is a numbered migration. This applies them, and most of
|
||||
// what follows is refusals rather than application — the applying is four lines.
|
||||
//
|
||||
// The refusals are the point. A migration runner that only moves forward is easy; what makes a
|
||||
// schema trustworthy months later is that it will not run when the record and the files disagree,
|
||||
// because at that moment nobody knows what the database actually contains and the honest thing is
|
||||
// to stop.
|
||||
|
||||
// migrationFile is `0001-name.sql` — the number, then what it does.
|
||||
var migrationFile = regexp.MustCompile(`^(\d{4})-([a-z0-9-]+)\.sql$`)
|
||||
|
||||
// Migration is one numbered change, read from the files the binary carries.
|
||||
type Migration struct {
|
||||
Number int
|
||||
Name string
|
||||
SQL string
|
||||
Checksum string
|
||||
}
|
||||
|
||||
// Applied is one row of the record — what ran, and what it looked like when it did.
|
||||
type Applied struct {
|
||||
Number int
|
||||
Name string
|
||||
Checksum string
|
||||
}
|
||||
|
||||
// ledger is where the record lives, inside the context's own database.
|
||||
//
|
||||
// In the same database as the schema it describes, deliberately: a record kept anywhere else can
|
||||
// be restored separately from the thing it describes, and then it is not a record of anything.
|
||||
const ledger = `
|
||||
create table if not exists migration (
|
||||
number integer primary key,
|
||||
name text not null,
|
||||
checksum text not null,
|
||||
applied timestamptz not null default now()
|
||||
)`
|
||||
|
||||
// lockKey is the advisory lock every migration run takes.
|
||||
//
|
||||
// One control plane runs (novox/hq ADR 0006), so a race needs two copies of it — which is exactly
|
||||
// what a restart during a slow migration produces, and is not rare enough to leave to chance. An
|
||||
// arbitrary constant; it only has to be the same in every copy of this binary.
|
||||
const lockKey int64 = 6_845_121_074
|
||||
|
||||
// LoadMigrations reads the migrations a context carries.
|
||||
//
|
||||
// From an embedded filesystem rather than from disk. The binary and its schema travel as one
|
||||
// thing, so a container image cannot be running one version of the code against a directory of
|
||||
// migrations from another.
|
||||
func LoadMigrations(files fs.FS, dir string) ([]Migration, error) {
|
||||
entries, err := fs.ReadDir(files, dir)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot read the migrations in %s: %w", dir, err)
|
||||
}
|
||||
|
||||
var migrations []Migration
|
||||
seen := map[int]string{}
|
||||
for _, entry := range entries {
|
||||
if entry.IsDir() {
|
||||
continue
|
||||
}
|
||||
match := migrationFile.FindStringSubmatch(entry.Name())
|
||||
if match == nil {
|
||||
// Not skipped quietly. A file that was meant to be a migration and is misnamed would
|
||||
// otherwise be ignored in silence, and the schema would simply lack it — the failure
|
||||
// arriving later as a missing column, a long way from its cause.
|
||||
return nil, fmt.Errorf(
|
||||
"%s is not a migration filename; it must be NNNN-what-it-does.sql, and a file "+
|
||||
"in this directory that is not a migration cannot be told apart from one "+
|
||||
"that was misnamed", path.Join(dir, entry.Name()))
|
||||
}
|
||||
number, _ := strconv.Atoi(match[1])
|
||||
if other, clash := seen[number]; clash {
|
||||
return nil, fmt.Errorf(
|
||||
"two migrations are numbered %04d: %s and %s. Order is the whole guarantee, and "+
|
||||
"two files with one number have none", number, other, entry.Name())
|
||||
}
|
||||
seen[number] = entry.Name()
|
||||
|
||||
body, err := fs.ReadFile(files, path.Join(dir, entry.Name()))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if strings.TrimSpace(string(body)) == "" {
|
||||
return nil, fmt.Errorf(
|
||||
"%s is empty. An empty migration records that something happened and changes "+
|
||||
"nothing, which is the one state that cannot be told from a mistake",
|
||||
entry.Name())
|
||||
}
|
||||
sum := sha256.Sum256(body)
|
||||
migrations = append(migrations, Migration{
|
||||
Number: number,
|
||||
Name: match[2],
|
||||
SQL: string(body),
|
||||
Checksum: hex.EncodeToString(sum[:]),
|
||||
})
|
||||
}
|
||||
|
||||
sort.Slice(migrations, func(i, j int) bool { return migrations[i].Number < migrations[j].Number })
|
||||
return migrations, nil
|
||||
}
|
||||
|
||||
// AppliedMigrations reads the record of what has run.
|
||||
func (s *Store) AppliedMigrations(ctx context.Context) ([]Applied, error) {
|
||||
if _, err := s.pool.Exec(ctx, ledger); err != nil {
|
||||
return nil, fmt.Errorf("cannot create the migration record in %s: %w", s.context, err)
|
||||
}
|
||||
rows, err := s.pool.Query(ctx,
|
||||
`select number, name, checksum from migration order by number`)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var applied []Applied
|
||||
for rows.Next() {
|
||||
var a Applied
|
||||
if err := rows.Scan(&a.Number, &a.Name, &a.Checksum); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
applied = append(applied, a)
|
||||
}
|
||||
return applied, rows.Err()
|
||||
}
|
||||
|
||||
// Pending is what has not run, having first established that what has run still matches.
|
||||
//
|
||||
// Two refusals, and both describe a database whose contents are no longer known:
|
||||
//
|
||||
// - a migration that ran and whose file has since changed. The database holds the old version
|
||||
// and the repository holds the new one, and nothing anywhere holds the difference.
|
||||
// - a migration numbered below one that already ran, which has not run itself. Almost always
|
||||
// two branches picking the same next number, merged in the order they happened to land. The
|
||||
// file is fine; applying it now would run the schema in an order nobody tested.
|
||||
//
|
||||
// Both are stops rather than warnings. There is no correct guess about a schema.
|
||||
func Pending(migrations []Migration, applied []Applied) ([]Migration, error) {
|
||||
record := map[int]Applied{}
|
||||
highest := 0
|
||||
for _, a := range applied {
|
||||
record[a.Number] = a
|
||||
if a.Number > highest {
|
||||
highest = a.Number
|
||||
}
|
||||
}
|
||||
|
||||
var pending []Migration
|
||||
var problems []string
|
||||
for _, m := range migrations {
|
||||
ran, wasApplied := record[m.Number]
|
||||
if wasApplied {
|
||||
if ran.Checksum != m.Checksum {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"migration %04d-%s ran against this database, and the file has changed since. "+
|
||||
"The database holds what the old file said and nothing holds the "+
|
||||
"difference. A migration that has run is finished — write a new one",
|
||||
m.Number, m.Name))
|
||||
}
|
||||
continue
|
||||
}
|
||||
if m.Number < highest {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"migration %04d-%s has never run, but %04d has. This is usually two branches "+
|
||||
"taking the same next number. Applying it now would run this schema in an "+
|
||||
"order that was never tested — renumber it above %04d",
|
||||
m.Number, m.Name, highest, highest))
|
||||
continue
|
||||
}
|
||||
pending = append(pending, m)
|
||||
}
|
||||
|
||||
if len(problems) > 0 {
|
||||
return nil, errors.New("this database and these migrations disagree:\n - " +
|
||||
strings.Join(problems, "\n - "))
|
||||
}
|
||||
return pending, nil
|
||||
}
|
||||
|
||||
// Migrate applies everything outstanding, in order, and reports what it did.
|
||||
//
|
||||
// Each migration runs in its own transaction together with the row recording it, so the two
|
||||
// cannot come apart: PostgreSQL runs DDL transactionally, so a migration that fails half way
|
||||
// leaves neither the change nor a claim that the change was made.
|
||||
func (s *Store) Migrate(ctx context.Context, migrations []Migration) ([]Migration, error) {
|
||||
conn, err := s.pool.Acquire(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer conn.Release()
|
||||
|
||||
// Held on one connection for the whole run, and released when it goes back to the pool. A
|
||||
// second copy of this process waits here rather than interleaving with the first.
|
||||
if _, err := conn.Exec(ctx, `select pg_advisory_lock($1)`, lockKey); err != nil {
|
||||
return nil, fmt.Errorf("cannot take the migration lock on %s: %w", s.context, err)
|
||||
}
|
||||
defer func() {
|
||||
_, _ = conn.Exec(context.WithoutCancel(ctx), `select pg_advisory_unlock($1)`, lockKey)
|
||||
}()
|
||||
|
||||
// Read *after* the lock. Reading before it would mean deciding what is outstanding from a
|
||||
// picture taken before another process was known to be finished with it.
|
||||
applied, err := s.AppliedMigrations(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
pending, err := Pending(migrations, applied)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var done []Migration
|
||||
for _, m := range pending {
|
||||
err := func() error {
|
||||
tx, err := conn.Begin(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = tx.Rollback(context.WithoutCancel(ctx)) }()
|
||||
|
||||
if _, err := tx.Exec(ctx, m.SQL); err != nil {
|
||||
return fmt.Errorf("migration %04d-%s failed, and nothing it did was kept: %w",
|
||||
m.Number, m.Name, err)
|
||||
}
|
||||
if _, err := tx.Exec(ctx,
|
||||
`insert into migration (number, name, checksum) values ($1, $2, $3)`,
|
||||
m.Number, m.Name, m.Checksum); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Commit(ctx)
|
||||
}()
|
||||
if err != nil {
|
||||
// What succeeded stays applied and stays recorded, which is why they are reported
|
||||
// alongside the failure: a person deciding what to do next needs to know the schema
|
||||
// moved, not only that the run did not finish.
|
||||
return done, err
|
||||
}
|
||||
done = append(done, m)
|
||||
}
|
||||
return done, nil
|
||||
}
|
||||
@@ -0,0 +1,206 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
"testing/fstest"
|
||||
)
|
||||
|
||||
// The refusals in this file are the ones that decide whether a schema can be trusted months
|
||||
// later. Each has a wrong answer that looks helpful — skip it, apply it anyway, warn and carry on
|
||||
// — and each of those turns "this database disagrees with these files" into a silence.
|
||||
|
||||
func load(t *testing.T, files fstest.MapFS) ([]Migration, error) {
|
||||
t.Helper()
|
||||
return LoadMigrations(files, "migrations")
|
||||
}
|
||||
|
||||
func file(body string) *fstest.MapFile { return &fstest.MapFile{Data: []byte(body)} }
|
||||
|
||||
func TestMigrationsAreReadInNumericOrder(t *testing.T) {
|
||||
// Read from a directory, which has no order of its own. Alphabetical happens to agree with
|
||||
// numeric while the numbers are the same width, which is exactly why this is asserted rather
|
||||
// than assumed.
|
||||
got, err := load(t, fstest.MapFS{
|
||||
"migrations/0010-tenth.sql": file("select 10"),
|
||||
"migrations/0002-second.sql": file("select 2"),
|
||||
"migrations/0001-first.sql": file("select 1"),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var order []int
|
||||
for _, m := range got {
|
||||
order = append(order, m.Number)
|
||||
}
|
||||
if len(order) != 3 || order[0] != 1 || order[1] != 2 || order[2] != 10 {
|
||||
t.Errorf("migrations came back in order %v; they run in the order they are returned", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAMisnamedFileIsAnErrorRatherThanSkipped(t *testing.T) {
|
||||
// The tempting behaviour is to ignore anything that does not match, so that a README can sit
|
||||
// in the directory. The cost is that a migration named `001-thing.sql` or `0002_thing.sql` is
|
||||
// then ignored in silence, and the schema simply lacks it — surfacing later as a missing
|
||||
// column, a long way from the file that was misnamed.
|
||||
_, err := load(t, fstest.MapFS{
|
||||
"migrations/0001-first.sql": file("select 1"),
|
||||
"migrations/0002_second.sql": file("select 2"),
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatal("a misnamed migration was skipped silently; it would never run and nothing would say so")
|
||||
}
|
||||
}
|
||||
|
||||
func TestTwoMigrationsWithOneNumberAreRefused(t *testing.T) {
|
||||
// Order is the entire guarantee. Two files with one number have none, and whichever the
|
||||
// filesystem returned first would win.
|
||||
_, err := load(t, fstest.MapFS{
|
||||
"migrations/0001-first.sql": file("select 1"),
|
||||
"migrations/0001-also-first.sql": file("select 2"),
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatal("two migrations numbered 0001 were accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAnEmptyMigrationIsRefused(t *testing.T) {
|
||||
// An empty migration records that something happened and changes nothing — the one state that
|
||||
// cannot be told apart from a mistake, and it is recorded as done for ever.
|
||||
_, err := load(t, fstest.MapFS{"migrations/0001-nothing.sql": file(" \n\t\n")})
|
||||
if err == nil {
|
||||
t.Fatal("an empty migration was accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAChangedMigrationIsRefused(t *testing.T) {
|
||||
// The one that matters most. The database holds what the old file said; the repository holds
|
||||
// the new one; nothing anywhere holds the difference. Applying it again would be wrong and
|
||||
// skipping it silently leaves the two permanently out of step.
|
||||
migrations := []Migration{{Number: 1, Name: "nodes", Checksum: "aaaa"}}
|
||||
applied := []Applied{{Number: 1, Name: "nodes", Checksum: "bbbb"}}
|
||||
|
||||
_, err := Pending(migrations, applied)
|
||||
if err == nil {
|
||||
t.Fatal("a migration whose file changed after it ran was accepted")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "0001-nodes") {
|
||||
t.Errorf("the refusal does not name the migration: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAnUnchangedMigrationIsNotPending(t *testing.T) {
|
||||
// The other half, and the one that makes the run idempotent. Without it every start would
|
||||
// re-apply the whole schema.
|
||||
migrations := []Migration{{Number: 1, Name: "nodes", Checksum: "aaaa"}}
|
||||
applied := []Applied{{Number: 1, Name: "nodes", Checksum: "aaaa"}}
|
||||
|
||||
pending, err := Pending(migrations, applied)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(pending) != 0 {
|
||||
t.Errorf("an already-applied migration came back as pending; every start would re-run it")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAMigrationArrivingBelowTheHighWaterMarkIsRefused(t *testing.T) {
|
||||
// Two branches take the same next number; they merge in whatever order they landed. The file
|
||||
// is fine and applying it now would run the schema in an order nobody tested — which is the
|
||||
// same class of fault as applying them out of order deliberately, arriving by accident.
|
||||
migrations := []Migration{
|
||||
{Number: 1, Name: "nodes", Checksum: "aaaa"},
|
||||
{Number: 2, Name: "late", Checksum: "cccc"},
|
||||
{Number: 3, Name: "third", Checksum: "bbbb"},
|
||||
}
|
||||
applied := []Applied{
|
||||
{Number: 1, Name: "nodes", Checksum: "aaaa"},
|
||||
{Number: 3, Name: "third", Checksum: "bbbb"},
|
||||
}
|
||||
|
||||
_, err := Pending(migrations, applied)
|
||||
if err == nil {
|
||||
t.Fatal("a migration numbered below one that already ran was applied out of order")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "0002-late") {
|
||||
t.Errorf("the refusal does not name the migration: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEveryDisagreementIsReportedNotOnlyTheFirst(t *testing.T) {
|
||||
// A person looking at this is deciding what to do about a schema. Being told one problem,
|
||||
// fixing it, and being told the next is how a single decision becomes four.
|
||||
migrations := []Migration{
|
||||
{Number: 1, Name: "one", Checksum: "aaaa"},
|
||||
{Number: 2, Name: "two", Checksum: "cccc"},
|
||||
}
|
||||
applied := []Applied{
|
||||
{Number: 1, Name: "one", Checksum: "changed"},
|
||||
{Number: 3, Name: "three", Checksum: "dddd"},
|
||||
}
|
||||
|
||||
_, err := Pending(migrations, applied)
|
||||
if err == nil {
|
||||
t.Fatal("expected refusals")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "0001-one") || !strings.Contains(err.Error(), "0002-two") {
|
||||
t.Errorf("only some problems were reported: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAContextNameThatCannotBeADatabaseIsRefused(t *testing.T) {
|
||||
// The name becomes an environment variable and a database name. One that is valid in one and
|
||||
// not the other fails at bootstrap, on a machine with no mesh on it and nobody watching.
|
||||
//
|
||||
// Asserted on *which* refusal fired, not merely that one did. Open has a second reason to
|
||||
// fail a line later — no credential — and an unusable name would have produced an error
|
||||
// either way, so a test asking only "was there an error" passes with this check deleted.
|
||||
// It was written that way first and confirmed to defend nothing.
|
||||
for _, name := range []string{"", "Inventory", "my-context", "9lives", "drop table"} {
|
||||
_, err := Open(t.Context(), name)
|
||||
if err == nil {
|
||||
t.Errorf("%q was accepted as a context name", name)
|
||||
continue
|
||||
}
|
||||
if !strings.Contains(err.Error(), "not a usable context name") {
|
||||
t.Errorf("%q was refused for the wrong reason: %v", name, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAContextWithoutItsCredentialIsRefused(t *testing.T) {
|
||||
// And the message has to say that reusing another context's connection is not the remedy,
|
||||
// because it is the obvious one and it is how ADR 0008 gets quietly undone.
|
||||
t.Setenv(Variable("inventory"), "")
|
||||
_, err := Open(t.Context(), "inventory")
|
||||
if err == nil {
|
||||
t.Fatal("a context with no credential opened a store")
|
||||
}
|
||||
if !strings.Contains(err.Error(), Variable("inventory")) {
|
||||
t.Errorf("the refusal does not name the variable that is missing: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAMalformedConnectionStringIsNotQuotedBack(t *testing.T) {
|
||||
// The value carries a password. An error message that quotes what it could not parse puts it
|
||||
// into a log, on the one machine where the bootstrap output is being watched by a person.
|
||||
secret := "hunter2-this-must-not-appear"
|
||||
t.Setenv(Variable("inventory"), "postgres://user:"+secret+"@host:notaport/db")
|
||||
|
||||
_, err := Open(t.Context(), "inventory")
|
||||
if err == nil {
|
||||
t.Fatal("a malformed connection string was accepted")
|
||||
}
|
||||
if strings.Contains(err.Error(), secret) {
|
||||
t.Errorf("the password appeared in the error: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheVariableAndDatabaseFollowTheContextName(t *testing.T) {
|
||||
if got := Variable("inventory"); got != "MESH_STORE_INVENTORY" {
|
||||
t.Errorf("Variable(inventory) = %q", got)
|
||||
}
|
||||
if got := Database("inventory"); got != "inventory" {
|
||||
t.Errorf("Database(inventory) = %q", got)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,135 @@
|
||||
// Package store is how a context reaches the database it exclusively owns.
|
||||
//
|
||||
// novox/hq ADR 0008: a context is granted only what it exclusively owns — no shared writes, no
|
||||
// read-only role on another context's store. That is a rule about credentials, so this package
|
||||
// makes it a rule about credentials rather than a rule about intentions.
|
||||
//
|
||||
// There is no mesh-wide connection string and no way to ask for one. A store is opened by naming
|
||||
// a context, and the settings for that context come from an environment variable named after it.
|
||||
// A control plane process that has been granted `inventory` holds MESH_STORE_INVENTORY and
|
||||
// nothing else, so reaching another context's store is not a matter of restraint — the process
|
||||
// has no address for it and no credential to present.
|
||||
//
|
||||
// Which is also how the rule is *checked*: what a context can reach is visible in the
|
||||
// declaration that runs it, as the list of variables it was given.
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"regexp"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// contextName is what a context may be called.
|
||||
//
|
||||
// Constrained because the name becomes part of an environment variable and part of a database
|
||||
// name, and a name that is valid in one and not the other is a fault discovered at bootstrap on
|
||||
// a machine with no mesh on it.
|
||||
var contextName = regexp.MustCompile(`^[a-z][a-z0-9]*$`)
|
||||
|
||||
// Store is one context's database.
|
||||
type Store struct {
|
||||
context string
|
||||
pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
// Variable is the environment variable holding a context's connection settings.
|
||||
//
|
||||
// Exported because the bootstrap has to set it and the declaration has to name it, and both
|
||||
// should read it from here rather than spell it out again.
|
||||
func Variable(context string) string {
|
||||
return "MESH_STORE_" + strings.ToUpper(context)
|
||||
}
|
||||
|
||||
// Database is what a context's database is called.
|
||||
//
|
||||
// Named after the context, so that a person looking at a PostgreSQL server can see which
|
||||
// contexts exist without a map. There is deliberately no database named for the mesh as a whole:
|
||||
// novox/hq ADR 0006 records that *the mesh database* names a thing that will not exist.
|
||||
func Database(context string) string { return context }
|
||||
|
||||
// Open connects to the database a context owns.
|
||||
//
|
||||
// The settings are read from the environment rather than passed in, which is not indirection for
|
||||
// its own sake: it means no caller anywhere can hand a context a connection to something else.
|
||||
func Open(ctx context.Context, name string) (*Store, error) {
|
||||
if !contextName.MatchString(name) {
|
||||
return nil, fmt.Errorf(
|
||||
"%q is not a usable context name: it becomes an environment variable and a database "+
|
||||
"name, so it must be lower-case letters and digits, starting with a letter", name)
|
||||
}
|
||||
|
||||
dsn := os.Getenv(Variable(name))
|
||||
if strings.TrimSpace(dsn) == "" {
|
||||
return nil, fmt.Errorf(
|
||||
"this process has no %s, so it was not granted the %s store. A context reaches only "+
|
||||
"the store it exclusively owns (novox/hq ADR 0008), so this is either the wrong "+
|
||||
"context or a missing grant — it is never something to work around by reusing "+
|
||||
"another context's connection", Variable(name), name)
|
||||
}
|
||||
|
||||
config, err := pgxpool.ParseConfig(dsn)
|
||||
if err != nil {
|
||||
// Deliberately not wrapping the driver's error verbatim into a message that gets logged:
|
||||
// a malformed DSN often *is* the password, and the value is the one thing here that must
|
||||
// not be quoted back.
|
||||
return nil, fmt.Errorf(
|
||||
"the connection settings in %s could not be read; the value is not quoted here "+
|
||||
"because it carries a password", Variable(name))
|
||||
}
|
||||
|
||||
pool, err := pgxpool.NewWithConfig(ctx, config)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot open the %s store: %w", name, err)
|
||||
}
|
||||
return &Store{context: name, pool: pool}, nil
|
||||
}
|
||||
|
||||
// Context is which context this store belongs to.
|
||||
func (s *Store) Context() string { return s.context }
|
||||
|
||||
// Pool is the connection pool, for the context that owns it.
|
||||
func (s *Store) Pool() *pgxpool.Pool { return s.pool }
|
||||
|
||||
// Close releases the connections.
|
||||
func (s *Store) Close() {
|
||||
if s.pool != nil {
|
||||
s.pool.Close()
|
||||
}
|
||||
}
|
||||
|
||||
// Ready waits until the database answers, or gives up.
|
||||
//
|
||||
// A read-back rather than a connect: pgxpool connects lazily, so a Store that was opened without
|
||||
// error proves only that a string parsed. The bootstrap raises PostgreSQL and the control plane
|
||||
// moments later, and "the container is running" is not "the database will answer" — that
|
||||
// distinction has already cost a debugging session on this project once.
|
||||
func (s *Store) Ready(ctx context.Context, within time.Duration) error {
|
||||
deadline := time.Now().Add(within)
|
||||
var last error
|
||||
for {
|
||||
err := s.pool.Ping(ctx)
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
last = err
|
||||
if ctx.Err() != nil {
|
||||
return errors.Join(ctx.Err(), last)
|
||||
}
|
||||
if time.Now().After(deadline) {
|
||||
return fmt.Errorf(
|
||||
"the %s store did not answer within %s: %w", s.context, within, last)
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return errors.Join(ctx.Err(), last)
|
||||
case <-time.After(250 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user