#!/usr/bin/python3
# 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-<id>.json
# Nothing else on the command line is accepted. Python 3, standard library only (a JSON plan cannot be parsed
# safely in sh). Tests: configs/test_felhom_os_apply.py (a fake runner; nothing real is executed).
#
# THE TRUST MODEL. The plan is written by the agent, so a broken-into agent writes whatever plan it likes. The
# protection is therefore what this file REFUSES, not where the plan came from: no removal, no downgrade, no new
# package, no package outside the plan, only Debian origin in the fast lane, 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.
#
# LAYERS (agent v0.141.0): "guest" (the customer LXC, entered with `pct exec`) and "host" (this Proxmox host, run
# directly). LANE: "fast" for those two. Agent v0.142.0 adds the layer "docker" (the guest's Docker engine set, `11`
# §5.8), which is the SLOW lane: lane "slow" only, the six Docker packages only, origin "Docker CE" only, and only
# with an authority this file checks ITSELF (R3): a signed operator job verified with `ssh-keygen -Y verify` against
# the ROOT-OWNED signers file (TRUST_SIGNERS), bound to this host (TRUST_FILE host_id), unexpired and never replayed;
# or, for an unsigned ring-0 step, the root-owned TRUST_FILE saying `"ring0_slow_lane": true` (set by hand on the demo
# boxes only). The agent's own config is NOT trusted for either: the agent can write it. A Docker step also needs
# `live-restore` ON (R15) — without it every container restarts.
#
# Modes (plan field "mode"):
#   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).
#   facts       (v0.142.0) read-only versions for the hub's System page: host Debian, running and next-boot kernel,
#               held packages, kernel taint, the crash guard; guest Debian, Docker engine, containerd, live-restore.
#   live-restore-on  (v0.142.0, layer guest) the ONE-TIME step of `09` decision 87: merge `"live-restore": true`
#               into the guest's /etc/docker/daemon.json and `systemctl reload docker`. NEVER a restart (R-835).
# Output: log lines on stderr and the journal (tag felhom-os-apply); the LAST stdout line is
#   OSAPPLY-REPORT <one JSON object>
# 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 calendar
import json
import os
import re
import stat
import subprocess
import sys
import time

PLAN_DIR = "/var/lib/felhom-agent/os"
PLAN_RE = re.compile(r"^plan-[A-Za-z0-9._-]{1,80}\.json$")
AGENT_USER = "felhom-agent"
FAST_ORIGINS = ("Debian", "Debian-Security")
# Debian package name and version grammar (Debian policy §5.6.1, §5.6.12).
NAME_RE = re.compile(r"^[a-z0-9][a-z0-9+.-]+$")
VERSION_RE = re.compile(r"^(?:[0-9]+:)?[0-9][A-Za-z0-9.+~-]*$")
SNAP_RE = re.compile(r"^[0-9]{8}T[0-9]{6}Z$")
RESERVED_VMIDS = set(range(990000, 990010)) | {9999}
DRIVES_PARENT = "/mnt/felhom-drives"
SNAPSHOT_LIST = "/etc/apt/sources.list.d/felhom-os-snapshot.list"
APT_ENV = ["env", "DEBIAN_FRONTEND=noninteractive", "APT_LISTCHANGES_FRONTEND=none", "NEEDRESTART_MODE=l", "LC_ALL=C"]
DPKG_OPTS = ["-o", "Dpkg::Options::=--force-confold", "-o", "Dpkg::Options::=--force-confdef"]
MIN_FREE = 500 * 1024 * 1024
# 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)")
# restart_needed() leaves out processes whose cgroup line matches (grep basic regex). Host: the LXC guests' own
# processes (`0::/lxc/<vmid>/...`) -- NOT lxc-start itself, whose cgroup is `0::/lxc.monitor/<vmid>` (measured
# 2026-10-04 on demo-felhom: the old pattern "lxc" hid lxc-start with 20 deleted maps, so "reboot needed" stayed false
# after a libc6 update). Pinned by test_restart_skip_patterns_against_real_cgroups.
RESTART_SKIP_CGROUP = {"guest": "docker", "host": ":/lxc/"}
HOST_SERVICES = ["pveproxy", "pvedaemon", "pvestatd", "pve-cluster", "felhom-agent"]
# The Docker engine set (`11` §5.8): the only names the docker layer may touch, from the only origin it may use.
DOCKER_NAMES = ("containerd.io", "docker-buildx-plugin", "docker-ce", "docker-ce-cli", "docker-ce-rootless-extras",
                "docker-compose-plugin")
