b6810f14ff
gates / gates (push) Successful in 23s
The unit's data files are stamped with the versions that wrote them; the capture keeps the definition the data belongs to; a restore never starts data under another version's definition (unit restores refuse a mismatch; the off-site restore writes the snapshot's definition); every tier's time is its data's; the conversion-copy release needs a dump on the new engine. File-browser sync single-flight + no empty kept folder (R-695); the kept view joins the folder's owning group, language switch resyncs (R-691); a restore-generated login is not shown as the password (R-694). Red-proofs in felhom.eu/documentation/audits/version-travel-2026-09-26/. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
932 lines
36 KiB
Go
932 lines
36 KiB
Go
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) {
|
||
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
|
||
}
|
||
// 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)
|
||
}
|
||
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, 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)
|
||
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
|
||
}
|