stacks: the box converts a PostgreSQL major as a guarded-update step (09 6.4 part 10, decisions 35/37/38)
gates / gates (push) Successful in 25s

A step whose ladder entry carries engine_conversion {service, engine, from, to}
converts the database: the old engine alone, the check (owners, roles,
extensions, per-table row counts), pg_dumpall validated by its completion line,
the volume emptied only after the undo copy's marker is validated again, the new
engine alone, the load with ON_ERROR_STOP, the check again + PG_VERSION. Any
failure goes to the existing undo; a restart during converting is undone.
A PostgreSQL major move without the mark is refused before anything moves.
The old datadir's copy is kept until a backup is proven after the conversion.
17 tests, 9 red-proofs (audits/night-2026-09-26/B/).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
This commit is contained in:
2026-09-25 12:57:03 +02:00
parent 3d49df1e5a
commit 2caae38a71
13 changed files with 1523 additions and 12 deletions
+864
View File
@@ -0,0 +1,864 @@
package stacks
import (
"bufio"
"context"
"fmt"
"io"
"os"
"path/filepath"
"regexp"
"sort"
"strconv"
"strings"
"syscall"
"time"
"gitea.dooplex.hu/admin/felhom-controller/internal/dockerexec"
"gopkg.in/yaml.v3"
)
// ── The PostgreSQL major conversion (`09` §3 decisions 16, 35, 37, 38; §6.4 part 10; v0.273.0) ─────
//
// WHAT IT IS FOR. The postgres image converts nothing: it REFUSES to start on an older major's datadir
// (R-463, measured 2026-09-21: "database files are incompatible with server"), and 18 refuses even an
// EMPTY volume at the old mount point (measured 2026-09-25, audits/night-2026-09-26/A/). So a step that
// moves a PostgreSQL image across a major must rebuild the datadir: save everything from the OLD engine,
// start the NEW one empty, load it back, check.
//
// THE BOX NEVER GUESSES. It converts only when the ladder entry of the step carries the harness's mark
// (`engine_conversion: {service, engine, from, to}`), written by `upgrade-test.py --write-ladder` after
// the product's own conversion was proven on a scratch box. A step that moves a PostgreSQL image across a
// major WITHOUT the mark is refused before anything moves (defence in depth — the catalog gate should
// never let one through). Pinned by TestConvert_MajorWithoutMarkIsRefused.
//
// IT RIDES THE GUARDED UPDATE. Precondition backup → safety dump → pin → pull → undo copy of every named
// volume (the old datadir, with its finished-marker) → CONVERTING → start → verify with the new probe.
// Any failure after the copy — the old engine will not start, the dump is cut off, the load errors, the
// counts differ, the new datadir is the wrong major, the app is unhealthy — goes to the EXISTING undo:
// every volume back from its copy, old pin, old definition, old probe. `undone`, else HOLD.
//
// THE ORDER INSIDE `converting`, and the order is the point:
//
// 1. the OLD database container — stopped by the copy, never recreated — is started ALONE;
// 2. the check's "before": per database owner/encoding/collation, every role, every extension, every
// table's exact row count (decision 38 / A3);
// 3. `pg_dumpall` from the old engine into the stack dir, validated by its completion line;
// 4. the old engine stops;
// 5. THE DB VOLUME IS EMPTIED ONLY AFTER THE COPY'S FINISHED-MARKER IS VALIDATED AGAIN — here, and again
// inside the helper that empties it (the Restore shape). Pinned by TestConvert_NeverEmptiedWithoutMarker;
// 6. the NEW engine starts alone on the empty volume (`up -d --no-deps <svc>` with the new definition),
// its entrypoint initialising the bootstrap role and database from the app's own env;
// 7. the entrypoint's databases are dropped when they hold no table, the dump's `CREATE ROLE <x>;` line is
// skipped for a role that already exists (its `ALTER ROLE … PASSWORD` still runs), and the load runs
// with ON_ERROR_STOP — so ANY other error stops it (measured: raw, exactly two "already exists");
// 8. the check's "after" must EQUAL the before, and `$PGDATA/PG_VERSION` must be the mark's `to`.
//
// AFTER SUCCESS the old datadir's copy is KEPT (decision 35, B5) until the first backup of the converted
// app is proven on the new engine; ReleaseConversionCopies then removes it, logged by name both times.
// UpdatePhaseConverting is the conversion's phase, between copying and starting.
const UpdatePhaseConverting = "converting"
// preUpdateConvertDir holds the conversion's dump inside the stack dir. Kept on a HOLD (evidence),
// removed on done and undone.
const preUpdateConvertDir = "pre-update-convert"
// EngineConversion is the ladder entry's mark: this step converts `service` from major `from` to `to`.
type EngineConversion struct {
Service string `yaml:"service" json:"service"`
Engine string `yaml:"engine" json:"engine"`
From int `yaml:"from" json:"from"`
To int `yaml:"to" json:"to"`
}
// ConversionCopy is app.yaml's record of the old datadir's copy kept after a successful conversion.
type ConversionCopy struct {
Volume string `yaml:"volume" json:"volume"`
Copy string `yaml:"copy" json:"copy"`
At string `yaml:"at" json:"at"` // RFC3339 — a backup proven after this releases the copy
From int `yaml:"from" json:"from"`
To int `yaml:"to" json:"to"`
}
// conversionDumpMargin is A5's margin on the dump's bound (the DB volume's own size).
const conversionDumpMargin = 1.25
// ── recognising a PostgreSQL major move ──────────────────────────────────────────────────────────
// isPostgresImage is the backup side's predicate (appbackup.dbTypeForImage, R-484) for the Postgres
// family — kept equal by TestConvert_PostgresFamilyMatchesTheBackupSide.
func isPostgresImage(ref string) bool {
img := strings.ToLower(ref)
return strings.Contains(img, "postgres") || strings.Contains(img, "postgis") ||
strings.Contains(img, "pgvector") || strings.Contains(img, "timescaledb")
}
var leadingDigits = regexp.MustCompile(`^(\d+)`)
// postgresMajor reads the major from a Postgres-family tag: `16-alpine` → 16, `17.2` → 17,
// `16-3.5-alpine` (postgis) → 16, `pg16` (pgvector) → 16, `16-vectorchord0.4.3-…` → 16.
func postgresMajor(ref string) (int, bool) {
ref = stripDigest(ref)
_, tag := splitImageRef(ref)
tag = strings.TrimPrefix(strings.ToLower(tag), "pg")
m := leadingDigits.FindStringSubmatch(tag)
if m == nil {
return 0, false
}
n, err := strconv.Atoi(m[1])
return n, err == nil
}
func stripDigest(ref string) string {
if i := strings.Index(ref, "@"); i >= 0 {
return ref[:i]
}
return ref
}
// errNoConversionMark is the refusal of a PostgreSQL major move the ladder does not mark (B1).
type errNoConversionMark struct{ service, from, to, why string }
func (e *errNoConversionMark) Error() string {
return fmt.Sprintf("service %s moves PostgreSQL %s → %s across a major and %s — refusing before anything moves", e.service, e.from, e.to, e.why)
}
// planEngineConversion decides whether the step about to be pinned is a PostgreSQL major move, and if
// so whether its ladder entry marks it. Returns nil, nil when no Postgres major moves.
func planEngineConversion(pinned, stepImgs map[string]string, entry *LadderEntry) (*EngineConversion, error) {
var moved []string
for svc, newRef := range stepImgs {
oldRef, ok := pinned[svc]
if !ok || !isPostgresImage(newRef) {
continue
}
if stripDigest(oldRef) == stripDigest(newRef) {
continue
}
om, ok1 := postgresMajor(oldRef)
nm, ok2 := postgresMajor(newRef)
if ok1 && ok2 && om == nm {
continue // within a major — the engine reads its own datadir
}
moved = append(moved, svc)
}
if len(moved) == 0 {
return nil, nil
}
sort.Strings(moved)
svc := moved[0]
fail := func(why string) (*EngineConversion, error) {
return nil, &errNoConversionMark{service: svc, from: pinned[svc], to: stepImgs[svc], why: why}
}
if len(moved) > 1 {
return fail(fmt.Sprintf("so do %v — one conversion per step", moved[1:]))
}
om, ok1 := postgresMajor(pinned[svc])
nm, ok2 := postgresMajor(stepImgs[svc])
if !ok1 || !ok2 {
return fail("the major cannot be read from the tag")
}
if entry == nil || entry.EngineConversion == nil {
return fail("the step's ladder entry carries no engine_conversion mark (no test proved the conversion)")
}
c := entry.EngineConversion
if c.Service != svc || !strings.EqualFold(c.Engine, "postgres") || c.From != om || c.To != nm {
return fail(fmt.Sprintf("the mark says %s %s %d → %d, the step moves %s %d → %d", c.Service, c.Engine, c.From, c.To, svc, om, nm))
}
cp := *c
return &cp, nil
}
// ── the compose facts the conversion needs ──────────────────────────────────────────────────────
type composeSvcVolumesDoc struct {
Name string `yaml:"name"`
Services map[string]struct {
Volumes []interface{} `yaml:"volumes"`
} `yaml:"services"`
Volumes map[string]*composeVolDef `yaml:"volumes"`
}
// composeProject is the compose project name the manager runs a stack under: the file's top-level
// `name:`, else the stack directory's name (DeclaredVolumeNames' rule).
func composeProject(composePath string) string {
data, err := os.ReadFile(composePath)
if err == nil {
var doc composeSvcVolumesDoc
if yaml.Unmarshal(data, &doc) == nil && strings.TrimSpace(doc.Name) != "" {
return strings.TrimSpace(doc.Name)
}
}
return filepath.Base(filepath.Dir(composePath))
}
// serviceDataVolume returns the declared named volume a service mounts at /var/lib/postgresql or below:
// its compose KEY and its Docker name.
func serviceDataVolume(composePath, service string) (key, dockerName string, err error) {
data, err := os.ReadFile(composePath)
if err != nil {
return "", "", err
}
var doc composeSvcVolumesDoc
if err := yaml.Unmarshal(data, &doc); err != nil {
return "", "", fmt.Errorf("parsing %s: %w", composePath, err)
}
svc, ok := doc.Services[service]
if !ok {
return "", "", fmt.Errorf("%s declares no service %q", composePath, service)
}
project := strings.TrimSpace(doc.Name)
if project == "" {
project = filepath.Base(filepath.Dir(composePath))
}
for _, v := range svc.Volumes {
src, dst := "", ""
switch x := v.(type) {
case string:
parts := strings.Split(x, ":")
if len(parts) >= 2 {
src, dst = parts[0], parts[1]
}
case map[string]interface{}:
src, _ = x["source"].(string)
dst, _ = x["target"].(string)
}
if !strings.HasPrefix(dst, "/var/lib/postgresql") {
continue
}
def, declared := doc.Volumes[src]
if !declared {
return "", "", fmt.Errorf("service %s mounts %s at %s, which is not a named volume the file declares (a bind mount cannot be converted)", service, src, dst)
}
name := project + "_" + src
if def != nil && strings.TrimSpace(def.Name) != "" {
name = strings.TrimSpace(def.Name)
}
return src, name, nil
}
return "", "", fmt.Errorf("service %s mounts no named volume under /var/lib/postgresql", service)
}
// freeBytesAt is the free space of the filesystem holding path (the dump's home). Seam for tests.
func (m *Manager) freeBytesAt(path string) (int64, bool) {
if m.convertFreeFn != nil {
return m.convertFreeFn(path)
}
var s syscall.Statfs_t
if err := syscall.Statfs(path, &s); err != nil {
return 0, false
}
return int64(s.Bavail) * int64(s.Bsize), true
}
// convertSpaceError is the refusal of B6's „no-space" sentence.
type convertSpaceError struct{ need, free int64 }
func (e *convertSpaceError) Error() string {
return fmt.Sprintf("the conversion needs %s free for its dump (the database volume's size × %.2f + the %.0f GB floor) and %s is free", humanBytes(e.need), conversionDumpMargin, updateDiskFloorGiB, humanBytes(e.free))
}
func humanBytes(b int64) string {
switch {
case b >= 1<<30:
return fmt.Sprintf("%.1f GB", float64(b)/(1<<30))
case b >= 1<<20:
return fmt.Sprintf("%.0f MB", float64(b)/(1<<20))
}
return fmt.Sprintf("%d kB", b/1024)
}
// planConversionSpace checks, before anything moves, that the conversion's dump fits beside the stack
// dir (A5): the dump's bound is the DB volume's own size, with a margin, plus the update's 2 GB floor.
// The undo copy's own space is planUndoCopies' check, unchanged.
func (m *Manager) planConversionSpace(name, dir, dbVol string) error {
b, err := m.copier().VolumeBytes(dbVol)
if err != nil {
return fmt.Errorf("sizing the database volume %s: %w", dbVol, err)
}
need := int64(float64(b)*conversionDumpMargin) + int64(updateDiskFloorGiB*(1<<30))
free, known := m.freeBytesAt(dir)
if !known {
m.logger.Printf("[WARN] [stacks] update %s: free space beside the stack dir is unreadable — the conversion proceeds without its space check", name)
return nil
}
m.logger.Printf("[INFO] [stacks] update %s: conversion space — database volume %s is %s; the dump needs up to %s + the %.0f GB floor; %s free beside the stack dir", name, dbVol, humanBytes(b), humanBytes(int64(float64(b)*conversionDumpMargin)), updateDiskFloorGiB, humanBytes(free))
if free < need {
return &convertSpaceError{need: need, free: free}
}
return nil
}
// ── the process boundary ────────────────────────────────────────────────────────────────────────
// pgConverter is the conversion's process boundary. Production is dockerPGConverter; tests inject a
// fake and never touch docker (R-650).
type pgConverter interface {
ServiceContainer(project, service string) (string, error)
Start(container string) error
Stop(container string) error
Env(container, key string) string
WaitReady(ctx context.Context, container, user string, timeout time.Duration) error
Snapshot(container, user string) ([]string, error)
DumpAll(container, user, path string) error
// Empty empties vol ONLY when copyVol carries its finished-marker (checked inside the same helper).
Empty(copyVol, vol string) error
// DropEmptyDatabases drops every non-template database except `postgres` that holds no table, and
// refuses (error) if any holds one — only an entrypoint-made database may be dropped.
DropEmptyDatabases(container, user string) ([]string, error)
Roles(container, user string) ([]string, error)
Load(container, user, path string, skipLines map[string]bool) error
DataVersion(container string) (string, error)
}
func (m *Manager) converter() pgConverter {
if m.pgConv != nil {
return m.pgConv
}
return dockerPGConverter{m: m}
}
// ── the conversion ──────────────────────────────────────────────────────────────────────────────
// dumpAllCompleteLine is the last comment pg_dumpall writes; a dump without it is cut off.
const dumpAllCompleteLine = "PostgreSQL database cluster dump complete"
// validateDumpAll refuses a dump without its completion line (the lesson of 2026-09-23: a truncated
// PostgreSQL dump loads with exit 0).
func validateDumpAll(path string) error {
f, err := os.Open(path)
if err != nil {
return err
}
defer f.Close()
st, err := f.Stat()
if err != nil {
return err
}
off := st.Size() - 4096
if off < 0 {
off = 0
}
buf := make([]byte, st.Size()-off)
if _, err := f.ReadAt(buf, off); err != nil && err != io.EOF {
return err
}
if !strings.Contains(string(buf), dumpAllCompleteLine) {
return fmt.Errorf("the dump %s (%d bytes) has no completion line — it is cut off", path, st.Size())
}
return nil
}
// convertEngine is the `converting` phase. On error the caller runs the undo (failAndHold → tryUndo).
func (m *Manager) convertEngine(ctx context.Context, name, dir string, env []string, entry *updateJournalEntry) error {
c := entry.Convert
pc := m.converter()
t0 := m.now()
project := composeProject(ComposePathIn(dir))
dbVol := entry.ConvertVolume
var copyVol string
for _, u := range entry.UndoCopies {
if u.Volume == dbVol {
copyVol = u.Copy
}
}
if copyVol == "" {
return fmt.Errorf("the database volume %s has no undo copy — refusing to touch it", dbVol)
}
m.logger.Printf("[INFO] [stacks] update %s: CONVERTING %s PostgreSQL %d → %d (volume %s, kept copy %s)", name, c.Service, c.From, c.To, dbVol, copyVol)
// 1. the OLD engine, alone. The copy stopped the app's containers without recreating them, so the
// service's container is still the old image.
oldC, err := pc.ServiceContainer(project, c.Service)
if err != nil {
return fmt.Errorf("finding the old database container: %w", err)
}
if err := pc.Start(oldC); err != nil {
return fmt.Errorf("starting the old engine %s: %w", oldC, err)
}
user := pc.Env(oldC, "POSTGRES_USER")
if user == "" {
user = "postgres"
}
if err := pc.WaitReady(ctx, oldC, user, 3*time.Minute); err != nil {
_ = pc.Stop(oldC)
return fmt.Errorf("the old engine did not become ready: %w", err)
}
// 2. the check's "before"
before, err := pc.Snapshot(oldC, user)
if err != nil {
_ = pc.Stop(oldC)
return fmt.Errorf("reading the old engine's counts: %w", err)
}
// 3. the dump, validated
cdir := filepath.Join(dir, preUpdateConvertDir)
if err := os.MkdirAll(cdir, 0o700); err != nil {
_ = pc.Stop(oldC)
return fmt.Errorf("creating %s: %w", cdir, err)
}
dump := filepath.Join(cdir, c.Service+"-dumpall.sql")
td := time.Now()
if err := pc.DumpAll(oldC, user, dump); err != nil {
_ = pc.Stop(oldC)
return fmt.Errorf("pg_dumpall from the old engine: %w", err)
}
if err := validateDumpAll(dump); err != nil {
_ = pc.Stop(oldC)
return err
}
var dumpSize int64
if fi, err := os.Stat(dump); err == nil {
dumpSize = fi.Size()
}
m.logger.Printf("[INFO] [stacks] update %s: pg_dumpall from %d done in %s — %s, completion line present; before: %s", name, c.From, time.Since(td).Round(time.Millisecond), humanBytes(dumpSize), summariseSnapshot(before))
// 4. the old engine stops
if err := pc.Stop(oldC); err != nil {
return fmt.Errorf("stopping the old engine: %w", err)
}
// 5. the volume is emptied ONLY after the copy's marker is validated again
if !m.copier().Complete(copyVol) {
return fmt.Errorf("the kept copy %s has no finished-marker — the database volume is NOT emptied", copyVol)
}
entry.ConvertTouched = true
if !m.enterUpdatePhase(name, entry, UpdatePhaseConverting) {
return fmt.Errorf("could not journal that the database volume is about to be emptied")
}
if err := pc.Empty(copyVol, dbVol); err != nil {
return fmt.Errorf("emptying %s: %w", dbVol, err)
}
m.logger.Printf("[INFO] [stacks] update %s: database volume %s emptied (its copy %s holds the old datadir, marker checked twice)", name, dbVol, copyVol)
// 6. the NEW engine, alone, on the empty volume
if _, err := m.updateCompose(dir, env, "up", "-d", "--no-deps", c.Service); err != nil {
return fmt.Errorf("starting the new engine: %w", err)
}
newC, err := pc.ServiceContainer(project, c.Service)
if err != nil {
return fmt.Errorf("finding the new database container: %w", err)
}
if err := pc.WaitReady(ctx, newC, user, 5*time.Minute); err != nil {
return fmt.Errorf("the new engine did not become ready: %w", err)
}
// 7. make room for exactly the two objects the entrypoint made, then load with ON_ERROR_STOP
dropped, err := pc.DropEmptyDatabases(newC, user)
if err != nil {
return fmt.Errorf("preparing the new engine: %w", err)
}
roles, err := pc.Roles(newC, user)
if err != nil {
return fmt.Errorf("reading the new engine's roles: %w", err)
}
skip := map[string]bool{}
for _, r := range roles {
skip["CREATE ROLE "+quoteIdent(r)+";"] = true
}
tl := time.Now()
if err := pc.Load(newC, user, dump, skip); err != nil {
return fmt.Errorf("loading the dump into %d: %w", c.To, err)
}
m.logger.Printf("[INFO] [stacks] update %s: loaded into %d in %s (dropped the entrypoint's empty database(s) %v; skipped CREATE ROLE for %v)", name, c.To, time.Since(tl).Round(time.Millisecond), dropped, roles)
// 8. the check's "after" must equal the before; PG_VERSION must be the mark's `to`
after, err := pc.Snapshot(newC, user)
if err != nil {
return fmt.Errorf("reading the new engine's counts: %w", err)
}
if diff := diffSnapshots(before, after); diff != "" {
return fmt.Errorf("the check differs after the load: %s", diff)
}
ver, err := pc.DataVersion(newC)
if err != nil {
return fmt.Errorf("reading PG_VERSION of the new datadir: %w", err)
}
if strings.TrimSpace(ver) != strconv.Itoa(c.To) {
return fmt.Errorf("the new datadir's PG_VERSION is %q, the mark says %d", strings.TrimSpace(ver), c.To)
}
m.logger.Printf("[INFO] [stacks] update %s: CONVERTED %d → %d in %s — the check is equal (%s), PG_VERSION %s", name, c.From, c.To, m.now().Sub(t0).Round(time.Millisecond), summariseSnapshot(after), strings.TrimSpace(ver))
return nil
}
func quoteIdent(s string) string {
if regexp.MustCompile(`^[a-z_][a-z0-9_$]*$`).MatchString(s) && !pgReserved[s] {
return s
}
return `"` + strings.ReplaceAll(s, `"`, `""`) + `"`
}
// pgReserved is the handful of reserved words a role could plausibly be named (quote_ident quotes them).
var pgReserved = map[string]bool{"user": true, "all": true, "default": true, "group": true, "order": true, "table": true}
// summariseSnapshot is the log's one-line form: databases, tables, rows.
func summariseSnapshot(lines []string) string {
dbs, tables, rows := 0, 0, int64(0)
for _, l := range lines {
switch {
case strings.HasPrefix(l, "db:"):
dbs++
case strings.HasPrefix(l, "rows:"):
tables++
if i := strings.LastIndex(l, "="); i >= 0 {
n, _ := strconv.ParseInt(l[i+1:], 10, 64)
rows += n
}
}
}
return fmt.Sprintf("%d database(s), %d table(s), %d row(s)", dbs, tables, rows)
}
// diffSnapshots names the first differences, "" when equal.
func diffSnapshots(before, after []string) string {
a, b := map[string]bool{}, map[string]bool{}
for _, l := range before {
a[l] = true
}
for _, l := range after {
b[l] = true
}
var d []string
for _, l := range before {
if !b[l] {
d = append(d, "missing after: "+l)
}
}
for _, l := range after {
if !a[l] {
d = append(d, "new after: "+l)
}
}
if len(d) == 0 {
return ""
}
if len(d) > 6 {
d = append(d[:6], fmt.Sprintf("… %d more", len(d)-6))
}
return strings.Join(d, "; ")
}
// ── keeping, then releasing, the old datadir's copy (B5) ────────────────────────────────────────
func (m *Manager) recordConversionCopy(name, dir string, cc *ConversionCopy) {
cfg := LoadAppConfig(dir)
if cfg == nil || (cc == nil && cfg.ConversionCopy == nil) {
return
}
cfg.ConversionCopy = cc
meta := LoadMetadata(dir)
if err := SaveAppConfig(dir, cfg, m.encKey, SensitiveEnvVars(&meta)); err != nil {
m.logger.Printf("[ERROR] [stacks] update %s: recording conversion_copy failed: %v", name, err)
return
}
m.mu.Lock()
if st, ok := m.stacks[name]; ok && st.AppConfig != nil {
st.AppConfig.ConversionCopy = cc
}
m.mu.Unlock()
}
// ReleaseConversionCopies removes each kept pre-conversion datadir copy whose app now has a backup
// PROVEN after the conversion (any tier — the backup side returns only proven copies). Returns the
// names released. Run periodically from main.go (TestConvert_ReleaseIsWiredAtStartup).
func (m *Manager) ReleaseConversionCopies(ctx context.Context) []string {
g := m.guards()
if g == nil {
return nil
}
var released []string
for _, st := range m.GetStacks() {
if st.AppConfig == nil || st.AppConfig.ConversionCopy == nil {
continue
}
cc := st.AppConfig.ConversionCopy
at, err := time.Parse(time.RFC3339, cc.At)
if err != nil {
m.logger.Printf("[WARN] [stacks] %s: conversion_copy has an unreadable time %q — the copy %s is kept", st.Name, cc.At, cc.Copy)
continue
}
rp, ok, _ := g.RestorePoints(ctx, st.Name, func(p UpdateRestorePoint) bool { return p.ProvenAt.After(at) })
if !ok {
continue
}
if err := m.copier().Remove(cc.Copy); err != nil {
m.logger.Printf("[WARN] [stacks] %s: could not remove the pre-conversion copy %s: %v — kept, tried again later", st.Name, cc.Copy, err)
continue
}
m.recordConversionCopy(st.Name, filepath.Dir(st.ComposePath), nil)
m.logger.Printf("[INFO] [stacks] %s: REMOVED the pre-conversion datadir copy %s (PostgreSQL %d) — the converted app has a backup proven on %d: %s at %s", st.Name, cc.Copy, cc.From, cc.To, updateTierName(rp.Tier), rp.ProvenAt.UTC().Format(time.RFC3339))
released = append(released, st.Name)
}
return released
}
// ── production boundary ─────────────────────────────────────────────────────────────────────────
type dockerPGConverter struct{ m *Manager }
func (d dockerPGConverter) ServiceContainer(project, service string) (string, error) {
out, err := d.m.execCommand("docker", "ps", "-a", "--filter", "label=com.docker.compose.project="+project,
"--filter", "label=com.docker.compose.service="+service, "--format", "{{.Names}}")
if err != nil {
return "", err
}
var names []string
for _, l := range strings.Split(out, "\n") {
if l = strings.TrimSpace(l); l != "" {
names = append(names, l)
}
}
if len(names) != 1 {
return "", fmt.Errorf("project %s service %s has %d container(s) %v, want exactly one", project, service, len(names), names)
}
return names[0], nil
}
func (d dockerPGConverter) Start(c string) error {
_, err := d.m.execCommand("docker", "start", c)
return err
}
func (d dockerPGConverter) Stop(c string) error {
_, err := d.m.execCommand("docker", "stop", "-t", "60", c)
return err
}
func (d dockerPGConverter) Env(c, key string) string {
out, err := d.m.execCommand("docker", "inspect", c, "--format", "{{range .Config.Env}}{{println .}}{{end}}")
if err != nil {
return ""
}
for _, l := range strings.Split(out, "\n") {
if strings.HasPrefix(l, key+"=") {
return strings.TrimSpace(strings.TrimPrefix(l, key+"="))
}
}
return ""
}
// WaitReady asks over TCP on 127.0.0.1: the entrypoint's temporary init server listens on the socket
// only, so a TCP answer is the FINAL server (loading into the temporary one would be lost on its stop).
func (d dockerPGConverter) WaitReady(ctx context.Context, c, user string, timeout time.Duration) error {
deadline := time.Now().Add(timeout)
last := ""
for time.Now().Before(deadline) {
out, err := d.m.execCommand("docker", "exec", c, "psql", "-h", "127.0.0.1", "-U", user, "-d", "postgres", "-Atc", "select 1")
if err == nil && strings.TrimSpace(out) == "1" {
return nil
}
if err != nil {
last = err.Error()
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(2 * time.Second):
}
}
return fmt.Errorf("no answer within %s (last: %s)", timeout, firstLine(last))
}
func firstLine(s string) string {
if i := strings.IndexByte(s, '\n'); i >= 0 {
return s[:i]
}
return s
}
func (d dockerPGConverter) psql(c, user, db, sql string) (string, error) {
return d.m.execCommand("docker", "exec", c, "psql", "-h", "127.0.0.1", "-U", user, "-d", db, "-v", "ON_ERROR_STOP=1", "-Atc", sql)
}
const snapshotCountsSQL = `SELECT n.nspname||'.'||c.relname||'='||(xpath('/row/c/text()', query_to_xml(format('select count(*) as c from %I.%I', n.nspname, c.relname), false, true, '')))[1]::text FROM pg_class c JOIN pg_namespace n ON n.oid=c.relnamespace WHERE c.relkind IN ('r','p') AND n.nspname NOT IN ('pg_catalog','information_schema') AND n.nspname NOT LIKE 'pg_toast%' ORDER BY 1`
// Snapshot is A3's check, as measured on 9202.
func (d dockerPGConverter) Snapshot(c, user string) ([]string, error) {
var out []string
add := func(prefix, s string) {
for _, l := range strings.Split(s, "\n") {
if l = strings.TrimSpace(l); l != "" {
out = append(out, prefix+l)
}
}
}
s, err := d.psql(c, user, "postgres", "select datname||' owner='||pg_get_userbyid(datdba)||' enc='||pg_encoding_to_char(encoding)||' coll='||datcollate from pg_database where not datistemplate order by 1")
if err != nil {
return nil, err
}
add("db:", s)
s, err = d.psql(c, user, "postgres", "select rolname||' super='||rolsuper||' login='||rolcanlogin||' pw='||(rolpassword is not null) from pg_roles where rolname !~ '^pg_' order by 1")
if err != nil {
return nil, err
}
add("role:", s)
dbs, err := d.psql(c, user, "postgres", "select datname from pg_database where not datistemplate and datname<>'postgres' order by 1")
if err != nil {
return nil, err
}
for _, db := range strings.Split(dbs, "\n") {
if db = strings.TrimSpace(db); db == "" {
continue
}
s, err := d.psql(c, user, db, "select extname from pg_extension order by 1")
if err != nil {
return nil, err
}
add("ext:"+db+":", s)
s, err = d.psql(c, user, db, snapshotCountsSQL)
if err != nil {
return nil, err
}
add("rows:"+db+":", s)
}
sort.Strings(out)
return out, nil
}
// DumpAll streams pg_dumpall's stdout straight into the file — never through the controller's memory.
func (d dockerPGConverter) DumpAll(c, user, path string) error {
f, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0o600)
if err != nil {
return err
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Hour)
defer cancel()
cmd := dockerexec.CommandContext(ctx, "docker", "exec", c, "pg_dumpall", "-h", "127.0.0.1", "-U", user)
var stderr strings.Builder
cmd.Stdout, cmd.Stderr = f, &stderr
runErr := cmd.Run()
if err := f.Sync(); err != nil && runErr == nil {
runErr = err
}
if err := f.Close(); err != nil && runErr == nil {
runErr = err
}
if runErr != nil {
return fmt.Errorf("%w (stderr: %s)", runErr, firstLine(stderr.String()))
}
return nil
}
func (d dockerPGConverter) Empty(copyVol, vol string) error {
_, err := d.m.execCommand("docker", "run", "--rm", "-v", copyVol+":/from:ro", "-v", vol+":/to", undoHelperImage,
"sh", "-c", "test -f /from/"+undoCopyMarker+" && find /to -mindepth 1 -delete && sync")
return err
}
func (d dockerPGConverter) DropEmptyDatabases(c, user string) ([]string, error) {
dbs, err := d.psql(c, user, "postgres", "select datname from pg_database where not datistemplate and datname<>'postgres' order by 1")
if err != nil {
return nil, err
}
var dropped []string
for _, db := range strings.Split(dbs, "\n") {
if db = strings.TrimSpace(db); db == "" {
continue
}
n, err := d.psql(c, user, db, "select count(*) from pg_class c join pg_namespace n on n.oid=c.relnamespace where c.relkind in ('r','p') and n.nspname not in ('pg_catalog','information_schema')")
if err != nil {
return nil, err
}
if strings.TrimSpace(n) != "0" {
return nil, fmt.Errorf("database %s on the NEW engine already holds %s table(s) — it was not made empty by the entrypoint; refusing to drop it", db, strings.TrimSpace(n))
}
if _, err := d.psql(c, user, "postgres", "DROP DATABASE "+quoteIdent(db)); err != nil {
return nil, err
}
dropped = append(dropped, db)
}
return dropped, nil
}
func (d dockerPGConverter) Roles(c, user string) ([]string, error) {
s, err := d.psql(c, user, "postgres", "select rolname from pg_roles where rolname !~ '^pg_' order by 1")
if err != nil {
return nil, err
}
var out []string
for _, l := range strings.Split(s, "\n") {
if l = strings.TrimSpace(l); l != "" {
out = append(out, l)
}
}
return out, nil
}
// Load streams the dump into psql with ON_ERROR_STOP, skipping exactly the lines named.
func (d dockerPGConverter) Load(c, user, path string, skip map[string]bool) error {
f, err := os.Open(path)
if err != nil {
return err
}
defer f.Close()
pr, pw := io.Pipe()
go func() {
pw.CloseWithError(filterLines(f, pw, skip))
}()
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Hour)
defer cancel()
cmd := dockerexec.CommandContext(ctx, "docker", "exec", "-i", c, "psql", "-h", "127.0.0.1", "-U", user, "-d", "postgres", "-v", "ON_ERROR_STOP=1", "-q")
var stderr strings.Builder
cmd.Stdin, cmd.Stderr, cmd.Stdout = pr, &stderr, io.Discard
if err := cmd.Run(); err != nil {
return fmt.Errorf("%w (psql: %s)", err, firstLine(strings.TrimSpace(stderr.String())))
}
return nil
}
// filterLines copies r to w, dropping each line that equals a key of skip exactly.
func filterLines(r io.Reader, w io.Writer, skip map[string]bool) error {
br := bufio.NewReaderSize(r, 1<<20)
for {
line, err := br.ReadString('\n')
if len(line) > 0 && !skip[strings.TrimRight(line, "\r\n")] {
if _, werr := io.WriteString(w, line); werr != nil {
return werr
}
}
if err == io.EOF {
return nil
}
if err != nil {
return err
}
}
}
func (d dockerPGConverter) DataVersion(c string) (string, error) {
return d.m.execCommand("docker", "exec", c, "sh", "-c", `cat "$PGDATA/PG_VERSION"`)
}
// appLabel is the app's display name for a household sentence, else its stack name.
func appLabel(st *Stack) string {
if st != nil && strings.TrimSpace(st.Meta.DisplayName) != "" {
return st.Meta.DisplayName
}
if st == nil {
return ""
}
return st.Name
}
// conversionFor is B1's decision for one step, shared by the preflight and the job: the conversion to
// run (nil = none), the DB volume's Docker name, and on refusal the bundle key that names it.
func (m *Manager) conversionFor(st *Stack, pinned map[string]string, step LadderStep) (*EngineConversion, string, string, error) {
imgs, err := ParseComposeImages(step.Source)
if err != nil {
// Not this check's refusal: advancePinTo reads the same file and refuses with its own sentence.
return nil, "", "", nil
}
conv, err := planEngineConversion(pinned, imgs, step.Entry)
if err != nil {
return nil, "", "err.stacks.update_engine_no_test", err
}
if conv == nil {
return nil, "", "", nil
}
oldKey, vol, err := serviceDataVolume(st.ComposePath, conv.Service)
if err != nil {
return nil, "", "", fmt.Errorf("the running definition: %w", err)
}
newKey, _, err := serviceDataVolume(step.Source, conv.Service)
if err != nil {
return nil, "", "", fmt.Errorf("the step's definition: %w", err)
}
if oldKey != newKey {
return nil, "", "", fmt.Errorf("the step mounts volume %q for %s, the running definition %q — the conversion rebuilds the datadir IN its volume and the undo copies that one", newKey, conv.Service, oldKey)
}
return conv, vol, "", nil
}