From ca78c17b293cb4b43e95767ebf0e638359870bf6 Mon Sep 17 00:00:00 2001 From: kisfenyo Date: Mon, 5 Oct 2026 07:14:32 +0200 Subject: [PATCH] v0.144.0 code: R8 measures the real download (R-865); an OS pass's report survives a killed agent (R-868); the debug pass runs from the saved block when the hub is away (R-866) Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS --- REUSE.md | 2 + cmd/felhom-agent/main.go | 37 +++- cmd/felhom-agent/r866_selftest_block_test.go | 55 ++++++ configs/felhom-os-apply | 49 ++++- configs/test_felhom_os_apply.py | 125 ++++++++++++- internal/osupdate/dockerjob.go | 3 + internal/osupdate/leg.go | 100 +++++++++-- internal/osupdate/leg_test.go | 18 ++ internal/osupdate/unsent.go | 178 +++++++++++++++++++ internal/osupdate/unsent_test.go | 123 +++++++++++++ 10 files changed, 662 insertions(+), 28 deletions(-) create mode 100644 cmd/felhom-agent/r866_selftest_block_test.go create mode 100644 internal/osupdate/unsent.go create mode 100644 internal/osupdate/unsent_test.go diff --git a/REUSE.md b/REUSE.md index 2a54b8e..dc2bb41 100644 --- a/REUSE.md +++ b/REUSE.md @@ -17,6 +17,8 @@ | `stageTemp` | internal/localapi/intermediary.go | `stageTemp(pattern, content) (path, err)` | random-named temp before a root `install` (audit B1) | Fixed /tmp names are a TOCTOU — sudoers globs expect `/tmp/felhom-*-*.ext` | | `BUNDLE_FILES` + `Bundle` (mode `bundle`, `--install-bundle`) | configs/felhom-os-apply | the ONE table of root-owned paths + the installer of them | ANY new root-owned file the installer writes (sudoers line, wrapper, unit) — add it to the table, never a new installer fetch (R-840) | The builder (`scripts/build-config-bundle.py`) and the installer read the same table; `test_every_root_file_the_installer_writes_is_in_the_bundle` fails on a path the bundle lacks. Trust files (`/etc/felhom/os-trust.json`, `operator-signers`) are NEVER bundle paths (R17) | | `osupdate.ConfigUpdateExecutor` | internal/osupdate/bundle.go | signed op `agent_config_update` {agent_version, bundle_sha256} | delivering the bundle to an installed box | a courier only: the root wrapper re-verifies signature, host, nonce and sha itself | +| `osupdate.Leg.SendUnsent` / `lockPass` (v0.144.0, R-868) | internal/osupdate/unsent.go | `(ctx) int` | an OS-pass report the agent never sent (killed mid-pass): the wrapper keeps `report---apply.json` beside the plan; the agent deletes it once the hub has it | any new caller that runs an apply pass must hold `lockPass` (flock, across processes) — the sender must never take a running pass's copy | +| `osupdate.LoadSavedBlock` (v0.144.0, R-866) | internal/osupdate/leg.go | `(planDir) (block, savedAt, ok)` | the hub's newest os_update block as the daemon last received it (`os-update-block.json`) | the debug pass uses it ONLY when the hub cannot be reached, and says so in its header; no saved block → no pass | | `guesthook.InstallSnippet` / `Register` | internal/guesthook/install.go | `InstallSnippet(ctx, runner) error` | pre-start self-heal hook install (C1 net) | Same random-temp+install pattern; snippet delegates to the agent binary (no shell logic). Issues `mkdir -p /var/lib/vz/snippets` FIRST (v0.63.0, B2 — fresh boxes lack the dir; sudoers grants exactly that argv) | ### Disk / format safety (role gates, durable IDs, format guards) diff --git a/cmd/felhom-agent/main.go b/cmd/felhom-agent/main.go index e6e46b3..4021080 100644 --- a/cmd/felhom-agent/main.go +++ b/cmd/felhom-agent/main.go @@ -868,6 +868,13 @@ func runDaemon(cfg config.Config, logger *slog.Logger, logRing *applog.Ring) int logger.Info("felhom-agent daemon starting", "version", version, "host_id", cfg.Hub.HostID, "hub_url", cfg.Hub.URL, "interval_s", hcfg.PollSeconds) // hub key intentionally not logged + // R-868 (v0.144.0): an OS pass whose agent was killed kept its report on disk — send it now (a pass that starts + // first sends it itself; the pass lock keeps the two apart). + go func() { + if n := osLeg.SendUnsent(ctx); n > 0 { + logger.Info("osupdate: sent kept report(s) at start", "count", n) + } + }() // Reconcile (slice 4) runs alongside the hub loop, sharing the per-guest queue // (doc 03 §10). At slice 4 the desired-state provider is empty (no hub serving @@ -3575,15 +3582,15 @@ func runSelftestOSUpdate(ctx context.Context, cfg config.Config, logger *slog.Lo 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) + blk, source, ok := selftestOSBlock(ctx, client, leg.PlanDir) + if !ok { + fmt.Fprintln(os.Stderr, "selftest=os-update:", source) return 1 } - leg.SetBlock(resp.DesiredState.OSUpdate) + leg.SetBlock(blk) b := leg.Block() - 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) + fmt.Printf("=== felhom-agent %s selftest=os-update vmid=%d ring=%d enabled=%v guest-release=%v host-release=%v appliance=%v block=%s ===\n", + version, vmid, b.Ring, b.Enabled, b.Release != nil, b.HostRelease != nil, leg.Appliance, source) start := time.Now() pass := leg.Run(ctx, vmid, "debug") worst := pass.Guest @@ -3609,6 +3616,24 @@ func runSelftestOSUpdate(ctx context.Context, cfg config.Config, logger *slog.Lo return 1 } +// selftestOSBlock is the debug pass's os_update block (R-866, v0.144.0): the hub's, fetched fresh; when the hub cannot +// be reached, the block the daemon last saved — named in the selftest's first line, so a pass with the hub away can +// be exercised by hand (the daemon's own leg already ran from its last block; the selftest stopped). No saved block +// and no hub → not run (ok=false), never a guessed block. +func selftestOSBlock(ctx context.Context, f interface { + FetchDesiredState(context.Context) (*hub.DesiredStateResponse, error) +}, planDir string) (*hub.WireOSUpdate, string, bool) { + resp, err := f.FetchDesiredState(ctx) + if err == nil { + return resp.DesiredState.OSUpdate, "hub", true + } + saved, at, ok := osupdate.LoadSavedBlock(planDir) + if !ok { + return nil, fmt.Sprintf("desired state: %v — and no saved block (the daemon saves one when the hub sends it)", err), false + } + return saved, fmt.Sprintf("SAVED(%s; hub unreachable: %v)", at.UTC().Format(time.RFC3339), err), true +} + // runSelftestFacts prints the versions the host report carries (R-852, agent v0.142.0) — read-only. // // sudo -u felhom-agent felhom-agent --config … --selftest=os-facts -vmid 9201 diff --git a/cmd/felhom-agent/r866_selftest_block_test.go b/cmd/felhom-agent/r866_selftest_block_test.go new file mode 100644 index 0000000..169b3d4 --- /dev/null +++ b/cmd/felhom-agent/r866_selftest_block_test.go @@ -0,0 +1,55 @@ +package main + +import ( + "context" + "errors" + "strings" + "testing" + + "gitea.dooplex.hu/admin/felhom-agent/internal/hub" + "gitea.dooplex.hu/admin/felhom-agent/internal/osupdate" +) + +type r866Fetcher struct{ err error } + +func (f r866Fetcher) FetchDesiredState(context.Context) (*hub.DesiredStateResponse, error) { + if f.err != nil { + return nil, f.err + } + r := &hub.DesiredStateResponse{} + r.DesiredState.OSUpdate = &hub.WireOSUpdate{Ring: 1, Enabled: true} + return r, nil +} + +// R-866 (v0.144.0). THE NIGHT'S SHAPE (A3, Tester 1 box, hub blocked): `selftest=os-update: desired state: hub: +// transport error … connect: invalid argument` — the debug pass could not run at all. Now it runs from the block the +// daemon saved, and its header says so. +// COMPANION RED-PROOF: return at once on a fetch error in selftestOSBlock → "the pass did not run from the saved block". +func TestR866_DebugPassUsesTheSavedBlockWhenTheHubIsAway(t *testing.T) { + dir := t.TempDir() + leg := &osupdate.Leg{PlanDir: dir} + r := &hub.DesiredStateResponse{} + r.DesiredState.OSUpdate = &hub.WireOSUpdate{Ring: 0, Enabled: true} + leg.OnDesiredState(context.Background(), r) // the daemon received a block and saved it + away := r866Fetcher{err: errors.New("hub: transport error: connect: invalid argument")} + b, src, ok := selftestOSBlock(context.Background(), away, dir) + if !ok || b == nil || b.Ring != 0 { + t.Fatalf("the pass did not run from the saved block: ok=%v block=%+v src=%q", ok, b, src) + } + if !strings.HasPrefix(src, "SAVED(") || !strings.Contains(src, "hub unreachable") { + t.Fatalf("the header must say the block is the saved one: %q", src) + } + // the hub reachable: its block wins, and the header says "hub" + b, src, ok = selftestOSBlock(context.Background(), r866Fetcher{}, dir) + if !ok || b.Ring != 1 || src != "hub" { + t.Fatalf("hub block not used: %+v %q", b, src) + } +} + +// No hub and nothing saved: the pass does not run on a guessed block. +func TestR866_NoHubNoSavedBlockDoesNotRun(t *testing.T) { + _, why, ok := selftestOSBlock(context.Background(), r866Fetcher{err: errors.New("down")}, t.TempDir()) + if ok || !strings.Contains(why, "no saved block") { + t.Fatalf("ok=%v why=%q", ok, why) + } +} diff --git a/configs/felhom-os-apply b/configs/felhom-os-apply index 29661a9..1a410c0 100755 --- a/configs/felhom-os-apply +++ b/configs/felhom-os-apply @@ -200,6 +200,32 @@ class Runner: if rc != 0: raise Refused("R7", f"could not write {path} in the guest") + def save_report(self, plan_path, report): + """R-868: keep an apply pass's report on disk until the agent has sent it (the agent deletes it). Written + INTO the agent's own plan dir as root, so: the dir is opened with O_NOFOLLOW and must be a real directory + owned by the agent (a symlink swapped in for it is refused); the file is created O_EXCL|O_NOFOLLOW after + removing an old one, then handed to the agent (0600). Any failure only loses the copy — never the run.""" + base = os.path.basename(plan_path) + name = "report-" + base[len("plan-"):] + dfd = os.open(PLAN_DIR, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + try: + st = os.fstat(dfd) + if not stat.S_ISDIR(st.st_mode) or st.st_uid != self.agent_uid(): + raise OSError(f"{PLAN_DIR} is not the agent's own directory") + try: + os.unlink(name, dir_fd=dfd) + except FileNotFoundError: + pass + fd = os.open(name, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600, dir_fd=dfd) + try: + os.write(fd, (json.dumps(report, sort_keys=True) + "\n").encode()) + os.fchown(fd, self.agent_uid(), -1) + finally: + os.close(fd) + finally: + os.close(dfd) + return os.path.join(PLAN_DIR, name) + def now(self): return time.time() @@ -808,6 +834,16 @@ class Apply: 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) + # R-868 (v0.144.0): the agent's run id, trigger and ring travel in the report, so a report the agent never + # received (it was killed mid-pass) can be sent later from the saved copy. Plain ids only; anything else + # is dropped, never refused (the plan's other checks decide). + rid, trig, ring = plan.get("run_id"), plan.get("trigger"), plan.get("ring") + if isinstance(rid, str) and re.match(r"^[A-Za-z0-9._-]{1,80}$", rid): + self.report["run_id"] = rid + if isinstance(trig, str) and re.match(r"^[a-z0-9_-]{1,20}$", trig): + self.report["trigger"] = trig + if ring in (0, 1) and not isinstance(ring, bool): + self.report["ring"] = ring if self.mode == "facts": return self.facts() if self.mode == "bundle": @@ -1027,7 +1063,11 @@ class Apply: 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) + # R-865 (v0.144.0): NO `-s`. With `-s` apt prints the simulation ("Inst …") and no URI list, so this summed + # 0 B and R8 only ever applied its 500 MB floor (measured 2026-10-04: 0 URIs with -s, 3 URIs without). + # `--print-uris` alone downloads nothing — measured on 9202 2026-10-05: the archive cache and the versions + # unchanged. Pinned by test_R8_counts_the_real_download / test_download_bytes_never_simulates. + rc, out, _ = self.x(APT_ENV + ["apt-get", "-o", "Debug::NoLocking=1", "--print-uris", "-q"] + args) total = 0 for l in out.splitlines(): m = re.match(r"^'[^']+' \S+ ([0-9]+) ", l) @@ -1406,6 +1446,13 @@ def main(argv, runner=None, environ=None): a.report["failed"] = {"rc": 124, "timeout": str(e.cmd)[:200]} rc = 3 a.report["pass_seconds"] = round(time.time() - t0, 1) + if a.report.get("mode") == "apply": + # R-868: BEFORE the stdout line — a killed agent never reads stdout, and this copy is how its report still + # reaches the hub (the agent sends an unsent copy when it starts, and deletes it once sent). + try: + r.save_report(argv[2], a.report) + except Exception as e: + r.log(f"os-apply: the report copy could not be saved (the run is unaffected): {e}") 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 4f2841d..2e96ce1 100644 --- a/configs/test_felhom_os_apply.py +++ b/configs/test_felhom_os_apply.py @@ -73,6 +73,7 @@ class Fake: self.sig_rc = 0 self.nonces = {} self.clock = 1791115200.0 # 2026-10-04T12:00:00Z + self.saved_reports = [] # R-868: (plan path, report) the wrapper kept on disk self.files[osapply.TRUST_FILE] = json.dumps({"host_id": "demo-hp-bb76ea", "ring0_slow_lane": False}) self.stats[osapply.TRUST_FILE] = St(mode=statmod.S_IFREG | 0o644, uid=0) self.files[osapply.TRUST_SIGNERS] = 'felhom-op-1 namespaces="felhom-op-v1" ssh-ed25519 AAAA\n' @@ -113,6 +114,10 @@ class Fake: def log(self, line): self.logs.append(line) + def save_report(self, plan_path, report): + self.saved_reports.append((plan_path, json.loads(json.dumps(report)))) + return plan_path.replace("/plan-", "/report-") + def host(self, argv, timeout=600, stdin=None): self.calls.append(("host", argv)) if argv[0] == "/usr/sbin/pct" and argv[1] == "status": @@ -168,8 +173,8 @@ class Fake: if "-f" in a: self.dpkg_audit = "" return 0, "Setting up x (1) ...\n" if getattr(self, "repaired", False) else "", "" - if "-s" in a: - return self.sim(a) + if "--print-uris" in a or "-s" in a: + return self.sim(a) # --print-uris prints and installs nothing, with or without -s (9202, 2026-10-05) if "install" in a: if self.install_rc: return self.install_rc, "", "E: boom" @@ -232,7 +237,14 @@ class Fake: def sim(self, a): if "--print-uris" in a: - return 0, "'http://x/libc6.deb' libc6.deb 4000000 SHA256:x\n", "" + if "-s" in a: + # real apt (measured 9202 2026-10-05): with -s it prints the SIMULATION, no URI list + return 0, "Inst libc6 [2.41-12+deb13u4] (2.41-12+deb13u4 Debian:13.7/stable [amd64])\n", "" + # real apt's line shape, verbatim from 9202 2026-10-05 (audits/night-fixes-2026-10-05/partC/) + return 0, ("Need to get 4347 kB of archives.\n" + "'http://deb.debian.org/debian/pool/main/b/bash/bash_5.2.37-2%2bb10_amd64.deb' bash_5.2.37-2+b10_amd64.deb 1500792 MD5Sum:27b11721fea83d73b96e0f7023863771\n" + "'http://deb.debian.org/debian/pool/main/g/glibc/libc6_2.41-12%2bdeb13u4_amd64.deb' libc6_2.41-12+deb13u4_amd64.deb 2846580 MD5Sum:5559581916477ef1f57ea9f82cecf22e\n" + + getattr(self, "extra_uris", "")), "" 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 @@ -273,7 +285,7 @@ class Happy(unittest.TestCase): self.assertEqual(rep["pending"][0]["name"], "bash") self.assertTrue(any(l.startswith("os-apply: REPAIR ") for l in f.logs), "the repair line must always print") self.assertTrue(any(l.startswith("os-apply: DONE rc=0") for l in f.logs)) - inst = [c for c in f.calls if c[0] == "guest" and "install" in c[2] and "-s" not in c[2] and "-f" not in c[2]] + inst = [c for c in f.calls if c[0] == "guest" and "install" in c[2] and "-s" not in c[2] and "-f" not in c[2] and "--print-uris" not in c[2]] self.assertTrue(inst and "Dpkg::Options::=--force-confold" in inst[0][2], "must keep existing config files") def test_already_current_is_a_no_op(self): @@ -343,7 +355,7 @@ class Refusals(unittest.TestCase): self.assertEqual(rc, 2, rep) self.assertEqual(rep["refused"]["code"], code, rep) self.assertTrue(any(l.startswith(f"os-apply: REFUSED: {code} ") for l in f.logs), f.logs) - inst = [c for c in f.calls if c[0] == "guest" and "install" in c[2] and "-s" not in c[2] and "-f" not in c[2]] + inst = [c for c in f.calls if c[0] == "guest" and "install" in c[2] and "-s" not in c[2] and "-f" not in c[2] and "--print-uris" not in c[2]] self.assertEqual(inst, [], "a refusal must install nothing") return rep @@ -454,6 +466,24 @@ class Refusals(unittest.TestCase): f.free = 100 * 1024 * 1024 self.refused(f, "R8") + # R-865: the download is the real one. 2 GB of URIs, 5 GB free: 5 GB < 3 x 2 GB -> R8, with the size in the line. + # COMPANION RED-PROOF: put "-s" back into download_bytes -> 0 B -> no refusal -> this test fails. + def test_R8_counts_the_real_download(self): + f = Fake() + f.free = 5 * 1024 ** 3 + f.extra_uris = "'http://deb.debian.org/debian/pool/main/b/big/big_1_amd64.deb' big_1_amd64.deb 2000000000 MD5Sum:x\n" + rep = self.refused(f, "R8") + self.assertIn("download 2004347372 B", str(rep)) + + def test_download_bytes_never_simulates(self): + f = Fake() + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + calls = [c[2] for c in f.calls if c[0] == "guest" and "--print-uris" in c[2]] + self.assertTrue(calls, "download_bytes was never called") + for c in calls: + self.assertNotIn("-s", c, f"--print-uris with -s prints no URIs: {c}") + def test_R9_guest_locked_by_a_backup(self): f = Fake() f.files["/etc/pve/lxc/9201.conf"] = CONF_OK + "lock: backup\n" @@ -962,3 +992,88 @@ class RealSignatureCheck(unittest.TestCase): if __name__ == "__main__": unittest.main() + + +class UnsentReport(unittest.TestCase): + """R-868 (v0.144.0): an apply pass keeps its report on disk until the agent has sent it. + COMPANION RED-PROOF: drop the r.save_report call in main() -> test_apply_keeps_a_copy fails.""" + + def test_apply_keeps_a_copy_with_the_agents_ids(self): + f = Fake() + f.plan.update(run_id="20261005T0257-ab12", trigger="debug", ring=0) + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + self.assertEqual(len(f.saved_reports), 1, "an apply pass must keep its report on disk") + path, saved = f.saved_reports[0] + self.assertEqual(path, PLAN) + self.assertEqual(saved, rep, "the copy is the report the agent would have read") + self.assertEqual((saved["run_id"], saved["trigger"], saved["ring"]), ("20261005T0257-ab12", "debug", 0)) + + def test_a_refusal_is_kept_too(self): + f = Fake() + f.free = 1 + rc, rep = run(f) + self.assertEqual(rc, 2) + self.assertEqual(f.saved_reports[0][1]["refused"]["code"], "R8") + + def test_other_modes_keep_nothing(self): + f = Fake() + f.plan["mode"] = "health" + run(f) + self.assertEqual(f.saved_reports, []) + + def test_odd_ids_are_dropped_not_trusted(self): + f = Fake() + f.plan.update(run_id="../../etc/x", trigger="Night; rm", ring=True) + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + for k in ("run_id", "trigger", "ring"): + self.assertNotIn(k, rep) + + +class SaveReportOnDisk(unittest.TestCase): + """The real Runner.save_report on a temp dir: root writes into the AGENT's directory, so a symlink must never + be followed — neither for the directory nor for the file name.""" + + def setUp(self): + import tempfile + self.tmp = tempfile.mkdtemp() + self.dir = os.path.join(self.tmp, "os") + os.mkdir(self.dir) + self.prev = osapply.PLAN_DIR + osapply.PLAN_DIR = self.dir + self.r = osapply.Runner() + self.r.agent_uid = lambda: os.getuid() + + def tearDown(self): + import shutil + osapply.PLAN_DIR = self.prev + shutil.rmtree(self.tmp) + + def test_writes_0600_next_to_the_plan(self): + p = self.r.save_report(os.path.join(self.dir, "plan-r1-guest-apply.json"), {"mode": "apply"}) + self.assertEqual(p, os.path.join(self.dir, "report-r1-guest-apply.json")) + st = os.stat(p) + self.assertEqual(statmod.S_IMODE(st.st_mode), 0o600) + with open(p) as fh: + self.assertEqual(json.load(fh), {"mode": "apply"}) + + def test_a_symlink_at_the_name_is_replaced_not_followed(self): + victim = os.path.join(self.tmp, "victim") + with open(victim, "w") as fh: + fh.write("untouched") + os.symlink(victim, os.path.join(self.dir, "report-r2-guest-apply.json")) + self.r.save_report(os.path.join(self.dir, "plan-r2-guest-apply.json"), {"mode": "apply"}) + with open(victim) as fh: + self.assertEqual(fh.read(), "untouched") + self.assertFalse(os.path.islink(os.path.join(self.dir, "report-r2-guest-apply.json"))) + + def test_a_symlinked_directory_is_refused(self): + real = os.path.join(self.tmp, "elsewhere") + os.mkdir(real) + link = os.path.join(self.tmp, "linked") + os.symlink(real, link) + osapply.PLAN_DIR = link + with self.assertRaises(OSError): + self.r.save_report(os.path.join(link, "plan-r3-guest-apply.json"), {"mode": "apply"}) + self.assertEqual(os.listdir(real), []) diff --git a/internal/osupdate/dockerjob.go b/internal/osupdate/dockerjob.go index 237074c..4d9f04f 100644 --- a/internal/osupdate/dockerjob.go +++ b/internal/osupdate/dockerjob.go @@ -73,6 +73,9 @@ func (e DockerStepExecutor) Execute(ctx context.Context, op string, params json. // RunDockerSigned is one signed Docker step (ring 1 or an undo): live-restore first (decision 87, a no-op when on), // then the docker layer with the signed envelope, which the wrapper verifies itself. func (l *Leg) RunDockerSigned(ctx context.Context, vmid int, p DockerStepParams, blob []byte, sig string) Report { + unlock := l.lockPass(true) + defer unlock() + l.sendUnsentLocked(ctx) // R-868 runID := l.now().UTC().Format("20060102T150405Z") lg := l.log().With("run", runID, "vmid", vmid, "trigger", "signed", "release", p.ReleaseID, "undo", p.Undo) if err := l.EnsureLiveRestore(ctx, runID, vmid); err != nil { diff --git a/internal/osupdate/leg.go b/internal/osupdate/leg.go index 3e33b18..9c675a2 100644 --- a/internal/osupdate/leg.go +++ b/internal/osupdate/leg.go @@ -71,16 +71,16 @@ type Container struct { // 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"` // ControllerDockerOK: the controller reaches the engine from INSIDE its container (R-858, wrapper ≥ v0.142.1; // nil from an older wrapper = not checked). Its own health check stayed "healthy" while it was blind. - ControllerDockerOK *bool `json:"controller_docker_ok,omitempty"` - HostServices map[string]string `json:"host_services,omitempty"` - GuestRunning *bool `json:"guest_running,omitempty"` - Guest *Health `json:"guest,omitempty"` + ControllerDockerOK *bool `json:"controller_docker_ok,omitempty"` + 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. @@ -106,6 +106,13 @@ type WrapperReport struct { LiveRestore json.RawMessage `json:"live_restore"` Facts json.RawMessage `json:"facts"` Bundle json.RawMessage `json:"bundle"` // the config bundle's result (R-840, mode "bundle") + // R-868 (v0.144.0): the agent's ids, echoed from the plan, so a report kept on disk can be sent without the + // agent process that started the pass. ReleaseID / VMID were always in the report. + RunID string `json:"run_id"` + Trigger string `json:"trigger"` + Ring *int `json:"ring"` + ReleaseID string `json:"release_id"` + VMID int `json:"vmid"` } func (w WrapperReport) refused() bool { return len(w.Refused) > 0 && string(w.Refused) != "null" } @@ -136,6 +143,8 @@ type Report struct { DockerEngine string `json:"docker_engine,omitempty"` // docker layer: the engine after the step Authority string `json:"authority,omitempty"` // docker layer: ring0 | signed Undo bool `json:"undo,omitempty"` // docker layer: a signed undo (downgrade) + + unsent string // R-868: the wrapper's kept copy of this pass's report — deleted once the hub has it } // Reporter posts a report to the hub (*hub.Client). @@ -162,14 +171,65 @@ type Leg struct { block *hub.WireOSUpdate } +// planFile / reportFile: the plan the agent writes and the copy of the report the wrapper keeps beside it (R-868). +func planFile(dir, runID, layer, mode string) string { + return filepath.Join(dir, fmt.Sprintf("plan-%s-%s-%s.json", runID, layer, mode)) +} + +func reportFile(dir, runID, layer, mode string) string { + return filepath.Join(dir, fmt.Sprintf("report-%s-%s-%s.json", runID, layer, mode)) +} + +func (l *Leg) planDir() string { + if l.PlanDir == "" { + return DefaultPlanDir + } + return l.PlanDir +} + // OnDesiredState stores the hub's os_update block (desired.RawConsumer — store only, never block). func (l *Leg) OnDesiredState(_ context.Context, resp *hub.DesiredStateResponse) { if resp == nil { return } l.mu.Lock() - defer l.mu.Unlock() l.block = resp.DesiredState.OSUpdate + l.mu.Unlock() + l.saveBlock(resp.DesiredState.OSUpdate) +} + +// SavedBlockFile is the hub's newest os_update block as the daemon last received it (R-866, v0.144.0): the debug +// pass falls back to it when the hub cannot be reached, and says so. +const SavedBlockFile = "os-update-block.json" + +type savedBlock struct { + SavedAt time.Time `json:"saved_at"` + Block *hub.WireOSUpdate `json:"block"` +} + +func (l *Leg) saveBlock(b *hub.WireOSUpdate) { + dir := l.planDir() + if err := os.MkdirAll(dir, 0o700); err != nil { + return + } + body, _ := json.Marshal(savedBlock{SavedAt: l.now().UTC(), Block: b}) + tmp := filepath.Join(dir, SavedBlockFile+".tmp") + if err := os.WriteFile(tmp, body, 0o600); err == nil { + _ = os.Rename(tmp, filepath.Join(dir, SavedBlockFile)) + } +} + +// LoadSavedBlock reads the block the daemon saved (R-866). ok=false: none saved yet. +func LoadSavedBlock(dir string) (b *hub.WireOSUpdate, savedAt time.Time, ok bool) { + raw, err := os.ReadFile(filepath.Join(dir, SavedBlockFile)) + if err != nil { + return nil, time.Time{}, false + } + var s savedBlock + if json.Unmarshal(raw, &s) != nil { + return nil, time.Time{}, false + } + return s.Block, s.SavedAt, true } // Block returns the newest os_update block. No block (an older hub, or nothing fetched yet) = ring 1, ON, no @@ -352,15 +412,12 @@ func DockerHealthVerdict(before, after *Health, wantEngine, gotEngine string) (b // call writes the plan and runs the wrapper once. func (l *Leg) call(ctx context.Context, runID string, plan map[string]any) (WrapperReport, error) { - dir := l.PlanDir - if dir == "" { - dir = DefaultPlanDir - } + dir := l.planDir() if err := os.MkdirAll(dir, 0o700); err != nil { return WrapperReport{}, fmt.Errorf("osupdate: plan dir: %w", err) } b, _ := json.Marshal(plan) - path := filepath.Join(dir, fmt.Sprintf("plan-%s-%s-%s.json", runID, plan["layer"], plan["mode"])) + path := planFile(dir, runID, fmt.Sprint(plan["layer"]), fmt.Sprint(plan["mode"])) if err := os.WriteFile(path, b, 0o600); err != nil { return WrapperReport{}, fmt.Errorf("osupdate: write plan: %w", err) } @@ -394,6 +451,9 @@ type Pass struct { // Run is one pass: the guest layer, then (on an appliance, after a good guest step) the host layer, then (ring 0 // only, after good earlier steps) the Docker engine set. trigger is "night" or "debug". func (l *Leg) Run(ctx context.Context, vmid int, trigger string) Pass { + unlock := l.lockPass(true) + defer unlock() + l.sendUnsentLocked(ctx) // R-868: a report a killed agent never sent goes first g, h := l.runFast(ctx, vmid, trigger) p := Pass{Guest: g, Host: h} if g.Outcome == "skipped" { @@ -521,7 +581,8 @@ func (l *Leg) runLayer(ctx context.Context, runID, layer string, vmid int, trigg lg.Info("osupdate: START", "enabled", blk.Enabled, "release", rel.ID) plan := map[string]any{"release_id": rel.ID, "layer": layer, "lane": lane, "vmid": vmid, "snapshot": rel.Snapshot, - "packages": []Package{}, "mode": "apply", "select": "listed"} + "packages": []Package{}, "mode": "apply", "select": "listed", + "run_id": runID, "trigger": trigger, "ring": blk.Ring} // R-868: echoed into the wrapper's kept copy if rel.ID == "" { plan["release_id"] = "none" } @@ -554,6 +615,9 @@ func (l *Leg) runLayer(ctx context.Context, runID, layer string, vmid int, trigg } rep.Mode = plan["mode"].(string) wr, err := l.call(ctx, runID, plan) + if rep.Mode == "apply" { + rep.unsent = reportFile(l.planDir(), runID, layer, rep.Mode) // the wrapper kept a copy (R-868) + } switch { case err != nil: rep.Outcome, rep.HealthReason = "failed", err.Error() @@ -684,7 +748,11 @@ func (l *Leg) finish(ctx context.Context, lg *slog.Logger, rep Report) Report { rctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), time.Minute) defer cancel() if err := l.Hub.PostOSReport(rctx, body); err != nil { - lg.Warn("osupdate: reporting to the hub failed (the run itself is done)", "err", err) + lg.Warn("osupdate: reporting to the hub failed (the run itself is done; the kept copy is sent at the next start or pass)", "err", err) + } else if rep.unsent != "" { + if rerr := os.Remove(rep.unsent); rerr != nil && !os.IsNotExist(rerr) { + lg.Warn("osupdate: could not delete the sent report's kept copy (it may be sent twice)", "path", rep.unsent, "err", rerr) + } } } return rep diff --git a/internal/osupdate/leg_test.go b/internal/osupdate/leg_test.go index 6fec4a0..c42fb28 100644 --- a/internal/osupdate/leg_test.go +++ b/internal/osupdate/leg_test.go @@ -3,6 +3,7 @@ package osupdate import ( "context" "encoding/json" + "fmt" "io" "os" "os/exec" @@ -21,6 +22,7 @@ type fakeWrapper struct { applyRep map[string]WrapperReport // per layer healthSeq map[string][]*Health // per layer: answers to successive "health" calls plans []map[string]any + keep bool // R-868: like the real wrapper, keep an apply report beside the plan } func yes() *bool { b := true; return &b } @@ -72,6 +74,22 @@ func (f *fakeWrapper) Run(_ context.Context, name string, args ...string) ([]byt rep.Health = ok } } + if f.keep && plan["mode"] == "apply" { + kept := rep + kept.Layer, kept.RunID, _ = layer, fmt.Sprint(plan["run_id"]), 0 + if tr, ok := plan["trigger"].(string); ok { + kept.Trigger = tr + } + if r, ok := plan["ring"].(float64); ok { + ri := int(r) + kept.Ring = &ri + } + kb, _ := json.Marshal(kept) + dst := filepath.Join(filepath.Dir(args[1]), "report-"+strings.TrimPrefix(filepath.Base(args[1]), "plan-")) + if err := os.WriteFile(dst, kb, 0o600); err != nil { + f.t.Fatal(err) + } + } out, _ := json.Marshal(rep) return []byte("OSAPPLY-REPORT " + string(out) + "\n"), []byte("os-apply: DONE rc=0\n"), nil } diff --git a/internal/osupdate/unsent.go b/internal/osupdate/unsent.go new file mode 100644 index 0000000..2ce49b0 --- /dev/null +++ b/internal/osupdate/unsent.go @@ -0,0 +1,178 @@ +package osupdate + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "strings" + "syscall" + + "gitea.dooplex.hu/admin/felhom-agent/internal/hub" +) + +// ── R-868 (v0.144.0): a pass whose agent was killed still reports ───────────────────────────────────── +// +// MEASURED 2026-10-05 02:57 UTC on demo-hp (night drill A5): the debug pass and the agent daemon were kill -9-ed while +// apt-get ran. The root wrapper (its own process under sudo) finished all six packages, but the agent that would have +// read its stdout and posted the report was gone; the hub never learned what the pass installed. +// +// THE MECHANISM: the wrapper writes its apply report to /report---apply.json BEFORE printing +// it (configs/felhom-os-apply save_report). The agent deletes that copy once the hub has the report (finish). A copy +// still on disk is a report nobody sent: SendUnsent posts it — at the agent's start, and before every pass — and then +// deletes it. A pass lock (flock on /pass.lock, released by the kernel when a process dies) keeps the +// sender from picking up the copy of a pass that is still running, also across the daemon and a selftest process. +// Pinned by TestR868_* (unsent_test.go). + +// lockPass takes the pass lock. block=false returns ok=false at once when another pass holds it. A lock that cannot +// be opened at all (no plan dir yet) does not stop a pass: the unlock is then a no-op. +func (l *Leg) lockPass(block bool) (unlock func()) { + u, _ := l.tryLockPass(block) + return u +} + +func (l *Leg) tryLockPass(block bool) (unlock func(), ok bool) { + dir := l.planDir() + _ = os.MkdirAll(dir, 0o700) + f, err := os.OpenFile(filepath.Join(dir, "pass.lock"), os.O_CREATE|os.O_RDWR, 0o600) + if err != nil { + l.log().Warn("osupdate: pass lock unavailable — continuing without it", "err", err) + return func() {}, true + } + how := syscall.LOCK_EX + if !block { + how |= syscall.LOCK_NB + } + if err := syscall.Flock(int(f.Fd()), how); err != nil { + f.Close() + return func() {}, false + } + return func() { _ = syscall.Flock(int(f.Fd()), syscall.LOCK_UN); f.Close() }, true +} + +// SendUnsent posts every report a pass kept on disk and nobody sent (R-868). It skips when a pass runs now (that +// pass sends them first). Called at the agent's start. +func (l *Leg) SendUnsent(ctx context.Context) int { + unlock, ok := l.tryLockPass(false) + if !ok { + l.log().Info("osupdate: a pass is running — its start sends any kept report") + return 0 + } + defer unlock() + return l.sendUnsentLocked(ctx) +} + +func (l *Leg) sendUnsentLocked(ctx context.Context) int { + if l.Hub == nil { + return 0 // nobody to send to: keep the copies for a process that has the hub + } + files, _ := filepath.Glob(filepath.Join(l.planDir(), "report-*.json")) + sent := 0 + for _, f := range files { + b, err := os.ReadFile(f) + if err != nil { + l.log().Warn("osupdate: a kept report cannot be read", "path", f, "err", err) + continue + } + var wr WrapperReport + if err := json.Unmarshal(b, &wr); err != nil || wr.Layer == "" { + l.log().Warn("osupdate: a kept report is not a report — moved aside", "path", f, "err", err) + _ = os.Rename(f, f+".bad") + continue + } + rep := l.reportFromKept(ctx, wr, f) + lg := l.log().With("run", rep.RunID, "layer", rep.Layer, "vmid", rep.VMID, "ring", rep.Ring, "trigger", rep.Trigger) + lg.Info("osupdate: sending a report the agent never sent (the agent stopped mid-pass, R-868)", "path", f) + before := rep.unsent + _ = l.finish(ctx, lg, rep) + if _, err := os.Stat(before); os.IsNotExist(err) { + sent++ + // the pass's plan file is left behind too when the agent was killed inside call() + _ = os.Remove(filepath.Join(filepath.Dir(f), "plan-"+strings.TrimPrefix(filepath.Base(f), "report-"))) + } + } + return sent +} + +// reportFromKept builds the hub report from a kept wrapper report, as runLayer would have. Health: the copy's own +// before/after reading; when that is not healthy after an install, one fresh reading (services restart after an +// install, and the pass that would have waited for them is gone). +func (l *Leg) reportFromKept(ctx context.Context, wr WrapperReport, path string) Report { + ring := 1 + if wr.Ring != nil { + ring = *wr.Ring + } + runID, trigger := wr.RunID, wr.Trigger + if runID == "" { + runID = strings.TrimSuffix(strings.TrimPrefix(filepath.Base(path), "report-"), ".json") + } + if trigger == "" { + trigger = "unknown" + } + rep := Report{RunID: runID, Layer: wr.Layer, Trigger: trigger, Mode: wr.Mode, Ring: ring, ReleaseID: wr.ReleaseID, + VMID: wr.VMID, unsent: path} + prefix := "sent after the agent stopped mid-pass (R-868)" + switch { + case wr.refused(): + rep.Outcome, rep.Refused, rep.HealthReason = "refused", wr.Refused, prefix + return rep + case wr.failed(): + rep.Outcome, rep.Refused = "failed", wr.Failed + case len(wr.Upgraded) == 0: + rep.Outcome = "nothing" + default: + rep.Outcome = "applied" + } + rep.Upgraded, rep.PassSeconds = wr.Upgraded, wr.PassSeconds + rep.DockerEngine, rep.Authority, rep.Undo = wr.DockerEngine, wr.Authority, wr.Undo + wantEngine := "" + for _, u := range wr.Upgraded { + if u.Name == "docker-ce" { + wantEngine = EngineOf(u.Version) + } + } + verdict := func(h *Health) (bool, string) { + switch wr.Layer { + case LayerDocker: + return DockerHealthVerdict(wr.HealthBefore, h, wantEngine, wr.DockerEngine) + case LayerHost: + t := hub.TunnelUnknown + if l.Tunnel != nil { + t, _ = l.Tunnel.Status(ctx) + } + return HostHealthVerdict(wr.HealthBefore, h, t) + } + return HealthVerdict(wr.HealthBefore, h) + } + ok, why := verdict(wr.HealthAfter) + if !ok && len(wr.Upgraded) > 0 && wr.VMID > 0 { + lane := "fast" + if wr.Layer == LayerDocker { + lane = "slow" + } + if hr, err := l.call(ctx, "kept-"+runID, map[string]any{"release_id": "kept", "layer": wr.Layer, "lane": lane, + "vmid": wr.VMID, "mode": "health", "packages": []Package{}}); err == nil && hr.Health != nil { + ok, why = verdict(hr.Health) + } + } + rep.Healthy, rep.HealthReason = ok, prefix + if why != "" { + rep.HealthReason = prefix + ": " + why + } + if !ok && rep.Outcome == "applied" { + rep.Outcome = "health_failed" + } + rep.Installed, rep.Pending = wr.Installed, wr.Pending + rep.RestartNeeded, rep.DockerRestartNeeded, rep.RebootNeeded = wr.RestartNeeded, wr.DockerRestartNeeded, wr.RebootNeeded + rep.RebootScanned = wr.RebootScanned + if wr.Layer == LayerDocker { + rep.Installed, rep.Pending = onlyDocker(wr.Installed), onlyDockerPending(wr.Pending) + } else { + planned := map[string]bool{} + for _, u := range wr.Upgraded { + planned[u.Name] = true + } + rep.NotCovered = notCovered(wr.Pending, ring, planned) + } + return rep +} diff --git a/internal/osupdate/unsent_test.go b/internal/osupdate/unsent_test.go new file mode 100644 index 0000000..d744235 --- /dev/null +++ b/internal/osupdate/unsent_test.go @@ -0,0 +1,123 @@ +package osupdate + +import ( + "context" + "encoding/json" + "errors" + "os" + "path/filepath" + "testing" + + "gitea.dooplex.hu/admin/felhom-agent/internal/hub" +) + +// R-868 (v0.144.0). THE NIGHT'S SHAPE (A5, demo-hp 2026-10-05 02:57 UTC): the wrapper finished six packages, the agent +// was kill -9-ed before it read the report; the hub never got it. Here: the wrapper's kept copy and the plan file are +// on disk, a NEW agent process starts — it must send exactly one `applied` report and delete both files. +// COMPANION RED-PROOF: drop the SendUnsent body (return 0) → "the hub got no report". +func TestR868_KilledPassIsReportedAtStart(t *testing.T) { + w := &fakeWrapper{t: t} + l, h := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + ring := 0 + kept := WrapperReport{Mode: "apply", Layer: LayerGuest, RunID: "20261005T025700Z", Trigger: "debug", Ring: &ring, VMID: 9201, + ReleaseID: "ring0-20261005T025700Z", HealthBefore: guestOK(), HealthAfter: guestOK(), + Upgraded: []Package{{Name: "libc6", Version: "u4"}, {Name: "openssl", Version: "u3"}}} + b, _ := json.Marshal(kept) + rp := reportFile(l.PlanDir, kept.RunID, LayerGuest, "apply") + pp := planFile(l.PlanDir, kept.RunID, LayerGuest, "apply") + must(t, os.WriteFile(rp, b, 0o600)) + must(t, os.WriteFile(pp, []byte("{}"), 0o600)) + if n := l.SendUnsent(context.Background()); n != 1 { + t.Fatalf("sent %d, want 1", n) + } + if len(h.reports) != 1 { + t.Fatalf("the hub got no report (or several): %+v", h.reports) + } + r := h.reports[0] + if r.Outcome != "applied" || !r.Healthy || r.Trigger != "debug" || r.RunID != kept.RunID || r.Ring != 0 || len(r.Upgraded) != 2 || r.VMID != 9201 { + t.Fatalf("report = %+v", r) + } + for _, p := range []string{rp, pp} { + if _, err := os.Stat(p); !os.IsNotExist(err) { + t.Fatalf("%s still on disk after the hub got it", filepath.Base(p)) + } + } + // a second start sends nothing again — no duplicate report + if n := l.SendUnsent(context.Background()); n != 0 || len(h.reports) != 1 { + t.Fatalf("sent again: %d, reports %d", n, len(h.reports)) + } +} + +// A normal pass: the hub gets ONE report per layer and the kept copies are gone after it (nothing resent later). +// COMPANION RED-PROOF: drop the os.Remove(rep.unsent) in finish → the next pass resends → "2 guest reports". +func TestR868_NormalPassLeavesNoCopyAndNoDuplicate(t *testing.T) { + w := &fakeWrapper{t: t, keep: true, applyRep: map[string]WrapperReport{LayerGuest: {Upgraded: []Package{{Name: "libc6", Version: "u4"}}}}} + l, h := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + l.Run(context.Background(), 9201, "debug") + left, _ := filepath.Glob(filepath.Join(l.PlanDir, "report-*.json")) + if len(left) != 0 { + t.Fatalf("kept copies left after the hub got them: %v", left) + } + l.Run(context.Background(), 9201, "debug") // the next pass sends kept copies first + guest := 0 + for _, r := range h.reports { + if r.Layer == LayerGuest && r.Outcome == "applied" { + guest++ + } + } + if guest != 2 { + t.Fatalf("%d guest reports for 2 passes (a duplicate or a loss)", guest) + } +} + +type failingHub struct{ n int } + +func (h *failingHub) PostOSReport(context.Context, []byte) error { + h.n++ + return errors.New("hub away") +} + +// The hub away: the copy stays, and goes at the next chance. +func TestR868_HubAwayKeepsTheCopy(t *testing.T) { + w := &fakeWrapper{t: t, keep: true, applyRep: map[string]WrapperReport{LayerGuest: {Upgraded: []Package{{Name: "libc6", Version: "u4"}}}}} + l, _ := newLeg(t, w, &hub.WireOSUpdate{Ring: 0, Enabled: true}) + l.Appliance = false + l.Hub = &failingHub{} + l.Run(context.Background(), 9201, "night") + left, _ := filepath.Glob(filepath.Join(l.PlanDir, "report-*.json")) + if len(left) != 2 { // ring 0: the guest step and the Docker step each kept one + t.Fatalf("the copies must stay while the hub is away: %v", left) + } + h := &fakeHub{} + l.Hub = h + if n := l.SendUnsent(context.Background()); n != 2 { + t.Fatalf("sent %d: %+v", n, h.reports) + } + for _, r := range h.reports { + if r.Trigger != "night" || r.Ring != 0 { + t.Fatalf("the kept report lost its ids: %+v", r) + } + } +} + +// A pass in progress holds the lock: the sender at start must not take that pass's copy (it would be sent twice). +func TestR868_RunningPassKeepsTheSenderOff(t *testing.T) { + w := &fakeWrapper{t: t} + l, h := newLeg(t, w, nil) + must(t, os.WriteFile(reportFile(l.PlanDir, "r1", LayerGuest, "apply"), []byte(`{"mode":"apply","layer":"guest"}`), 0o600)) + unlock := l.lockPass(true) + if n := l.SendUnsent(context.Background()); n != 0 || len(h.reports) != 0 { + t.Fatalf("sent while a pass ran: %d", n) + } + unlock() + if n := l.SendUnsent(context.Background()); n != 1 { + t.Fatalf("not sent after the pass: %d", n) + } +} + +func must(t *testing.T, err error) { + t.Helper() + if err != nil { + t.Fatal(err) + } +}