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 }