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) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
This commit is contained in:
@@ -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-<run>-<layer>-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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
+48
-1
@@ -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
|
||||
|
||||
|
||||
@@ -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), [])
|
||||
|
||||
@@ -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 {
|
||||
|
||||
+84
-16
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 <plan dir>/report-<run>-<layer>-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 <plan dir>/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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user