// Package osupdate is the agent's OS-update leg (`11-os-updates.md` §8 steps 2–3, §5.8): the customer GUEST's Debian // fast lane (agent v0.140.0), after it in the same pass the HOST's (agent v0.141.0), and then — ring 0 only — the // guest's DOCKER engine set, the slow lane (agent v0.142.0; a ring-1 box takes a Docker step only inside a signed // operator job, DockerStepExecutor). It also reads the box's versions for the hub's System page (Facts, R-852). // // It runs right after the night's successful whole-guest backup, while the backup goroutine still holds the host-wide // heavy-op gate (so it never overlaps a backup or a restore-test, `11` C10), at most once per night. All root work is // the wrapper `felhom-os-apply` (configs/, its own tests); this package only builds plans, calls the wrapper through // sudo, judges health and reports to the hub — one report per layer. // // NO AUTOMATIC UNDO (R-837 measured; `09` §3 decision 81): a failed health check stops, reports `health_failed` and the // hub mails the operator; the whole-guest backup taken minutes earlier is the guest's undo, by hand; a host package is // put back by hand from the previous release's snapshot (runbook). The host is NEVER rebooted by this package. package osupdate import ( "context" "encoding/json" "fmt" "log/slog" "os" "path/filepath" "sort" "strings" "sync" "time" "gitea.dooplex.hu/admin/felhom-agent/internal/hub" "gitea.dooplex.hu/admin/felhom-agent/internal/proxmox" ) // WrapperPath is the pinned sudoers vector (configs/felhom-agent.sudoers FELHOM_OSAPPLY). const WrapperPath = "/usr/local/sbin/felhom-os-apply" // DefaultPlanDir is where plans are written (the sudoers glob names it). const DefaultPlanDir = "/var/lib/felhom-agent/os" // Layers. const ( LayerGuest = "guest" LayerHost = "host" LayerDocker = "docker" // the guest's Docker engine set — slow lane (`11` §5.8) ) // DockerNames are the six packages of the Docker engine set (the wrapper's DOCKER_NAMES). var DockerNames = map[string]bool{"containerd.io": true, "docker-buildx-plugin": true, "docker-ce": true, "docker-ce-cli": true, "docker-ce-rootless-extras": true, "docker-compose-plugin": true} // Package is one name=version with its origin. type Package struct { Name string `json:"name"` Version string `json:"version"` Origin string `json:"origin"` } // Pending is one update the sources offer (origin as apt names it, possibly several). type Pending struct { Name string `json:"name"` From string `json:"from"` To string `json:"to"` Origin []string `json:"origin"` } // Container is one container's state as the wrapper saw it. type Container struct { State string `json:"state"` Health string `json:"health"` // healthy | unhealthy | starting | none ID string `json:"id,omitempty"` } // 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"` HostServices map[string]string `json:"host_services,omitempty"` GuestRunning *bool `json:"guest_running,omitempty"` Guest *Health `json:"guest,omitempty"` } // WrapperReport is the wrapper's OSAPPLY-REPORT object. type WrapperReport struct { Mode string `json:"mode"` Layer string `json:"layer"` Refused json.RawMessage `json:"refused"` Failed json.RawMessage `json:"failed"` Upgraded []Package `json:"upgraded"` Installed []Package `json:"installed"` Pending []Pending `json:"pending"` RestartNeeded []string `json:"restart_needed"` DockerRestartNeeded bool `json:"docker_restart_needed"` RebootNeeded bool `json:"reboot_needed"` RebootScanned bool `json:"reboot_scanned"` HealthBefore *Health `json:"health_before"` HealthAfter *Health `json:"health_after"` Health *Health `json:"health"` PassSeconds float64 `json:"pass_seconds"` DockerEngine string `json:"docker_engine"` Authority string `json:"authority"` Undo bool `json:"undo"` LiveRestore json.RawMessage `json:"live_restore"` Facts json.RawMessage `json:"facts"` } func (w WrapperReport) refused() bool { return len(w.Refused) > 0 && string(w.Refused) != "null" } func (w WrapperReport) failed() bool { return len(w.Failed) > 0 && string(w.Failed) != "null" } // Report is what the hub receives per layer (hub osupdates.Report — field-exact). type Report struct { RunID string `json:"run_id"` Layer string `json:"layer"` Trigger string `json:"trigger"` Mode string `json:"mode"` Ring int `json:"ring"` ReleaseID string `json:"release_id"` Outcome string `json:"outcome"` Healthy bool `json:"healthy"` HealthReason string `json:"health_reason,omitempty"` VMID int `json:"vmid"` Upgraded []Package `json:"upgraded,omitempty"` Installed []Package `json:"installed,omitempty"` Pending []Pending `json:"pending,omitempty"` NotCovered []string `json:"not_covered,omitempty"` RestartNeeded []string `json:"restart_needed,omitempty"` DockerRestartNeeded bool `json:"docker_restart_needed,omitempty"` RebootNeeded bool `json:"reboot_needed,omitempty"` RebootScanned bool `json:"reboot_scanned,omitempty"` // the pass looked (host: every pass) — a false RebootNeeded then means "not needed" Refused json.RawMessage `json:"refused,omitempty"` PassSeconds float64 `json:"pass_seconds,omitempty"` 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) } // Reporter posts a report to the hub (*hub.Client). type Reporter interface { PostOSReport(ctx context.Context, body []byte) error } // Leg runs one OS-update pass for the customer guest and then the host. type Leg struct { Runner proxmox.Runner Hub Reporter Tunnel hub.CloudflaredProber // the host health rule needs the tunnel `running` (R-841) Appliance bool // agent.json deployment_mode; the wrapper re-checks the ROOT-owned record (R12) Logger *slog.Logger PlanDir string StatePath string // last night run (once per night) HealthWait time.Duration // how long health may take to come back (default 5 min) HealthPoll time.Duration // default 15 s MinGap time.Duration // between night runs (default 20 h) Now func() time.Time Sleep func(context.Context, time.Duration) mu sync.Mutex block *hub.WireOSUpdate } // 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 } // Block returns the newest os_update block. No block (an older hub, or nothing fetched yet) = ring 1, ON, no // release: the box reports and installs nothing. func (l *Leg) Block() hub.WireOSUpdate { l.mu.Lock() defer l.mu.Unlock() if l.block == nil { return hub.WireOSUpdate{Ring: 1, Enabled: true} } return *l.block } // SetBlock sets the block directly (the selftest fetches the desired state itself). func (l *Leg) SetBlock(b *hub.WireOSUpdate) { l.mu.Lock() defer l.mu.Unlock() l.block = b } func (l *Leg) now() time.Time { if l.Now != nil { return l.Now() } return time.Now() } func (l *Leg) log() *slog.Logger { if l.Logger != nil { return l.Logger } return slog.Default() } func (l *Leg) sleep(ctx context.Context, d time.Duration) { if l.Sleep != nil { l.Sleep(ctx, d) return } select { case <-ctx.Done(): case <-time.After(d): } } // IsFast reports whether every origin apt names is Debian / Debian-Security (the fast lane, `11` C3). func IsFast(origins []string) bool { if len(origins) == 0 { return false } for _, o := range origins { if o != "Debian" && o != "Debian-Security" { return false } } return true } // HealthVerdict is THE guest health rule (`11` §8.1; pinned by TestHealthVerdict*): docker answers, the guest's // network resolves, the controller's own health check is `healthy`, and every container that was running at the // start of the pass runs again — and healthy again if it was. "starting" is not yet healthy. func HealthVerdict(before, after *Health) (bool, string) { if after == nil { return false, "no health reading" } if !after.DockerOK { return false, "docker does not answer" } if !after.NetworkOK { return false, "the guest cannot resolve deb.debian.org" } if after.Controller != "healthy" { return false, "the controller is " + after.Controller } if before == nil { return true, "" } names := make([]string, 0, len(before.Containers)) for n := range before.Containers { names = append(names, n) } sort.Strings(names) for _, n := range names { b := before.Containers[n] if b.State != "running" { continue } a, ok := after.Containers[n] if !ok || a.State != "running" { return false, n + " was running and is not" } if b.Health == "healthy" && a.Health != "healthy" { return false, n + " was healthy and is " + a.Health } } return true, "" } // HostHealthVerdict is THE host health rule (`11` §8.2; pinned by TestHostHealthVerdict): the Proxmox daemons and the // agent are active, the customer guest still runs, the guest's own rule still passes against the start of the pass, // and the tunnel is `running` (R-841). func HostHealthVerdict(before, after *Health, tunnel string) (bool, string) { if after == nil { return false, "no health reading" } svcs := make([]string, 0, len(after.HostServices)) for s := range after.HostServices { svcs = append(svcs, s) } sort.Strings(svcs) if len(svcs) == 0 { return false, "no host service reading" } for _, s := range svcs { if after.HostServices[s] != "active" { return false, s + " is " + after.HostServices[s] } } if after.GuestRunning == nil || !*after.GuestRunning { return false, "the customer guest is not running" } var gb *Health if before != nil { gb = before.Guest } if ok, why := HealthVerdict(gb, after.Guest); !ok { return false, "guest: " + why } if tunnel != hub.TunnelRunning { return false, "the tunnel is " + tunnel } return true, "" } // EngineOf is the engine version `docker version` prints for a docker-ce package version: "5:29.8.2-1~debian.13~trixie" // → "29.8.2" (no epoch, no Debian revision). func EngineOf(pkgVersion string) string { v := pkgVersion if i := strings.Index(v, ":"); i >= 0 { v = v[i+1:] } if i := strings.Index(v, "-"); i >= 0 { v = v[:i] } return v } // DockerHealthVerdict is THE Docker-step health rule (`11` §5.8; pinned by TestDockerHealthVerdict): the guest rule, // plus every container running at the start still runs as the SAME container (same id — a changed id means the // household's apps restarted, which `live-restore` exists to prevent), plus the engine now reports the version the step // installed (wantEngine "" = no engine change expected). func DockerHealthVerdict(before, after *Health, wantEngine, gotEngine string) (bool, string) { if ok, why := HealthVerdict(before, after); !ok { return false, why } if before != nil { names := make([]string, 0, len(before.Containers)) for n := range before.Containers { names = append(names, n) } sort.Strings(names) for _, n := range names { b := before.Containers[n] if b.State != "running" || b.ID == "" { continue } if a := after.Containers[n]; a.ID != b.ID { return false, n + " is a new container (id changed) — the engine step restarted it" } } } if wantEngine != "" && gotEngine != wantEngine { return false, "the engine is " + gotEngine + ", not " + wantEngine } return true, "" } // 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 } 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"])) if err := os.WriteFile(path, b, 0o600); err != nil { return WrapperReport{}, fmt.Errorf("osupdate: write plan: %w", err) } defer os.Remove(path) stdout, stderr, err := l.Runner.Run(ctx, WrapperPath, "--plan", path) for _, line := range strings.Split(strings.TrimSpace(string(stderr)), "\n") { if strings.HasPrefix(line, "os-apply: ") { l.log().Info("osupdate: wrapper", "line", line) } } var rep WrapperReport found := false for _, line := range strings.Split(string(stdout), "\n") { if strings.HasPrefix(line, "OSAPPLY-REPORT ") { if jerr := json.Unmarshal([]byte(strings.TrimPrefix(line, "OSAPPLY-REPORT ")), &rep); jerr == nil { found = true } } } if !found { return rep, fmt.Errorf("osupdate: wrapper gave no report (err %v): %s", err, strings.TrimSpace(string(stderr))) } return rep, nil // a refusal / failure is IN the report (exit 2 / 3), not an error here } // Pass is one leg's reports; an empty Layer means the step did not run. type Pass struct { Guest, Host, Docker Report } // 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 { g, h := l.runFast(ctx, vmid, trigger) p := Pass{Guest: g, Host: h} if g.Outcome == "skipped" { return p } blk := l.Block() okStep := func(r Report) bool { return (r.Outcome == "applied" || r.Outcome == "nothing" || r.Outcome == "inventory") && r.Healthy } lg := l.log().With("run", g.RunID, "vmid", vmid, "trigger", trigger) switch { case blk.Ring != 0 || !blk.Enabled: lg.Info("osupdate: docker step skipped — ring 1 takes an engine set only inside a signed operator job (`11` §5.8)", "ring", blk.Ring, "enabled", blk.Enabled) case !okStep(g) || (h.Layer != "" && !okStep(h)): lg.Warn("osupdate: docker step skipped — an earlier step did not end healthy") default: if err := l.EnsureLiveRestore(ctx, g.RunID, vmid); err != nil { p.Docker = l.finish(ctx, lg, Report{RunID: g.RunID, Layer: LayerDocker, Trigger: trigger, Ring: 0, VMID: vmid, Mode: "apply", Outcome: "failed", HealthReason: "live-restore could not be turned on: " + err.Error()}) return p } p.Docker = l.runLayer(ctx, g.RunID, LayerDocker, vmid, trigger, blk, dockerOpts{}) } return p } // dockerOpts is a signed Docker step (DockerStepExecutor); the zero value is ring 0's unsigned "pending-docker". type dockerOpts struct { releaseID string packages []Package undo bool signed map[string]string // blob_b64, sig — the wrapper verifies them ITSELF } // EnsureLiveRestore is the ONE-TIME step of `09` decision 87: the wrapper merges `"live-restore": true` into the guest's // daemon.json and RELOADS docker (never a restart, R-835). A no-op when it is already on. func (l *Leg) EnsureLiveRestore(ctx context.Context, runID string, vmid int) error { wr, err := l.call(ctx, runID, map[string]any{"release_id": "live-restore", "layer": LayerGuest, "lane": "fast", "vmid": vmid, "mode": "live-restore-on", "packages": []Package{}}) if err != nil { return err } if wr.refused() { return fmt.Errorf("refused: %s", wr.Refused) } if wr.failed() { return fmt.Errorf("failed: %s", wr.Failed) } l.log().Info("osupdate: live-restore", "vmid", vmid, "result", string(wr.LiveRestore)) return nil } // Facts reads the box's versions through the wrapper's read-only facts mode (R-852): host Debian, kernels, held // packages, taint, the crash guard; guest Debian, Docker engine, containerd, live-restore. Raw JSON, the wrapper's shape. func (l *Leg) Facts(ctx context.Context, vmid int) (json.RawMessage, error) { wr, err := l.call(ctx, "facts"+l.now().UTC().Format("150405"), map[string]any{"release_id": "facts", "layer": LayerHost, "lane": "fast", "vmid": vmid, "mode": "facts", "packages": []Package{}}) if err != nil { return nil, err } if wr.refused() { return nil, fmt.Errorf("facts refused: %s", wr.Refused) } if len(wr.Facts) == 0 { return nil, fmt.Errorf("facts: the wrapper returned none (an older wrapper?)") } return wr.Facts, nil } // runFast is the guest + host fast lane (agent v0.141.x behaviour). func (l *Leg) runFast(ctx context.Context, vmid int, trigger string) (guest Report, host Report) { runID := l.now().UTC().Format("20060102T150405Z") lg := l.log().With("run", runID, "vmid", vmid, "trigger", trigger) if trigger == "night" && l.StatePath != "" { gap := l.MinGap if gap == 0 { gap = 20 * time.Hour } if b, err := os.ReadFile(l.StatePath); err == nil { if last, perr := time.Parse(time.RFC3339, strings.TrimSpace(string(b))); perr == nil && l.now().Sub(last) < gap { lg.Info("osupdate: skipped — already ran tonight", "last", last.UTC().Format(time.RFC3339)) return Report{RunID: runID, Layer: LayerGuest, Outcome: "skipped"}, Report{} } } } blk := l.Block() guest = l.runLayer(ctx, runID, LayerGuest, vmid, trigger, blk, dockerOpts{}) if trigger == "night" && l.StatePath != "" { _ = os.WriteFile(l.StatePath, []byte(l.now().UTC().Format(time.RFC3339)), 0o600) } switch { case !l.Appliance: lg.Info("osupdate: host step skipped — not an appliance install (a BYO host belongs to its owner, `11` §1)") case !(guest.Outcome == "applied" || guest.Outcome == "nothing" || guest.Outcome == "inventory") || !guest.Healthy: lg.Warn("osupdate: host step skipped — the guest step did not end healthy", "guest_outcome", guest.Outcome, "reason", guest.HealthReason) default: host = l.runLayer(ctx, runID, LayerHost, vmid, trigger, blk, dockerOpts{}) } return guest, host } func (l *Leg) runLayer(ctx context.Context, runID, layer string, vmid int, trigger string, blk hub.WireOSUpdate, do dockerOpts) Report { rel := hub.WireOSRelease{ID: "ring0-" + runID} var wire *hub.WireOSRelease switch layer { case LayerGuest: wire = blk.Release case LayerHost: wire = blk.HostRelease } lane := "fast" if layer == LayerDocker { lane = "slow" if do.signed != nil { rel = hub.WireOSRelease{ID: do.releaseID} } } else if blk.Ring == 1 { rel = hub.WireOSRelease{} if wire != nil { rel = *wire } } rep := Report{RunID: runID, Layer: layer, Trigger: trigger, Ring: blk.Ring, VMID: vmid, ReleaseID: rel.ID} lg := l.log().With("run", runID, "layer", layer, "vmid", vmid, "ring", blk.Ring, "trigger", trigger) lg.Info("osupdate: START", "enabled", blk.Enabled, "release", rel.ID) plan := map[string]any{"release_id": rel.ID, "layer": layer, "lane": lane, "vmid": vmid, "snapshot": rel.Snapshot, "packages": []Package{}, "mode": "apply", "select": "listed"} if rel.ID == "" { plan["release_id"] = "none" } planned := map[string]bool{} switch { case layer == LayerDocker && do.signed != nil: plan["packages"], plan["signed"] = do.packages, do.signed if do.undo { plan["undo"] = true } for _, p := range do.packages { planned[p.Name] = true } case layer == LayerDocker: plan["select"] = "pending-docker" // ring 0: the wrapper checks the box's ROOT-OWNED ring-0 mark itself case !blk.Enabled: plan["mode"] = "inventory" lg.Info("osupdate: switched OFF for this box — reporting only") case blk.Ring == 0: plan["select"] = "pending-fast" // the wrapper picks every pending Debian / Debian-Security upgrade case len(rel.Packages) == 0: plan["mode"] = "inventory" // ring 1 with no approved release for this layer: nothing to install default: var pk []Package for _, p := range rel.Packages { pk = append(pk, Package{Name: p.Name, Version: p.Version, Origin: p.Origin}) planned[p.Name] = true } plan["packages"] = pk } rep.Mode = plan["mode"].(string) wr, err := l.call(ctx, runID, plan) switch { case err != nil: rep.Outcome, rep.HealthReason = "failed", err.Error() return l.finish(ctx, lg, rep) case wr.refused(): rep.Outcome, rep.Refused = "refused", wr.Refused return l.finish(ctx, lg, rep) case wr.failed(): rep.Outcome, rep.Refused = "failed", wr.Failed } rep.Upgraded, rep.PassSeconds = wr.Upgraded, wr.PassSeconds rep.DockerEngine, rep.Authority, rep.Undo = wr.DockerEngine, wr.Authority, wr.Undo if rep.Outcome == "" { switch { case rep.Mode == "inventory" && !blk.Enabled: rep.Outcome = "inventory" case len(wr.Upgraded) == 0: rep.Outcome = "nothing" default: rep.Outcome = "applied" } } if blk.Ring == 0 { for _, u := range wr.Upgraded { planned[u.Name] = true } } // Health: compare with the start of the pass; give restarted services time (only after an install). cur := wr.HealthAfter wantEngine := "" for _, u := range wr.Upgraded { if u.Name == "docker-ce" { wantEngine = EngineOf(u.Version) } } verdict := func(h *Health) (bool, string) { if layer == LayerDocker { return DockerHealthVerdict(wr.HealthBefore, h, wantEngine, wr.DockerEngine) } if layer == LayerHost { t := hub.TunnelUnknown if l.Tunnel != nil { t, _ = l.Tunnel.Status(ctx) } return HostHealthVerdict(wr.HealthBefore, h, t) } return HealthVerdict(wr.HealthBefore, h) } if len(wr.Upgraded) > 0 { wait, poll := l.HealthWait, l.HealthPoll if wait == 0 { wait = 5 * time.Minute } if poll == 0 { poll = 15 * time.Second } deadline := l.now().Add(wait) for { ok, why := verdict(cur) rep.Healthy, rep.HealthReason = ok, why if ok || !l.now().Before(deadline) || ctx.Err() != nil { break } l.sleep(ctx, poll) hp := map[string]any{"release_id": plan["release_id"], "layer": layer, "lane": lane, "vmid": vmid, "mode": "health", "packages": []Package{}} hr, herr := l.call(ctx, runID, hp) if herr == nil && hr.Health != nil { cur = hr.Health } } if !rep.Healthy && rep.Outcome == "applied" { rep.Outcome = "health_failed" } } else { rep.Healthy, rep.HealthReason = verdict(cur) } rep.Installed, rep.Pending = wr.Installed, wr.Pending rep.RestartNeeded, rep.DockerRestartNeeded, rep.RebootNeeded = wr.RestartNeeded, wr.DockerRestartNeeded, wr.RebootNeeded rep.RebootScanned = wr.RebootScanned if layer == LayerDocker { // the docker report carries the engine set only (the guest report already carries the Debian packages) rep.Installed, rep.Pending = onlyDocker(wr.Installed), onlyDockerPending(wr.Pending) rep.NotCovered = nil } else { rep.NotCovered = notCovered(wr.Pending, blk.Ring, planned) } return l.finish(ctx, lg, rep) } func onlyDocker(in []Package) []Package { var out []Package for _, p := range in { if DockerNames[p.Name] { out = append(out, p) } } return out } func onlyDockerPending(in []Pending) []Pending { var out []Pending for _, p := range in { if DockerNames[p.Name] { out = append(out, p) } } return out } // notCovered lists pending updates no approved release covers: in ring 0 everything outside the fast lane; in ring 1 // also every fast-lane update the release did not name. func notCovered(pending []Pending, ring int, planned map[string]bool) []string { var out []string for _, p := range pending { if !IsFast(p.Origin) || (ring == 1 && !planned[p.Name]) { out = append(out, p.Name) } } return out } func (l *Leg) finish(ctx context.Context, lg *slog.Logger, rep Report) Report { lg.Info("osupdate: DONE", "outcome", rep.Outcome, "healthy", rep.Healthy, "reason", rep.HealthReason, "upgraded", len(rep.Upgraded), "pending", len(rep.Pending), "not_covered", len(rep.NotCovered), "restart_needed", len(rep.RestartNeeded), "reboot_needed", rep.RebootNeeded, "wrapper_seconds", rep.PassSeconds) if l.Hub != nil { body, _ := json.Marshal(rep) rctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), time.Minute) 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) } } return rep }