Files
felhom-controller/controller/internal/stacks/pgconvert.go
T
admin 820e8efde1
gates / gates (push) Successful in 26s
v0.276.0: a restore and a drive move keep the app's records (R-697, R-700)
A drive move persisted through the restore's fresh app.yaml write and dropped the pin: the syncer
then copied the catalog verbatim and the next start jumped the app past its ladder (R-700).
persistDriveFlip now changes HDD_PATH and nothing else. The restore's write carries the life
records (conversion copies, desired_state, update history) from the app.yaml it replaces, and a
second conversion no longer overwrites the first kept copy's record (R-697).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
2026-09-27 17:59:04 +02:00

981 lines
38 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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"`
// Service (v0.275.0) is the converted compose service — the release looks for a dump written by ITS
// engine at major To. "" on a record written before v0.275.0: then every Postgres-family image in the
// dump's recorded set must be at To.
Service string `yaml:"service,omitempty" json:"service,omitempty"`
}
// DataDumpStamp is one database dump in the app's own recovery unit, with what the backup side recorded
// when it was WRITTEN (v0.275.0, R-696): its time and the running images (service -> ref@digest).
type DataDumpStamp struct {
File string
At time.Time
Images map[string]string
}
// DumpStampSource is the backup side's record of the app's own unit's database dumps (v0.275.0). An
// OPTIONAL extension of UpdateGuards: guards without it never release a conversion copy (fail closed —
// a copy outliving its backup costs disk, never data).
type DumpStampSource interface {
DumpStamps(name string) []DataDumpStamp
}
// convertedDumpAt is the release's second condition (A4): a database dump written AFTER the conversion
// whose recorded engine is the NEW major. The first condition — a copy of the app proven after the
// conversion on any tier — says the DATA is newer; this one says it was written by the converted engine,
// which a unit's refresh time or a snapshot's time cannot say (R-696: the demo-hp release cited a unit
// re-captured over a PostgreSQL 16 dump).
func convertedDumpAt(stamps []DataDumpStamp, cc *ConversionCopy, after time.Time) (DataDumpStamp, bool) {
for _, st := range stamps {
if !st.At.After(after) || len(st.Images) == 0 {
continue
}
if cc.Service != "" {
if ref, ok := st.Images[cc.Service]; ok {
if mj, ok := postgresMajor(ref); ok && mj == cc.To {
return st, true
}
}
continue
}
seen, all := 0, true
for _, ref := range st.Images {
if !isPostgresImage(ref) {
continue
}
seen++
if mj, ok := postgresMajor(ref); !ok || mj != cc.To {
all = false
}
}
if seen > 0 && all {
return st, true
}
}
return DataDumpStamp{}, false
}
// 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) {
set := func(cfg *AppConfig) {
// R-697 (v0.276.0): a newer conversion never overwrites an older kept copy's record — the older one
// moves to the earlier list and is released by the same rule.
if cc != nil && cfg.ConversionCopy != nil && cfg.ConversionCopy.Copy != cc.Copy && !hasCopy(cfg.EarlierConversionCopies, cfg.ConversionCopy.Copy) {
cfg.EarlierConversionCopies = append(cfg.EarlierConversionCopies, *cfg.ConversionCopy)
}
cfg.ConversionCopy = cc
}
cfg := LoadAppConfig(dir)
if cfg == nil || (cc == nil && cfg.ConversionCopy == nil) {
return
}
set(cfg)
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 = cfg.ConversionCopy
st.AppConfig.EarlierConversionCopies = cfg.EarlierConversionCopies
}
m.mu.Unlock()
}
func hasCopy(list []ConversionCopy, copyVol string) bool {
for _, c := range list {
if c.Copy == copyVol {
return true
}
}
return false
}
// forgetEarlierConversionCopy drops one released copy from the earlier list.
func (m *Manager) forgetEarlierConversionCopy(name, dir, copyVol string) {
m.mutateAppConfig(name, dir, "earlier_conversion_copies", func(cfg *AppConfig) bool {
var keep []ConversionCopy
for _, c := range cfg.EarlierConversionCopies {
if c.Copy != copyVol {
keep = append(keep, c)
}
}
if len(keep) == len(cfg.EarlierConversionCopies) {
return false
}
cfg.EarlierConversionCopies = keep
return true
})
}
// 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 {
continue
}
// v0.276.0 (R-697): the current record and every earlier one a later conversion superseded.
for _, e := range st.AppConfig.EarlierConversionCopies {
e := e
if m.releaseOneConversionCopy(ctx, g, st, &e) {
m.forgetEarlierConversionCopy(st.Name, filepath.Dir(st.ComposePath), e.Copy)
released = append(released, st.Name)
}
}
if cc := st.AppConfig.ConversionCopy; cc != nil && m.releaseOneConversionCopy(ctx, g, st, cc) {
m.recordConversionCopy(st.Name, filepath.Dir(st.ComposePath), nil)
released = append(released, st.Name)
}
}
return released
}
// releaseOneConversionCopy removes one kept copy when its release conditions hold; true = removed.
func (m *Manager) releaseOneConversionCopy(ctx context.Context, g UpdateGuards, st Stack, cc *ConversionCopy) bool {
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)
return false
}
rp, ok, _ := g.RestorePoints(ctx, st.Name, func(p UpdateRestorePoint) bool { return p.ProvenAt.After(at) })
if !ok {
return false
}
// v0.275.0 (A4): AND a dump the converted engine wrote. Without the stamps (older guards, or a unit
// whose data is unstamped) the copy is KEPT — logged, retried at the next pass.
src, hasStamps := g.(DumpStampSource)
var dump DataDumpStamp
if hasStamps {
dump, ok = convertedDumpAt(src.DumpStamps(st.Name), cc, at)
}
if !hasStamps || !ok {
if m.isDebug() {
m.logger.Printf("[DEBUG] [stacks] %s: the pre-conversion copy %s is KEPT — a copy proven at %s exists, but no database dump written after %s by PostgreSQL %d is recorded yet", st.Name, cc.Copy, rp.ProvenAt.UTC().Format(time.RFC3339), cc.At, cc.To)
}
return false
}
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)
return false
}
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, its dump %s written %s by %v", st.Name, cc.Copy, cc.From, cc.To, updateTierName(rp.Tier), rp.ProvenAt.UTC().Format(time.RFC3339), dump.File, dump.At.UTC().Format(time.RFC3339), dump.Images)
return true
}
// ── 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
}