diff --git a/configs/felhom-os-apply b/configs/felhom-os-apply index e44b3dd..848aedd 100755 --- a/configs/felhom-os-apply +++ b/configs/felhom-os-apply @@ -63,6 +63,9 @@ 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 +# R-876: dpkg's state in ONE call — `--audit`, a marker line, then the update journal's file names. +JOURNAL_MARK = "@@FELHOM-DPKG-JOURNAL@@" +DPKG_STATE_SCRIPT = "dpkg --audit; echo " + JOURNAL_MARK + "; ls -A /var/lib/dpkg/updates 2>/dev/null; true" # 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 @@ -902,20 +905,34 @@ class Apply: self.report["health_after"] = self.health() return 0 - def repair(self): - rc, before, _ = self.x(["dpkg", "--audit"]) + def dpkg_state(self): + """`dpkg --audit` AND dpkg's update journal, in ONE call (R-876, agent v0.145.0). A crash in the middle of an + install can leave `/var/lib/dpkg/updates/` non-empty while `--audit` reads clean — measured on demo-hp + 2026-10-05 — and that journal is exactly what apt refuses on ("dpkg was interrupted"). One `sh -c` with a + constant script keeps R-845's speed: a clean pass still costs one call here, as before.""" + rc, out, _ = self.x(["sh", "-c", DPKG_STATE_SCRIPT]) + audit, _, journal = out.partition(JOURNAL_MARK + "\n") + return audit, [l for l in journal.split() if l] + + def repair(self, force=False): + before, journal = self.dpkg_state() 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) + after, journal_after = "", [] + # nothing half-done and no update journal → nothing to run (R-845: two calls saved on every clean pass). + # R-876: the JOURNAL counts too, and `force` (apt said "dpkg was interrupted") always repairs. + if before.strip() or journal or force: 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"]) + after, journal_after = self.dpkg_state() 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}") + self.report["repair"] = {"half_configured_before": configured, "journal_before": len(journal), "fixed": fixed, + "clean_after": after.strip() == "" and not journal_after} + self.r.log(f"os-apply: REPAIR configured={configured} journal={len(journal)} fixed={fixed}" + (" forced" if force else "")) if after.strip(): raise Refused("R13", "dpkg is still broken after the repair: " + after.strip().splitlines()[0]) + if journal_after: + raise Refused("R13", f"dpkg's update journal is still not empty after the repair ({len(journal_after)} file(s))") def pending_fast(self): """Ring 0 (select pending-fast): every pending upgrade of an INSTALLED package whose every origin is Debian / @@ -1032,6 +1049,12 @@ class Apply: 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) + if rc != 0 and "dpkg was interrupted" in (out + err): + # R-876 (belt): apt says dpkg was interrupted although the repair found nothing — repair and + # retry ONCE. Never a loop. + self.r.log("os-apply: INTERRUPTED apt says dpkg was interrupted — repairing and retrying once") + self.repair(force=True) + 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 diff --git a/configs/test_felhom_os_apply.py b/configs/test_felhom_os_apply.py index 6228f8b..318dd33 100644 --- a/configs/test_felhom_os_apply.py +++ b/configs/test_felhom_os_apply.py @@ -154,6 +154,8 @@ class Fake: if cmd == "dpkg" and a[1] == "--audit": return 0, self.dpkg_audit, "" if cmd == "dpkg" and a[1] == "--configure": + self.configured_calls = getattr(self, "configured_calls", 0) + 1 + self.dpkg_journal = [] # `dpkg --configure -a` replays and empties the update journal return 0, "", "" if cmd == "fuser": return (0, " 123", "") if self.lock_held else (1, "", "") @@ -176,6 +178,9 @@ class Fake: 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 getattr(self, "dpkg_journal", []): + # real apt (demo-hp 2026-10-05): a non-empty update journal refuses every install + return 100, "", "E: dpkg was interrupted, you must manually run 'sudo dpkg --configure -a' to correct the problem.\n" if self.install_rc: return self.install_rc, "", "E: boom" for x in a: @@ -220,6 +225,8 @@ class Fake: return 0, "", "" if cmd == "getent": return 0, "1.2.3.4 deb.debian.org\n", "" + if cmd == "sh" and a[2] == osapply.DPKG_STATE_SCRIPT: + return 0, self.dpkg_audit + osapply.JOURNAL_MARK + "\n" + "".join(j + "\n" for j in getattr(self, "dpkg_journal", [])), "" if cmd == "sh": if "vmlinuz" in a[2]: return 0, "/boot/vmlinuz-7.0.2-6-pve\n/boot/vmlinuz-7.0.14-20-pve\n", "" @@ -1123,3 +1130,52 @@ class AgentDiesMidPass(unittest.TestCase): self.assertEqual(len(f.saved_reports), 1, "the kept copy is the only way this report reaches the hub") self.assertEqual(len(f.saved_reports[0][1]["upgraded"]), 2) self.assertTrue(any(l.startswith("os-apply: DONE") for l in journal), "the journal must still get the DONE line") + + +class CrashLeftTheJournal(unittest.TestCase): + """R-876 — THE MEASURED SHAPE (demo-hp 2026-10-05, a crash while dpkg unpacked): `dpkg --audit` CLEAN, but + /var/lib/dpkg/updates holds 3 files; apt refuses every install ("dpkg was interrupted"). v0.144.1 logged + `REPAIR configured=0 fixed=0`, then `FAILED rc=100`, every pass, until a person ran `dpkg --configure -a`. + COMPANION RED-PROOFS: drop `or journal` from repair()'s condition -> test 1 fails; drop the interrupted-retry + -> test 3 fails; make repair() always run -> test 2 fails (R-845's speed).""" + + def test_the_next_pass_repairs_by_itself_and_finishes(self): + f = Fake() + f.dpkg_audit = "" + f.dpkg_journal = ["0000", "0001", "0002"] + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + self.assertEqual(rep["repair"]["journal_before"], 3) + self.assertTrue(rep["repair"]["clean_after"]) + self.assertEqual(len(rep["upgraded"]), 2) + cfg = next(i for i, c in enumerate(f.calls) if c[0] == "guest" and "--configure" in c[2]) + inst = next(i for i, c in enumerate(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.assertLess(cfg, inst, "the repair must run before the install") + self.assertFalse(any("INTERRUPTED" in l for l in f.logs), "the journal check must catch it BEFORE apt refuses") + + def test_a_clean_pass_still_costs_one_state_call(self): + f = Fake() + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + state = [c for c in f.calls if c[0] == "guest" and c[2][:2] == ["sh", "-c"] and c[2][2] == osapply.DPKG_STATE_SCRIPT] + self.assertEqual(len(state), 1, "R-845: a clean pass reads dpkg's state ONCE and runs no repair") + self.assertEqual(getattr(f, "configured_calls", 0), 0) + + def test_apt_interrupted_is_repaired_and_retried_once(self): + """The belt: the journal probe saw nothing (a race, an odd layout), apt still says interrupted.""" + f = Fake() + orig = f.emulate + state = {"first": True} + + def emulate(argv): + a = [x for x in argv if not re.match(r"^[A-Z_]+=", x) and x != "env"] + if a[0] == "apt-get" and "install" in a and "-s" not in a and "-f" not in a and "--print-uris" not in a and state["first"]: + state["first"] = False + return 100, "", "E: dpkg was interrupted, you must manually run 'sudo dpkg --configure -a' to correct the problem.\n" + return orig(argv) + f.emulate = emulate + rc, rep = run(f) + self.assertEqual(rc, 0, rep) + self.assertTrue(any("INTERRUPTED" in l for l in f.logs), f.logs) + self.assertTrue(any(l.startswith("os-apply: REPAIR ") and l.endswith("forced") for l in f.logs), f.logs) diff --git a/internal/backup/r874_first_eval_test.go b/internal/backup/r874_first_eval_test.go new file mode 100644 index 0000000..8fa69a9 --- /dev/null +++ b/internal/backup/r874_first_eval_test.go @@ -0,0 +1,63 @@ +package backup + +import ( + "context" + "fmt" + "sync/atomic" + "testing" + "time" + + "gitea.dooplex.hu/admin/felhom-agent/internal/reconcile" +) + +// R-874 (v0.145.0). THE MEASURED SHAPE (Part F spike, Tester 2): power-on sessions of ~1.5 h and ~5 min against a +// 6 h evaluation ticker that restarts at every start — no restore-test ever evaluated. Now the first evaluation runs +// FirstEval after start. +// COMPANION RED-PROOF: drop the first-evaluation timer in Run (back to the bare ticker) → "no evaluation within". +func TestR874_FirstEvaluationAfterStart(t *testing.T) { + var picks int32 + s := NewScheduler(SchedulerOptions{ + Runner: &fakeRTRunner{res: reconcile.RestoreTestResult{Pass: true, Verified: "boot+running"}}, + Pick: func(context.Context) (string, error) { + atomic.AddInt32(&picks, 1) + return fmt.Sprintf("local:backup/vzdump-lxc-9201-%d.tar.zst", atomic.LoadInt32(&picks)), nil + }, + Store: NewStore(), Spec: (&specSpy{}).build, + Cadence: 6 * time.Hour, FirstEval: 30 * time.Millisecond, Logger: quiet(), + }) + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { _ = s.Run(ctx); close(done) }() + deadline := time.Now().Add(3 * time.Second) + for atomic.LoadInt32(&picks) == 0 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + cancel() + <-done + if atomic.LoadInt32(&picks) == 0 { + t.Fatal("no evaluation within 3 s of start (FirstEval 30 ms) — a box with short sessions never gets a restore-test") + } +} + +// The earned restraint stays: an agent that restarts before FirstEval never evaluates (a crash loop does not +// hammer a failing tier). +func TestR874_CrashLoopNeverEvaluates(t *testing.T) { + var picks int32 + for i := 0; i < 5; i++ { // five quick "restarts" + s := NewScheduler(SchedulerOptions{ + Runner: &fakeRTRunner{res: reconcile.RestoreTestResult{Pass: true}}, + Pick: func(context.Context) (string, error) { atomic.AddInt32(&picks, 1); return "x", nil }, + Store: NewStore(), Spec: (&specSpy{}).build, + Cadence: 6 * time.Hour, FirstEval: 200 * time.Millisecond, Logger: quiet(), + }) + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + _ = s.Run(ctx) + cancel() + } + if n := atomic.LoadInt32(&picks); n != 0 { + t.Fatalf("a restart before FirstEval evaluated %d time(s)", n) + } + if DefaultFirstEval != 30*time.Minute { + t.Fatalf("DefaultFirstEval = %s, the documented 30 min", DefaultFirstEval) + } +} diff --git a/internal/backup/schedule.go b/internal/backup/schedule.go index 5b7d07a..ea6f942 100644 --- a/internal/backup/schedule.go +++ b/internal/backup/schedule.go @@ -60,12 +60,21 @@ type Scheduler struct { // R-85 tier rotation. All optional: without them the scheduler behaves exactly as before // (single tier via `pick`), which keeps every existing caller and test working untouched. - tiers []string // configured tier target ids, primary first - tierPick TierPicker // newest archive on a named tier - rtState *RestoreTestState // persisted last-successful-per-tier (drives oldest-first) - inFlight *InFlight // shared with the backup path — Scenario F + tiers []string // configured tier target ids, primary first + tierPick TierPicker // newest archive on a named tier + rtState *RestoreTestState // persisted last-successful-per-tier (drives oldest-first) + inFlight *InFlight // shared with the backup path — Scenario F + firstEval time.Duration // R-874: the first evaluation after start } +// DefaultFirstEval (R-874): the first due-ness evaluation runs 30 minutes after the agent starts, then every +// cadence. MEASURED need (2026-10-05 Part F spike): a box whose power-on sessions are all shorter than the 6 h +// interval (Tester 2: ~1.5 h and ~5 min) NEVER evaluated, because the ticker restarts at each start. 30 minutes +// keeps the earned restraint below — a crash-looping agent restarts far more often than that and still never +// evaluates — while a box that stays on for half an hour gets its due test. Pinned by +// TestR874_FirstEvaluationAfterStart and TestR874_CrashLoopNeverEvaluates. +const DefaultFirstEval = 30 * time.Minute + // SchedulerOptions configures a Scheduler. type SchedulerOptions struct { Runner RestoreTestRunner @@ -81,6 +90,8 @@ type SchedulerOptions struct { // 0 → no settle requirement (any archive is a candidate). Settle time.Duration Logger *slog.Logger + // FirstEval (R-874, v0.145.0) is when the FIRST evaluation runs after start; 0 → DefaultFirstEval. + FirstEval time.Duration // R-85 (all optional — omit for the pre-R-85 single-tier behaviour): // Tiers are the configured tier target ids (primary first); TierPick resolves an archive on a @@ -110,6 +121,12 @@ func NewScheduler(opts SchedulerOptions) *Scheduler { tierPick: opts.TierPick, rtState: opts.State, inFlight: opts.InFlight, + firstEval: func() time.Duration { + if opts.FirstEval > 0 { + return opts.FirstEval + } + return DefaultFirstEval + }(), } } @@ -120,8 +137,9 @@ func NewScheduler(opts SchedulerOptions) *Scheduler { // trigger any more: its phase is the process's uptime, and agent deploys reset it, which is exactly // the defect R-86 removes. What decides that a test happens is `EvaluateDue`. // -// It still does NOT evaluate immediately on start — the first evaluation is one interval in. That -// is an EARNED restraint, kept deliberately: a restore is heavy, agent restarts are routine, and a +// It still does NOT evaluate immediately on start. v0.145.0 (R-874): the first evaluation is +// firstEval (30 min) in, then every interval — it was one full interval in, which a box with short +// power-on sessions never reached. The restraint itself is EARNED and kept: a restore is heavy, agent restarts are routine, and a // crash-loop that evaluated at start would hammer a permanently-failing tier as fast as it could // restart. Due-ness does not expire while we wait, so the only cost is up to one interval of // latency on a tier that just became due. On-demand runs use `--selftest=restore-test`. @@ -135,6 +153,16 @@ func (s *Scheduler) Run(ctx context.Context) error { } s.logger.Info("backup: restore-test scheduler starting (per-archive due-check)", "eval_interval", s.cadence, "settle", s.settle) + first := time.NewTimer(s.firstEval) + defer first.Stop() + select { + case <-ctx.Done(): + s.logger.Info("backup: restore-test scheduler shutting down", "reason", ctx.Err()) + return nil + case <-first.C: + s.logger.Info("backup: restore-test first evaluation after start (R-874)", "after", s.firstEval) + s.tick(ctx) + } t := time.NewTicker(s.cadence) defer t.Stop() for { diff --git a/internal/osupdate/unsent.go b/internal/osupdate/unsent.go index e476b3d..c647a25 100644 --- a/internal/osupdate/unsent.go +++ b/internal/osupdate/unsent.go @@ -101,7 +101,7 @@ func (l *Leg) sendUnsentLocked(ctx context.Context) int { } 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) + lg.Info("osupdate: sending a kept report late (the agent was stopped mid-pass or the hub was away, R-868)", "path", f) before := rep.unsent _ = l.finish(ctx, lg, rep) if _, err := os.Stat(before); os.IsNotExist(err) { @@ -130,7 +130,7 @@ func (l *Leg) reportFromKept(ctx context.Context, wr WrapperReport, path string) } 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)" + prefix := "sent late — kept on the box until the hub could take it (R-868, R-875)" switch { case wr.refused(): rep.Outcome, rep.Refused, rep.HealthReason = "refused", wr.Refused, prefix diff --git a/internal/osupdate/unsent_test.go b/internal/osupdate/unsent_test.go index 7a70756..774e44d 100644 --- a/internal/osupdate/unsent_test.go +++ b/internal/osupdate/unsent_test.go @@ -6,6 +6,7 @@ import ( "errors" "os" "path/filepath" + "strings" "testing" "time" @@ -38,6 +39,10 @@ func TestR868_KilledPassIsReportedAtStart(t *testing.T) { 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) } + // R-875 (v0.145.0): neutral — the copy cannot tell a killed agent from an absent hub. + if !strings.HasPrefix(r.HealthReason, "sent late") || strings.Contains(r.HealthReason, "stopped mid-pass") { + t.Fatalf("health_reason = %q, want the neutral \"sent late …\"", r.HealthReason) + } 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))