Files
felhom-agent/internal/osupdate/leg.go
T

460 lines
15 KiB
Go

// Package osupdate is the agent's OS-update leg for the customer GUEST (`11-os-updates.md` §8 step 2, agent v0.140.0).
//
// 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.
//
// NO AUTOMATIC UNDO in this release (R-837, measured 2026-10-04): a customer guest cannot be snapshotted — PVE
// refuses any snapshot not named `vzdump` when the guest has host-path binds (mp8, mp9). A failed health check
// therefore stops, reports `health_failed` and the hub mails the operator; the whole-guest backup taken minutes
// earlier is the undo, by hand.
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"
// 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 guest's 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
}
// Health is the guest's health snapshot (the wrapper's `health` object).
type Health struct {
DockerOK bool `json:"docker_ok"`
NetworkOK bool `json:"network_ok"`
Controller string `json:"controller"`
Containers map[string]Container `json:"containers"`
}
// WrapperReport is the wrapper's OSAPPLY-REPORT object.
type WrapperReport struct {
Mode string `json:"mode"`
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"`
HealthBefore *Health `json:"health_before"`
HealthAfter *Health `json:"health_after"`
Health *Health `json:"health"`
}
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 (hub osupdates.Report — field-exact).
type Report struct {
RunID string `json:"run_id"`
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"`
Refused json.RawMessage `json:"refused,omitempty"`
}
// 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.
type Leg struct {
Runner proxmox.Runner
Hub Reporter
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
have bool
}
// 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, l.have = resp.DesiredState.OSUpdate, true
}
// 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, l.have = b, true
}
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
}
func fastOrigin(origins []string) string {
for _, o := range origins {
if o == "Debian-Security" {
return o
}
}
return "Debian"
}
// HealthVerdict is THE health rule (`11` §5.4.1; pinned by TestHealthVerdict_*): after the run, docker answers, the
// guest's network resolves, the controller's own health check is `healthy`, and every container that was running
// before is running again — and healthy again if it was healthy before. "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, ""
}
// MergeBaseline is the health baseline from two readings: a container counts as running (and healthy) if EITHER
// reading saw it so. nil-safe. Pinned by TestHealth_BaselineIsTheStartOfTheLeg.
func MergeBaseline(a, b *Health) *Health {
if a == nil {
return b
}
if b == nil {
return a
}
out := &Health{DockerOK: a.DockerOK && b.DockerOK, NetworkOK: a.NetworkOK && b.NetworkOK, Controller: b.Controller,
Containers: map[string]Container{}}
for _, h := range []*Health{a, b} {
for n, c := range h.Containers {
cur, seen := out.Containers[n]
if !seen || (cur.State != "running" && c.State == "running") {
out.Containers[n] = c
continue
}
if c.State == "running" && c.Health == "healthy" {
out.Containers[n] = c
}
}
}
return out
}
// call writes the plan and runs the wrapper once.
func (l *Leg) call(ctx context.Context, runID, mode string, vmid int, rel hub.WireOSRelease, pkgs []Package) (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)
}
plan := map[string]any{"release_id": rel.ID, "layer": "guest", "lane": "fast", "vmid": vmid, "mode": mode,
"snapshot": rel.Snapshot, "packages": pkgs}
if pkgs == nil {
plan["packages"] = []Package{}
}
b, _ := json.Marshal(plan)
path := filepath.Join(dir, "plan-"+runID+"-"+mode+".json")
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
}
// Run is one pass: inventory → (switch, ring, plan) → apply → health → report. trigger is "night" or "debug".
func (l *Leg) Run(ctx context.Context, vmid int, trigger string) Report {
runID := l.now().UTC().Format("20060102T150405Z")
blk := l.Block()
rel := hub.WireOSRelease{ID: "ring0-" + runID}
if blk.Ring == 1 {
rel = hub.WireOSRelease{}
if blk.Release != nil {
rel = *blk.Release
}
}
rep := Report{RunID: runID, Trigger: trigger, Ring: blk.Ring, VMID: vmid, ReleaseID: rel.ID, Mode: "inventory"}
lg := l.log().With("run", runID, "vmid", vmid, "ring", blk.Ring, "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))
rep.Outcome = "skipped"
return rep
}
}
}
lg.Info("osupdate: START", "enabled", blk.Enabled, "release", rel.ID)
inv, err := l.call(ctx, runID, "inventory", vmid, rel, nil)
if err != nil || inv.refused() || inv.failed() {
rep.Outcome, rep.Refused = "refused", inv.Refused
if err != nil {
rep.Outcome, rep.HealthReason = "failed", err.Error()
}
return l.finish(ctx, lg, rep)
}
var plan []Package
switch {
case !blk.Enabled:
rep.Outcome = "inventory"
lg.Info("osupdate: switched OFF for this box — reporting only")
case blk.Ring == 0:
for _, p := range inv.Pending {
if IsFast(p.Origin) {
plan = append(plan, Package{Name: p.Name, Version: p.To, Origin: fastOrigin(p.Origin)})
}
}
default:
for _, p := range rel.Packages {
plan = append(plan, Package{Name: p.Name, Version: p.Version, Origin: p.Origin})
}
}
pkgNames := map[string]bool{}
for _, p := range plan {
pkgNames[p.Name] = true
}
if blk.Enabled && len(plan) == 0 {
rep.Outcome = "nothing"
}
final := inv
if blk.Enabled && len(plan) > 0 {
rep.Mode = "apply"
ap, err := l.call(ctx, runID, "apply", vmid, rel, plan)
switch {
case err != nil:
rep.Outcome, rep.HealthReason = "failed", err.Error()
return l.finish(ctx, lg, rep)
case ap.refused():
rep.Outcome, rep.Refused = "refused", ap.Refused
return l.finish(ctx, lg, rep)
case ap.failed():
rep.Outcome, rep.Refused = "failed", ap.Failed
}
final = ap
rep.Upgraded = ap.Upgraded
if rep.Outcome == "" {
if len(ap.Upgraded) == 0 {
rep.Outcome = "nothing"
} else {
rep.Outcome = "applied"
}
}
// Health: compare with what the guest looked like BEFORE the run; give restarted services time.
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)
cur := ap.HealthAfter
// The baseline is the guest as it was at the START of the leg (the inventory's reading) merged with the
// apply's own "before": an app that stops at any point during the run counts. Found live 2026-10-04: an app
// stopped between the inventory and the apply's own reading was taken as "stopped before" and ignored.
base := MergeBaseline(inv.HealthAfter, ap.HealthBefore)
for {
ok, why := HealthVerdict(base, cur)
rep.Healthy, rep.HealthReason = ok, why
if ok || !l.now().Before(deadline) || ctx.Err() != nil {
break
}
l.sleep(ctx, poll)
hr, herr := l.call(ctx, runID, "health", vmid, rel, nil)
if herr == nil && hr.Health != nil {
cur = hr.Health
}
}
if !rep.Healthy && rep.Outcome == "applied" {
rep.Outcome = "health_failed"
}
} else {
ok, why := HealthVerdict(nil, inv.HealthAfter)
rep.Healthy, rep.HealthReason = ok, why
}
rep.Installed, rep.Pending = final.Installed, final.Pending
rep.RestartNeeded, rep.DockerRestartNeeded, rep.RebootNeeded = final.RestartNeeded, final.DockerRestartNeeded, final.RebootNeeded
rep.NotCovered = notCovered(final.Pending, blk.Ring, pkgNames)
if trigger == "night" && l.StatePath != "" {
_ = os.WriteFile(l.StatePath, []byte(l.now().UTC().Format(time.RFC3339)), 0o600)
}
return l.finish(ctx, lg, rep)
}
// 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))
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
}