diff --git a/cmd/felhom-agent/main.go b/cmd/felhom-agent/main.go index 29dca44..2d6f087 100644 --- a/cmd/felhom-agent/main.go +++ b/cmd/felhom-agent/main.go @@ -785,7 +785,7 @@ func runDaemon(cfg config.Config, logger *slog.Logger, logRing *applog.Ring) int pbsStore := pbs.NewSnapshotStore() pbsTargets := pbsTargetsFromPVE(cfg, px, logger) pbsReporter := pbs.NewLiveSnapshotReporter(pbsTargets, pbsStore, pbs.DefaultLiveSnapshotTimeout, logger) - collector := hub.NewCollector(px, hub.SystemctlProber{}, observer, backupStore, backupStore, pbsReporter, cfg.Hub.HostID, version, logger) + collector := hub.NewCollector(px, newTunnelProber(cfg, px), observer, backupStore, backupStore, pbsReporter, cfg.Hub.HostID, version, logger) collector.SetBackupTargetResolver(primaryBackupTargetOf(cfg)) // R-109: the recipe names the live target // Privileged-capability self-check (v0.44.0): probe the sudoers grants the non-root agent // depends on. The probe runs `sudo -n -l` LITERALLY (a policy LIST, never executing the @@ -844,7 +844,7 @@ func runDaemon(cfg config.Config, logger *slog.Logger, logRing *applog.Ring) int 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) + osLeg := newOSLeg(cfg, client, px, 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. @@ -1127,7 +1127,7 @@ func runDaemon(cfg config.Config, logger *slog.Logger, logRing *applog.Ring) int return case <-time.After(90 * time.Second): } - osLeg.Run(ctx, vmid, "night") + _, _ = osLeg.Run(ctx, vmid, "night") }) } if localTokens != nil { @@ -2026,7 +2026,7 @@ func runSelftestHub(ctx context.Context, cfg config.Config, logger *slog.Logger) // pbs coord. The live reporter lists snapshots directly (fresh store, last-known-good fallback) so // the selftest reflects exactly what a freshly-restarted daemon's first collect emits. pbsReporter := pbs.NewLiveSnapshotReporter(pbsTargetsFromPVE(cfg, px, logger), pbs.NewSnapshotStore(), pbs.DefaultLiveSnapshotTimeout, logger) - collector := hub.NewCollector(px, hub.SystemctlProber{}, observer, nil, nil, pbsReporter, cfg.Hub.HostID, version, logger) + collector := hub.NewCollector(px, newTunnelProber(cfg, px), observer, nil, nil, pbsReporter, cfg.Hub.HostID, version, logger) // R-109: wire the backup-target resolver here TOO. Without it selftest=hub would print a recipe whose // backup_target reads unknown/agent_backup_config_unavailable while the daemon's is resolved — and // this one-shot exists precisely so "the report it would send" can be trusted to match. @@ -3508,7 +3508,7 @@ func (f *selftestFlag) Set(v string) error { // 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 { +func newOSLeg(cfg config.Config, client *hub.Client, px *proxmox.Client, logger *slog.Logger) *osupdate.Leg { mode := proxmox.RunnerMode(cfg.Privileged.Mode) if mode == "" { mode = proxmox.RunnerSudo @@ -3518,6 +3518,9 @@ func newOSLeg(cfg config.Config, client *hub.Client, logger *slog.Logger) *osupd Logger: logger, PlanDir: osupdate.DefaultPlanDir, StatePath: filepath.Join(osupdate.DefaultPlanDir, "last-night-run"), + // The host step (agent v0.141.0) runs only on an appliance install; the wrapper re-checks the ROOT-owned record. + Appliance: cfg.IsAppliance(), + Tunnel: newTunnelProber(cfg, px), } if client != nil { l.Hub = client @@ -3539,7 +3542,12 @@ func runSelftestOSUpdate(ctx context.Context, cfg config.Config, logger *slog.Lo fmt.Fprintln(os.Stderr, "selftest=os-update: hub client:", err) return 1 } - leg := newOSLeg(cfg, client, logger) + px, perr := newProxmoxClient(cfg) + if perr != nil { + fmt.Fprintln(os.Stderr, "selftest=os-update: proxmox client:", perr) + return 1 + } + leg := newOSLeg(cfg, client, px, logger) resp, err := client.FetchDesiredState(ctx) if err != nil { fmt.Fprintln(os.Stderr, "selftest=os-update: desired state:", err) @@ -3547,15 +3555,67 @@ func runSelftestOSUpdate(ctx context.Context, cfg config.Config, logger *slog.Lo } 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}) + fmt.Printf("=== felhom-agent %s selftest=os-update vmid=%d ring=%d enabled=%v guest-release=%v host-release=%v appliance=%v ===\n", + version, vmid, b.Ring, b.Enabled, b.Release != nil, b.HostRelease != nil, leg.Appliance) + start := time.Now() + g, h := leg.Run(ctx, vmid, "debug") + for _, rep := range []osupdate.Report{g, h} { + if rep.Layer == "" { + fmt.Println(" host step: skipped (see the log line above)") + continue + } + printJSON("os-update report ("+rep.Layer+")", 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, "reboot_needed": rep.RebootNeeded, "wrapper_seconds": rep.PassSeconds, "refused": rep.Refused}) + } + fmt.Printf(" pass took %s\n", time.Since(start).Round(100*time.Millisecond)) + rep := g + if h.Layer != "" && !(h.Outcome == "applied" || h.Outcome == "nothing" || h.Outcome == "inventory") { + rep = h + } switch rep.Outcome { case "applied", "nothing", "inventory", "skipped": return 0 } return 1 } + +// newTunnelProber reads the box's REAL tunnel (R-841, agent v0.141.0): the cloudflared container in each running +// customer guest — a guest that binds /mnt/felhom-drives, the same rule the OS wrapper's R10 uses — through the +// existing `pct exec [0-9]* -- docker inspect -f *` sudoers line. +func newTunnelProber(cfg config.Config, px *proxmox.Client) hub.CloudflaredProber { + mode := proxmox.RunnerMode(cfg.Privileged.Mode) + if mode == "" { + mode = proxmox.RunnerSudo + } + return hub.GuestTunnelProber{ + Runner: &proxmox.ExecRunner{Mode: mode, SudoPath: cfg.Privileged.SudoPath}, + Guests: func(ctx context.Context) ([]int, error) { + if px == nil { + return nil, fmt.Errorf("no proxmox client") + } + gs, err := px.ListLXC(ctx) + if err != nil { + return nil, err + } + var out []int + for _, g := range gs { + if g.Status != "running" { + continue + } + gc, err := px.GuestConfig(ctx, g.VMID) + if err != nil { + continue + } + for _, v := range gc.MountPoints() { + if src, _, _ := strings.Cut(v, ","); src == "/mnt/felhom-drives" { + out = append(out, g.VMID) + break + } + } + } + return out, nil + }, + } +} diff --git a/configs/felhom-os-apply b/configs/felhom-os-apply index ab51c46..a1a5d32 100755 --- a/configs/felhom-os-apply +++ b/configs/felhom-os-apply @@ -1,5 +1,5 @@ #!/usr/bin/python3 -# felhom-os-apply — the ROOT half of the agent's operating-system update leg (`11-os-updates.md` §5.4.1). +# felhom-os-apply — the ROOT half of the agent's operating-system update leg (`11-os-updates.md` §5.4.1, §8.1–8.2). # # 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 @@ -8,22 +8,28 @@ # # 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. +# package, no package outside the plan, only Debian origin in the fast lane, no kernel / boot package on the host, +# only the box's own customer guest, and the host layer only on a box whose ROOT-OWNED install record says +# "appliance" (a BYO host belongs to its owner, `11` §1). Package signatures stay Debian's: apt checks every Release +# file, 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. +# LAYERS (agent v0.141.0): "guest" (the customer LXC, entered with `pct exec`) and "host" (this Proxmox host, run +# directly). LANE: "fast" only — the slow lane (kernel, Proxmox, Docker) is REFUSED (R3, R14) until `11` §8 steps 5–6. # # 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. +# inventory `apt-get update`, then report what is installed (with origin), what is pending, and health. +# apply repair first, pick the packages (select "listed": the plan's name=version list; "pending-fast": every +# pending Debian / Debian-Security upgrade, for ring 0), check every refusal on an `apt-get -s` +# simulation of EXACTLY name=version, install, clean, scan for restart-needed, report 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. +# +# SPEED (R-845, agent v0.141.0). Every `pct exec` costs ~0.9 s (measured on demo-hp), and v0.140.0 made one per +# package for version comparisons — 272 packages ≈ 4 minutes. Versions are now compared with the HOST's dpkg (the same +# Debian algorithm), madison/policy run once per pass for all packages, the restart scan runs only after an install, +# and one wrapper call does the whole pass (no separate inventory call before an apply). import json import os import re @@ -46,6 +52,13 @@ 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 +# The installer's ROOT-OWNED record (felhom-host-install.sh `state_set mode`); the agent cannot write it. +INSTALL_STATE = "/var/lib/felhom-install/state.json" +# Kernel, boot and firmware packages are the SLOW lane on the host whatever their origin (`11` C3, §5.2): a host +# reboot is needed for them to take effect, and a bad one can stop the box from booting. +HOST_SLOW_RE = re.compile(r"^(linux-(image|headers|kbuild|modules|base)|proxmox-kernel|proxmox-default-kernel|pve-kernel|" + r"pve-firmware|firmware-|grub|shim|systemd-boot|intel-microcode|amd64-microcode|efibootmgr)") +HOST_SERVICES = ["pveproxy", "pvedaemon", "pvestatd", "pve-cluster", "felhom-agent"] class Refused(Exception): @@ -57,8 +70,8 @@ class Refused(Exception): 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) + def host(self, argv, timeout=600, stdin=None): + p = subprocess.run(argv, capture_output=True, text=True, timeout=timeout, input=stdin) return p.returncode, p.stdout, p.stderr def guest(self, vmid, argv, timeout=1800): @@ -75,6 +88,16 @@ class Runner: import pwd return pwd.getpwnam(AGENT_USER).pw_uid + def write_file(self, layer, vmid, path, body): + """Write a small text file in the target layer — never via a shell string.""" + if layer == "host": + with open(path, "w") as f: + f.write(body) + return + rc, _, _ = self.host(["/usr/sbin/pct", "exec", str(vmid), "--", "tee", path], 60, stdin=body) + if rc != 0: + raise Refused("R7", f"could not write {path} in the guest") + def log(self, line): print(line, file=sys.stderr, flush=True) try: @@ -117,8 +140,9 @@ class Apply: 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)") + layer = plan.get("layer") + if layer not in ("guest", "host"): + raise Refused("R12", f"layer {layer!r} is not guest or host") if plan.get("lane", "fast") != "fast": raise Refused("R3", "the slow lane is refused in this release") vmid = plan.get("vmid") @@ -129,9 +153,16 @@ class Apply: 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") + select = plan.get("select", "listed") + if select not in ("listed", "pending-fast"): + raise Refused("R11", f"unknown select {select!r}") 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") + if not isinstance(pk, list): + raise Refused("R11", "packages must be a list") + if mode == "apply" and select == "listed" and not pk: + raise Refused("R11", "packages must be a non-empty list in apply mode (select listed)") + if select == "pending-fast" and pk: + raise Refused("R11", "select pending-fast takes no package list") seen = set() for e in pk: if not isinstance(e, dict): @@ -146,10 +177,27 @@ class Apply: 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)") + if layer == "host" and HOST_SLOW_RE.match(n): + raise Refused("R14", f"{n} is a kernel / boot / firmware package — the host's slow lane") 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 + return mode, layer, vmid, select + + def check_appliance(self): + """R12: the host layer only on a box whose ROOT-OWNED install record says appliance (`11` §1: never BYO).""" + try: + st = self.r.stat(INSTALL_STATE) + except OSError: + raise Refused("R12", f"no install record ({INSTALL_STATE}) — this box cannot prove it is an appliance") + if st.st_uid != 0 or (st.st_mode & 0o022): + raise Refused("R12", f"{INSTALL_STATE} is not root-owned and root-only-writable — it proves nothing") + try: + mode = json.loads(self.r.read_file(INSTALL_STATE)).get("mode") + except (OSError, ValueError, AttributeError): + raise Refused("R12", f"{INSTALL_STATE} is unreadable — this box cannot prove it is an appliance") + if mode != "appliance": + raise Refused("R12", f"this box was installed as {mode!r}, not appliance — its host belongs to its owner") def check_guest(self, vmid): if vmid in RESERVED_VMIDS: @@ -169,12 +217,24 @@ class Apply: 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): + # ---------- target helpers ---------- + def x(self, argv, timeout=1800): + """Run in the TARGET layer: the guest via pct exec, or the host directly.""" + if self.layer == "host": + return self.r.host(argv, timeout) return self.r.guest(self.vmid, argv, timeout) + def g(self, argv, timeout=1800): + """Run in the customer GUEST whatever the layer (its health).""" + return self.r.guest(self.vmid, argv, timeout) + + def dpkg_cmp(self, a, op, b): + # The HOST's dpkg: the same Debian version algorithm, and no `pct exec` (0.9 s) per comparison (R-845). + rc, _, _ = self.r.host(["dpkg", "--compare-versions", a, op, b], 30) + return rc == 0 + def installed(self): - rc, out, _ = self.g(["dpkg-query", "-W", "-f", "${Package}\t${Version}\t${db:Status-Abbrev}\n"]) + rc, out, _ = self.x(["dpkg-query", "-W", "-f", "${Package}\t${Version}\t${db:Status-Abbrev}\n"]) res = {} for l in out.splitlines(): parts = l.split("\t") @@ -182,21 +242,20 @@ class Apply: 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() + def madison_all(self, names): + """name -> set of downloadable versions, ONE call for all names.""" + res = {n: set() for n in names} + if not names: + return res + rc, out, _ = self.x(["apt-cache", "madison"] + sorted(names)) 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 + if len(f) >= 3 and f[0] in res: + res[f[0]].add(f[1]) + return res def simulate(self, args): - rc, out, err = self.g(APT_ENV + ["apt-get", "-s", "-q"] + args) + rc, out, err = self.x(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) @@ -213,17 +272,17 @@ class Apply: 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", "/"]) + rc, out, _ = self.x(["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"]) + rc, out, _ = self.x(["fuser", "/var/lib/dpkg/lock-frontend", "/var/lib/dpkg/lock"]) return rc == 0 and out.strip() != "" - def health(self): + def guest_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 = {} @@ -238,21 +297,34 @@ class Apply: "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 health(self): + if self.layer == "guest": + return self.guest_health() + rc, out, _ = self.r.host(["systemctl", "is-active"] + HOST_SERVICES, 30) + states = out.split() + svc = {s: (states[i] if i < len(states) else "unknown") for i, s in enumerate(HOST_SERVICES)} + src, sout, _ = self.r.host(["/usr/sbin/pct", "status", str(self.vmid)], 30) + running = src == 0 and "running" in sout + return {"host_services": svc, "guest_running": running, "guest": self.guest_health() if running else None} - def inventory(self): - inst = self.installed() + def restart_needed(self): + """Processes still mapping deleted files, OUTSIDE containers (C11). Guest: outside docker; host: outside the + LXC guests (the host's /proc shows guest processes too).""" + skip = "docker" if self.layer == "guest" else "lxc" + script = ('for p in /proc/[0-9]*; do grep -q "(deleted)" $p/maps 2>/dev/null || continue; ' + 'grep -q "%s" $p/cgroup 2>/dev/null && continue; echo "${p#/proc/} $(cat $p/comm 2>/dev/null)"; done' % skip) + rc, out, _ = self.x(["sh", "-c", script], timeout=120) + lines = [l for l in out.splitlines() if " " in l] + procs = sorted({l.split(" ", 1)[1] for l in lines}) + pid1 = any(l.split(" ", 1)[0] == "1" for l in lines) + return procs, pid1 or "lxc-start" in procs + + def inventory(self, inst=None): + inst = inst if inst is not None else 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]) + if names: + rc, out, _ = self.x(["apt-cache", "policy"] + names) # ONE call (R-845) cur, star = None, False for l in out.splitlines(): if not l.startswith(" "): @@ -270,10 +342,12 @@ class Apply: 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 "proxmox" in src: + return "Proxmox" if "security" in src and "debian" in src: return "Debian-Security" if "docker.com" in src: @@ -282,6 +356,7 @@ class Apply: return "Debian" return "other" rc, pend, remv, _ = self.simulate(["dist-upgrade"]) + self._pending = pend 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"], @@ -291,51 +366,70 @@ class Apply: # ---------- 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.mode, self.layer, self.vmid, self.select = self.check_plan(plan) + self.report.update(mode=self.mode, layer=self.layer, release_id=plan.get("release_id"), vmid=self.vmid) + if self.layer == "host": + self.check_appliance() 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', []))}") + log(f"os-apply: START release={plan.get('release_id')} layer={self.layer}" + + (f":{self.vmid}" if self.layer == "guest" else "") + + f" lane=fast mode={self.mode} select={self.select} packages={len(plan.get('packages', []))}") if self.apt_lock_held(): - raise Refused("R9", "another apt/dpkg holds the lock in the guest") + raise Refused("R9", f"another apt/dpkg holds the lock on the {self.layer}") 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) + rc, out, err = self.x(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:]}") + raise Refused("R7", f"apt-get update failed on the {self.layer}: {(out + err).strip().splitlines()[-1:]}") + installed_after = None if self.mode == "apply": - rc = self.apply(plan) + rc, installed_after = 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.update(self.inventory(installed_after)) 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"]) + rc, before, _ = self.x(["dpkg", "--audit"]) configured = len([l for l in before.splitlines() if l.startswith(" ")]) - fixed = len(re.findall(r"^Setting up ", out, re.M)) + fixed = 0 + after = "" + if before.strip(): # nothing half-done → nothing to run (R-845: two calls saved on every clean pass) + self.x(APT_ENV + ["dpkg", "--configure", "-a", "--force-confold"]) + rc2, out, err = self.x(APT_ENV + ["apt-get", "-f", "install", "-y", "-q"] + DPKG_OPTS) + _, after, _ = self.x(["dpkg", "--audit"]) + 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 pending_fast(self): + """Ring 0 (select pending-fast): every pending upgrade of an INSTALLED package whose every origin is Debian / + Debian-Security — and, on the host, not a kernel / boot / firmware package.""" + rc, pend, remv, _ = self.simulate(["dist-upgrade"]) + out = [] + for p in pend: + o = self.origin_name(p["origin"]) + if p["from"] is None or not o or not o <= set(FAST_ORIGINS): + continue + if self.layer == "host" and HOST_SLOW_RE.match(p["name"]): + continue + out.append({"name": p["name"], "version": p["to"], "origin": "Debian-Security" if "Debian-Security" in o else "Debian"}) + return out + def apply(self, plan): log = self.r.log + packages = plan["packages"] if self.select == "listed" else self.pending_fast() inst = self.installed() upgrade, already, notinst = [], 0, 0 - for e in plan["packages"]: + for e in packages: n, v = e["name"], e["version"] if n not in inst: notinst += 1 @@ -345,13 +439,18 @@ class Apply: continue upgrade.append((n, v)) from_snap = 0 - missing = [(n, v) for n, v in upgrade if v not in self.madison(n)] + if upgrade: + avail = self.madison_all([n for n, _ in upgrade]) + missing = [(n, v) for n, v in upgrade if v not in avail[n]] + else: + missing = [] 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)] + avail = self.madison_all([n for n, _ in missing]) + still = [(n, v) for n, v in missing if v not in avail[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}") @@ -362,7 +461,7 @@ class Apply: if not upgrade: self.report["upgraded"] = [] log("os-apply: DONE rc=0 seconds=0 upgraded=0 (nothing to do)") - return 0 + return 0, inst 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: @@ -382,12 +481,14 @@ class Apply: 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") + if self.layer == "host" and HOST_SLOW_RE.match(p["name"]): + raise Refused("R14", f"{p['name']} is a kernel / boot / firmware package — the host's slow lane") 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) + rc, out, err = self.x(APT_ENV + ["apt-get", "-y", "-q"] + DPKG_OPTS + args) secs = time.time() - t0 # dpkg says "Installing new version of config file X" when X was NOT changed locally (the package's new # version is taken), and "Configuration file 'X'" + "Keeping old config file" when it was (--force-confold @@ -404,23 +505,27 @@ class Apply: log(f"os-apply: CONFFILE kept {conflict} (changed locally; the package's version is {conflict}.dpkg-dist)") self.report.setdefault("conffiles_kept", []).append(conflict) conflict = None - self.g(["apt-get", "clean"]) + self.x(["apt-get", "clean"]) if rc != 0: - _, aud, _ = self.g(["dpkg", "--audit"]) + _, aud, _ = self.x(["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 + return 3, None 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 + procs, reboot = self.restart_needed() # only after an install (R-845) + self.report["restart_needed"] = procs + self.report["docker_restart_needed"] = any(p in ("dockerd", "containerd") for p in procs) + self.report["reboot_needed"] = reboot + log(f"os-apply: DONE rc=0 seconds={secs:.1f} upgraded={len(upgrade)} restart-needed={','.join(procs) or '-'} reboot-needed={'yes' if reboot else 'no'}") + return 0, None 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) + rc, out, _ = self.x(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) @@ -429,33 +534,22 @@ class Apply: return total def add_snapshot_sources(self, snap): - rc, out, _ = self.g(["sh", "-c", ". /etc/os-release && echo $VERSION_CODENAME"]) + rc, out, _ = self.x(["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})") + raise Refused("R7", f"cannot read the {self.layer}'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.write_file(self.layer, 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) + rc, out, err = self.x(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 + self.x(["rm", "-f", SNAPSHOT_LIST]) + self.x(APT_ENV + ["apt-get", "-q", "update"], timeout=600) def main(argv, runner=None): @@ -465,6 +559,7 @@ def main(argv, runner=None): print("OSAPPLY-REPORT " + json.dumps({"refused": {"code": "R1", "reason": "usage"}})) return 2 a = Apply(r, argv[2]) + t0 = time.time() try: rc = a.run() except Refused as e: @@ -475,6 +570,7 @@ def main(argv, runner=None): r.log(f"os-apply: FAILED rc=124 step=timeout — {e.cmd}") a.report["failed"] = {"rc": 124, "timeout": str(e.cmd)[:200]} rc = 3 + a.report["pass_seconds"] = round(time.time() - t0, 1) print("OSAPPLY-REPORT " + json.dumps(a.report, sort_keys=True)) return rc diff --git a/configs/test_felhom_os_apply.py b/configs/test_felhom_os_apply.py index 3ed3d28..500514e 100644 --- a/configs/test_felhom_os_apply.py +++ b/configs/test_felhom_os_apply.py @@ -61,6 +61,9 @@ class Fake: self.written = {} self.status = "status: running" self.lock_held = False + self.services = {} + self.files[osapply.INSTALL_STATE] = json.dumps({"mode": "appliance"}) + self.stats[osapply.INSTALL_STATE] = St(mode=statmod.S_IFREG | 0o644, uid=0) # Runner interface def read_file(self, p): @@ -81,14 +84,19 @@ class Fake: def log(self, line): self.logs.append(line) - def host(self, argv, timeout=600): + def host(self, argv, timeout=600, stdin=None): self.calls.append(("host", argv)) - if argv[1] == "status": + if argv[0] == "/usr/sbin/pct" and argv[1] == "status": return 0, self.status + "\n", "" - return 1, "", "unexpected host call" + if argv[0] == "dpkg" and argv[1] == "--compare-versions": + return (0 if dpkg_cmp(argv[2], argv[3], argv[4]) else 1), "", "" + if argv[0] == "systemctl" and argv[1] == "is-active": + return 0, "\n".join(self.services.get(s, "active") for s in argv[2:]) + "\n", "" + return self.emulate(argv) - def guest_write(self, vmid, path, body): + def write_file(self, layer, vmid, path, body): self.written[path] = body + self.write_layer = layer if path == osapply.SNAPSHOT_LIST: self.snap_active = True @@ -100,6 +108,9 @@ class Fake: def guest(self, vmid, argv, timeout=1800): self.calls.append(("guest", vmid, argv)) + return self.emulate(argv) + + def emulate(self, 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": @@ -113,7 +124,8 @@ class Fake: 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])), "" + self.madison_calls = getattr(self, "madison_calls", 0) + 1 + return 0, "".join(f" {n} | {v} | http://deb.debian.org trixie/main amd64 Packages\n" for n in a[2:] for v in self.avail(n)), "" if cmd == "apt-cache" and a[1] == "policy": out = "" for n in a[2:]: @@ -146,6 +158,8 @@ class Fake: if cmd == "sh": if "os-release" in a[2]: return 0, "trixie\n", "" + if "(deleted)" in a[2]: + return 0, getattr(self, "restart_out", ""), "" return 0, "", "" if cmd == "rm": self.snap_active = False @@ -156,6 +170,9 @@ class Fake: if "--print-uris" in a: return 0, "'http://x/libc6.deb' libc6.deb 4000000 SHA256:x\n", "" if "dist-upgrade" in a: + if getattr(self, "pending_sim", None) is not None and not getattr(self, "_pending_used", False): + self._pending_used = True + return 0, "\n".join(self.pending_sim) + "\n", "" 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: @@ -209,7 +226,7 @@ class Happy(unittest.TestCase): self.assertEqual(rc, 0, rep) self.assertEqual(f.installed["libc6"], "2.41-12+deb13u3") self.assertIn("installed", rep) - self.assertIn("restart_needed", rep) + self.assertNotIn("restart_needed", rep, "the restart scan runs only after an install (R-845)") def test_health_mode(self): f = Fake() @@ -419,11 +436,44 @@ class Refusals(unittest.TestCase): f.plan["packages"][0]["name"] = "--purge" self.refused(f, "R11") - def test_R12_host_layer(self): + def test_R12_unknown_layer(self): + f = Fake() + f.plan["layer"] = "vm" + self.refused(f, "R12") + + def test_R12_host_on_a_byo_box(self): f = Fake() f.plan["layer"] = "host" + f.files[osapply.INSTALL_STATE] = json.dumps({"mode": "byo"}) self.refused(f, "R12") + def test_R12_host_without_an_install_record(self): + f = Fake() + f.plan["layer"] = "host" + del f.stats[osapply.INSTALL_STATE] + self.refused(f, "R12") + + def test_R12_host_record_not_root_owned(self): + # the agent can write agent.json's deployment_mode; only a ROOT-owned record proves anything + f = Fake() + f.plan["layer"] = "host" + f.stats[osapply.INSTALL_STATE] = St(mode=statmod.S_IFREG | 0o644, uid=999) + self.refused(f, "R12") + + def test_R14_kernel_package_in_a_host_plan(self): + f = Fake() + f.plan["layer"] = "host" + f.plan["packages"].append({"name": "linux-image-amd64", "version": "6.12.1-1", "origin": "Debian"}) + self.refused(f, "R14") + + def test_R14_kernel_package_pulled_by_the_simulation(self): + f = Fake() + f.plan["layer"] = "host" + f.extra_sim = ["Inst grub-common [2.12-9] (2.12-10 Debian:13.7/stable [amd64])"] + f.installed["grub-common"] = "2.12-9" + f.plan["packages"].append({"name": "grub-common", "version": "2.12-10", "origin": "Debian"}) + self.refused(f, "R14") + def test_R13_repair_does_not_fix_it(self): f = Fake() f.dpkg_audit = "The following packages are broken\n perl\n" @@ -453,6 +503,79 @@ class Conffiles(unittest.TestCase): self.assertFalse(any("kept /etc/debian_version" in l for l in f.logs), "an updated file must not be reported as kept") +class HostLayer(unittest.TestCase): + def test_host_runs_on_the_host_not_in_the_guest(self): + f = Fake() + f.plan["layer"] = "host" + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + inst = [c for c in f.calls if c[0] == "host" and "install" in c[1] and "-s" not in c[1] and "-f" not in c[1]] + self.assertTrue(inst, "the host install must run on the host") + self.assertFalse([c for c in f.calls if c[0] == "guest" and "install" in c[2]], "nothing installed in the guest") + self.assertEqual(sorted(rep["health_after"]["host_services"]), sorted(osapply.HOST_SERVICES)) + self.assertTrue(rep["health_after"]["guest_running"]) + + def test_pending_fast_skips_proxmox_docker_and_kernel(self): + f = Fake() + f.plan["layer"] = "host" + f.plan["select"] = "pending-fast" + f.plan["packages"] = [] + f.installed.update({"pve-manager": "9.2.2", "linux-image-amd64": "6.12.1", "docker-ce": "29.7"}) + f.pending_sim = [ + "Inst libc6 [2.41-12+deb13u3] (2.41-12+deb13u4 Debian:13.7/stable [amd64])", + "Inst openssl [3.5.6-1~deb13u1] (3.5.7-1~deb13u3 Debian:13.7/stable, Debian-Security:13/stable-security [amd64])", + "Inst pve-manager [9.2.2] (9.2.21 Proxmox Debian Repository:stable [amd64])", + "Inst linux-image-amd64 [6.12.1] (6.12.9 Debian:13.7/stable [amd64])", + "Inst docker-ce [29.7] (29.8 Docker CE:trixie [amd64])", + "Inst brand-new (1.0 Debian:13.7/stable [amd64])", + ] + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + got = sorted(u["name"] for u in rep["upgraded"]) + self.assertEqual(got, ["libc6", "openssl"], "pending-fast must take only installed, Debian-origin, non-kernel packages") + + def test_reboot_needed_when_pid1_or_lxc_start(self): + f = Fake() + f.plan["layer"] = "host" + f.restart_out = "1 systemd\n2101 lxc-start\n530 sshd\n" + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + self.assertTrue(rep["reboot_needed"]) + self.assertIn("lxc-start", rep["restart_needed"]) + + def test_reboot_needed_for_lxc_start_alone(self): + # lxc-start runs the guest; only a guest restart (or a host reboot) replaces it + f = Fake() + f.plan["layer"] = "host" + f.restart_out = "2101 lxc-start\n530 sshd\n" + rc, rep = run(f) + self.assertTrue(rep["reboot_needed"], rep) + + def test_no_reboot_for_ordinary_daemons(self): + f = Fake() + f.plan["layer"] = "host" + f.restart_out = "530 sshd\n611 cron\n" + rc, rep = run(f) + self.assertFalse(rep["reboot_needed"], rep) + + +class Speed(unittest.TestCase): + # R-845: no `pct exec` per package — version checks on the host, madison once, the restart scan only after an install. + def test_no_per_package_guest_calls(self): + f = Fake() + for i in range(40): + f.installed[f"pkg{i}"] = "1.0-1" + f.live[f"pkg{i}"] = {"1.0-2"} + f.plan["packages"].append({"name": f"pkg{i}", "version": "1.0-2", "origin": "Debian"}) + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + guest_cmp = [c for c in f.calls if c[0] == "guest" and "--compare-versions" in c[2]] + self.assertEqual(guest_cmp, [], "version comparisons must run on the host") + self.assertEqual(f.madison_calls, 1, "madison must run once for all packages") + guest_calls = len([c for c in f.calls if c[0] == "guest"]) + self.assertLess(guest_calls, 30, f"{guest_calls} guest calls for 42 packages — something is per-package again") + + class Failure(unittest.TestCase): def test_install_failure_is_rc3_with_dpkg_state(self): f = Fake() diff --git a/internal/hub/cloudflared.go b/internal/hub/cloudflared.go index 0579610..94c0f99 100644 --- a/internal/hub/cloudflared.go +++ b/internal/hub/cloudflared.go @@ -2,45 +2,111 @@ package hub import ( "context" - "os/exec" + "fmt" "strings" + + "gitea.dooplex.hu/admin/felhom-agent/internal/proxmox" ) -// CloudflaredProber reports the cloudflared tunnel service health. It is a -// READ-ONLY probe: the agent does NOT manage or restart cloudflared in this slice -// (that is the tunnel-management slice — this is the seam for it). Injectable so -// tests use a fake and never exec. +// Tunnel states the agent reports (R-841, agent v0.141.0). THREE, never two: a probe that could not ask is +// `unknown`, which the hub never shows as up or down and never alarms on (R-96 rule 3). +const ( + TunnelRunning = "running" // the cloudflared container runs AND its readiness check says CONNECTED + TunnelNotRunning = "not_running" // stopped / exited / absent, or running but NOT connected (Detail says which) + TunnelUnknown = "unknown" // the probe could not ask (guest down, pct/sudo error, health still starting) +) + +// CloudflaredProber reports the box's tunnel. Injectable so tests use a fake and never exec. type CloudflaredProber interface { - // Status returns one of: "active" | "inactive" | "failed" | "unknown". - Status(ctx context.Context) (string, error) + // Status returns one of the Tunnel* states and a short detail (why not_running / why unknown). + Status(ctx context.Context) (status, detail string) } -// SystemctlProber runs `systemctl is-active cloudflared`. This is NOT a Privileged -// (root-CLI) op — `is-active` is non-root readable and is not one of the three -// proven root exceptions, so it does not go through internal/proxmox.Privileged. -type SystemctlProber struct { - Unit string // defaults to "cloudflared" +// GuestTunnelProber reads the REAL tunnel: the `cloudflared` container in the box's own customer guest. +// +// Before v0.141.0 the agent ran `systemctl is-active cloudflared` on the HOST — a unit that does not exist (cloudflared +// is a guest container, `11-os-updates.md` C8), so every box reported `inactive` (R-841). +// +// It uses ONLY the existing sudoers line `pct exec [0-9]* -- docker inspect -f *` (03 §3): the container's state, exit +// code and Docker health status. The health status comes from the compose health check controller v0.292.0 adds +// (`cloudflared tunnel --metrics localhost:20241 ready` → /ready: 200 only with ≥ 1 connection). Measured 2026-10-04: +// with a wrong token the container stays "running" while /ready answers 503 — so the container state alone would lie. +// A container with no health check (an older controller) is judged on its state alone, and Detail says so. +type GuestTunnelProber struct { + Runner proxmox.Runner + // Guests returns the box's customer guest vmids (running pool guests that bind /mnt/felhom-drives). + Guests func(ctx context.Context) ([]int, error) } -// Status maps `systemctl is-active` output to the report vocabulary. systemctl -// exits non-zero for inactive/failed, so the output string is authoritative over -// the exit code; any exec error (binary missing, etc.) maps to "unknown". -func (p SystemctlProber) Status(ctx context.Context) (string, error) { - unit := p.Unit - if unit == "" { - unit = "cloudflared" +const tunnelInspect = `{{.State.Status}}|{{.State.ExitCode}}|{{if .State.Health}}{{.State.Health.Status}}{{else}}none{{end}}` + +// Status probes every customer guest and reports the worst state (normally there is exactly one guest). +func (p GuestTunnelProber) Status(ctx context.Context) (string, string) { + if p.Runner == nil || p.Guests == nil { + return TunnelUnknown, "no probe wired" } - out, _ := exec.CommandContext(ctx, "systemctl", "is-active", unit).Output() - switch strings.TrimSpace(string(out)) { - case "active": - return "active", nil - case "failed": - return "failed", nil - case "inactive", "deactivating", "activating": - return "inactive", nil - case "": - return "unknown", nil // no output → systemctl/exec problem - default: - return "unknown", nil + vmids, err := p.Guests(ctx) + if err != nil { + return TunnelUnknown, "could not list the customer guest: " + err.Error() } + if len(vmids) == 0 { + return TunnelUnknown, "no running customer guest" + } + worst, wdetail := "", "" + rank := map[string]int{TunnelRunning: 0, TunnelUnknown: 1, TunnelNotRunning: 2} + for _, v := range vmids { + out, errOut, err := p.Runner.Run(ctx, "/usr/sbin/pct", "exec", fmt.Sprint(v), "--", "docker", "inspect", "-f", tunnelInspect, "cloudflared") + st, d := ClassifyTunnel(string(out), string(errOut), err) + if len(vmids) > 1 { + d = fmt.Sprintf("guest %d: %s", v, d) + } + if worst == "" || rank[st] > rank[worst] { + worst, wdetail = st, d + } + } + return worst, wdetail +} + +// ClassifyTunnel maps one `docker inspect` answer to a state. Pure; pinned by TestClassifyTunnel. +func ClassifyTunnel(stdout, stderr string, err error) (string, string) { + out := strings.TrimSpace(stdout) + if err != nil || out == "" { + if strings.Contains(stderr, "No such object") || strings.Contains(stderr, "No such container") { + return TunnelNotRunning, "no cloudflared container in the guest" + } + return TunnelUnknown, "could not ask the guest: " + firstLine(stderr, err) + } + parts := strings.Split(out, "|") + if len(parts) != 3 { + return TunnelUnknown, "unreadable docker answer: " + out + } + state, code, health := parts[0], parts[1], parts[2] + if state != "running" { + return TunnelNotRunning, fmt.Sprintf("container %s, exit code %s", state, code) + } + switch health { + case "healthy": + return TunnelRunning, "connected" + case "unhealthy": + return TunnelNotRunning, "container running but the tunnel is NOT connected (cloudflared /ready fails)" + case "starting": + return TunnelUnknown, "container running, readiness check still starting" + case "none": + return TunnelRunning, "container running (no readiness check on this controller — connection not checked)" + } + return TunnelUnknown, "unknown health state " + health +} + +func firstLine(stderr string, err error) string { + s := strings.TrimSpace(stderr) + if i := strings.IndexByte(s, '\n'); i >= 0 { + s = s[:i] + } + if s == "" && err != nil { + s = err.Error() + } + if len(s) > 160 { + s = s[:160] + } + return s } diff --git a/internal/hub/cloudflared_test.go b/internal/hub/cloudflared_test.go new file mode 100644 index 0000000..261f775 --- /dev/null +++ b/internal/hub/cloudflared_test.go @@ -0,0 +1,67 @@ +package hub + +import ( + "context" + "errors" + "io" + "strings" + "testing" +) + +// R-841: the three states from one `docker inspect` answer. Red-proof: map "unhealthy" to running (the container +// state alone — what a plain "is it running" probe would say) and the "running but not connected" case fails. +func TestClassifyTunnel(t *testing.T) { + cases := []struct { + name, out, errOut string + err error + want string + detail string + }{ + {"connected", "running|0|healthy\n", "", nil, TunnelRunning, "connected"}, + {"running but not connected", "running|0|unhealthy\n", "", nil, TunnelNotRunning, "NOT connected"}, + {"stopped", "exited|137|unhealthy\n", "", nil, TunnelNotRunning, "exit code 137"}, + {"absent", "", "Error: No such object: cloudflared", errors.New("exit status 1"), TunnelNotRunning, "no cloudflared container"}, + {"still starting", "running|0|starting\n", "", nil, TunnelUnknown, "starting"}, + {"no health check (older controller)", "running|0|none\n", "", nil, TunnelRunning, "connection not checked"}, + {"guest not running", "", "CT 9201 not running", errors.New("exit status 255"), TunnelUnknown, "could not ask"}, + {"sudo refused", "", "sudo: a password is required", errors.New("exit status 1"), TunnelUnknown, "could not ask"}, + } + for _, c := range cases { + st, d := ClassifyTunnel(c.out, c.errOut, c.err) + if st != c.want || !strings.Contains(d, c.detail) { + t.Errorf("%s: got %q (%s), want %q (…%s…)", c.name, st, d, c.want, c.detail) + } + } +} + +type tunnelRunner struct { + calls []string + out map[string]string +} + +func (r *tunnelRunner) Run(_ context.Context, name string, args ...string) ([]byte, []byte, error) { + line := name + " " + strings.Join(args, " ") + r.calls = append(r.calls, line) + return []byte(r.out[args[1]]), nil, nil +} +func (r *tunnelRunner) RunStdin(ctx context.Context, _ io.Reader, name string, args ...string) ([]byte, []byte, error) { + return r.Run(ctx, name, args...) +} + +// The probe uses EXACTLY the existing sudoers shape `pct exec -- docker inspect -f cloudflared`, and +// with no customer guest it is unknown, never down. +func TestGuestTunnelProber(t *testing.T) { + r := &tunnelRunner{out: map[string]string{"9201": "running|0|healthy"}} + p := GuestTunnelProber{Runner: r, Guests: func(context.Context) ([]int, error) { return []int{9201}, nil }} + if st, d := p.Status(context.Background()); st != TunnelRunning || d != "connected" { + t.Fatalf("got %q %q", st, d) + } + want := "/usr/sbin/pct exec 9201 -- docker inspect -f " + tunnelInspect + " cloudflared" + if len(r.calls) != 1 || r.calls[0] != want { + t.Fatalf("command = %q, want %q", r.calls, want) + } + none := GuestTunnelProber{Runner: r, Guests: func(context.Context) ([]int, error) { return nil, nil }} + if st, _ := none.Status(context.Background()); st != TunnelUnknown { + t.Fatalf("no guest → %q, want unknown", st) + } +} diff --git a/internal/hub/collect.go b/internal/hub/collect.go index 496f176..0936c03 100644 --- a/internal/hub/collect.go +++ b/internal/hub/collect.go @@ -280,7 +280,7 @@ func (c *Collector) Collect(ctx context.Context) (*HostReport, error) { PBSSnapshots: c.collectPBSSnapshots(ctx), AuditTail: []AuditEntry{}, - Cloudflared: Cloudflared{Status: c.cloudflaredStatus(ctx)}, + Cloudflared: c.cloudflared(ctx), Capabilities: c.capabilities(ctx), LeafFingerprint: c.leafFP, Addresses: c.collectAddresses(), @@ -555,16 +555,18 @@ func (c *Collector) collectPBSSnapshots(ctx context.Context) []PBSSnapshot { return []PBSSnapshot{} } -func (c *Collector) cloudflaredStatus(ctx context.Context) string { +func (c *Collector) cloudflared(ctx context.Context) Cloudflared { if c.cf == nil { - return "unknown" + return Cloudflared{Status: TunnelUnknown, Detail: "no probe wired"} } - st, err := c.cf.Status(ctx) - if err != nil || st == "" { - c.logger.Warn("hub: cloudflared probe failed", "err", err) - return "unknown" + st, d := c.cf.Status(ctx) + if st == "" { + st = TunnelUnknown } - return st + if st == TunnelUnknown { + c.logger.Debug("hub: tunnel probe could not decide", "detail", d) + } + return Cloudflared{Status: st, Detail: d} } func percent(used, total int64) float64 { diff --git a/internal/hub/collect_guestnet_test.go b/internal/hub/collect_guestnet_test.go index 0b76de1..3c2910e 100644 --- a/internal/hub/collect_guestnet_test.go +++ b/internal/hub/collect_guestnet_test.go @@ -16,7 +16,7 @@ func (f fakeGuestNet) GuestNetStatus(context.Context) *GuestNetStatus { return f func TestCollect_GuestNetOmittedWhenReporterNil(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, fakeObserver{}, nil, nil, nil, "h", "0.92.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, fakeObserver{}, nil, nil, nil, "h", "0.92.0", quietLogger()) r, err := c.Collect(context.Background()) if err != nil { t.Fatalf("Collect: %v", err) @@ -38,7 +38,7 @@ func TestCollect_GuestNetOmittedWhenReporterNil(t *testing.T) { func TestCollect_GuestNetPopulatedWhenWired(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, fakeObserver{}, nil, nil, nil, "h", "0.92.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, fakeObserver{}, nil, nil, nil, "h", "0.92.0", quietLogger()) c.SetGuestNetReporter(fakeGuestNet{st: &GuestNetStatus{ CheckedAt: "2026-07-21T10:00:00Z", Guests: []GuestNetGuest{{ diff --git a/internal/hub/collect_mgmtplane_test.go b/internal/hub/collect_mgmtplane_test.go index 3e36a13..776bcdb 100644 --- a/internal/hub/collect_mgmtplane_test.go +++ b/internal/hub/collect_mgmtplane_test.go @@ -12,7 +12,7 @@ func (f fakeMgmtPlane) MgmtPlaneStatus(context.Context) *MgmtPlaneStatus { retur func TestCollect_MgmtPlaneOmittedWhenReporterNil(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, fakeObserver{}, nil, nil, nil, "h", "0.71.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, fakeObserver{}, nil, nil, nil, "h", "0.71.0", quietLogger()) r, err := c.Collect(context.Background()) if err != nil { t.Fatalf("Collect: %v", err) @@ -24,7 +24,7 @@ func TestCollect_MgmtPlaneOmittedWhenReporterNil(t *testing.T) { func TestCollect_MgmtPlanePopulatedWhenWired(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, fakeObserver{}, nil, nil, nil, "h", "0.71.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, fakeObserver{}, nil, nil, nil, "h", "0.71.0", quietLogger()) c.SetMgmtPlaneReporter(fakeMgmtPlane{st: &MgmtPlaneStatus{ PrivsepDirOK: true, SshdReachable: true, HealedRecently: true, PrivsepHealedAt: "2026-07-05T16:42:17Z", }}) @@ -47,7 +47,7 @@ func (f fakeOOB) OOBStatus(context.Context) *OOBStatus { return f.st } func TestCollect_OOBOmittedWhenNil(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, fakeObserver{}, nil, nil, nil, "h", "0.72.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, fakeObserver{}, nil, nil, nil, "h", "0.72.0", quietLogger()) r, _ := c.Collect(context.Background()) if r.OOB != nil { t.Fatalf("no reporter → oob omitted, got %+v", r.OOB) @@ -56,7 +56,7 @@ func TestCollect_OOBOmittedWhenNil(t *testing.T) { func TestCollect_OOBPopulatedWhenWired(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, fakeObserver{}, nil, nil, nil, "h", "0.72.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, fakeObserver{}, nil, nil, nil, "h", "0.72.0", quietLogger()) c.SetOOBReporter(fakeOOB{st: &OOBStatus{FelhomSshdActive: true, FelhomSshdPort: 8822, Reachable: true}}) r, _ := c.Collect(context.Background()) if r.OOB == nil || r.OOB.FelhomSshdPort != 8822 || !r.OOB.Reachable { diff --git a/internal/hub/collect_test.go b/internal/hub/collect_test.go index 6d1398f..03c4b09 100644 --- a/internal/hub/collect_test.go +++ b/internal/hub/collect_test.go @@ -33,7 +33,7 @@ func TestCollect_StorageTargetsFromObserver(t *testing.T) { obs := fakeObserver{targets: []StorageTarget{ {Name: "local-lvm", Type: StorageTypeLVMThin, State: StorageStateAttached, Reachable: true}, }} - c := NewCollector(px, fakeProber{status: "active"}, obs, nil, nil, nil, "h", "0.5.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, obs, nil, nil, nil, "h", "0.5.0", quietLogger()) r, err := c.Collect(context.Background()) if err != nil { t.Fatalf("Collect: %v", err) @@ -45,7 +45,7 @@ func TestCollect_StorageTargetsFromObserver(t *testing.T) { func TestCollect_StorageObserverErrorDegradesToEmpty(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, fakeObserver{err: errors.New("proxmox down")}, nil, nil, nil, "h", "0.5.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, fakeObserver{err: errors.New("proxmox down")}, nil, nil, nil, "h", "0.5.0", quietLogger()) r, err := c.Collect(context.Background()) if err != nil { t.Fatalf("a storage observe error must not sink the heartbeat: %v", err) @@ -64,7 +64,7 @@ func TestCollect_HostAndGuests(t *testing.T) { }, cfg: map[int]proxmox.GuestConfig{100: {Cores: 2, Memory: 2048}}, } - c := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "demo-host-01", "0.3.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "demo-host-01", "0.3.0", quietLogger()) r, err := c.Collect(context.Background()) if err != nil { t.Fatalf("Collect: %v", err) @@ -88,7 +88,7 @@ func TestCollect_HostAndGuests(t *testing.T) { if g.Spec.Cores != 2 || g.Spec.MemoryBytes != 2147483648 || g.Spec.DiskBytes != 21474836480 { t.Errorf("spec = %+v", g.Spec) } - if r.Cloudflared.Status != "active" { + if r.Cloudflared.Status != "running" || r.Cloudflared.Detail != "connected" { t.Errorf("cloudflared = %q", r.Cloudflared.Status) } } @@ -104,7 +104,7 @@ func TestCollect_GuestConfigFailureKeepsStatusOmitsSpec(t *testing.T) { cfg: map[int]proxmox.GuestConfig{100: {Cores: 2}}, cfgErr: map[int]error{200: errors.New("config read failed")}, } - c := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "h", "0.3.1", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "h", "0.3.1", quietLogger()) r, err := c.Collect(context.Background()) if err != nil { t.Fatalf("a per-guest failure must NOT fail the whole report: %v", err) @@ -125,7 +125,7 @@ func TestCollect_GuestConfigFailureKeepsStatusOmitsSpec(t *testing.T) { func TestCollect_NodeStatusFailureIsHardError(t *testing.T) { px := &fakePx{node: "n", nsErr: errors.New("proxmox down")} - c := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "h", "0.3.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "h", "0.3.0", quietLogger()) if _, err := c.Collect(context.Background()); err == nil { t.Fatal("NodeStatus failure must be a hard error (no useful report)") } @@ -133,7 +133,7 @@ func TestCollect_NodeStatusFailureIsHardError(t *testing.T) { func TestCollect_CloudflaredProbeErrorIsUnknown(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{err: errors.New("no systemctl")}, nil, nil, nil, nil, "h", "0.3.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "", detail: "could not ask"}, nil, nil, nil, nil, "h", "0.3.0", quietLogger()) r, err := c.Collect(context.Background()) if err != nil { t.Fatalf("cloudflared failure must not be fatal: %v", err) @@ -153,7 +153,7 @@ func TestCollect_LeafFingerprint(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} const fp = "60b5974d586f5f3c8ec41eb998d0f07406178219c36bf6d3ff377570279d8245" - c := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "h", "0.48.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "h", "0.48.0", quietLogger()) c.SetLeafFingerprint(fp) r, err := c.Collect(context.Background()) if err != nil { @@ -164,7 +164,7 @@ func TestCollect_LeafFingerprint(t *testing.T) { } // Companion: no SetLeafFingerprint (local API disabled) → empty, never a fabricated value. - c2 := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "h", "0.48.0", quietLogger()) + c2 := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "h", "0.48.0", quietLogger()) r2, _ := c2.Collect(context.Background()) if r2.LeafFingerprint != "" { t.Fatalf("unset leaf_fingerprint = %q, want empty", r2.LeafFingerprint) diff --git a/internal/hub/dr_recipe_test.go b/internal/hub/dr_recipe_test.go index 574349c..c7640e4 100644 --- a/internal/hub/dr_recipe_test.go +++ b/internal/hub/dr_recipe_test.go @@ -384,7 +384,7 @@ func TestCollectDRRecipe_ProductionPath(t *testing.T) { obs := fakeObserver{targets: capturedDemoFelhomTargets()} pbsRep := fakePBSReporter{snaps: capturedDemoFelhomSnapshots()} - c := NewCollector(px, fakeProber{status: "active"}, obs, nil, nil, pbsRep, "h", "0.118.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, obs, nil, nil, pbsRep, "h", "0.118.0", quietLogger()) c.SetBackupTargetResolver(func() ConfiguredBackupTarget { return ConfiguredBackupTarget{StorageID: "felhom-backup", Known: true} }) @@ -409,7 +409,7 @@ func TestCollectDRRecipe_ProductionPath(t *testing.T) { // test that would have caught shipping the seam without wiring it. func TestCollectDRRecipe_UnwiredSeamReportsUnknown(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, fakeObserver{targets: capturedDemoFelhomTargets()}, + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, fakeObserver{targets: capturedDemoFelhomTargets()}, nil, nil, nil, "h", "0.118.0", quietLogger()) r, err := c.Collect(context.Background()) diff --git a/internal/hub/hostmetrics_test.go b/internal/hub/hostmetrics_test.go index 63d0f21..fb9c799 100644 --- a/internal/hub/hostmetrics_test.go +++ b/internal/hub/hostmetrics_test.go @@ -16,7 +16,7 @@ func intp(v int) *int { return &v } // HostMetricsNow returns a fresh host block with cpu% from NodeStatus and the temp from the reader. func TestHostMetricsNow_PopulatesTemp(t *testing.T) { px := &fakePx{node: "demo-felhom", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "h", "0.14.0", quietLogger()). + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "h", "0.14.0", quietLogger()). SetTempReader(fakeTemp{c: intp(46)}) h, err := c.HostMetricsNow(context.Background()) if err != nil { @@ -36,7 +36,7 @@ func TestHostMetricsNow_PopulatesTemp(t *testing.T) { // A missing temp sensor gracefully nulls cpu_temp_c without failing the host read. func TestHostMetricsNow_GracefulNullTemp(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "h", "0.14.0", quietLogger()). + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "h", "0.14.0", quietLogger()). SetTempReader(fakeTemp{c: nil}) h, err := c.HostMetricsNow(context.Background()) if err != nil { @@ -50,7 +50,7 @@ func TestHostMetricsNow_GracefulNullTemp(t *testing.T) { // A NodeStatus failure is a hard error (no useful host view). func TestHostMetricsNow_NodeStatusErrorIsHard(t *testing.T) { px := &fakePx{node: "n", nsErr: errors.New("proxmox down")} - c := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "h", "0.14.0", quietLogger()) + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "h", "0.14.0", quietLogger()) if _, err := c.HostMetricsNow(context.Background()); err == nil { t.Fatal("NodeStatus failure must be a hard error") } @@ -59,7 +59,7 @@ func TestHostMetricsNow_NodeStatusErrorIsHard(t *testing.T) { // Collect() (the hub report) also carries the temp now — the operator freebie. func TestCollect_HostReportCarriesTemp(t *testing.T) { px := &fakePx{node: "n", ns: newTestNodeStatus()} - c := NewCollector(px, fakeProber{status: "active"}, nil, nil, nil, nil, "h", "0.14.0", quietLogger()). + c := NewCollector(px, fakeProber{status: "running", detail: "connected"}, nil, nil, nil, nil, "h", "0.14.0", quietLogger()). SetTempReader(fakeTemp{c: intp(51)}) r, err := c.Collect(context.Background()) if err != nil { diff --git a/internal/hub/mock_test.go b/internal/hub/mock_test.go index 627ea6f..e09cca8 100644 --- a/internal/hub/mock_test.go +++ b/internal/hub/mock_test.go @@ -56,7 +56,7 @@ func (f *fakePx) GuestConfig(ctx context.Context, vmid int) (proxmox.GuestConfig // fakeProber is a fake CloudflaredProber. type fakeProber struct { status string - err error + detail string } -func (p fakeProber) Status(ctx context.Context) (string, error) { return p.status, p.err } +func (p fakeProber) Status(ctx context.Context) (string, string) { return p.status, p.detail } diff --git a/internal/hub/report.go b/internal/hub/report.go index 04329e1..47b72a0 100644 --- a/internal/hub/report.go +++ b/internal/hub/report.go @@ -289,9 +289,10 @@ type GuestSpec struct { DiskBytes int64 `json:"disk_bytes"` } -// Cloudflared is the tunnel service health (read-only probe this slice). +// Cloudflared is the box's tunnel (R-841, agent v0.141.0): the cloudflared container in the customer guest. type Cloudflared struct { - Status string `json:"status"` // active | inactive | failed | unknown + Status string `json:"status"` // running | not_running | unknown (TunnelRunning …) + Detail string `json:"detail,omitempty"` // why not_running / unknown, or "connected" } // The following element types are declared now so the empty collections above are @@ -559,6 +560,9 @@ type WireOSUpdate struct { Ring int `json:"ring"` Enabled bool `json:"enabled"` Release *WireOSRelease `json:"release,omitempty"` + // HostRelease is the newest approved HOST release (hub v0.131.0, `11` §8 step 3) — a separate set: a version + // approved for the guest is not approved for the host by that fact alone. + HostRelease *WireOSRelease `json:"host_release,omitempty"` } // WireOSRelease is an approved version set; Snapshot is the approval time (YYYYMMDDTHHMMSSZ) the wrapper uses diff --git a/internal/osupdate/leg.go b/internal/osupdate/leg.go index dfe78eb..58b0197 100644 --- a/internal/osupdate/leg.go +++ b/internal/osupdate/leg.go @@ -1,14 +1,14 @@ -// Package osupdate is the agent's OS-update leg for the customer GUEST (`11-os-updates.md` §8 step 2, agent v0.140.0). +// Package osupdate is the agent's OS-update leg (`11-os-updates.md` §8 steps 2–3): the customer GUEST's Debian fast +// lane (agent v0.140.0) and, after it in the same pass, the HOST's (agent v0.141.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. +// sudo, judges health and reports to the hub — one report per layer. // -// 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. +// NO AUTOMATIC UNDO (R-837 measured; `09` §3 decision 81): a failed health check stops, reports `health_failed` and the +// hub mails the operator; the whole-guest backup taken minutes earlier is the guest's undo, by hand; a host package is +// put back by hand from the previous release's snapshot (runbook). The host is NEVER rebooted by this package. package osupdate import ( @@ -33,6 +33,12 @@ 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" +// Layers. +const ( + LayerGuest = "guest" + LayerHost = "host" +) + // Package is one name=version with its origin. type Package struct { Name string `json:"name"` @@ -40,7 +46,7 @@ type Package struct { Origin string `json:"origin"` } -// Pending is one update the guest's sources offer (origin as apt names it, possibly several). +// Pending is one update the sources offer (origin as apt names it, possibly several). type Pending struct { Name string `json:"name"` From string `json:"from"` @@ -54,17 +60,22 @@ type Container struct { Health string `json:"health"` // healthy | unhealthy | starting | none } -// Health is the guest's health snapshot (the wrapper's `health` object). +// Health is one health reading. Guest layer: DockerOK..Containers. Host layer: HostServices, GuestRunning and the +// guest's own reading in Guest. type Health struct { - DockerOK bool `json:"docker_ok"` - NetworkOK bool `json:"network_ok"` - Controller string `json:"controller"` - Containers map[string]Container `json:"containers"` + DockerOK bool `json:"docker_ok"` + NetworkOK bool `json:"network_ok"` + Controller string `json:"controller"` + Containers map[string]Container `json:"containers"` + HostServices map[string]string `json:"host_services,omitempty"` + GuestRunning *bool `json:"guest_running,omitempty"` + Guest *Health `json:"guest,omitempty"` } // WrapperReport is the wrapper's OSAPPLY-REPORT object. type WrapperReport struct { Mode string `json:"mode"` + Layer string `json:"layer"` Refused json.RawMessage `json:"refused"` Failed json.RawMessage `json:"failed"` Upgraded []Package `json:"upgraded"` @@ -76,14 +87,16 @@ type WrapperReport struct { HealthBefore *Health `json:"health_before"` HealthAfter *Health `json:"health_after"` Health *Health `json:"health"` + PassSeconds float64 `json:"pass_seconds"` } 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). +// Report is what the hub receives per layer (hub osupdates.Report — field-exact). type Report struct { RunID string `json:"run_id"` + Layer string `json:"layer"` Trigger string `json:"trigger"` Mode string `json:"mode"` Ring int `json:"ring"` @@ -100,6 +113,7 @@ type Report struct { DockerRestartNeeded bool `json:"docker_restart_needed,omitempty"` RebootNeeded bool `json:"reboot_needed,omitempty"` Refused json.RawMessage `json:"refused,omitempty"` + PassSeconds float64 `json:"pass_seconds,omitempty"` } // Reporter posts a report to the hub (*hub.Client). @@ -107,10 +121,12 @@ type Reporter interface { PostOSReport(ctx context.Context, body []byte) error } -// Leg runs one OS-update pass for the customer guest. +// Leg runs one OS-update pass for the customer guest and then the host. type Leg struct { Runner proxmox.Runner Hub Reporter + Tunnel hub.CloudflaredProber // the host health rule needs the tunnel `running` (R-841) + Appliance bool // agent.json deployment_mode; the wrapper re-checks the ROOT-owned record (R12) Logger *slog.Logger PlanDir string StatePath string // last night run (once per night) @@ -122,7 +138,6 @@ type Leg struct { mu sync.Mutex block *hub.WireOSUpdate - have bool } // OnDesiredState stores the hub's os_update block (desired.RawConsumer — store only, never block). @@ -132,7 +147,7 @@ func (l *Leg) OnDesiredState(_ context.Context, resp *hub.DesiredStateResponse) } l.mu.Lock() defer l.mu.Unlock() - l.block, l.have = resp.DesiredState.OSUpdate, true + l.block = resp.DesiredState.OSUpdate } // Block returns the newest os_update block. No block (an older hub, or nothing fetched yet) = ring 1, ON, no @@ -150,7 +165,7 @@ func (l *Leg) Block() hub.WireOSUpdate { func (l *Leg) SetBlock(b *hub.WireOSUpdate) { l.mu.Lock() defer l.mu.Unlock() - l.block, l.have = b, true + l.block = b } func (l *Leg) now() time.Time { @@ -191,18 +206,9 @@ func IsFast(origins []string) bool { 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. +// HealthVerdict is THE guest health rule (`11` §8.1; pinned by TestHealthVerdict*): docker answers, the guest's +// network resolves, the controller's own health check is `healthy`, and every container that was running at the +// start of the pass runs again — and healthy again if it was. "starting" is not yet healthy. func HealthVerdict(before, after *Health) (bool, string) { if after == nil { return false, "no health reading" @@ -240,34 +246,44 @@ func HealthVerdict(before, after *Health) (bool, string) { return true, "" } -// MergeBaseline is the health baseline from two readings: a container counts as running (and healthy) if EITHER -// reading saw it so. nil-safe. Pinned by TestHealth_BaselineIsTheStartOfTheLeg. -func MergeBaseline(a, b *Health) *Health { - if a == nil { - return b +// HostHealthVerdict is THE host health rule (`11` §8.2; pinned by TestHostHealthVerdict): the Proxmox daemons and the +// agent are active, the customer guest still runs, the guest's own rule still passes against the start of the pass, +// and the tunnel is `running` (R-841). +func HostHealthVerdict(before, after *Health, tunnel string) (bool, string) { + if after == nil { + return false, "no health reading" } - if b == nil { - return a + svcs := make([]string, 0, len(after.HostServices)) + for s := range after.HostServices { + svcs = append(svcs, s) } - out := &Health{DockerOK: a.DockerOK && b.DockerOK, NetworkOK: a.NetworkOK && b.NetworkOK, Controller: b.Controller, - Containers: map[string]Container{}} - for _, h := range []*Health{a, b} { - for n, c := range h.Containers { - cur, seen := out.Containers[n] - if !seen || (cur.State != "running" && c.State == "running") { - out.Containers[n] = c - continue - } - if c.State == "running" && c.Health == "healthy" { - out.Containers[n] = c - } + sort.Strings(svcs) + if len(svcs) == 0 { + return false, "no host service reading" + } + for _, s := range svcs { + if after.HostServices[s] != "active" { + return false, s + " is " + after.HostServices[s] } } - return out + if after.GuestRunning == nil || !*after.GuestRunning { + return false, "the customer guest is not running" + } + var gb *Health + if before != nil { + gb = before.Guest + } + if ok, why := HealthVerdict(gb, after.Guest); !ok { + return false, "guest: " + why + } + if tunnel != hub.TunnelRunning { + return false, "the tunnel is " + tunnel + } + 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) { +func (l *Leg) call(ctx context.Context, runID string, plan map[string]any) (WrapperReport, error) { dir := l.PlanDir if dir == "" { dir = DefaultPlanDir @@ -275,13 +291,8 @@ func (l *Leg) call(ctx context.Context, runID, mode string, vmid int, rel hub.Wi 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") + path := filepath.Join(dir, fmt.Sprintf("plan-%s-%s-%s.json", runID, plan["layer"], plan["mode"])) if err := os.WriteFile(path, b, 0o600); err != nil { return WrapperReport{}, fmt.Errorf("osupdate: write plan: %w", err) } @@ -307,20 +318,11 @@ func (l *Leg) call(ctx context.Context, runID, mode string, vmid int, rel hub.Wi 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 { +// Run is one pass: the guest layer, then (on an appliance, after a good guest step) the host layer. Returns both +// reports (host empty when skipped). trigger is "night" or "debug". +func (l *Leg) Run(ctx context.Context, vmid int, trigger string) (guest Report, host 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) - + lg := l.log().With("run", runID, "vmid", vmid, "trigger", trigger) if trigger == "night" && l.StatePath != "" { gap := l.MinGap if gap == 0 { @@ -329,68 +331,107 @@ func (l *Leg) Run(ctx context.Context, vmid int, trigger string) Report { 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 + return Report{RunID: runID, Layer: LayerGuest, Outcome: "skipped"}, Report{} } } } + blk := l.Block() + guest = l.runLayer(ctx, runID, LayerGuest, vmid, trigger, blk) + if trigger == "night" && l.StatePath != "" { + _ = os.WriteFile(l.StatePath, []byte(l.now().UTC().Format(time.RFC3339)), 0o600) + } + switch { + case !l.Appliance: + lg.Info("osupdate: host step skipped — not an appliance install (a BYO host belongs to its owner, `11` §1)") + case !(guest.Outcome == "applied" || guest.Outcome == "nothing" || guest.Outcome == "inventory") || !guest.Healthy: + lg.Warn("osupdate: host step skipped — the guest step did not end healthy", "guest_outcome", guest.Outcome, "reason", guest.HealthReason) + default: + host = l.runLayer(ctx, runID, LayerHost, vmid, trigger, blk) + } + return guest, host +} + +func (l *Leg) runLayer(ctx context.Context, runID, layer string, vmid int, trigger string, blk hub.WireOSUpdate) Report { + rel := hub.WireOSRelease{ID: "ring0-" + runID} + var wire *hub.WireOSRelease + if layer == LayerGuest { + wire = blk.Release + } else { + wire = blk.HostRelease + } + if blk.Ring == 1 { + rel = hub.WireOSRelease{} + if wire != nil { + rel = *wire + } + } + rep := Report{RunID: runID, Layer: layer, Trigger: trigger, Ring: blk.Ring, VMID: vmid, ReleaseID: rel.ID} + lg := l.log().With("run", runID, "layer", layer, "vmid", vmid, "ring", blk.Ring, "trigger", trigger) 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) + plan := map[string]any{"release_id": rel.ID, "layer": layer, "lane": "fast", "vmid": vmid, "snapshot": rel.Snapshot, + "packages": []Package{}, "mode": "apply", "select": "listed"} + if rel.ID == "" { + plan["release_id"] = "none" } - var plan []Package + planned := map[string]bool{} switch { case !blk.Enabled: - rep.Outcome = "inventory" + plan["mode"] = "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)}) - } - } + plan["select"] = "pending-fast" // the wrapper picks every pending Debian / Debian-Security upgrade + case len(rel.Packages) == 0: + plan["mode"] = "inventory" // ring 1 with no approved release for this layer: nothing to install default: + var pk []Package for _, p := range rel.Packages { - plan = append(plan, Package{Name: p.Name, Version: p.Version, Origin: p.Origin}) + pk = append(pk, Package{Name: p.Name, Version: p.Version, Origin: p.Origin}) + planned[p.Name] = true } + plan["packages"] = pk } - pkgNames := map[string]bool{} - for _, p := range plan { - pkgNames[p.Name] = true + rep.Mode = plan["mode"].(string) + wr, err := l.call(ctx, runID, plan) + switch { + case err != nil: + rep.Outcome, rep.HealthReason = "failed", err.Error() + return l.finish(ctx, lg, rep) + case wr.refused(): + rep.Outcome, rep.Refused = "refused", wr.Refused + return l.finish(ctx, lg, rep) + case wr.failed(): + rep.Outcome, rep.Refused = "failed", wr.Failed } - 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) + rep.Upgraded, rep.PassSeconds = wr.Upgraded, wr.PassSeconds + if rep.Outcome == "" { 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 + case rep.Mode == "inventory" && !blk.Enabled: + rep.Outcome = "inventory" + case len(wr.Upgraded) == 0: + rep.Outcome = "nothing" + default: + rep.Outcome = "applied" } - final = ap - rep.Upgraded = ap.Upgraded - if rep.Outcome == "" { - if len(ap.Upgraded) == 0 { - rep.Outcome = "nothing" - } else { - rep.Outcome = "applied" + } + if blk.Ring == 0 { + for _, u := range wr.Upgraded { + planned[u.Name] = true + } + } + // Health: compare with the start of the pass; give restarted services time (only after an install). + cur := wr.HealthAfter + verdict := func(h *Health) (bool, string) { + if layer == LayerHost { + t := hub.TunnelUnknown + if l.Tunnel != nil { + t, _ = l.Tunnel.Status(ctx) } + return HostHealthVerdict(wr.HealthBefore, h, t) } - // Health: compare with what the guest looked like BEFORE the run; give restarted services time. + return HealthVerdict(wr.HealthBefore, h) + } + if len(wr.Upgraded) > 0 { wait, poll := l.HealthWait, l.HealthPoll if wait == 0 { wait = 5 * time.Minute @@ -399,19 +440,15 @@ func (l *Leg) Run(ctx context.Context, vmid int, trigger string) Report { poll = 15 * time.Second } deadline := l.now().Add(wait) - cur := ap.HealthAfter - // The baseline is the guest as it was at the START of the leg (the inventory's reading) merged with the - // apply's own "before": an app that stops at any point during the run counts. Found live 2026-10-04: an app - // stopped between the inventory and the apply's own reading was taken as "stopped before" and ignored. - base := MergeBaseline(inv.HealthAfter, ap.HealthBefore) for { - ok, why := HealthVerdict(base, cur) + ok, why := verdict(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) + hp := map[string]any{"release_id": plan["release_id"], "layer": layer, "lane": "fast", "vmid": vmid, "mode": "health", "packages": []Package{}} + hr, herr := l.call(ctx, runID, hp) if herr == nil && hr.Health != nil { cur = hr.Health } @@ -420,20 +457,16 @@ func (l *Leg) Run(ctx context.Context, vmid int, trigger string) Report { 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) + rep.Healthy, rep.HealthReason = verdict(cur) } + rep.Installed, rep.Pending = wr.Installed, wr.Pending + rep.RestartNeeded, rep.DockerRestartNeeded, rep.RebootNeeded = wr.RestartNeeded, wr.DockerRestartNeeded, wr.RebootNeeded + rep.NotCovered = notCovered(wr.Pending, blk.Ring, planned) 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. +// 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 { @@ -446,7 +479,8 @@ func notCovered(pending []Pending, ring int, planned map[string]bool) []string { 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)) + "upgraded", len(rep.Upgraded), "pending", len(rep.Pending), "not_covered", len(rep.NotCovered), + "restart_needed", len(rep.RestartNeeded), "reboot_needed", rep.RebootNeeded, "wrapper_seconds", rep.PassSeconds) if l.Hub != nil { body, _ := json.Marshal(rep) rctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), time.Minute) diff --git a/internal/osupdate/leg_test.go b/internal/osupdate/leg_test.go index 2be35e0..6b0876e 100644 --- a/internal/osupdate/leg_test.go +++ b/internal/osupdate/leg_test.go @@ -14,15 +14,27 @@ import ( "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. +// fakeWrapper plays /usr/local/sbin/felhom-os-apply: it reads the plan the leg wrote and answers per layer and mode. type fakeWrapper struct { t *testing.T pending []Pending - applyRep WrapperReport - healthSeq []*Health // answers to successive "health" calls + applyRep map[string]WrapperReport // per layer + healthSeq map[string][]*Health // per layer: answers to successive "health" calls plans []map[string]any } +func yes() *bool { b := true; return &b } + +func guestOK() *Health { + return &Health{DockerOK: true, NetworkOK: true, Controller: "healthy", Containers: map[string]Container{ + "felhom-controller": {State: "running", Health: "healthy"}, "app": {State: "running", Health: "healthy"}}} +} + +func hostOK() *Health { + return &Health{HostServices: map[string]string{"pveproxy": "active", "pvedaemon": "active", "pvestatd": "active", + "pve-cluster": "active", "felhom-agent": "active"}, GuestRunning: yes(), Guest: guestOK()} +} + 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) @@ -34,27 +46,30 @@ func (f *fakeWrapper) Run(_ context.Context, name string, args ...string) ([]byt var plan map[string]any json.Unmarshal(b, &plan) f.plans = append(f.plans, plan) + layer := plan["layer"].(string) + ok := guestOK() + if layer == LayerHost { + ok = hostOK() + } 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, + rep = WrapperReport{Mode: "inventory", Pending: f.pending, HealthBefore: ok, HealthAfter: ok, Installed: []Package{{Name: "libc6", Version: "u3", Origin: "Debian"}}} case "apply": - rep = f.applyRep + rep = f.applyRep[layer] rep.Mode = "apply" if rep.HealthBefore == nil { - rep.HealthBefore = healthy + rep.HealthBefore = ok } if rep.HealthAfter == nil { - rep.HealthAfter = healthy + rep.HealthAfter = ok } case "health": - if len(f.healthSeq) > 0 { - rep.Health, f.healthSeq = f.healthSeq[0], f.healthSeq[1:] + if seq := f.healthSeq[layer]; len(seq) > 0 { + rep.Health, f.healthSeq[layer] = seq[0], seq[1:] } else { - rep.Health = healthy + rep.Health = ok } } out, _ := json.Marshal(rep) @@ -74,11 +89,21 @@ func (h *fakeHub) PostOSReport(_ context.Context, body []byte) error { return nil } +type fakeTunnel struct{ st string } + +func (t fakeTunnel) Status(context.Context) (string, string) { return t.st, "" } + 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) + if w.applyRep == nil { + w.applyRep = map[string]WrapperReport{} + } + if w.healthSeq == nil { + w.healthSeq = map[string][]*Health{} + } l := &Leg{Runner: w, Hub: h, PlanDir: t.TempDir(), StatePath: filepath.Join(t.TempDir(), "last"), - HealthWait: time.Minute, HealthPoll: 10 * time.Second, + HealthWait: time.Minute, HealthPoll: 10 * time.Second, Appliance: true, Tunnel: fakeTunnel{hub.TunnelRunning}, Now: func() time.Time { return now }, Sleep: func(_ context.Context, d time.Duration) { now = now.Add(d) }} if blk != nil { @@ -93,66 +118,73 @@ var pend = []Pending{ {Name: "docker-ce", From: "29.7", To: "29.8", Origin: []string{"Docker CE"}}, } -func modes(w *fakeWrapper) string { +func calls(w *fakeWrapper) string { var m []string for _, p := range w.plans { - m = append(m, p["mode"].(string)) + m = append(m, p["layer"].(string)+":"+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:]}} +// Ring 0: ONE wrapper call per layer (R-845), select pending-fast (the wrapper picks every Debian / Debian-Security +// upgrade), from live sources; the guest step first, then the host step. +func TestRing0_OneCallPerLayer(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: map[string]WrapperReport{ + LayerGuest: {Upgraded: []Package{{Name: "libc6", Version: "u4"}, {Name: "openssl", Version: "u3"}}, Pending: pend[2:]}, + LayerHost: {Upgraded: []Package{{Name: "openssl", Version: "u3"}}}, + }} 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) + g, ho := l.Run(context.Background(), 9201, "night") + if g.Outcome != "applied" || !g.Healthy || ho.Outcome != "applied" || !ho.Healthy { + t.Fatalf("guest %+v\nhost %+v", g, ho) } - pk := w.plans[1]["packages"].([]any) - if len(pk) != 2 { - t.Fatalf("ring-0 plan = %v, want libc6 + openssl only", pk) + if calls(w) != "guest:apply,host:apply" { + t.Fatalf("calls = %s, want one apply per layer, guest first", calls(w)) } - if o := pk[1].(map[string]any)["origin"]; o != "Debian-Security" { - t.Fatalf("openssl origin = %v", o) + for _, p := range w.plans { + if p["select"] != "pending-fast" || p["snapshot"] != "" || len(p["packages"].([]any)) != 0 { + t.Fatalf("ring-0 plan = %v", p) + } } - if w.plans[1]["snapshot"] != "" { - t.Fatalf("ring 0 installs from live sources, snapshot = %v", w.plans[1]["snapshot"]) + if len(g.NotCovered) != 1 || g.NotCovered[0] != "docker-ce" { + t.Fatalf("not covered = %v", g.NotCovered) } - 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" { + if len(h.reports) != 2 || h.reports[0].Layer != LayerGuest || h.reports[1].Layer != LayerHost { 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) +// Ring 1 installs EXACTLY each layer's own approved release (a guest release is not a host release). +func TestRing1_EachLayerItsOwnRelease(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: map[string]WrapperReport{ + LayerGuest: {Upgraded: []Package{{Name: "libc6", Version: "g-u4"}}, Pending: pend[1:]}, + LayerHost: {Upgraded: []Package{{Name: "openssl", Version: "h-u3"}}}, + }} + gr := &hub.WireOSRelease{ID: "os-g", Snapshot: "20261004T080000Z", Packages: []hub.WireOSPackage{{Name: "libc6", Version: "g-u4", Origin: "Debian"}}} + hr := &hub.WireOSRelease{ID: "os-h", Snapshot: "20261004T090000Z", Packages: []hub.WireOSPackage{{Name: "openssl", Version: "h-u3", Origin: "Debian-Security"}}} + l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 1, Enabled: true, Release: gr, HostRelease: hr}) + g, ho := l.Run(context.Background(), 9201, "night") + if g.ReleaseID != "os-g" || ho.ReleaseID != "os-h" { + t.Fatalf("release ids %q %q", g.ReleaseID, ho.ReleaseID) } - 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) + gp, hp := w.plans[0], w.plans[1] + if gp["snapshot"] != "20261004T080000Z" || gp["packages"].([]any)[0].(map[string]any)["version"] != "g-u4" { + t.Fatalf("guest plan %v", gp) } - // 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) + if hp["layer"] != LayerHost || hp["snapshot"] != "20261004T090000Z" || hp["packages"].([]any)[0].(map[string]any)["name"] != "openssl" { + t.Fatalf("host plan %v", hp) + } + if strings.Join(g.NotCovered, ",") != "openssl,docker-ce" { + t.Fatalf("guest not covered = %v", g.NotCovered) } } -func TestRing1_NoReleaseInstallsNothing(t *testing.T) { +func TestRing1_NoReleaseIsInventory(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)) + g, ho := l.Run(context.Background(), 9201, "night") + if g.Outcome != "nothing" || ho.Outcome != "nothing" || calls(w) != "guest:inventory,host:inventory" { + t.Fatalf("g=%+v h=%+v calls=%s", g, ho, calls(w)) } } @@ -160,31 +192,53 @@ func TestRing1_NoReleaseInstallsNothing(t *testing.T) { 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)) + if g, _ := l.Run(context.Background(), 9201, "night"); g.Outcome != "nothing" || g.Ring != 1 { + t.Fatalf("g=%+v calls=%s", g, calls(w)) } } -// Switched OFF: the box reports but installs nothing (the brief, Part D 2). +// Switched OFF: the box reports but installs nothing, on both layers. 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)) + g, ho := l.Run(context.Background(), 9201, "night") + if g.Outcome != "inventory" || ho.Outcome != "inventory" || calls(w) != "guest:inventory,host:inventory" || len(h.reports) != 2 { + t.Fatalf("g=%+v h=%+v calls=%s", g, ho, calls(w)) } } -// Unhealthy after the run, and still unhealthy at the end of the wait → health_failed (the hub mails the operator). +// Not an appliance (BYO, `11` §1): the host step never runs — no host plan at all. +func TestBYO_NoHostPlan(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: map[string]WrapperReport{LayerGuest: {Upgraded: []Package{{Name: "libc6"}}}}} + l, h := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + l.Appliance = false + _, ho := l.Run(context.Background(), 9201, "night") + if ho.Outcome != "" || calls(w) != "guest:apply" || len(h.reports) != 1 { + t.Fatalf("a BYO box got a host step: host=%+v calls=%s", ho, calls(w)) + } +} + +// A failed guest step skips the host step that night. +func TestGuestFailure_SkipsTheHost(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: map[string]WrapperReport{ + LayerGuest: {Refused: json.RawMessage(`{"code":"R6","reason":"x"}`)}}} + l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + g, ho := l.Run(context.Background(), 9201, "night") + if g.Outcome != "refused" || ho.Outcome != "" || calls(w) != "guest:apply" { + t.Fatalf("g=%+v h=%+v calls=%s", g, ho, calls(w)) + } +} + +// Unhealthy after the run, and still unhealthy at the end of the wait → health_failed; the host step is skipped. 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}} + w := &fakeWrapper{t: t, pending: pend, applyRep: map[string]WrapperReport{LayerGuest: {Upgraded: []Package{{Name: "libc6"}}, HealthAfter: bad}}, + healthSeq: map[string][]*Health{LayerGuest: {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) + g, ho := l.Run(context.Background(), 9201, "night") + if g.Outcome != "health_failed" || g.Healthy || !strings.Contains(g.HealthReason, "app was running") || ho.Outcome != "" { + t.Fatalf("g=%+v h=%+v", g, ho) } if h.reports[0].Outcome != "health_failed" { t.Fatal("the hub was not told") @@ -194,11 +248,22 @@ func TestHealth_FailsAfterTheWait(t *testing.T) { // 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}} + w := &fakeWrapper{t: t, pending: pend, applyRep: map[string]WrapperReport{LayerGuest: {Upgraded: []Package{{Name: "libc6"}}, HealthAfter: starting}}, + healthSeq: map[string][]*Health{LayerGuest: {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) + if g, _ := l.Run(context.Background(), 9201, "night"); g.Outcome != "applied" || !g.Healthy { + t.Fatalf("g = %+v", g) + } +} + +// The host step judged unhealthy when the tunnel is down after it. +func TestHost_TunnelDownFailsTheHostStep(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: map[string]WrapperReport{LayerHost: {Upgraded: []Package{{Name: "openssl"}}}}} + l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + l.Tunnel = fakeTunnel{hub.TunnelNotRunning} + _, ho := l.Run(context.Background(), 9201, "night") + if ho.Outcome != "health_failed" || !strings.Contains(ho.HealthReason, "tunnel") { + t.Fatalf("host = %+v", ho) } } @@ -227,24 +292,49 @@ func TestHealthVerdict(t *testing.T) { } } -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) +// The host rule: every listed daemon active, the guest running and passing its own rule, the tunnel running. +// Red-proofs: drop any one check and its case fails. +func TestHostHealthVerdict(t *testing.T) { + no := false + svcDown := hostOK() + svcDown.HostServices["pveproxy"] = "failed" + guestDown := hostOK() + guestDown.GuestRunning = &no + guestApp := hostOK() + guestApp.Guest = &Health{DockerOK: true, NetworkOK: true, Controller: "healthy", Containers: map[string]Container{"felhom-controller": {State: "running", Health: "healthy"}}} + cases := []struct { + name string + after *Health + tunnel string + want bool + why string + }{ + {"all good", hostOK(), hub.TunnelRunning, true, ""}, + {"a daemon down", svcDown, hub.TunnelRunning, false, "pveproxy"}, + {"the guest stopped", guestDown, hub.TunnelRunning, false, "guest is not running"}, + {"an app in the guest gone", guestApp, hub.TunnelRunning, false, "app was running"}, + {"the tunnel down", hostOK(), hub.TunnelNotRunning, false, "tunnel"}, + {"the tunnel unknown", hostOK(), hub.TunnelUnknown, false, "tunnel"}, + {"no services read", &Health{GuestRunning: yes(), Guest: guestOK()}, hub.TunnelRunning, false, "no host service"}, } - if rep := l.Run(context.Background(), 9201, "debug"); rep.Outcome == "skipped" { - t.Fatal("the debug action must not be throttled") + for _, c := range cases { + got, why := HostHealthVerdict(hostOK(), c.after, c.tunnel) + if got != c.want || !strings.Contains(why, c.why) { + t.Errorf("%s: got %v (%s), want %v (…%s…)", c.name, got, why, c.want, c.why) + } } } -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) +func TestOncePerNight(t *testing.T) { + w := &fakeWrapper{t: t, pending: pend, applyRep: map[string]WrapperReport{LayerGuest: {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 g, _ := l.Run(context.Background(), 9201, "night"); g.Outcome != "skipped" || len(w.plans) != n { + t.Fatalf("a second night run in the same night ran: %+v", g) + } + if g, _ := l.Run(context.Background(), 9201, "debug"); g.Outcome == "skipped" { + t.Fatal("the debug action must not be throttled") } } @@ -254,7 +344,7 @@ func TestWrapperSuite(t *testing.T) { if err != nil { t.Skip("python3 not available") } - cmd := exec.Command(py, "../../configs/test_felhom_os_apply.py") + cmd := exec.Command(py, "-B", "../../configs/test_felhom_os_apply.py") out, err := cmd.CombinedOutput() if err != nil { t.Fatalf("wrapper suite failed: %v\n%s", err, out) @@ -263,19 +353,3 @@ func TestWrapperSuite(t *testing.T) { t.Fatalf("wrapper suite did not report OK:\n%s", out) } } - -// An app that stops BETWEEN the start of the leg and the apply's own "before" reading still fails the run. -// Measured live 2026-10-04 on demo-hp (privatebin stopped 1 s after the apply plan was written: the old rule -// passed). Red-proof: use ap.HealthBefore alone as the baseline and this fails. -func TestHealth_BaselineIsTheStartOfTheLeg(t *testing.T) { - stoppedEarly := &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"}}, - HealthBefore: stoppedEarly, HealthAfter: stoppedEarly}, - healthSeq: []*Health{stoppedEarly, stoppedEarly, stoppedEarly, stoppedEarly, stoppedEarly, stoppedEarly, stoppedEarly}} - l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) - rep := l.Run(context.Background(), 9201, "night") - if rep.Outcome != "health_failed" || !strings.Contains(rep.HealthReason, "app was running") { - t.Fatalf("an app that stopped during the run passed: %+v", rep) - } -}