DOCKER_ORIGIN = "Docker CE"
# ROOT-OWNED trust anchors (the installer writes them; the demo boxes got them by hand, R-840). Never the agent's config.
TRUST_FILE = "/etc/felhom/os-trust.json"          # {"host_id": "...", "ring0_slow_lane": false}
TRUST_SIGNERS = "/etc/felhom/operator-signers"    # ssh allowed_signers: <key_id> namespaces="felhom-op-v1" <key>
SIG_NAMESPACE = "felhom-op-v1"
SIGNED_OP = "os_docker_step"
NONCE_FILE = "/var/lib/felhom-os-apply/nonces.json"
DAEMON_JSON = "/etc/docker/daemon.json"
CRASH_GUARD_STATE = "/var/lib/felhom-crash-guard/state.json"


class Refused(Exception):
    def __init__(self, code, reason):
        super().__init__(f"{code} {reason}")
        self.code, self.reason = code, reason


class Runner:
    """Runs commands for real. Tests replace it with a fake. `guest` runs inside the container via pct exec."""

    def host(self, argv, timeout=600, 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):
        return self.host(["/usr/sbin/pct", "exec", str(vmid), "--"] + argv, timeout)

    def read_file(self, path):
        with open(path) as f:
            return f.read()

    def stat(self, path):
        return os.lstat(path)

    def agent_uid(self):
        import pwd
        return pwd.getpwnam(AGENT_USER).pw_uid

    def 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 now(self):
        return time.time()

    def sleep(self, s):
        time.sleep(s)

    def verify_sig(self, signers, key_id, namespace, blob, sig):
        """`ssh-keygen -Y verify` over the EXACT signed bytes. Files in a root-only temp dir; nothing via a shell."""
        import tempfile
        with tempfile.TemporaryDirectory(prefix="felhom-os-apply-") as d:
            sp = os.path.join(d, "sig")
            with open(sp, "w") as f:
                f.write(sig)
            p = subprocess.run(["ssh-keygen", "-Y", "verify", "-f", signers, "-I", key_id, "-n", namespace, "-s", sp],
                               input=blob, capture_output=True, timeout=30)
            return p.returncode

    def read_nonces(self):
        try:
            with open(NONCE_FILE) as f:
                d = json.load(f)
            return d if isinstance(d, dict) else {}
        except (OSError, ValueError):
            return {}

    def write_nonces(self, d):
        os.makedirs(os.path.dirname(NONCE_FILE), mode=0o700, exist_ok=True)
        tmp = NONCE_FILE + ".tmp"
        with open(tmp, "w") as f:
            json.dump(d, f)
        os.replace(tmp, NONCE_FILE)

    def log(self, line):
        print(line, file=sys.stderr, flush=True)
        try:
            subprocess.run(["logger", "-t", "felhom-os-apply", line], timeout=10)
        except Exception:
            pass


