Files
mesh-controller/internal/store/migrate.go
jschoubben 306c4ca13b 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.
2026-08-29 02:44:09 +02:00

258 lines
8.8 KiB
Go

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
}