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 ` 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 ;` 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 }