class Apply:
    def __init__(self, runner, plan_path):
        self.r = runner
        self.plan_path = plan_path
        self.report = {"refused": None, "mode": None}

    # ---------- checks ----------
    def load_plan(self):
        p = self.plan_path
        d, base = os.path.dirname(p), os.path.basename(p)
        if d != PLAN_DIR or not PLAN_RE.match(base) or ".." in p:
            raise Refused("R1", f"the plan must be {PLAN_DIR}/plan-<id>.json, got {p!r}")
        try:
            st = self.r.stat(p)
        except OSError as e:
            raise Refused("R1", f"cannot stat the plan: {e}")
        if not stat.S_ISREG(st.st_mode):
            raise Refused("R1", "the plan is not a regular file (a symlink or a device is refused)")
        if st.st_uid != self.r.agent_uid():
            raise Refused("R1", f"the plan is not owned by {AGENT_USER}")
        if st.st_size > 2 * 1024 * 1024:
            raise Refused("R1", "the plan is larger than 2 MB")
        try:
            plan = json.loads(self.r.read_file(p))
        except (OSError, ValueError) as e:
            raise Refused("R1", f"the plan is not valid JSON: {e}")
        if not isinstance(plan, dict):
            raise Refused("R1", "the plan is not a JSON object")
        return plan

    def check_plan(self, plan):
        mode = plan.get("mode", "apply")
        if mode not in ("apply", "inventory", "health", "facts", "live-restore-on"):
            raise Refused("R11", f"unknown mode {mode!r}")
        layer = plan.get("layer")
        if layer not in ("guest", "host", "docker"):
            raise Refused("R12", f"layer {layer!r} is not guest, host or docker")
        lane = plan.get("lane", "fast")
        if layer == "docker" and lane != "slow":
            raise Refused("R3", "the Docker engine is the slow lane (`11` §5.8); a fast-lane Docker plan is refused")
        if layer != "docker" and lane != "fast":
            raise Refused("R3", f"the {layer} layer has no slow lane in this release (kernel, Proxmox: `11` §8 step 6)")
        if mode == "facts" and layer != "host":
            raise Refused("R11", "facts is a host-layer mode (it reads the host and the guest)")
        if mode == "live-restore-on" and layer != "guest":
            raise Refused("R11", "live-restore-on is a guest-layer mode")
        if plan.get("undo") and layer != "docker":
            raise Refused("R5", "an undo (downgrade) exists only for the Docker layer, inside a signed job")
        vmid = plan.get("vmid")
        if not isinstance(vmid, int) or isinstance(vmid, bool) or vmid <= 0:
            raise Refused("R11", f"vmid must be a positive integer, got {vmid!r}")
        rid = plan.get("release_id", "")
        if not isinstance(rid, str) or not re.match(r"^[A-Za-z0-9._:-]{1,80}$", rid):
            raise Refused("R11", f"release_id {rid!r} is not a plain id")
        if plan.get("allow_new"):
            raise Refused("R6", "allow_new is a slow-lane field; the fast lane never adds a package")
        select = plan.get("select", "listed")
        if select not in ("listed", "pending-fast", "pending-docker"):
            raise Refused("R11", f"unknown select {select!r}")
        if (select == "pending-docker") != (layer == "docker" and select != "listed"):
            if select == "pending-docker" or layer == "docker":
                raise Refused("R11", f"select {select!r} does not fit layer {layer!r}")
        pk = plan.get("packages", [])
        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 in ("pending-fast", "pending-docker") and pk:
            raise Refused("R11", f"select {select} takes no package list")
        seen = set()
        for e in pk:
            if not isinstance(e, dict):
                raise Refused("R11", "every package entry must be an object")
            n, v, o = e.get("name"), e.get("version"), e.get("origin")
            if not isinstance(n, str) or not NAME_RE.match(n):
                raise Refused("R11", f"package name {n!r} is not a Debian package name")
            if not isinstance(v, str) or not VERSION_RE.match(v):
                raise Refused("R11", f"version {v!r} of {n} is not a Debian version string")
            if n in seen:
                raise Refused("R11", f"package {n} is named twice")
            seen.add(n)
            if layer == "docker":
                if n not in DOCKER_NAMES or o != DOCKER_ORIGIN:
                    raise Refused("R2", f"{n} ({o!r}) is not one of the six Docker packages from {DOCKER_ORIGIN!r}")
                continue
            if n in DOCKER_NAMES:
                raise Refused("R2", f"{n} is a Docker package — the slow lane (`11` §5.8), never in a {layer} plan")
            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, 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 load_trust(self):
        """The ROOT-OWNED trust record. Absent or agent-writable → no slow-lane authority at all (R3)."""
        try:
            st = self.r.stat(TRUST_FILE)
        except OSError:
            raise Refused("R3", f"no {TRUST_FILE} — this box has no slow-lane trust anchor")
        if st.st_uid != 0 or (st.st_mode & 0o022):
            raise Refused("R3", f"{TRUST_FILE} is not root-owned and root-only-writable — it proves nothing")
        try:
            t = json.loads(self.r.read_file(TRUST_FILE))
        except (OSError, ValueError):
            raise Refused("R3", f"{TRUST_FILE} is unreadable")
        if not isinstance(t, dict) or not isinstance(t.get("host_id"), str) or not t["host_id"]:
            raise Refused("R3", f"{TRUST_FILE} names no host_id")
        return t

    def verify_signed(self, signed, trust):
        """R3: an operator-signed os_docker_step, checked HERE (not by the agent): signature against the root-owned
        signers file, op, host binding, time window, and a root-owned nonce record (no replay)."""
        import base64
        if not isinstance(signed, dict) or not isinstance(signed.get("blob_b64"), str) or not isinstance(signed.get("sig"), str):
            raise Refused("R3", "the signed job is malformed")
        try:
            blob = base64.b64decode(signed["blob_b64"], validate=True)
            op = json.loads(blob)
        except (ValueError, TypeError):
            raise Refused("R3", "the signed blob is not base64 JSON")
        key_id = op.get("key_id", "")
        if not isinstance(key_id, str) or not re.match(r"^[A-Za-z0-9._-]{1,64}$", key_id):
            raise Refused("R3", "the signed blob names no plain key_id")
        try:
            st = self.r.stat(TRUST_SIGNERS)
        except OSError:
            raise Refused("R3", f"no {TRUST_SIGNERS} — no operator key to check a signed job against")
        if st.st_uid != 0 or (st.st_mode & 0o022):
            raise Refused("R3", f"{TRUST_SIGNERS} is not root-owned and root-only-writable")
        rc = self.r.verify_sig(TRUST_SIGNERS, key_id, SIG_NAMESPACE, blob, signed["sig"])
        if rc != 0:
            raise Refused("R3", f"the operator signature does not verify (ssh-keygen rc={rc})")
        if op.get("op") != SIGNED_OP:
            raise Refused("R3", f"the signed op is {op.get('op')!r}, not {SIGNED_OP}")
        if (op.get("target") or {}).get("host_id") != trust["host_id"]:
            raise Refused("R3", "the signed job is for another host")
        now = self.r.now()
        try:
            exp = calendar.timegm(time.strptime(op["expires_at"], "%Y-%m-%dT%H:%M:%SZ"))
            iss = calendar.timegm(time.strptime(op["issued_at"], "%Y-%m-%dT%H:%M:%SZ"))
        except (KeyError, ValueError, TypeError):
            raise Refused("R3", "the signed job has no readable time window")
        if now > exp or now < iss - 120:
            raise Refused("R3", "the signed job is expired or not yet valid")
        nonce = op.get("nonce")
        if not isinstance(nonce, str) or not nonce:
            raise Refused("R3", "the signed job has no nonce")
        seen = self.r.read_nonces()
        if nonce in seen:
            raise Refused("R3", "the signed job was already used (replay)")
        seen[nonce] = exp
        self.r.write_nonces({k: v for k, v in seen.items() if v > now})
        return op.get("params") or {}

    def docker_authority(self, plan):
        """R3 for the docker layer: returns (who, undo). A signed job binds the EXACT package list and the undo flag."""
        trust = self.load_trust()
        signed = plan.get("signed")
        if signed:
            params = self.verify_signed(signed, trust)
            want = sorted(f"{e.get('name')}={e.get('version')}" for e in params.get("packages") or [])
            got = sorted(f"{e['name']}={e['version']}" for e in plan.get("packages", []))
            if not want or want != got:
                raise Refused("R3", "the plan's packages are not exactly the signed job's packages")
            if bool(params.get("undo")) != bool(plan.get("undo")):
                raise Refused("R3", "the plan's undo flag is not the signed job's")
            if params.get("vmid") not in (None, self.vmid):
                raise Refused("R3", "the signed job names another guest")
            return "signed", bool(plan.get("undo"))
        if plan.get("undo"):
            raise Refused("R3", "an undo (downgrade) needs a signed operator job")
        if trust.get("ring0_slow_lane") is True:
            return "ring0", False
        raise Refused("R3", "a Docker step needs a signed operator job (ring 1) or this box's root-owned ring-0 mark")

    def live_restore(self):
        rc, out, _ = self.g(["docker", "info", "--format", "{{.LiveRestoreEnabled}}"], timeout=60)
        return out.strip() if rc == 0 and out.strip() in ("true", "false") else "unknown"

    def container_ids(self):
        rc, out, _ = self.g(["docker", "ps", "-q", "--no-trunc"], timeout=60)
        return sorted(out.split()) if rc == 0 else None

    def live_restore_on(self):
        """`09` decision 87: merge live-restore into daemon.json and RELOAD (C5: a reload turns it on, no restart)."""
        log = self.r.log
        if self.live_restore() == "true":
            self.report["live_restore"] = {"result": "already on"}
            log("os-apply: LIVE-RESTORE already on")
            return 0
        rc, cur, _ = self.g(["cat", DAEMON_JSON], timeout=30)
        try:
            conf = json.loads(cur) if rc == 0 and cur.strip() else {}
        except ValueError:
            raise Refused("R16", f"{DAEMON_JSON} in the guest is not valid JSON — not touched")
        if not isinstance(conf, dict):
            raise Refused("R16", f"{DAEMON_JSON} is not a JSON object — not touched")
        before = self.container_ids()
        conf["live-restore"] = True
        self.r.write_file("guest", self.vmid, DAEMON_JSON, json.dumps(conf, indent=2, sort_keys=True) + "\n")
        rrc, _, rerr = self.g(["systemctl", "reload", "docker"], timeout=120)
        state = "unknown"
        for _ in range(10):
            state = self.live_restore()
            if state == "true":
                break
            self.r.sleep(1)
        after = self.container_ids()
        same = before is not None and before == after
        self.report["live_restore"] = {"result": "on" if state == "true" else "failed", "reload_rc": rrc,
                                       "containers_before": len(before or []), "same_ids": same}
        log(f"os-apply: LIVE-RESTORE reload_rc={rrc} state={state} containers={len(before or [])} same-ids={'yes' if same else 'NO'}")
        if state != "true":
            # put the old file back (and reload again) — still never a restart
            self.r.write_file("guest", self.vmid, DAEMON_JSON, cur if rc == 0 else "{}\n")
            self.g(["systemctl", "reload", "docker"], timeout=120)
            self.report["failed"] = {"rc": 3, "step": "live-restore", "reason": (rerr or "").strip()[-200:]}
            return 3
        return 0

    def kernel_next_boot(self):
        """Which kernel GRUB boots next, read without root-only files (grubenv + /etc/default/grub + /boot)."""
        try:
            dflt = re.search(r'^GRUB_DEFAULT=["\']?([^"\'\n]*)', self.r.read_file("/etc/default/grub"), re.M)
            dflt = dflt.group(1) if dflt else "0"
        except OSError:
            dflt = "0"
        env = {}
        try:
            for l in self.r.read_file("/boot/grub/grubenv").splitlines():
                if "=" in l and not l.startswith("#"):
                    k, v = l.split("=", 1)
                    env[k] = v
        except OSError:
            pass

        def ver(entry):
            m = re.search(r"gnulinux-([0-9][^>\s]*?-pve)-(?:advanced|recovery)", entry)
            return m.group(1) if m else "unknown"
        if env.get("next_entry"):
            return ver(env["next_entry"]), "next_entry (a one-shot GRUB cannot clear on LVM /boot)"
        if dflt == "saved":
            return (ver(env["saved_entry"]), "saved default") if env.get("saved_entry") else ("unknown", "saved default unset")
        if dflt == "0":
            rc, out, _ = self.r.host(["sh", "-c", "ls /boot/vmlinuz-* 2>/dev/null"], 30)
            vers = [l.split("vmlinuz-", 1)[1] for l in out.split() if "vmlinuz-" in l]
            best = None
            for v in vers:
                if best is None or self.dpkg_cmp(v, "gt", best):
                    best = v
            return (best or "unknown"), "GRUB_DEFAULT=0 (the newest installed)"
        return "unknown", f"GRUB_DEFAULT={dflt}"

    def facts(self):
        """Read-only versions for the hub's System page (R-852). A value that cannot be read is "unknown"."""
        def first(cmd, timeout=30):
            rc, out, _ = self.r.host(cmd, timeout)
            v = out.strip().splitlines()[0].strip() if rc == 0 and out.strip() else ""
            return v or "unknown"
        h = {"debian": first(["cat", "/etc/debian_version"]), "kernel_running": first(["uname", "-r"])}
        h["kernel_next_boot"], h["kernel_next_boot_source"] = self.kernel_next_boot()
        rc, out, _ = self.r.host(["apt-mark", "showhold"], 60)
        h["held"] = sorted(out.split()) if rc == 0 else None
        try:
            t = int(self.r.read_file("/proc/sys/kernel/tainted").strip())
            h["tainted"], h["oops_this_boot"], h["warn_this_boot"] = t, bool(t & 128), bool(t & 512)
        except (OSError, ValueError):
            h["tainted"], h["oops_this_boot"], h["warn_this_boot"] = None, None, None
        try:
            h["kernel_panic"] = int(self.r.read_file("/proc/sys/kernel/panic").strip())
        except (OSError, ValueError):
            h["kernel_panic"] = None
        try:
            h["crash_guard"] = json.loads(self.r.read_file(CRASH_GUARD_STATE))
        except (OSError, ValueError):
            h["crash_guard"] = None
        g = {"debian": "unknown", "docker_engine": "unknown", "containerd": "unknown", "live_restore": "unknown"}
        try:
            self.check_guest(self.vmid)
            running = True
        except Refused as e:
            running, g["unknown_reason"] = False, f"{e.code} {e.reason}"
        if running:
            script = ('echo "debian=$(cat /etc/debian_version 2>/dev/null)"; '
                      'echo "engine=$(docker version --format \'{{.Server.Version}}\' 2>/dev/null)"; '
                      'echo "containerd=$(dpkg-query -W -f \'${Version}\' containerd.io 2>/dev/null)"; '
                      'echo "live=$(docker info --format \'{{.LiveRestoreEnabled}}\' 2>/dev/null)"')
            rc, out, _ = self.g(["sh", "-c", script], timeout=60)
            kv = dict(l.split("=", 1) for l in out.splitlines() if "=" in l)
            for k, src in (("debian", "debian"), ("docker_engine", "engine"), ("containerd", "containerd")):
                g[k] = kv.get(src, "").strip() or "unknown"
            g["live_restore"] = {"true": "on", "false": "off"}.get(kv.get("live", "").strip(), "unknown")
        self.report["facts"] = {"host": h, "guest": g}
        return 0

    def check_guest(self, vmid):
        if vmid in RESERVED_VMIDS:
            raise Refused("R10", f"vmid {vmid} is a reserved scratch vmid")
        try:
            conf = self.r.read_file(f"/etc/pve/lxc/{vmid}.conf")
        except OSError:
            raise Refused("R10", f"vmid {vmid} is not a container on this host")
        cur = conf.split("\n[", 1)[0]  # the current config, not a snapshot section
        binds = [l for l in cur.splitlines() if re.match(r"^mp[0-9]+: " + re.escape(DRIVES_PARENT) + r",", l)]
        if not binds:
            raise Refused("R10", f"vmid {vmid} does not bind {DRIVES_PARENT} — it is not this box's customer guest")
        lock = [l for l in cur.splitlines() if l.startswith("lock:")]
        if lock:
            raise Refused("R9", f"vmid {vmid} is locked ({lock[0].split(':', 1)[1].strip()}) — a backup or restore is running")
        rc, out, _ = self.r.host(["/usr/sbin/pct", "status", str(vmid)])
        if rc != 0 or "running" not in out:
            raise Refused("R10", f"vmid {vmid} is not running")

    # ---------- 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)  # guest and docker both live in the customer guest

    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.x(["dpkg-query", "-W", "-f", "${Package}\t${Version}\t${db:Status-Abbrev}\n"])
        res = {}
        for l in out.splitlines():
            parts = l.split("\t")
            if len(parts) == 3 and parts[2].startswith("ii"):
                res[parts[0]] = parts[1]
        return res

    def 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] in res:
                res[f[0]].add(f[1])
        return res

    def simulate(self, 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)
            if m:
                inst.append({"name": m.group(1), "from": m.group(2), "to": m.group(3), "origin": m.group(4)})
            m = re.match(r"^Remv (\S+)", l)
            if m:
                remv.append(m.group(1))
        return rc, inst, remv, out + err

    @staticmethod
    def origin_name(origin):
        # "Debian:13.7/stable, Debian-Security:13/stable-security" -> {"Debian", "Debian-Security"}
        return {o.strip().split(":")[0] for o in origin.split(",") if o.strip()}

    def free_bytes(self):
        rc, out, _ = self.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.x(["fuser", "/var/lib/dpkg/lock-frontend", "/var/lib/dpkg/lock"])
        return rc == 0 and out.strip() != ""

    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", "--no-trunc", "--format", "{{.Names}}\t{{.State}}\t{{.Status}}\t{{.ID}}"], timeout=60)
        cont = {}
        for l in out.splitlines():
            p = l.split("\t")
            if len(p) >= 3:
                h = "healthy" if "(healthy)" in p[2] else "unhealthy" if "(unhealthy)" in p[2] else \
                    "starting" if "(health: starting)" in p[2] else "none"
                cont[p[0]] = {"state": p[1], "health": h}
                if len(p) >= 4 and p[3]:
                    cont[p[0]]["id"] = p[3]
        nrc, _, _ = self.g(["getent", "hosts", "deb.debian.org"], timeout=30)
        return {"docker_ok": rc == 0, "containers": cont,
                "controller": cont.get("felhom-controller", {}).get("health", "absent"),
                "network_ok": nrc == 0}

    def health(self):
        if self.layer in ("guest", "docker"):
            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 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 = RESTART_SKIP_CGROUP["host" if self.layer == "host" else "guest"]
        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 = {}
        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(" "):
                    cur, star = l.rstrip(":"), False
                    continue
                s = l.strip()
                if s.startswith("*** "):
                    star = True
                    continue
                if star and cur and re.match(r"^[0-9-]+ ", s):
                    if "/var/lib/dpkg/status" in s:
                        origins.setdefault(cur, "local")
                    else:
                        origins[cur] = s
                    continue
                if star and not re.match(r"^[0-9-]+ ", s):
                    star = False

        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:
                return "Docker"
            if "debian" in src:
                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"],
                         "origin": sorted(self.origin_name(p["origin"]))} for p in pend],
        }

    # ---------- the run ----------
    def run(self):
        plan = self.load_plan()
        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.mode == "facts":
            return self.facts()
        if self.layer == "host":
            self.check_appliance()
        self.check_guest(self.vmid)
        log = self.r.log
        if self.mode == "live-restore-on":
            return self.live_restore_on()
        self.who, self.allow_downgrade = ("fast", False)
        if self.layer == "docker" and self.mode == "apply":
            self.who, self.allow_downgrade = self.docker_authority(plan)
            if self.live_restore() != "true":
                raise Refused("R15", "live-restore is not ON in the guest — a Docker step would restart every container")
            self.report["authority"] = self.who
            self.report["undo"] = self.allow_downgrade
        if self.mode == "health":
            self.report["health"] = self.health()
            return 0
        log(f"os-apply: START release={plan.get('release_id')} layer={self.layer}" +
            (f":{self.vmid}" if self.layer != "host" else "") +
            f" lane={plan.get('lane', 'fast')} mode={self.mode} select={self.select} packages={len(plan.get('packages', []))}" +
            (f" authority={self.who}{' UNDO' if self.allow_downgrade else ''}" if self.layer == "docker" else ""))
        if self.apt_lock_held():
            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.x(APT_ENV + ["apt-get", "-q", "update"], timeout=600)
        if rc != 0:
            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, installed_after = self.apply(plan)
            if rc:
                return rc
        self.report.update(self.inventory(installed_after))
        if "reboot_needed" not in self.report:
            # EVERY layer is scanned on EVERY pass: a reboot (host) or a restart (guest) must CLEAR "restart needed",
            # or the fleet view keeps a stale date (R-849, v0.142.0; the host since v0.141.1). One pct exec, ~1 s.
            # Pinned by test_every_layer_scans_every_pass.
            self.report["restart_needed"], self.report["reboot_needed"] = self.restart_needed()
            self.report["docker_restart_needed"] = any(p in ("dockerd", "containerd") for p in self.report["restart_needed"])
        if self.layer == "docker":
            rc_v, out_v, _ = self.g(["docker", "version", "--format", "{{.Server.Version}}"], timeout=60)
            self.report["docker_engine"] = out_v.strip() if rc_v == 0 and out_v.strip() else "unknown"
        self.report["reboot_scanned"] = "reboot_needed" in self.report
        self.report["health_after"] = self.health()
        return 0

    def repair(self):
        rc, before, _ = self.x(["dpkg", "--audit"])
        configured = len([l for l in before.splitlines() if l.startswith(" ")])
        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 pending_docker(self):
        """Ring 0 (select pending-docker): the newest pending version of each INSTALLED Docker package, Docker origin."""
        rc, pend, remv, _ = self.simulate(["dist-upgrade"])
        return [{"name": p["name"], "version": p["to"], "origin": DOCKER_ORIGIN} for p in pend
                if p["from"] is not None and p["name"] in DOCKER_NAMES and self.origin_name(p["origin"]) == {DOCKER_ORIGIN}]

    def origin_ok(self, origin):
        o = self.origin_name(origin)
        if self.layer == "docker":
            return o == {DOCKER_ORIGIN}
        return bool(o & set(FAST_ORIGINS))

    def apply(self, plan):
        log = self.r.log
        if self.select == "listed":
            packages = plan["packages"]
        elif self.select == "pending-docker":
            packages = self.pending_docker()
        else:
            packages = self.pending_fast()
        cmp_op = "ne" if self.allow_downgrade else "gt"
        inst = self.installed()
        upgrade, already, notinst = [], 0, 0
        for e in packages:
            n, v = e["name"], e["version"]
            if n not in inst:
                notinst += 1
                continue
            if not self.dpkg_cmp(v, cmp_op, inst[n]):
                already += 1
                continue
            upgrade.append((n, v))
        from_snap = 0
        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)
            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}")
            from_snap = len(missing)
        try:
            log(f"os-apply: PLAN upgrade={len(upgrade)} already={already} not-installed={notinst} from-snapshot={from_snap}")
            self.report["plan"] = {"upgrade": len(upgrade), "already": already, "not_installed": notinst, "from_snapshot": from_snap}
            if not upgrade:
                self.report["upgraded"] = []
                log("os-apply: DONE rc=0 seconds=0 upgraded=0 (nothing to do)")
                return 0, inst
            args = ["install", "--only-upgrade", "--no-install-recommends"] + \
                (["--allow-downgrades"] if self.allow_downgrade else []) + [f"{n}={v}" for n, v in upgrade]
            rc, sim, remv, text = self.simulate(args)
            if rc != 0:
                tail = text.strip().splitlines()[-1] if text.strip() else ""
                raise Refused("R7", "the simulation failed: " + tail)
            if remv:
                raise Refused("R4", f"the plan would remove {', '.join(remv[:5])}")
            want = dict(upgrade)
            for p in sim:
                if p["from"] is None:
                    raise Refused("R6", f"the plan would add a package that is not installed: {p['name']}")
                if p["name"] not in want:
                    raise Refused("R6", f"the plan would touch {p['name']}, which is not in the plan")
                if p["to"] != want[p["name"]]:
                    raise Refused("R6", f"{p['name']} would go to {p['to']}, not the approved {want[p['name']]}")
                if not self.allow_downgrade and not self.dpkg_cmp(p["to"], "gt", p["from"]):
                    raise Refused("R5", f"{p['name']} would be downgraded {p['from']} -> {p['to']}")
                if not self.origin_ok(p["origin"]):
                    raise Refused("R2", f"{p['name']} would come from {p['origin']}, not the {self.layer} layer's origin")
                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.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
            # keeps the local one; the package's version lands as X.dpkg-dist). Measured live 2026-10-04 (debian_version).
            conflict = None
            for l in (out + err).splitlines():
                m = re.search(r"Installing new version of config file (\S+?)\s*\.\.\.", l)
                if m:
                    log(f"os-apply: CONFFILE updated {m.group(1)} (it was not changed locally)")
                m = re.search(r"Configuration file '([^']+)'", l)
                if m:
                    conflict = m.group(1)
                if conflict and "Keeping old config file" in l:
                    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.x(["apt-get", "clean"])
            if rc != 0:
                _, 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, None
            self.report["upgraded"] = [{"name": n, "version": v} for n, v in upgrade]
            self.report["seconds"] = round(secs, 1)
            procs, reboot = self.restart_needed()
            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.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)
            if m:
                total += int(m.group(1))
        return total

    def add_snapshot_sources(self, snap):
        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 {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.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.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.x(["rm", "-f", SNAPSHOT_LIST])
        self.x(APT_ENV + ["apt-get", "-q", "update"], timeout=600)


def main(argv, runner=None):
    r = runner or Runner()
    if len(argv) != 3 or argv[1] != "--plan":
        r.log("os-apply: REFUSED: R1 usage: felhom-os-apply --plan /var/lib/felhom-agent/os/plan-<id>.json")
        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:
        r.log(f"os-apply: REFUSED: {e.code} {e.reason}")
        a.report["refused"] = {"code": e.code, "reason": e.reason}
        rc = 2
    except subprocess.TimeoutExpired as e:
        r.log(f"os-apply: FAILED rc=124 step=timeout — {e.cmd}")
        a.report["failed"] = {"rc": 124, "timeout": str(e.cmd)[:200]}
        rc = 3
    a.report["pass_seconds"] = round(time.time() - t0, 1)
    print("OSAPPLY-REPORT " + json.dumps(a.report, sort_keys=True))
    return rc


if __name__ == "__main__":
    if os.geteuid() != 0:
        print("felhom-os-apply: must run as root (via sudo)", file=sys.stderr)
        sys.exit(2)
    sys.exit(main(sys.argv))
