From 23a8ef3de4a7ae30ca105f7127508c932a6cdeae Mon Sep 17 00:00:00 2001 From: kisfenyo Date: Sun, 4 Oct 2026 10:44:29 +0200 Subject: [PATCH] =?UTF-8?q?OS=20updates,=20guest=20fast=20lane=20(11=20?= =?UTF-8?q?=C2=A78=20step=202):=20felhom-os-apply=20wrapper=20(R1-R13=20re?= =?UTF-8?q?fusals,=20repair=20first,=20snapshot.debian.org=20fallback),=20?= =?UTF-8?q?FELHOM=5FOSAPPLY=20sudoers,=20the=20OS=20leg=20after=20the=20pr?= =?UTF-8?q?imary=20backup,=20hub=20os=5Fupdate=20block=20+=20os-report,=20?= =?UTF-8?q?--selftest=3Dos-update?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit No automatic undo: a customer guest cannot be snapshotted (R-837, measured). Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS --- cmd/felhom-agent/main.go | 72 +++ configs/felhom-agent.sudoers | 9 +- configs/felhom-os-apply | 475 ++++++++++++++++++ configs/test_felhom_os_apply.py | 452 +++++++++++++++++ internal/capability/manifest.go | 3 + internal/hub/client.go | 24 + internal/hub/osupdate_contract_test.go | 30 ++ internal/hub/report.go | 26 + .../desired-state-osupdate.golden.json | 17 + internal/localapi/afterbackup_test.go | 63 +++ internal/localapi/server.go | 13 + internal/osupdate/leg.go | 429 ++++++++++++++++ internal/osupdate/leg_test.go | 265 ++++++++++ 13 files changed, 1877 insertions(+), 1 deletion(-) create mode 100755 configs/felhom-os-apply create mode 100644 configs/test_felhom_os_apply.py create mode 100644 internal/hub/osupdate_contract_test.go create mode 100644 internal/hub/testdata/desired-state-osupdate.golden.json create mode 100644 internal/localapi/afterbackup_test.go create mode 100644 internal/osupdate/leg.go create mode 100644 internal/osupdate/leg_test.go diff --git a/cmd/felhom-agent/main.go b/cmd/felhom-agent/main.go index cfb677c..7140f02 100644 --- a/cmd/felhom-agent/main.go +++ b/cmd/felhom-agent/main.go @@ -48,6 +48,7 @@ import ( "gitea.dooplex.hu/admin/felhom-agent/internal/pbsdr" "gitea.dooplex.hu/admin/felhom-agent/internal/poke" "gitea.dooplex.hu/admin/felhom-agent/internal/provision" + "gitea.dooplex.hu/admin/felhom-agent/internal/osupdate" "gitea.dooplex.hu/admin/felhom-agent/internal/proxmox" "gitea.dooplex.hu/admin/felhom-agent/internal/reconcile" "gitea.dooplex.hu/admin/felhom-agent/internal/restorespace" @@ -242,6 +243,8 @@ func main() { os.Exit(runSelftestRestoreTest(context.Background(), cfg, logger, archive)) case "restore-test-due": os.Exit(runSelftestRestoreTestDue(context.Background(), cfg, logger)) + case "os-update": + os.Exit(runSelftestOSUpdate(context.Background(), cfg, logger, vmid)) case "pbs-verify": os.Exit(runSelftestPBSVerify(context.Background(), cfg, logger)) case "lanresolver": @@ -839,6 +842,10 @@ func runDaemon(cfg config.Config, logger *slog.Logger, logRing *applog.Ring) int // The "Down" channel sync hook: on each heartbeat, fetch desired-state when the generation // advances. The loop calls it via the EnvelopeObserver seam (hub does not import desired). desiredSyncer := desired.NewSyncer(client, desiredProvider, logger) + // OS updates, guest fast lane (agent v0.140.0, `11-os-updates.md` §8 step 2): the leg consumes the hub's + // os_update block and runs after each successful primary whole-guest backup (wired on the local API below). + osLeg := newOSLeg(cfg, client, logger) + desiredSyncer.AddConsumer(osLeg) // S5: consume a host_loss restore_directive into an inspectable restore PLAN (derive + surface, // execute nothing). The recipe is fetched on-demand (rare directive) via a fresh Collect. desiredSyncer.AddConsumer(dr.NewConsumer(func(ctx context.Context) *hub.DRRecipeHostHalf { @@ -1112,6 +1119,17 @@ func runDaemon(cfg config.Config, logger *slog.Logger, logRing *applog.Ring) int }, } localSrv := buildLocalAPIServer(cfg, px, backupStore, heavyOps, observer, driveKnown, hostOps, gate, collector, client, intentRec, guestBindStore, formatJobStore, logRing, escrowCeremonyCfg, logger, &localTokens) + if localSrv != nil { + localSrv.SetAfterPrimaryBackup(func(ctx context.Context, vmid int) { + // Let the controller finish bringing its apps back after the backup, then run (still under the gate). + select { + case <-ctx.Done(): + return + case <-time.After(90 * time.Second): + } + osLeg.Run(ctx, vmid, "night") + }) + } if localTokens != nil { defer localTokens.Close() } @@ -3483,3 +3501,57 @@ func (f *selftestFlag) Set(v string) error { } return nil } + +// newOSLeg builds the OS-update leg (agent v0.140.0). The wrapper runs through sudo (FELHOM_OSAPPLY); the plan and +// the once-per-night marker live in the agent's own os/ dir. +func newOSLeg(cfg config.Config, client *hub.Client, logger *slog.Logger) *osupdate.Leg { + mode := proxmox.RunnerMode(cfg.Privileged.Mode) + if mode == "" { + mode = proxmox.RunnerSudo + } + l := &osupdate.Leg{ + Runner: &proxmox.ExecRunner{Mode: mode, SudoPath: cfg.Privileged.SudoPath}, + Logger: logger, + PlanDir: osupdate.DefaultPlanDir, + StatePath: filepath.Join(osupdate.DefaultPlanDir, "last-night-run"), + } + if client != nil { + l.Hub = client + } + return l +} + +// runSelftestOSUpdate is the OS leg's DEBUG ACTION (agent v0.140.0): one pass for -vmid, now, exactly as the night +// runs it after a backup — the hub's os_update block (fetched fresh), the wrapper via sudo, the health wait, the +// report to the hub — with trigger "debug" (never throttled, and it does NOT count as a night run for approval). +// Run it as the agent user: sudo -u felhom-agent felhom-agent --config … --selftest=os-update -vmid 9201 +func runSelftestOSUpdate(ctx context.Context, cfg config.Config, logger *slog.Logger, vmid int) int { + if vmid <= 0 { + fmt.Fprintln(os.Stderr, "selftest=os-update: -vmid is required") + return 2 + } + client, err := hub.NewClient(cfg.Hub, logger) + if err != nil { + fmt.Fprintln(os.Stderr, "selftest=os-update: hub client:", err) + return 1 + } + leg := newOSLeg(cfg, client, logger) + resp, err := client.FetchDesiredState(ctx) + if err != nil { + fmt.Fprintln(os.Stderr, "selftest=os-update: desired state:", err) + return 1 + } + leg.SetBlock(resp.DesiredState.OSUpdate) + b := leg.Block() + fmt.Printf("=== felhom-agent %s selftest=os-update vmid=%d ring=%d enabled=%v release=%v ===\n", version, vmid, b.Ring, b.Enabled, b.Release != nil) + rep := leg.Run(ctx, vmid, "debug") + printJSON("os-update report", map[string]any{"run_id": rep.RunID, "ring": rep.Ring, "release_id": rep.ReleaseID, + "mode": rep.Mode, "outcome": rep.Outcome, "healthy": rep.Healthy, "health_reason": rep.HealthReason, + "upgraded": rep.Upgraded, "pending": len(rep.Pending), "not_covered": rep.NotCovered, + "restart_needed": rep.RestartNeeded, "docker_restart_needed": rep.DockerRestartNeeded, "refused": rep.Refused}) + switch rep.Outcome { + case "applied", "nothing", "inventory", "skipped": + return 0 + } + return 1 +} diff --git a/configs/felhom-agent.sudoers b/configs/felhom-agent.sudoers index 07bc3d6..4d82a1a 100644 --- a/configs/felhom-agent.sudoers +++ b/configs/felhom-agent.sudoers @@ -299,10 +299,17 @@ Cmnd_Alias FELHOM_SELFHEAL = \ # argument after the numeric vmid is a literal, so the grant cannot be widened by anything the guest or # the hub says. The address read is deliberately NOT duplicated here — it is already FELHOM_DNSMASQ's, # and the same command must not be granted twice under two names. +# OS updates, guest fast lane (`11-os-updates.md` §5.4.1, agent v0.140.0). The ONLY entry: the root wrapper with +# one plan file in the agent's own os/ dir. Every safety rule (no removal, no downgrade, no new or unlisted package, +# Debian origin only, the box's own customer guest only) lives in the wrapper, red-proved per rule +# (configs/test_felhom_os_apply.py). The agent gets NO apt grant of its own. +Cmnd_Alias FELHOM_OSAPPLY = \ + /usr/local/sbin/felhom-os-apply --plan /var/lib/felhom-agent/os/plan-*.json + Cmnd_Alias FELHOM_GUESTNET = \ /usr/sbin/pct exec [0-9]* -- ip route show default, \ /usr/sbin/pct exec [0-9]* -- cat /etc/network/interfaces, \ /usr/sbin/pct exec [0-9]* -- pgrep -x dhclient, \ /usr/sbin/pct exec [0-9]* -- dhclient -pf /run/dhclient.eth0.pid -lf /var/lib/dhcp/dhclient.eth0.leases eth0 -felhom-agent ALL=(root) NOPASSWD: FELHOM_MOUNT, FELHOM_DISK, FELHOM_PROVISION, FELHOM_FORMAT, FELHOM_DNSMASQ, FELHOM_GUESTHOOK, FELHOM_INTERMEDIARY, FELHOM_CONTROLLERSWAP, FELHOM_STALELOCK, FELHOM_NETMOUNT, FELHOM_WG, FELHOM_SELFUPDATE, FELHOM_SSHD, FELHOM_OOB, FELHOM_PBSDR, FELHOM_BACKUPTARGET, FELHOM_SELFHEAL, FELHOM_ESCROW, FELHOM_GUESTNET, FELHOM_SCRATCH_TEARDOWN +felhom-agent ALL=(root) NOPASSWD: FELHOM_MOUNT, FELHOM_DISK, FELHOM_PROVISION, FELHOM_FORMAT, FELHOM_DNSMASQ, FELHOM_GUESTHOOK, FELHOM_INTERMEDIARY, FELHOM_CONTROLLERSWAP, FELHOM_STALELOCK, FELHOM_NETMOUNT, FELHOM_WG, FELHOM_SELFUPDATE, FELHOM_SSHD, FELHOM_OOB, FELHOM_PBSDR, FELHOM_BACKUPTARGET, FELHOM_SELFHEAL, FELHOM_ESCROW, FELHOM_GUESTNET, FELHOM_SCRATCH_TEARDOWN, FELHOM_OSAPPLY diff --git a/configs/felhom-os-apply b/configs/felhom-os-apply new file mode 100755 index 0000000..016022d --- /dev/null +++ b/configs/felhom-os-apply @@ -0,0 +1,475 @@ +#!/usr/bin/python3 +# felhom-os-apply — the ROOT half of the agent's operating-system update leg (`11-os-updates.md` §5.4.1). +# +# Install as /usr/local/sbin/felhom-os-apply (0755 root:root). The non-root agent invokes it via `sudo -n` +# (FELHOM_OSAPPLY alias) with EXACTLY: felhom-os-apply --plan /var/lib/felhom-agent/os/plan-.json +# Nothing else on the command line is accepted. Python 3, standard library only (a JSON plan cannot be parsed +# safely in sh). Tests: configs/test_felhom_os_apply.py (a fake runner; nothing real is executed). +# +# THE TRUST MODEL. The plan is written by the agent, so a broken-into agent writes whatever plan it likes. The +# protection is therefore what this file REFUSES, not where the plan came from: no removal, no downgrade, no new +# package, no package outside the plan, only Debian origin in the fast lane, only the box's own customer guest. +# Package signatures stay Debian's: apt checks every Release file against the guest's keyring, including the +# snapshot.debian.org fallback (decision 79). Nothing here is overridable from the environment. +# +# THIS RELEASE: layer "guest", lane "fast" only. The host layer and the slow lane exist in the interface and are +# REFUSED (R3, R12) until `11` §8 steps 3 and 5 enable them. +# +# Modes (plan field "mode"): +# inventory read-only for packages: `apt-get update` in the guest, then report what is installed (with origin), +# what is pending, restart-needed and health. Installs nothing. +# apply repair first, check every refusal on an `apt-get -s` simulation of EXACTLY name=version, then +# install, clean, and report the same as inventory. +# health report health only (the agent polls it after a run). +# Output: log lines on stderr and the journal (tag felhom-os-apply); the LAST stdout line is +# OSAPPLY-REPORT +# which is what the agent parses. Exit 0 = done; 2 = refused (nothing changed); 3 = failed during install. +import json +import os +import re +import stat +import subprocess +import sys +import time + +PLAN_DIR = "/var/lib/felhom-agent/os" +PLAN_RE = re.compile(r"^plan-[A-Za-z0-9._-]{1,80}\.json$") +AGENT_USER = "felhom-agent" +FAST_ORIGINS = ("Debian", "Debian-Security") +# Debian package name and version grammar (Debian policy §5.6.1, §5.6.12). +NAME_RE = re.compile(r"^[a-z0-9][a-z0-9+.-]+$") +VERSION_RE = re.compile(r"^(?:[0-9]+:)?[0-9][A-Za-z0-9.+~-]*$") +SNAP_RE = re.compile(r"^[0-9]{8}T[0-9]{6}Z$") +RESERVED_VMIDS = set(range(990000, 990010)) | {9999} +DRIVES_PARENT = "/mnt/felhom-drives" +SNAPSHOT_LIST = "/etc/apt/sources.list.d/felhom-os-snapshot.list" +APT_ENV = ["env", "DEBIAN_FRONTEND=noninteractive", "APT_LISTCHANGES_FRONTEND=none", "NEEDRESTART_MODE=l", "LC_ALL=C"] +DPKG_OPTS = ["-o", "Dpkg::Options::=--force-confold", "-o", "Dpkg::Options::=--force-confdef"] +MIN_FREE = 500 * 1024 * 1024 + + +class Refused(Exception): + def __init__(self, code, reason): + super().__init__(f"{code} {reason}") + self.code, self.reason = code, reason + + +class Runner: + """Runs commands for real. Tests replace it with a fake. `guest` runs inside the container via pct exec.""" + + def host(self, argv, timeout=600): + p = subprocess.run(argv, capture_output=True, text=True, timeout=timeout) + return p.returncode, p.stdout, p.stderr + + def guest(self, vmid, argv, timeout=1800): + return self.host(["/usr/sbin/pct", "exec", str(vmid), "--"] + argv, timeout) + + def read_file(self, path): + with open(path) as f: + return f.read() + + def stat(self, path): + return os.lstat(path) + + def agent_uid(self): + import pwd + return pwd.getpwnam(AGENT_USER).pw_uid + + def log(self, line): + print(line, file=sys.stderr, flush=True) + try: + subprocess.run(["logger", "-t", "felhom-os-apply", line], timeout=10) + except Exception: + pass + + +class Apply: + def __init__(self, runner, plan_path): + self.r = runner + self.plan_path = plan_path + self.report = {"refused": None, "mode": None} + + # ---------- checks ---------- + def load_plan(self): + p = self.plan_path + d, base = os.path.dirname(p), os.path.basename(p) + if d != PLAN_DIR or not PLAN_RE.match(base) or ".." in p: + raise Refused("R1", f"the plan must be {PLAN_DIR}/plan-.json, got {p!r}") + try: + st = self.r.stat(p) + except OSError as e: + raise Refused("R1", f"cannot stat the plan: {e}") + if not stat.S_ISREG(st.st_mode): + raise Refused("R1", "the plan is not a regular file (a symlink or a device is refused)") + if st.st_uid != self.r.agent_uid(): + raise Refused("R1", f"the plan is not owned by {AGENT_USER}") + if st.st_size > 2 * 1024 * 1024: + raise Refused("R1", "the plan is larger than 2 MB") + try: + plan = json.loads(self.r.read_file(p)) + except (OSError, ValueError) as e: + raise Refused("R1", f"the plan is not valid JSON: {e}") + if not isinstance(plan, dict): + raise Refused("R1", "the plan is not a JSON object") + return plan + + def check_plan(self, plan): + mode = plan.get("mode", "apply") + if mode not in ("apply", "inventory", "health"): + raise Refused("R11", f"unknown mode {mode!r}") + if plan.get("layer") != "guest": + raise Refused("R12", f"layer {plan.get('layer')!r} is refused in this release (guest only)") + if plan.get("lane", "fast") != "fast": + raise Refused("R3", "the slow lane is refused in this release") + vmid = plan.get("vmid") + if not isinstance(vmid, int) or isinstance(vmid, bool) or vmid <= 0: + raise Refused("R11", f"vmid must be a positive integer, got {vmid!r}") + rid = plan.get("release_id", "") + if not isinstance(rid, str) or not re.match(r"^[A-Za-z0-9._:-]{1,80}$", rid): + raise Refused("R11", f"release_id {rid!r} is not a plain id") + if plan.get("allow_new"): + raise Refused("R6", "allow_new is a slow-lane field; the fast lane never adds a package") + pk = plan.get("packages", []) + if not isinstance(pk, list) or (mode == "apply" and not pk): + raise Refused("R11", "packages must be a non-empty list in apply mode") + seen = set() + for e in pk: + if not isinstance(e, dict): + raise Refused("R11", "every package entry must be an object") + n, v, o = e.get("name"), e.get("version"), e.get("origin") + if not isinstance(n, str) or not NAME_RE.match(n): + raise Refused("R11", f"package name {n!r} is not a Debian package name") + if not isinstance(v, str) or not VERSION_RE.match(v): + raise Refused("R11", f"version {v!r} of {n} is not a Debian version string") + if n in seen: + raise Refused("R11", f"package {n} is named twice") + seen.add(n) + if o not in FAST_ORIGINS: + raise Refused("R2", f"{n}: origin {o!r} is not Debian / Debian-Security (the fast lane, `11` C3)") + snap = plan.get("snapshot", "") + if snap and not SNAP_RE.match(snap): + raise Refused("R11", f"snapshot {snap!r} is not YYYYMMDDTHHMMSSZ") + return mode, vmid + + def check_guest(self, vmid): + if vmid in RESERVED_VMIDS: + raise Refused("R10", f"vmid {vmid} is a reserved scratch vmid") + try: + conf = self.r.read_file(f"/etc/pve/lxc/{vmid}.conf") + except OSError: + raise Refused("R10", f"vmid {vmid} is not a container on this host") + cur = conf.split("\n[", 1)[0] # the current config, not a snapshot section + binds = [l for l in cur.splitlines() if re.match(r"^mp[0-9]+: " + re.escape(DRIVES_PARENT) + r",", l)] + if not binds: + raise Refused("R10", f"vmid {vmid} does not bind {DRIVES_PARENT} — it is not this box's customer guest") + lock = [l for l in cur.splitlines() if l.startswith("lock:")] + if lock: + raise Refused("R9", f"vmid {vmid} is locked ({lock[0].split(':', 1)[1].strip()}) — a backup or restore is running") + rc, out, _ = self.r.host(["/usr/sbin/pct", "status", str(vmid)]) + if rc != 0 or "running" not in out: + raise Refused("R10", f"vmid {vmid} is not running") + + # ---------- guest helpers ---------- + def g(self, argv, timeout=1800): + return self.r.guest(self.vmid, argv, timeout) + + def installed(self): + rc, out, _ = self.g(["dpkg-query", "-W", "-f", "${Package}\t${Version}\t${db:Status-Abbrev}\n"]) + res = {} + for l in out.splitlines(): + parts = l.split("\t") + if len(parts) == 3 and parts[2].startswith("ii"): + res[parts[0]] = parts[1] + return res + + def dpkg_cmp(self, a, op, b): + rc, _, _ = self.g(["dpkg", "--compare-versions", a, op, b]) + return rc == 0 + + def madison(self, name): + rc, out, _ = self.g(["apt-cache", "madison", name]) + vs = set() + for l in out.splitlines(): + f = [x.strip() for x in l.split("|")] + if len(f) >= 3 and f[0] == name: + vs.add(f[1]) + return vs + + def simulate(self, args): + rc, out, err = self.g(APT_ENV + ["apt-get", "-s", "-q"] + args) + inst, remv = [], [] + for l in out.splitlines(): + m = re.match(r"^Inst (\S+) (?:\[([^]]*)\] )?\((\S+) (.*?) \[[a-z0-9]+\]\)", l) + if m: + inst.append({"name": m.group(1), "from": m.group(2), "to": m.group(3), "origin": m.group(4)}) + m = re.match(r"^Remv (\S+)", l) + if m: + remv.append(m.group(1)) + return rc, inst, remv, out + err + + @staticmethod + def origin_name(origin): + # "Debian:13.7/stable, Debian-Security:13/stable-security" -> {"Debian", "Debian-Security"} + return {o.strip().split(":")[0] for o in origin.split(",") if o.strip()} + + def free_bytes(self): + rc, out, _ = self.g(["df", "-B1", "--output=avail", "/"]) + try: + return int(out.strip().splitlines()[-1]) + except (ValueError, IndexError): + return -1 + + def apt_lock_held(self): + rc, out, _ = self.g(["fuser", "/var/lib/dpkg/lock-frontend", "/var/lib/dpkg/lock"]) + return rc == 0 and out.strip() != "" + + def health(self): + """The guest's signals: every container's state + health, the controller's own health, the network.""" + rc, out, _ = self.g(["docker", "ps", "-a", "--format", "{{.Names}}\t{{.State}}\t{{.Status}}"], timeout=60) + cont = {} + for l in out.splitlines(): + p = l.split("\t") + if len(p) == 3: + h = "healthy" if "(healthy)" in p[2] else "unhealthy" if "(unhealthy)" in p[2] else \ + "starting" if "(health: starting)" in p[2] else "none" + cont[p[0]] = {"state": p[1], "health": h} + nrc, _, _ = self.g(["getent", "hosts", "deb.debian.org"], timeout=30) + return {"docker_ok": rc == 0, "containers": cont, + "controller": cont.get("felhom-controller", {}).get("health", "absent"), + "network_ok": nrc == 0} + + def restart_needed(self): + """Processes still mapping deleted files, OUTSIDE docker containers (C11).""" + script = ('for p in /proc/[0-9]*; do grep -q "(deleted)" $p/maps 2>/dev/null || continue; ' + 'grep -q "docker" $p/cgroup 2>/dev/null && continue; echo "${p#/proc/} $(cat $p/comm 2>/dev/null)"; done') + rc, out, _ = self.g(["sh", "-c", script], timeout=120) + procs = sorted({l.split(" ", 1)[1] for l in out.splitlines() if " " in l}) + pid1 = any(l.split(" ", 1)[0] == "1" for l in out.splitlines()) + return procs, pid1 + + def inventory(self): + inst = self.installed() + names = sorted(inst) + origins = {} + for i in range(0, len(names), 200): + rc, out, _ = self.g(["apt-cache", "policy"] + names[i:i + 200]) + cur, star = None, False + for l in out.splitlines(): + if not l.startswith(" "): + cur, star = l.rstrip(":"), False + continue + s = l.strip() + if s.startswith("*** "): + star = True + continue + if star and cur and re.match(r"^[0-9-]+ ", s): + if "/var/lib/dpkg/status" in s: + origins.setdefault(cur, "local") + else: + origins[cur] = s + continue + if star and not re.match(r"^[0-9-]+ ", s): + star = False + # Map an index URL to an origin name the hub understands. + def oname(src): + if src in (None, "local"): + return "unknown" + if "security" in src and "debian" in src: + return "Debian-Security" + if "docker.com" in src: + return "Docker" + if "debian" in src: + return "Debian" + return "other" + rc, pend, remv, _ = self.simulate(["dist-upgrade"]) + return { + "installed": [{"name": n, "version": inst[n], "origin": oname(origins.get(n))} for n in names], + "pending": [{"name": p["name"], "from": p["from"], "to": p["to"], + "origin": sorted(self.origin_name(p["origin"]))} for p in pend], + } + + # ---------- the run ---------- + def run(self): + plan = self.load_plan() + self.mode, self.vmid = self.check_plan(plan) + self.report.update(mode=self.mode, release_id=plan.get("release_id"), vmid=self.vmid) + self.check_guest(self.vmid) + log = self.r.log + if self.mode == "health": + self.report["health"] = self.health() + return 0 + log(f"os-apply: START release={plan.get('release_id')} layer=guest:{self.vmid} lane=fast mode={self.mode} packages={len(plan.get('packages', []))}") + if self.apt_lock_held(): + raise Refused("R9", "another apt/dpkg holds the lock in the guest") + self.report["health_before"] = self.health() + if self.mode == "apply": + self.repair() + rc, out, err = self.g(APT_ENV + ["apt-get", "-q", "update"], timeout=600) + if rc != 0: + raise Refused("R7", f"apt-get update failed in the guest: {(out + err).strip().splitlines()[-1:]}") + if self.mode == "apply": + rc = self.apply(plan) + if rc: + return rc + procs, pid1 = self.restart_needed() + self.report.update(self.inventory()) + self.report["restart_needed"] = procs + self.report["docker_restart_needed"] = any(p in ("dockerd", "containerd") for p in procs) + self.report["reboot_needed"] = pid1 + self.report["health_after"] = self.health() + return 0 + + def repair(self): + rc, before, _ = self.g(["dpkg", "--audit"]) + self.g(APT_ENV + ["dpkg", "--configure", "-a", "--force-confold"]) + rc2, out, err = self.g(APT_ENV + ["apt-get", "-f", "install", "-y", "-q"] + DPKG_OPTS) + _, after, _ = self.g(["dpkg", "--audit"]) + configured = len([l for l in before.splitlines() if l.startswith(" ")]) + fixed = len(re.findall(r"^Setting up ", out, re.M)) + self.report["repair"] = {"half_configured_before": configured, "fixed": fixed, "clean_after": after.strip() == ""} + self.r.log(f"os-apply: REPAIR configured={configured} fixed={fixed}") + if after.strip(): + raise Refused("R13", "dpkg is still broken after the repair: " + after.strip().splitlines()[0]) + + def apply(self, plan): + log = self.r.log + inst = self.installed() + upgrade, already, notinst = [], 0, 0 + for e in plan["packages"]: + n, v = e["name"], e["version"] + if n not in inst: + notinst += 1 + continue + if not self.dpkg_cmp(v, "gt", inst[n]): + already += 1 + continue + upgrade.append((n, v)) + from_snap = 0 + missing = [(n, v) for n, v in upgrade if v not in self.madison(n)] + if missing: + snap = plan.get("snapshot", "") + if not snap: + raise Refused("R7", f"{missing[0][0]}={missing[0][1]} is not downloadable and the plan names no snapshot") + self.add_snapshot_sources(snap) + still = [(n, v) for n, v in missing if v not in self.madison(n)] + if still: + self.remove_snapshot_sources() + raise Refused("R7", f"{still[0][0]}={still[0][1]} is not downloadable, not even from snapshot {snap}") + from_snap = len(missing) + try: + log(f"os-apply: PLAN upgrade={len(upgrade)} already={already} not-installed={notinst} from-snapshot={from_snap}") + self.report["plan"] = {"upgrade": len(upgrade), "already": already, "not_installed": notinst, "from_snapshot": from_snap} + if not upgrade: + self.report["upgraded"] = [] + log("os-apply: DONE rc=0 seconds=0 upgraded=0 (nothing to do)") + return 0 + args = ["install", "--only-upgrade", "--no-install-recommends"] + [f"{n}={v}" for n, v in upgrade] + rc, sim, remv, text = self.simulate(args) + if rc != 0: + tail = text.strip().splitlines()[-1] if text.strip() else "" + raise Refused("R7", "the simulation failed: " + tail) + if remv: + raise Refused("R4", f"the plan would remove {', '.join(remv[:5])}") + want = dict(upgrade) + for p in sim: + if p["from"] is None: + raise Refused("R6", f"the plan would add a package that is not installed: {p['name']}") + if p["name"] not in want: + raise Refused("R6", f"the plan would touch {p['name']}, which is not in the plan") + if p["to"] != want[p["name"]]: + raise Refused("R6", f"{p['name']} would go to {p['to']}, not the approved {want[p['name']]}") + if not self.dpkg_cmp(p["to"], "gt", p["from"]): + raise Refused("R5", f"{p['name']} would be downgraded {p['from']} -> {p['to']}") + if not self.origin_name(p["origin"]) & set(FAST_ORIGINS): + raise Refused("R2", f"{p['name']} would come from {p['origin']}, not Debian") + need = self.download_bytes(args) + free = self.free_bytes() + if free >= 0 and free < max(MIN_FREE, 3 * need): + raise Refused("R8", f"free space {free} B is below max(500 MB, 3 x download {need} B)") + t0 = time.time() + rc, out, err = self.g(APT_ENV + ["apt-get", "-y", "-q"] + DPKG_OPTS + args) + secs = time.time() - t0 + for l in (out + err).splitlines(): + m = re.search(r"Installing new version of config file (\S+)|Configuration file '([^']+)'", l) + if m: + log(f"os-apply: CONFFILE kept {m.group(1) or m.group(2)}") + self.g(["apt-get", "clean"]) + if rc != 0: + _, aud, _ = self.g(["dpkg", "--audit"]) + first = aud.strip().splitlines()[0] if aud.strip() else "clean" + log(f"os-apply: FAILED rc={rc} step=install — dpkg state: {first}") + self.report["failed"] = {"rc": rc, "dpkg_audit": first, "tail": (out + err).strip().splitlines()[-3:]} + return 3 + self.report["upgraded"] = [{"name": n, "version": v} for n, v in upgrade] + self.report["seconds"] = round(secs, 1) + log(f"os-apply: DONE rc=0 seconds={secs:.1f} upgraded={len(upgrade)}") + return 0 + finally: + if from_snap: + self.remove_snapshot_sources() + + def download_bytes(self, args): + rc, out, _ = self.g(APT_ENV + ["apt-get", "-s", "-o", "Debug::NoLocking=1", "--print-uris", "-q"] + args) + total = 0 + for l in out.splitlines(): + m = re.match(r"^'[^']+' \S+ ([0-9]+) ", l) + if m: + total += int(m.group(1)) + return total + + def add_snapshot_sources(self, snap): + rc, out, _ = self.g(["sh", "-c", ". /etc/os-release && echo $VERSION_CODENAME"]) + code = out.strip() + if not re.match(r"^[a-z]+$", code): + raise Refused("R7", f"cannot read the guest's Debian codename ({code!r})") + body = (f"deb [check-valid-until=no] http://snapshot.debian.org/archive/debian/{snap} {code} main\n" + f"deb [check-valid-until=no] http://snapshot.debian.org/archive/debian-security/{snap} {code}-security main\n") + self.r.guest_write(self.vmid, SNAPSHOT_LIST, body) + self.r.log(f"os-apply: SNAPSHOT using snapshot.debian.org/{snap} for versions no longer published (decision 79)") + rc, out, err = self.g(APT_ENV + ["apt-get", "-q", "update"], timeout=600) + if rc != 0: + self.remove_snapshot_sources() + raise Refused("R7", "apt-get update against snapshot.debian.org failed") + + def remove_snapshot_sources(self): + self.g(["rm", "-f", SNAPSHOT_LIST]) + self.g(APT_ENV + ["apt-get", "-q", "update"], timeout=600) + + +def guest_write(self, vmid, path, body): + """Write a small text file inside the guest via `pct exec … tee` (stdin), never via a shell string.""" + p = subprocess.run(["/usr/sbin/pct", "exec", str(vmid), "--", "tee", path], input=body, + capture_output=True, text=True, timeout=60) + if p.returncode != 0: + raise Refused("R7", f"could not write {path} in the guest") + + +Runner.guest_write = guest_write + + +def main(argv, runner=None): + r = runner or Runner() + if len(argv) != 3 or argv[1] != "--plan": + r.log("os-apply: REFUSED: R1 usage: felhom-os-apply --plan /var/lib/felhom-agent/os/plan-.json") + print("OSAPPLY-REPORT " + json.dumps({"refused": {"code": "R1", "reason": "usage"}})) + return 2 + a = Apply(r, argv[2]) + try: + rc = a.run() + except Refused as e: + r.log(f"os-apply: REFUSED: {e.code} {e.reason}") + a.report["refused"] = {"code": e.code, "reason": e.reason} + rc = 2 + except subprocess.TimeoutExpired as e: + r.log(f"os-apply: FAILED rc=124 step=timeout — {e.cmd}") + a.report["failed"] = {"rc": 124, "timeout": str(e.cmd)[:200]} + rc = 3 + print("OSAPPLY-REPORT " + json.dumps(a.report, sort_keys=True)) + return rc + + +if __name__ == "__main__": + if os.geteuid() != 0: + print("felhom-os-apply: must run as root (via sudo)", file=sys.stderr) + sys.exit(2) + sys.exit(main(sys.argv)) diff --git a/configs/test_felhom_os_apply.py b/configs/test_felhom_os_apply.py new file mode 100644 index 0000000..13cb85e --- /dev/null +++ b/configs/test_felhom_os_apply.py @@ -0,0 +1,452 @@ +#!/usr/bin/env python3 +"""Tests for configs/felhom-os-apply (`11` §5.4.1). A fake runner plays the host and the guest: nothing is executed +for real except the local `dpkg --compare-versions` (pure, no network). Every refusal R1–R13 has a test; the red-proof +(each test fails when its rule is removed) is `audits/os-guest-lane-2026-10-04/partB/redproof.txt`. + +Run: python3 configs/test_felhom_os_apply.py (also run by internal/osupdate's Go test) +""" +import importlib.machinery +import importlib.util +import json +import os +import pathlib +import re +import stat as statmod +import subprocess +import unittest + +HERE = pathlib.Path(__file__).resolve().parent +_loader = importlib.machinery.SourceFileLoader("osapply", os.environ.get("OSAPPLY_UNDER_TEST", str(HERE / "felhom-os-apply"))) # red-proof seam +_spec = importlib.util.spec_from_loader("osapply", _loader) +osapply = importlib.util.module_from_spec(_spec) +_loader.exec_module(osapply) + +PLAN = "/var/lib/felhom-agent/os/plan-t1.json" +CONF_OK = ("arch: amd64\nmp0: local-lvm:vm-9201-disk-1,mp=/var/lib/felhom,backup=1,size=70G\n" + "mp8: /mnt/felhom-drives,mp=/mnt/felhom-drives\nrootfs: local-lvm:vm-9201-disk-0,size=32G\n") +DEB = "Debian:13.7/stable" +SEC = "Debian-Security:13/stable-security" + + +def dpkg_cmp(a, op, b): + return subprocess.run(["dpkg", "--compare-versions", a, op, b]).returncode == 0 + + +class St: + def __init__(self, mode=statmod.S_IFREG | 0o600, uid=999, size=100): + self.st_mode, self.st_uid, self.st_size = mode, uid, size + + +class Fake: + """The host + one guest. `installed` / `live` (name -> versions in the live archive) / `snapshot` (versions + the snapshot archive adds) / `extra_sim` (lines the simulation adds) / `dpkg_audit` / `free`.""" + + def __init__(self): + self.plan = {"release_id": "os-t1", "layer": "guest", "lane": "fast", "vmid": 9201, "mode": "apply", + "snapshot": "20261004T080000Z", + "packages": [{"name": "libc6", "version": "2.41-12+deb13u4", "origin": "Debian"}, + {"name": "openssl", "version": "3.5.7-1~deb13u3", "origin": "Debian-Security"}]} + self.files = {"/etc/pve/lxc/9201.conf": CONF_OK} + self.stats = {PLAN: St()} + self.installed = {"libc6": "2.41-12+deb13u3", "openssl": "3.5.6-1~deb13u1", "bash": "5.2.37-2+b9"} + self.live = {"libc6": {"2.41-12+deb13u4"}, "openssl": {"3.5.7-1~deb13u3"}} + self.snapshot = {} + self.snap_active = False + self.extra_sim = [] + self.dpkg_audit = "" + self.free = 10 * 1024 ** 3 + self.install_rc = 0 + self.calls = [] + self.logs = [] + self.written = {} + self.status = "status: running" + self.lock_held = False + + # Runner interface + def read_file(self, p): + if p == PLAN: + return json.dumps(self.plan) + if p not in self.files: + raise OSError("no such file") + return self.files[p] + + def stat(self, p): + if p not in self.stats: + raise OSError("no such file") + return self.stats[p] + + def agent_uid(self): + return 999 + + def log(self, line): + self.logs.append(line) + + def host(self, argv, timeout=600): + self.calls.append(("host", argv)) + if argv[1] == "status": + return 0, self.status + "\n", "" + return 1, "", "unexpected host call" + + def guest_write(self, vmid, path, body): + self.written[path] = body + if path == osapply.SNAPSHOT_LIST: + self.snap_active = True + + def avail(self, n): + v = set(self.live.get(n, set())) + if self.snap_active: + v |= self.snapshot.get(n, set()) + return v + + def guest(self, vmid, argv, timeout=1800): + self.calls.append(("guest", vmid, argv)) + a = [x for x in argv if not re.match(r"^[A-Z_]+=", x) and x != "env"] + cmd = a[0] + if cmd == "dpkg-query": + return 0, "".join(f"{n}\t{v}\tii \n" for n, v in self.installed.items()), "" + if cmd == "dpkg" and a[1] == "--compare-versions": + return (0 if dpkg_cmp(a[2], a[3], a[4]) else 1), "", "" + if cmd == "dpkg" and a[1] == "--audit": + return 0, self.dpkg_audit, "" + if cmd == "dpkg" and a[1] == "--configure": + return 0, "", "" + if cmd == "fuser": + return (0, " 123", "") if self.lock_held else (1, "", "") + if cmd == "apt-cache" and a[1] == "madison": + return 0, "".join(f" {a[2]} | {v} | http://deb.debian.org trixie/main amd64 Packages\n" for v in self.avail(a[2])), "" + if cmd == "apt-cache" and a[1] == "policy": + out = "" + for n in a[2:]: + out += f"{n}:\n Installed: {self.installed.get(n)}\n Version table:\n *** {self.installed.get(n)} 500\n 500 http://deb.debian.org/debian trixie/main amd64 Packages\n" + return 0, out, "" + if cmd == "apt-get": + if "update" in a: + return 0, "", "" + if "clean" in a: + return 0, "", "" + if "-f" in a: + self.dpkg_audit = "" + return 0, "Setting up x (1) ...\n" if getattr(self, "repaired", False) else "", "" + if "-s" in a: + return self.sim(a) + if "install" in a: + if self.install_rc: + return self.install_rc, "", "E: boom" + for x in a: + if "=" in x and not x.startswith("-") and "::" not in x: + n, v = x.split("=", 1) + self.installed[n] = v + return 0, "Setting up libc6 ...\n", "" + if cmd == "df": + return 0, f"Avail\n{self.free}\n", "" + if cmd == "docker": + return 0, "felhom-controller\trunning\tUp 1 hour (healthy)\napp\trunning\tUp 1 hour (healthy)\n", "" + if cmd == "getent": + return 0, "1.2.3.4 deb.debian.org\n", "" + if cmd == "sh": + if "os-release" in a[2]: + return 0, "trixie\n", "" + return 0, "", "" + if cmd == "rm": + self.snap_active = False + return 0, "", "" + return 1, "", f"unexpected guest call {a}" + + def sim(self, a): + if "--print-uris" in a: + return 0, "'http://x/libc6.deb' libc6.deb 4000000 SHA256:x\n", "" + if "dist-upgrade" in a: + return 0, "Inst bash [5.2.37-2+b9] (5.2.37-2+b10 Debian:13.7/stable [amd64])\n", "" + out = "" + for x in a: + if "=" in x and not x.startswith("-") and "::" not in x: + n, v = x.split("=", 1) + if v not in self.avail(n): + return 100, "", f"E: Version '{v}' for '{n}' was not found" + origin = SEC if n == "openssl" else DEB + out += f"Inst {n} [{self.installed[n]}] ({v} {origin} [amd64])\n" + out += "".join(l + "\n" for l in self.extra_sim) + return 0, out, "" + + +def run(f): + import io + import contextlib + buf = io.StringIO() + with contextlib.redirect_stdout(buf): + rc = osapply.main(["felhom-os-apply", "--plan", PLAN], runner=f) + line = [l for l in buf.getvalue().splitlines() if l.startswith("OSAPPLY-REPORT ")][-1] + return rc, json.loads(line[len("OSAPPLY-REPORT "):]) + + +class Happy(unittest.TestCase): + def test_apply_installs_exactly_the_plan(self): + f = Fake() + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + self.assertEqual(f.installed["libc6"], "2.41-12+deb13u4") + self.assertEqual(f.installed["openssl"], "3.5.7-1~deb13u3") + self.assertEqual(f.installed["bash"], "5.2.37-2+b9", "a package outside the plan was changed") + self.assertEqual(rep["plan"]["upgrade"], 2) + self.assertIn("installed", rep) + self.assertEqual(rep["pending"][0]["name"], "bash") + self.assertTrue(any(l.startswith("os-apply: REPAIR ") for l in f.logs), "the repair line must always print") + self.assertTrue(any(l.startswith("os-apply: DONE rc=0") for l in f.logs)) + inst = [c for c in f.calls if c[0] == "guest" and "install" in c[2] and "-s" not in c[2] and "-f" not in c[2]] + self.assertTrue(inst and "Dpkg::Options::=--force-confold" in inst[0][2], "must keep existing config files") + + def test_already_current_is_a_no_op(self): + f = Fake() + f.installed.update(libc6="2.41-12+deb13u4", openssl="3.5.7-1~deb13u3") + rc, rep = run(f) + self.assertEqual(rc, 0) + self.assertEqual(rep["plan"]["upgrade"], 0) + + def test_inventory_installs_nothing(self): + f = Fake() + f.plan["mode"] = "inventory" + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + self.assertEqual(f.installed["libc6"], "2.41-12+deb13u3") + self.assertIn("installed", rep) + self.assertIn("restart_needed", rep) + + def test_health_mode(self): + f = Fake() + f.plan["mode"] = "health" + rc, rep = run(f) + self.assertEqual(rc, 0) + self.assertEqual(rep["health"]["controller"], "healthy") + + +class Repair(unittest.TestCase): + def test_repair_runs_first_and_is_reported(self): + f = Fake() + f.dpkg_audit = "The following packages have been unpacked but not yet configured.\n perl Larry Wall's\n" + f.repaired = True + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + self.assertEqual(rep["repair"]["half_configured_before"], 1) + self.assertEqual(rep["repair"]["fixed"], 1) + order = [i for i, c in enumerate(f.calls) if c[0] == "guest" and c[2][-1:] != ["update"]] + first_cfg = next(i for i, c in enumerate(f.calls) if c[0] == "guest" and "--configure" in c[2]) + first_upd = next(i for i, c in enumerate(f.calls) if c[0] == "guest" and "update" in c[2]) + self.assertLess(first_cfg, first_upd, "the repair must run before anything else touches apt") + self.assertTrue(order) + + +class Snapshot(unittest.TestCase): + def test_a_replaced_version_comes_from_the_snapshot(self): + f = Fake() + f.live["openssl"] = {"3.5.7-1~deb13u4"} # Debian moved on + f.snapshot["openssl"] = {"3.5.7-1~deb13u3"} + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + self.assertEqual(rep["plan"]["from_snapshot"], 1) + self.assertEqual(f.installed["openssl"], "3.5.7-1~deb13u3", "must install the APPROVED version, not the newer one") + body = f.written[osapply.SNAPSHOT_LIST] + self.assertIn("snapshot.debian.org/archive/debian/20261004T080000Z trixie main", body) + self.assertIn("debian-security/20261004T080000Z trixie-security main", body) + self.assertFalse(f.snap_active, "the temporary snapshot sources must be removed after the run") + + def test_snapshot_does_not_have_it_either(self): + f = Fake() + f.live["openssl"] = set() + rc, rep = run(f) + self.assertEqual((rc, rep["refused"]["code"]), (2, "R7")) + self.assertFalse(f.snap_active) + + +class Refusals(unittest.TestCase): + def refused(self, f, code): + rc, rep = run(f) + self.assertEqual(rc, 2, rep) + self.assertEqual(rep["refused"]["code"], code, rep) + self.assertTrue(any(l.startswith(f"os-apply: REFUSED: {code} ") for l in f.logs), f.logs) + inst = [c for c in f.calls if c[0] == "guest" and "install" in c[2] and "-s" not in c[2] and "-f" not in c[2]] + self.assertEqual(inst, [], "a refusal must install nothing") + return rep + + def test_R1_usage(self): + import io + import contextlib + f = Fake() + with contextlib.redirect_stdout(io.StringIO()): + self.assertEqual(osapply.main(["felhom-os-apply", "--plan", PLAN, "--extra"], runner=f), 2) + self.assertEqual(osapply.main(["felhom-os-apply", "--plan"], runner=f), 2) + + def test_R1_path_outside_the_plan_dir(self): + import io + import contextlib + f = Fake() + with contextlib.redirect_stdout(io.StringIO()): + rc = osapply.main(["felhom-os-apply", "--plan", "/tmp/plan-x.json"], runner=f) + self.assertEqual(rc, 2) + self.assertTrue(any("R1" in l for l in f.logs)) + + def test_R1_symlink(self): + f = Fake() + f.stats[PLAN] = St(mode=statmod.S_IFLNK | 0o777) + self.refused(f, "R1") + + def test_R1_not_owned_by_the_agent(self): + f = Fake() + f.stats[PLAN] = St(uid=0) + self.refused(f, "R1") + + def test_R2_non_debian_origin_in_the_plan(self): + f = Fake() + f.plan["packages"][0]["origin"] = "Proxmox" + self.refused(f, "R2") + + def test_R2_non_debian_origin_in_the_simulation(self): + f = Fake() + f.installed["libc6"] = "2.41-12+deb13u3" + orig = f.sim + + def sim(a): + rc, out, err = orig(a) + return rc, out.replace("Debian:13.7/stable", "Proxmox Debian Repository:stable"), err + f.sim = sim + self.refused(f, "R2") + + def test_R3_slow_lane(self): + f = Fake() + f.plan["lane"] = "slow" + self.refused(f, "R3") + + def test_R4_removal(self): + f = Fake() + f.extra_sim = ["Remv bash [5.2.37-2+b9]"] + self.refused(f, "R4") + + def test_R5_downgrade(self): + f = Fake() + f.installed["libc6"] = "2.41-12+deb13u4" + f.plan["packages"] = [{"name": "openssl", "version": "3.5.7-1~deb13u3", "origin": "Debian-Security"}] + f.extra_sim = ["Inst openssl [3.5.6-1~deb13u1] (3.5.5-1 Debian:13.7/stable [amd64])"] + orig = f.sim + + def sim(a): # the simulation answers with a LOWER version than installed + rc, out, err = orig(a) + return rc, "\n".join(l for l in out.splitlines() if not l.startswith("Inst openssl [3.5.6-1~deb13u1] (3.5.7")) + "\n", err + f.sim = sim + f.plan["packages"][0]["version"] = "3.5.7-1~deb13u3" + rep = run(f)[1] + # The plan asks 3.5.7; the simulation goes to 3.5.5: that is BOTH a wrong version (R6) and a downgrade. + self.assertIn(rep["refused"]["code"], ("R5", "R6")) + + def test_R5_downgrade_exact(self): + f = Fake() + f.plan["packages"] = [{"name": "openssl", "version": "3.5.7-1~deb13u3", "origin": "Debian-Security"}] + f.installed["openssl"] = "3.5.6-1~deb13u1" + orig = f.sim + + def sim(a): + rc, out, err = orig(a) + return rc, out.replace("[3.5.6-1~deb13u1]", "[3.5.8-1]"), err + f.sim = sim + self.refused(f, "R5") + + def test_R6_new_package(self): + f = Fake() + f.extra_sim = ["Inst newthing (1.0 Debian:13.7/stable [amd64])"] + self.refused(f, "R6") + + def test_R6_unlisted_package(self): + f = Fake() + f.extra_sim = ["Inst bash [5.2.37-2+b9] (5.2.37-2+b10 Debian:13.7/stable [amd64])"] + self.refused(f, "R6") + + def test_R6_allow_new_is_slow_lane(self): + f = Fake() + f.plan["allow_new"] = ["proxmox-kernel-x"] + self.refused(f, "R6") + + def test_R7_not_downloadable_and_no_snapshot(self): + f = Fake() + f.live["openssl"] = set() + f.plan["snapshot"] = "" + self.refused(f, "R7") + + def test_R8_free_space(self): + f = Fake() + f.free = 100 * 1024 * 1024 + self.refused(f, "R8") + + def test_R9_guest_locked_by_a_backup(self): + f = Fake() + f.files["/etc/pve/lxc/9201.conf"] = CONF_OK + "lock: backup\n" + self.refused(f, "R9") + + def test_R9_apt_lock_held(self): + f = Fake() + f.lock_held = True + self.refused(f, "R9") + + def test_R10_not_the_boxs_own_guest(self): + f = Fake() + f.files["/etc/pve/lxc/9201.conf"] = CONF_OK.replace("mp8: /mnt/felhom-drives,", "mp8: /mnt/hdd_1/scratch,") + self.refused(f, "R10") + + def test_R10_reserved_vmid(self): + f = Fake() + f.plan["vmid"] = 990003 + self.refused(f, "R10") + + def test_R10_bind_only_in_a_snapshot_section(self): + f = Fake() + f.files["/etc/pve/lxc/9201.conf"] = "rootfs: x\n[snap1]\nmp8: /mnt/felhom-drives,mp=/mnt/felhom-drives\n" + self.refused(f, "R10") + + def test_R10_not_running(self): + f = Fake() + f.status = "status: stopped" + self.refused(f, "R10") + + def test_R11_duplicate(self): + f = Fake() + f.plan["packages"].append(dict(f.plan["packages"][0])) + self.refused(f, "R11") + + def test_R11_bad_version_string(self): + f = Fake() + f.plan["packages"][0]["version"] = "1.0; rm -rf /" + self.refused(f, "R11") + + def test_R11_bad_name(self): + f = Fake() + f.plan["packages"][0]["name"] = "--purge" + self.refused(f, "R11") + + def test_R12_host_layer(self): + f = Fake() + f.plan["layer"] = "host" + self.refused(f, "R12") + + def test_R13_repair_does_not_fix_it(self): + f = Fake() + f.dpkg_audit = "The following packages are broken\n perl\n" + orig = f.guest + + def guest(vmid, argv, timeout=1800): + rc, out, err = orig(vmid, argv, timeout) + if "-f" in argv: + f.dpkg_audit = "The following packages are broken\n perl\n" + return rc, out, err + f.guest = guest + self.refused(f, "R13") + + +class Failure(unittest.TestCase): + def test_install_failure_is_rc3_with_dpkg_state(self): + f = Fake() + f.install_rc = 100 + rc, rep = run(f) + self.assertEqual(rc, 3) + self.assertEqual(rep["failed"]["rc"], 100) + self.assertTrue(any(l.startswith("os-apply: FAILED rc=100 step=install") for l in f.logs)) + + +if __name__ == "__main__": + unittest.main() diff --git a/internal/capability/manifest.go b/internal/capability/manifest.go index 67b7acd..871d9a6 100644 --- a/internal/capability/manifest.go +++ b/internal/capability/manifest.go @@ -123,6 +123,9 @@ var manifest = []Capability{ {"dnsmasq-install", "dnsmasq package install", "/usr/bin/apt-get", []string{"install", "-y", "-q", "dnsmasq"}, false, ""}, {"dnsmasq-write", "dnsmasq drop-in write", "/usr/bin/install", []string{"-m", "0644", "/tmp/felhom-resolver-x.conf", "/etc/dnsmasq.d/felhom-x.conf"}, false, ""}, {"dnsmasq-enable", "dnsmasq enable", "/usr/bin/systemctl", []string{"enable", "--now", "dnsmasq"}, false, ""}, + + // ---- OS updates, guest fast lane (`11` §5.4.1; the wrapper holds every rule) ---- + {"osapply-run", "OS update wrapper (guest fast lane)", "/usr/local/sbin/felhom-os-apply", []string{"--plan", "/var/lib/felhom-agent/os/plan-x.json"}, false, ""}, {"dnsmasq-reload", "dnsmasq reload", "/usr/bin/systemctl", []string{"reload", "dnsmasq"}, false, ""}, {"dnsmasq-restart", "dnsmasq restart (LAN-DNS self-heal)", "/usr/bin/systemctl", []string{"restart", "dnsmasq"}, false, ""}, {"dnsmasq-rm", "dnsmasq drop-in remove (decommission)", "/usr/bin/rm", []string{"-f", "/etc/dnsmasq.d/felhom-x.conf"}, false, ""}, diff --git a/internal/hub/client.go b/internal/hub/client.go index 1859b57..49428be 100644 --- a/internal/hub/client.go +++ b/internal/hub/client.go @@ -425,3 +425,27 @@ func (c *Client) FetchRetainedIdentityEscrow(ctx context.Context) (*RetainedEscr } return &out, nil } + +// PostOSReport sends the OS-update leg's report after every run (hub v0.130.0): POST /api/v1/hosts/{id}/os-report. +// Per-host key, self-scoped on the hub. Errors are typed like RegisterWG's and never include the bearer. +func (c *Client) PostOSReport(ctx context.Context, body []byte) error { + if c.hostID == "" { + return fmt.Errorf("hub: PostOSReport requires a configured host_id") + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/api/v1/hosts/"+c.hostID+"/os-report", bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("hub: building os-report request: %w", err) + } + req.Header.Set("Authorization", "Bearer "+c.apiKey) + req.Header.Set("Content-Type", "application/json") + resp, err := c.hc.Do(req) + if err != nil { + return &TransportError{Err: err} + } + defer resp.Body.Close() + raw, _ := io.ReadAll(io.LimitReader(resp.Body, 64<<10)) + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return &HTTPError{StatusCode: resp.StatusCode, BodyTail: tail(raw, 256)} + } + return nil +} diff --git a/internal/hub/osupdate_contract_test.go b/internal/hub/osupdate_contract_test.go new file mode 100644 index 0000000..c5c94c7 --- /dev/null +++ b/internal/hub/osupdate_contract_test.go @@ -0,0 +1,30 @@ +package hub + +import ( + "encoding/json" + "os" + "testing" +) + +// The os_update block is a cross-repo contract: testdata/desired-state-osupdate.golden.json is byte-identical with +// felhom.eu/hub/internal/api/testdata (the hub's TestOSUpdate_DesiredBlockMatchesTheGolden proves the hub SERVES +// it). Here: the agent DECODES every field. A renamed json tag on either side fails one of the two tests. +func TestOSUpdateGolden_Decodes(t *testing.T) { + raw, err := os.ReadFile("testdata/desired-state-osupdate.golden.json") + if err != nil { + t.Fatal(err) + } + var resp DesiredStateResponse + if err := json.Unmarshal(raw, &resp); err != nil { + t.Fatal(err) + } + o := resp.DesiredState.OSUpdate + if o == nil || o.Ring != 1 || !o.Enabled || o.Release == nil { + t.Fatalf("os_update = %+v", o) + } + r := o.Release + if r.ID != "os-20261004-120000" || r.Snapshot != "20261004T120000Z" || len(r.Packages) != 2 || + r.Packages[1].Name != "openssl" || r.Packages[1].Version != "3.5.7-1~deb13u3" || r.Packages[1].Origin != "Debian-Security" { + t.Fatalf("release = %+v", r) + } +} diff --git a/internal/hub/report.go b/internal/hub/report.go index 017a924..04329e1 100644 --- a/internal/hub/report.go +++ b/internal/hub/report.go @@ -548,6 +548,32 @@ type WireDesiredState struct { RestoreDirective *WireRestoreDirective `json:"restore_directive,omitempty"` // slice 10D (forward-compat) Wireguard *WireWireguard `json:"wireguard,omitempty"` // S3 (doc 06 §3.2; golden-pinned) PBSDR *WirePBSDR `json:"pbs_dr,omitempty"` // PBS DR tier (slice 2 consumer) + OSUpdate *WireOSUpdate `json:"os_update,omitempty"` // OS updates, guest fast lane (agent v0.140.0) +} + +// WireOSUpdate is the hub-OWNED OS-update block (hub v0.130.0, `11-os-updates.md` §5.3), merged into the served +// document at read time. Ring 0 installs every pending Debian / Debian-Security fix; ring 1 installs exactly the +// newest approved release. Absent (older hub) → the agent treats the box as ring 1, ON, no release: it reports +// and installs nothing. Golden: testdata/desired-state-osupdate.golden.json (byte-identical with the hub's). +type WireOSUpdate struct { + Ring int `json:"ring"` + Enabled bool `json:"enabled"` + Release *WireOSRelease `json:"release,omitempty"` +} + +// WireOSRelease is an approved version set; Snapshot is the approval time (YYYYMMDDTHHMMSSZ) the wrapper uses +// for snapshot.debian.org when Debian has already replaced a version (decision 79). +type WireOSRelease struct { + ID string `json:"id"` + Snapshot string `json:"snapshot"` + Packages []WireOSPackage `json:"packages"` +} + +// WireOSPackage is one approved name=version and its origin ("Debian" | "Debian-Security"). +type WireOSPackage struct { + Name string `json:"name"` + Version string `json:"version"` + Origin string `json:"origin"` } // WirePBSDR is the hub's PBS-DR-tier descriptor (PBS DR slice 1, hub/internal/web/pbsdr.go diff --git a/internal/hub/testdata/desired-state-osupdate.golden.json b/internal/hub/testdata/desired-state-osupdate.golden.json new file mode 100644 index 0000000..da28c48 --- /dev/null +++ b/internal/hub/testdata/desired-state-osupdate.golden.json @@ -0,0 +1,17 @@ +{ + "generation": 1, + "desired_state": { + "os_update": { + "ring": 1, + "enabled": true, + "release": { + "id": "os-20261004-120000", + "snapshot": "20261004T120000Z", + "packages": [ + {"name": "libc6", "version": "2.41-12+deb13u4", "origin": "Debian"}, + {"name": "openssl", "version": "3.5.7-1~deb13u3", "origin": "Debian-Security"} + ] + } + } + } +} diff --git a/internal/localapi/afterbackup_test.go b/internal/localapi/afterbackup_test.go new file mode 100644 index 0000000..4c458c4 --- /dev/null +++ b/internal/localapi/afterbackup_test.go @@ -0,0 +1,63 @@ +package localapi + +import ( + "context" + "net/http" + "sync" + "testing" + "time" + + "gitea.dooplex.hu/admin/felhom-agent/internal/backup" +) + +// The OS leg (agent v0.140.0) runs after a SUCCESSFUL primary backup, and only then; and it runs BEFORE the +// host-wide heavy-op gate is released, so a restore-test cannot start in the middle of it (`11` C10). +// Red-proof: drop the `b.Success &&` guard and the failed-backup sub-case fails; move the call after release() +// and the gate sub-case fails. +func TestAfterPrimaryBackup(t *testing.T) { + run := func(t *testing.T, failErr string) (calls []int, gateHeld bool) { + gate := &backup.InFlight{} + b := &fakeBackups{failErr: failErr} + srv := newTestServerS(t, &fakeGuests{}, b, &fakeStore{}, nil) + srv.inFlight = gate + var mu sync.Mutex + done := make(chan struct{}, 1) + srv.SetAfterPrimaryBackup(func(_ context.Context, vmid int) { + rel, _, ok := gate.TryAcquire("probe") + mu.Lock() + calls = append(calls, vmid) + gateHeld = !ok + mu.Unlock() + if ok { + rel() + } + done <- struct{}{} + }) + h := srv.Handler() + if do(t, h, "POST", "/backup", "A", "").Code != http.StatusAccepted { + t.Fatal("POST /backup not accepted") + } + select { + case <-done: + case <-time.After(500 * time.Millisecond): + } + time.Sleep(20 * time.Millisecond) + mu.Lock() + defer mu.Unlock() + return calls, gateHeld + } + t.Run("success runs the leg under the gate", func(t *testing.T) { + calls, held := run(t, "") + if len(calls) != 1 { + t.Fatalf("the leg ran %d time(s), want 1", len(calls)) + } + if !held { + t.Fatal("the heavy-op gate was free while the leg ran — a restore-test could overlap it") + } + }) + t.Run("a failed backup runs nothing", func(t *testing.T) { + if calls, _ := run(t, "vzdump exploded"); len(calls) != 0 { + t.Fatalf("the leg ran after a FAILED backup: %v", calls) + } + }) +} diff --git a/internal/localapi/server.go b/internal/localapi/server.go index 6e81f95..456dbc8 100644 --- a/internal/localapi/server.go +++ b/internal/localapi/server.go @@ -160,6 +160,10 @@ type Options struct { // NetStorage is the privileged network-mount (NAS) surface (Part A1). OPTIONAL — when nil, the // /netstorage endpoints report "not configured". Satisfied by *storage.SudoHostOps. NetStorage NetworkStorageOps + // AfterPrimaryBackup (agent v0.140.0, `11-os-updates.md` §8 step 2) runs right after a SUCCESSFUL backup on the + // PRIMARY tier, inside the backup goroutine and BEFORE the host-wide heavy-op gate is released — so the OS leg + // that it starts can never overlap another backup or a restore-test (`11` C10). OPTIONAL — nil → nothing runs. + AfterPrimaryBackup func(ctx context.Context, vmid int) // Privileged runs the fenced root wrappers (E-2a: felhom-backup-target-apply). OPTIONAL — when // nil, POST /backup/target reports "not configured". Satisfied by *proxmox.ExecRunner. Privileged PrivilegedRunner @@ -276,6 +280,7 @@ type Server struct { tiers []BackupTier // inFlight (R-85) is shared with the restore-test scheduler so the two never run together. inFlight *backup.InFlight + afterPrimaryBackup func(ctx context.Context, vmid int) // the OS leg (agent v0.140.0); nil = none logger *slog.Logger now func() time.Time @@ -469,6 +474,7 @@ func NewServer(o Options) (*Server, error) { // the primary is always first, because that is what the untargeted endpoints act on. s.tiers = normalizeBackupTiers(o.BackupTiers, o.Backups, cadence) s.inFlight = o.InFlight + s.afterPrimaryBackup = o.AfterPrimaryBackup if s.backups == nil && len(s.tiers) > 0 { s.backups = s.tiers[0].Service } @@ -894,6 +900,10 @@ func (s *Server) handleBackup(w http.ResponseWriter, r *http.Request, vmid int) } s.store.RecordBackup(b) s.finishJob(key, jobID, b) + // OS leg (agent v0.140.0): after the night's whole-guest copy exists, still holding the heavy-op gate. + if b.Success && tier.Primary && s.afterPrimaryBackup != nil { + s.afterPrimaryBackup(base, vmid) + } }() writeStatus(w, http.StatusAccepted, true, BackupResponse{VMID: vmid, JobID: jobID, Phase: PhaseRunning}, "") } @@ -1465,3 +1475,6 @@ func writeStatus(w http.ResponseWriter, code int, ok bool, data any, errMsg stri w.WriteHeader(code) _ = json.NewEncoder(w).Encode(apiResponse{OK: ok, Data: data, Error: errMsg}) } + +// SetAfterPrimaryBackup wires the hook that runs after a successful primary-tier backup (the OS leg, agent v0.140.0). +func (s *Server) SetAfterPrimaryBackup(fn func(ctx context.Context, vmid int)) { s.afterPrimaryBackup = fn } diff --git a/internal/osupdate/leg.go b/internal/osupdate/leg.go new file mode 100644 index 0000000..90c483a --- /dev/null +++ b/internal/osupdate/leg.go @@ -0,0 +1,429 @@ +// Package osupdate is the agent's OS-update leg for the customer GUEST (`11-os-updates.md` §8 step 2, agent v0.140.0). +// +// It runs right after the night's successful whole-guest backup, while the backup goroutine still holds the host-wide +// heavy-op gate (so it never overlaps a backup or a restore-test, `11` C10), at most once per night. All root work is +// the wrapper `felhom-os-apply` (configs/, its own tests); this package only builds plans, calls the wrapper through +// sudo, judges health and reports to the hub. +// +// NO AUTOMATIC UNDO in this release (R-837, measured 2026-10-04): a customer guest cannot be snapshotted — PVE +// refuses any snapshot not named `vzdump` when the guest has host-path binds (mp8, mp9). A failed health check +// therefore stops, reports `health_failed` and the hub mails the operator; the whole-guest backup taken minutes +// earlier is the undo, by hand. +package osupdate + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "os" + "path/filepath" + "sort" + "strings" + "sync" + "time" + + "gitea.dooplex.hu/admin/felhom-agent/internal/hub" + "gitea.dooplex.hu/admin/felhom-agent/internal/proxmox" +) + +// WrapperPath is the pinned sudoers vector (configs/felhom-agent.sudoers FELHOM_OSAPPLY). +const WrapperPath = "/usr/local/sbin/felhom-os-apply" + +// DefaultPlanDir is where plans are written (the sudoers glob names it). +const DefaultPlanDir = "/var/lib/felhom-agent/os" + +// Package is one name=version with its origin. +type Package struct { + Name string `json:"name"` + Version string `json:"version"` + Origin string `json:"origin"` +} + +// Pending is one update the guest's sources offer (origin as apt names it, possibly several). +type Pending struct { + Name string `json:"name"` + From string `json:"from"` + To string `json:"to"` + Origin []string `json:"origin"` +} + +// Container is one container's state as the wrapper saw it. +type Container struct { + State string `json:"state"` + Health string `json:"health"` // healthy | unhealthy | starting | none +} + +// Health is the guest's health snapshot (the wrapper's `health` object). +type Health struct { + DockerOK bool `json:"docker_ok"` + NetworkOK bool `json:"network_ok"` + Controller string `json:"controller"` + Containers map[string]Container `json:"containers"` +} + +// WrapperReport is the wrapper's OSAPPLY-REPORT object. +type WrapperReport struct { + Mode string `json:"mode"` + Refused json.RawMessage `json:"refused"` + Failed json.RawMessage `json:"failed"` + Upgraded []Package `json:"upgraded"` + Installed []Package `json:"installed"` + Pending []Pending `json:"pending"` + RestartNeeded []string `json:"restart_needed"` + DockerRestartNeeded bool `json:"docker_restart_needed"` + RebootNeeded bool `json:"reboot_needed"` + HealthBefore *Health `json:"health_before"` + HealthAfter *Health `json:"health_after"` + Health *Health `json:"health"` +} + +func (w WrapperReport) refused() bool { return len(w.Refused) > 0 && string(w.Refused) != "null" } +func (w WrapperReport) failed() bool { return len(w.Failed) > 0 && string(w.Failed) != "null" } + +// Report is what the hub receives (hub osupdates.Report — field-exact). +type Report struct { + RunID string `json:"run_id"` + Trigger string `json:"trigger"` + Mode string `json:"mode"` + Ring int `json:"ring"` + ReleaseID string `json:"release_id"` + Outcome string `json:"outcome"` + Healthy bool `json:"healthy"` + HealthReason string `json:"health_reason,omitempty"` + VMID int `json:"vmid"` + Upgraded []Package `json:"upgraded,omitempty"` + Installed []Package `json:"installed,omitempty"` + Pending []Pending `json:"pending,omitempty"` + NotCovered []string `json:"not_covered,omitempty"` + RestartNeeded []string `json:"restart_needed,omitempty"` + DockerRestartNeeded bool `json:"docker_restart_needed,omitempty"` + RebootNeeded bool `json:"reboot_needed,omitempty"` + Refused json.RawMessage `json:"refused,omitempty"` +} + +// Reporter posts a report to the hub (*hub.Client). +type Reporter interface { + PostOSReport(ctx context.Context, body []byte) error +} + +// Leg runs one OS-update pass for the customer guest. +type Leg struct { + Runner proxmox.Runner + Hub Reporter + Logger *slog.Logger + PlanDir string + StatePath string // last night run (once per night) + HealthWait time.Duration // how long health may take to come back (default 5 min) + HealthPoll time.Duration // default 15 s + MinGap time.Duration // between night runs (default 20 h) + Now func() time.Time + Sleep func(context.Context, time.Duration) + + mu sync.Mutex + block *hub.WireOSUpdate + have bool +} + +// OnDesiredState stores the hub's os_update block (desired.RawConsumer — store only, never block). +func (l *Leg) OnDesiredState(_ context.Context, resp *hub.DesiredStateResponse) { + if resp == nil { + return + } + l.mu.Lock() + defer l.mu.Unlock() + l.block, l.have = resp.DesiredState.OSUpdate, true +} + +// Block returns the newest os_update block. No block (an older hub, or nothing fetched yet) = ring 1, ON, no +// release: the box reports and installs nothing. +func (l *Leg) Block() hub.WireOSUpdate { + l.mu.Lock() + defer l.mu.Unlock() + if l.block == nil { + return hub.WireOSUpdate{Ring: 1, Enabled: true} + } + return *l.block +} + +// SetBlock sets the block directly (the selftest fetches the desired state itself). +func (l *Leg) SetBlock(b *hub.WireOSUpdate) { + l.mu.Lock() + defer l.mu.Unlock() + l.block, l.have = b, true +} + +func (l *Leg) now() time.Time { + if l.Now != nil { + return l.Now() + } + return time.Now() +} + +func (l *Leg) log() *slog.Logger { + if l.Logger != nil { + return l.Logger + } + return slog.Default() +} + +func (l *Leg) sleep(ctx context.Context, d time.Duration) { + if l.Sleep != nil { + l.Sleep(ctx, d) + return + } + select { + case <-ctx.Done(): + case <-time.After(d): + } +} + +// IsFast reports whether every origin apt names is Debian / Debian-Security (the fast lane, `11` C3). +func IsFast(origins []string) bool { + if len(origins) == 0 { + return false + } + for _, o := range origins { + if o != "Debian" && o != "Debian-Security" { + return false + } + } + return true +} + +func fastOrigin(origins []string) string { + for _, o := range origins { + if o == "Debian-Security" { + return o + } + } + return "Debian" +} + +// HealthVerdict is THE health rule (`11` §5.4.1; pinned by TestHealthVerdict_*): after the run, docker answers, the +// guest's network resolves, the controller's own health check is `healthy`, and every container that was running +// before is running again — and healthy again if it was healthy before. "starting" is not yet healthy. +func HealthVerdict(before, after *Health) (bool, string) { + if after == nil { + return false, "no health reading" + } + if !after.DockerOK { + return false, "docker does not answer" + } + if !after.NetworkOK { + return false, "the guest cannot resolve deb.debian.org" + } + if after.Controller != "healthy" { + return false, "the controller is " + after.Controller + } + if before == nil { + return true, "" + } + names := make([]string, 0, len(before.Containers)) + for n := range before.Containers { + names = append(names, n) + } + sort.Strings(names) + for _, n := range names { + b := before.Containers[n] + if b.State != "running" { + continue + } + a, ok := after.Containers[n] + if !ok || a.State != "running" { + return false, n + " was running and is not" + } + if b.Health == "healthy" && a.Health != "healthy" { + return false, n + " was healthy and is " + a.Health + } + } + return true, "" +} + +// call writes the plan and runs the wrapper once. +func (l *Leg) call(ctx context.Context, runID, mode string, vmid int, rel hub.WireOSRelease, pkgs []Package) (WrapperReport, error) { + dir := l.PlanDir + if dir == "" { + dir = DefaultPlanDir + } + if err := os.MkdirAll(dir, 0o700); err != nil { + return WrapperReport{}, fmt.Errorf("osupdate: plan dir: %w", err) + } + plan := map[string]any{"release_id": rel.ID, "layer": "guest", "lane": "fast", "vmid": vmid, "mode": mode, + "snapshot": rel.Snapshot, "packages": pkgs} + if pkgs == nil { + plan["packages"] = []Package{} + } + b, _ := json.Marshal(plan) + path := filepath.Join(dir, "plan-"+runID+"-"+mode+".json") + if err := os.WriteFile(path, b, 0o600); err != nil { + return WrapperReport{}, fmt.Errorf("osupdate: write plan: %w", err) + } + defer os.Remove(path) + stdout, stderr, err := l.Runner.Run(ctx, WrapperPath, "--plan", path) + for _, line := range strings.Split(strings.TrimSpace(string(stderr)), "\n") { + if strings.HasPrefix(line, "os-apply: ") { + l.log().Info("osupdate: wrapper", "line", line) + } + } + var rep WrapperReport + found := false + for _, line := range strings.Split(string(stdout), "\n") { + if strings.HasPrefix(line, "OSAPPLY-REPORT ") { + if jerr := json.Unmarshal([]byte(strings.TrimPrefix(line, "OSAPPLY-REPORT ")), &rep); jerr == nil { + found = true + } + } + } + if !found { + return rep, fmt.Errorf("osupdate: wrapper gave no report (err %v): %s", err, strings.TrimSpace(string(stderr))) + } + return rep, nil // a refusal / failure is IN the report (exit 2 / 3), not an error here +} + +// Run is one pass: inventory → (switch, ring, plan) → apply → health → report. trigger is "night" or "debug". +func (l *Leg) Run(ctx context.Context, vmid int, trigger string) Report { + runID := l.now().UTC().Format("20060102T150405Z") + blk := l.Block() + rel := hub.WireOSRelease{ID: "ring0-" + runID} + if blk.Ring == 1 { + rel = hub.WireOSRelease{} + if blk.Release != nil { + rel = *blk.Release + } + } + rep := Report{RunID: runID, Trigger: trigger, Ring: blk.Ring, VMID: vmid, ReleaseID: rel.ID, Mode: "inventory"} + lg := l.log().With("run", runID, "vmid", vmid, "ring", blk.Ring, "trigger", trigger) + + if trigger == "night" && l.StatePath != "" { + gap := l.MinGap + if gap == 0 { + gap = 20 * time.Hour + } + if b, err := os.ReadFile(l.StatePath); err == nil { + if last, perr := time.Parse(time.RFC3339, strings.TrimSpace(string(b))); perr == nil && l.now().Sub(last) < gap { + lg.Info("osupdate: skipped — already ran tonight", "last", last.UTC().Format(time.RFC3339)) + rep.Outcome = "skipped" + return rep + } + } + } + lg.Info("osupdate: START", "enabled", blk.Enabled, "release", rel.ID) + + inv, err := l.call(ctx, runID, "inventory", vmid, rel, nil) + if err != nil || inv.refused() || inv.failed() { + rep.Outcome, rep.Refused = "refused", inv.Refused + if err != nil { + rep.Outcome, rep.HealthReason = "failed", err.Error() + } + return l.finish(ctx, lg, rep) + } + var plan []Package + switch { + case !blk.Enabled: + rep.Outcome = "inventory" + lg.Info("osupdate: switched OFF for this box — reporting only") + case blk.Ring == 0: + for _, p := range inv.Pending { + if IsFast(p.Origin) { + plan = append(plan, Package{Name: p.Name, Version: p.To, Origin: fastOrigin(p.Origin)}) + } + } + default: + for _, p := range rel.Packages { + plan = append(plan, Package{Name: p.Name, Version: p.Version, Origin: p.Origin}) + } + } + pkgNames := map[string]bool{} + for _, p := range plan { + pkgNames[p.Name] = true + } + if blk.Enabled && len(plan) == 0 { + rep.Outcome = "nothing" + } + final := inv + if blk.Enabled && len(plan) > 0 { + rep.Mode = "apply" + ap, err := l.call(ctx, runID, "apply", vmid, rel, plan) + switch { + case err != nil: + rep.Outcome, rep.HealthReason = "failed", err.Error() + return l.finish(ctx, lg, rep) + case ap.refused(): + rep.Outcome, rep.Refused = "refused", ap.Refused + return l.finish(ctx, lg, rep) + case ap.failed(): + rep.Outcome, rep.Refused = "failed", ap.Failed + } + final = ap + rep.Upgraded = ap.Upgraded + if rep.Outcome == "" { + if len(ap.Upgraded) == 0 { + rep.Outcome = "nothing" + } else { + rep.Outcome = "applied" + } + } + // Health: compare with what the guest looked like BEFORE the run; give restarted services time. + wait, poll := l.HealthWait, l.HealthPoll + if wait == 0 { + wait = 5 * time.Minute + } + if poll == 0 { + poll = 15 * time.Second + } + deadline := l.now().Add(wait) + cur := ap.HealthAfter + for { + ok, why := HealthVerdict(ap.HealthBefore, cur) + rep.Healthy, rep.HealthReason = ok, why + if ok || !l.now().Before(deadline) || ctx.Err() != nil { + break + } + l.sleep(ctx, poll) + hr, herr := l.call(ctx, runID, "health", vmid, rel, nil) + if herr == nil && hr.Health != nil { + cur = hr.Health + } + } + if !rep.Healthy && rep.Outcome == "applied" { + rep.Outcome = "health_failed" + } + } else { + ok, why := HealthVerdict(nil, inv.HealthAfter) + rep.Healthy, rep.HealthReason = ok, why + } + rep.Installed, rep.Pending = final.Installed, final.Pending + rep.RestartNeeded, rep.DockerRestartNeeded, rep.RebootNeeded = final.RestartNeeded, final.DockerRestartNeeded, final.RebootNeeded + rep.NotCovered = notCovered(final.Pending, blk.Ring, pkgNames) + if trigger == "night" && l.StatePath != "" { + _ = os.WriteFile(l.StatePath, []byte(l.now().UTC().Format(time.RFC3339)), 0o600) + } + return l.finish(ctx, lg, rep) +} + +// notCovered lists pending updates no approved release covers: in ring 0 everything outside the fast lane; in +// ring 1 also every fast-lane update the release did not name. +func notCovered(pending []Pending, ring int, planned map[string]bool) []string { + var out []string + for _, p := range pending { + if !IsFast(p.Origin) || (ring == 1 && !planned[p.Name]) { + out = append(out, p.Name) + } + } + return out +} + +func (l *Leg) finish(ctx context.Context, lg *slog.Logger, rep Report) Report { + lg.Info("osupdate: DONE", "outcome", rep.Outcome, "healthy", rep.Healthy, "reason", rep.HealthReason, + "upgraded", len(rep.Upgraded), "pending", len(rep.Pending), "not_covered", len(rep.NotCovered), "restart_needed", len(rep.RestartNeeded)) + if l.Hub != nil { + body, _ := json.Marshal(rep) + rctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), time.Minute) + defer cancel() + if err := l.Hub.PostOSReport(rctx, body); err != nil { + lg.Warn("osupdate: reporting to the hub failed (the run itself is done)", "err", err) + } + } + return rep +} diff --git a/internal/osupdate/leg_test.go b/internal/osupdate/leg_test.go new file mode 100644 index 0000000..8cf1922 --- /dev/null +++ b/internal/osupdate/leg_test.go @@ -0,0 +1,265 @@ +package osupdate + +import ( + "context" + "encoding/json" + "io" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + "gitea.dooplex.hu/admin/felhom-agent/internal/hub" +) + +// fakeWrapper plays /usr/local/sbin/felhom-os-apply: it reads the plan the leg wrote and answers per mode. +type fakeWrapper struct { + t *testing.T + pending []Pending + applyRep WrapperReport + healthSeq []*Health // answers to successive "health" calls + plans []map[string]any +} + +func (f *fakeWrapper) Run(_ context.Context, name string, args ...string) ([]byte, []byte, error) { + if name != WrapperPath || len(args) != 2 || args[0] != "--plan" { + f.t.Fatalf("unexpected command %s %v", name, args) + } + b, err := os.ReadFile(args[1]) + if err != nil { + f.t.Fatal(err) + } + var plan map[string]any + json.Unmarshal(b, &plan) + f.plans = append(f.plans, plan) + var rep WrapperReport + healthy := &Health{DockerOK: true, NetworkOK: true, Controller: "healthy", Containers: map[string]Container{ + "felhom-controller": {State: "running", Health: "healthy"}, "app": {State: "running", Health: "healthy"}}} + switch plan["mode"] { + case "inventory": + rep = WrapperReport{Mode: "inventory", Pending: f.pending, HealthAfter: healthy, + Installed: []Package{{Name: "libc6", Version: "u3", Origin: "Debian"}}} + case "apply": + rep = f.applyRep + rep.Mode = "apply" + if rep.HealthBefore == nil { + rep.HealthBefore = healthy + } + if rep.HealthAfter == nil { + rep.HealthAfter = healthy + } + case "health": + if len(f.healthSeq) > 0 { + rep.Health, f.healthSeq = f.healthSeq[0], f.healthSeq[1:] + } else { + rep.Health = healthy + } + } + out, _ := json.Marshal(rep) + return []byte("OSAPPLY-REPORT " + string(out) + "\n"), []byte("os-apply: DONE rc=0\n"), nil +} + +func (f *fakeWrapper) RunStdin(ctx context.Context, _ io.Reader, name string, args ...string) ([]byte, []byte, error) { + return f.Run(ctx, name, args...) +} + +type fakeHub struct{ reports []Report } + +func (h *fakeHub) PostOSReport(_ context.Context, body []byte) error { + var r Report + json.Unmarshal(body, &r) + h.reports = append(h.reports, r) + return nil +} + +func newLeg(t *testing.T, w *fakeWrapper, blk *hub.WireOSUpdate) (*Leg, *fakeHub) { + h := &fakeHub{} + now := time.Date(2026, 10, 4, 4, 0, 0, 0, time.UTC) + l := &Leg{Runner: w, Hub: h, PlanDir: t.TempDir(), StatePath: filepath.Join(t.TempDir(), "last"), + HealthWait: time.Minute, HealthPoll: 10 * time.Second, + Now: func() time.Time { return now }, + Sleep: func(_ context.Context, d time.Duration) { now = now.Add(d) }} + if blk != nil { + l.SetBlock(blk) + } + return l, h +} + +var pend = []Pending{ + {Name: "libc6", From: "u3", To: "u4", Origin: []string{"Debian"}}, + {Name: "openssl", From: "u1", To: "u3", Origin: []string{"Debian-Security", "Debian"}}, + {Name: "docker-ce", From: "29.7", To: "29.8", Origin: []string{"Docker CE"}}, +} + +func modes(w *fakeWrapper) string { + var m []string + for _, p := range w.plans { + m = append(m, p["mode"].(string)) + } + return strings.Join(m, ",") +} + +// Ring 0 plans every pending FAST-LANE update (Docker excluded, `11` C3) and reports what it now runs. +func TestRing0_PlansTheFastLaneOnly(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: WrapperReport{Upgraded: []Package{{Name: "libc6", Version: "u4"}, {Name: "openssl", Version: "u3"}}, Pending: pend[2:]}} + l, h := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + rep := l.Run(context.Background(), 9201, "night") + if rep.Outcome != "applied" || !rep.Healthy { + t.Fatalf("rep = %+v", rep) + } + pk := w.plans[1]["packages"].([]any) + if len(pk) != 2 { + t.Fatalf("ring-0 plan = %v, want libc6 + openssl only", pk) + } + if o := pk[1].(map[string]any)["origin"]; o != "Debian-Security" { + t.Fatalf("openssl origin = %v", o) + } + if w.plans[1]["snapshot"] != "" { + t.Fatalf("ring 0 installs from live sources, snapshot = %v", w.plans[1]["snapshot"]) + } + if len(rep.NotCovered) != 1 || rep.NotCovered[0] != "docker-ce" { + t.Fatalf("not covered = %v", rep.NotCovered) + } + if len(h.reports) != 1 || h.reports[0].Outcome != "applied" { + t.Fatalf("hub got %+v", h.reports) + } +} + +// Ring 1 installs EXACTLY the approved release (its versions, its snapshot), nothing it computed itself. +func TestRing1_InstallsExactlyTheRelease(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: WrapperReport{Upgraded: []Package{{Name: "libc6", Version: "u4-approved"}}, Pending: pend[1:]}} + rel := &hub.WireOSRelease{ID: "os-1", Snapshot: "20261004T080000Z", Packages: []hub.WireOSPackage{{Name: "libc6", Version: "u4-approved", Origin: "Debian"}}} + l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 1, Enabled: true, Release: rel}) + rep := l.Run(context.Background(), 9201, "night") + if rep.Outcome != "applied" || rep.ReleaseID != "os-1" { + t.Fatalf("rep = %+v", rep) + } + ap := w.plans[1] + pk := ap["packages"].([]any) + if len(pk) != 1 || pk[0].(map[string]any)["version"] != "u4-approved" || ap["snapshot"] != "20261004T080000Z" || ap["release_id"] != "os-1" { + t.Fatalf("ring-1 plan = %v", ap) + } + // openssl is pending and fast-lane but NOT in the release → not covered; docker-ce is never covered. + if strings.Join(rep.NotCovered, ",") != "openssl,docker-ce" { + t.Fatalf("not covered = %v", rep.NotCovered) + } +} + +func TestRing1_NoReleaseInstallsNothing(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend} + l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 1, Enabled: true}) + rep := l.Run(context.Background(), 9201, "night") + if rep.Outcome != "nothing" || modes(w) != "inventory" { + t.Fatalf("rep=%+v modes=%s", rep, modes(w)) + } +} + +// No block from the hub (an older hub): ring 1, ON, no release → reports, installs nothing. +func TestNoBlock_IsRing1Nothing(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend} + l, _ := newLeg(t, w, nil) + if rep := l.Run(context.Background(), 9201, "night"); rep.Outcome != "nothing" || rep.Ring != 1 || modes(w) != "inventory" { + t.Fatalf("rep=%+v modes=%s", rep, modes(w)) + } +} + +// Switched OFF: the box reports but installs nothing (the brief, Part D 2). +func TestSwitchOff_ReportsOnly(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend} + l, h := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: false}) + rep := l.Run(context.Background(), 9201, "night") + if rep.Outcome != "inventory" || modes(w) != "inventory" || len(h.reports) != 1 || len(rep.Pending) != 3 { + t.Fatalf("rep=%+v modes=%s", rep, modes(w)) + } +} + +// Unhealthy after the run, and still unhealthy at the end of the wait → health_failed (the hub mails the operator). +func TestHealth_FailsAfterTheWait(t *testing.T) { + bad := &Health{DockerOK: true, NetworkOK: true, Controller: "healthy", Containers: map[string]Container{ + "felhom-controller": {State: "running", Health: "healthy"}, "app": {State: "exited"}}} + w := &fakeWrapper{t: t, pending: pend, applyRep: WrapperReport{Upgraded: []Package{{Name: "libc6"}}, HealthAfter: bad}, + healthSeq: []*Health{bad, bad, bad, bad, bad, bad, bad, bad}} + l, h := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + rep := l.Run(context.Background(), 9201, "night") + if rep.Outcome != "health_failed" || rep.Healthy || !strings.Contains(rep.HealthReason, "app was running") { + t.Fatalf("rep = %+v", rep) + } + if h.reports[0].Outcome != "health_failed" { + t.Fatal("the hub was not told") + } +} + +// A service that takes a moment to come back is not a failure: the poll sees it recover inside the wait. +func TestHealth_RecoversInsideTheWait(t *testing.T) { + starting := &Health{DockerOK: true, NetworkOK: true, Controller: "starting"} + w := &fakeWrapper{t: t, pending: pend, applyRep: WrapperReport{Upgraded: []Package{{Name: "libc6"}}, HealthAfter: starting}, + healthSeq: []*Health{starting}} + l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + if rep := l.Run(context.Background(), 9201, "night"); rep.Outcome != "applied" || !rep.Healthy { + t.Fatalf("rep = %+v", rep) + } +} + +func TestHealthVerdict(t *testing.T) { + ok := &Health{DockerOK: true, NetworkOK: true, Controller: "healthy", Containers: map[string]Container{"a": {State: "running", Health: "healthy"}}} + cases := []struct { + name string + before *Health + after *Health + want bool + }{ + {"all good", ok, ok, true}, + {"no reading", ok, nil, false}, + {"docker down", ok, &Health{NetworkOK: true, Controller: "healthy"}, false}, + {"no network", ok, &Health{DockerOK: true, Controller: "healthy"}, false}, + {"controller starting", ok, &Health{DockerOK: true, NetworkOK: true, Controller: "starting"}, false}, + {"only the controller differs", ok, &Health{DockerOK: true, NetworkOK: true, Controller: "unhealthy", Containers: map[string]Container{"a": {State: "running", Health: "healthy"}}}, false}, + {"app gone", ok, &Health{DockerOK: true, NetworkOK: true, Controller: "healthy", Containers: map[string]Container{}}, false}, + {"app unhealthy", ok, &Health{DockerOK: true, NetworkOK: true, Controller: "healthy", Containers: map[string]Container{"a": {State: "running", Health: "unhealthy"}}}, false}, + {"stopped before stays stopped", &Health{Containers: map[string]Container{"x": {State: "exited"}}}, &Health{DockerOK: true, NetworkOK: true, Controller: "healthy"}, true}, + } + for _, c := range cases { + if got, why := HealthVerdict(c.before, c.after); got != c.want { + t.Errorf("%s: got %v (%s), want %v", c.name, got, why, c.want) + } + } +} + +func TestOncePerNight(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: WrapperReport{Upgraded: []Package{{Name: "libc6"}}}} + l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + l.Run(context.Background(), 9201, "night") + n := len(w.plans) + if rep := l.Run(context.Background(), 9201, "night"); rep.Outcome != "skipped" || len(w.plans) != n { + t.Fatalf("a second night run in the same night ran: %+v", rep) + } + if rep := l.Run(context.Background(), 9201, "debug"); rep.Outcome == "skipped" { + t.Fatal("the debug action must not be throttled") + } +} + +func TestRefusedIsReported(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: WrapperReport{Refused: json.RawMessage(`{"code":"R6","reason":"x"}`)}} + l, h := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + if rep := l.Run(context.Background(), 9201, "night"); rep.Outcome != "refused" || h.reports[0].Outcome != "refused" { + t.Fatalf("rep = %+v", rep) + } +} + +// The root wrapper's own suite (configs/test_felhom_os_apply.py) runs with `go test ./...` so CI covers it. +func TestWrapperSuite(t *testing.T) { + py, err := exec.LookPath("python3") + if err != nil { + t.Skip("python3 not available") + } + cmd := exec.Command(py, "../../configs/test_felhom_os_apply.py") + out, err := cmd.CombinedOutput() + if err != nil { + t.Fatalf("wrapper suite failed: %v\n%s", err, out) + } + if !strings.Contains(string(out), "OK") { + t.Fatalf("wrapper suite did not report OK:\n%s", out) + } +}