From 1fc1e74cc3c9702fb5819cdf243cb5bf9aebac88 Mon Sep 17 00:00:00 2001 From: jochen Date: Wed, 7 Oct 2026 01:51:59 +0200 Subject: [PATCH] Judge whether what a module runs stays up, and say it (hq ADR 0240, to-be 48 Phase A) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A container that crash-looped after its compose applied passed every check the gate had: nothing looked at what a module runs. The node-engine now judges every long-running resource on every look — one read of the runtime, one per service manager — keeps the restarts it counts across recreates and its own restarts, and says the state in every report and as an event on change, again every minute while not healthy. It reads only; nothing is restarted for being unhealthy. --- cmd/mesh-host/main.go | 99 +++++- internal/link/bus.go | 21 ++ internal/link/health_test.go | 52 +++ internal/link/hearing_nats.go | 4 + internal/link/messages.go | 68 ++++ internal/link/messages_test.go | 7 + internal/link/queue.go | 34 ++ internal/liveness/liveness.go | 508 +++++++++++++++++++++++++++++ internal/liveness/liveness_test.go | 332 +++++++++++++++++++ internal/liveness/replay_test.go | 91 ++++++ internal/liveness/runtime.go | 134 ++++++++ 11 files changed, 1346 insertions(+), 4 deletions(-) create mode 100644 internal/link/health_test.go create mode 100644 internal/liveness/liveness.go create mode 100644 internal/liveness/liveness_test.go create mode 100644 internal/liveness/replay_test.go create mode 100644 internal/liveness/runtime.go diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index ee2a401..e34ced0 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -34,6 +34,7 @@ import ( "github.com/novox/mesh-host/internal/identity" "github.com/novox/mesh-host/internal/inventory" "github.com/novox/mesh-host/internal/link" + "github.com/novox/mesh-host/internal/liveness" "github.com/novox/mesh-host/internal/outward" "github.com/novox/mesh-host/internal/profile" "github.com/novox/mesh-host/internal/reachable" @@ -1052,6 +1053,26 @@ func runLink(ctx context.Context, opts options) error { sched.RecordWindowsIn(apply.WindowsIn(filepath.Dir(opts.state))) go sched.Run(ctx) + // The one liveness judge on this machine (novox/hq ADR 0240), made before anything applies: every + // report this host makes carries what it says. + if j, err := liveness.Open(filepath.Join(filepath.Dir(opts.state), liveness.FileName), + &liveness.Exec{Run: apply.ExecRunner}); j != nil { + if err != nil { + say(err.Error()) + } + windows := apply.WindowsIn(filepath.Dir(opts.state)) + j.HeldNow = func(now time.Time) map[string]bool { + held := map[string]bool{} + for _, w := range windows.Open(now) { + for _, name := range w.Holds { + held[name] = true + } + } + return held + } + judging = j + } + // **Standing aside for a successor happens between reconciles and nowhere else** (novox/hq ADR // 0141). A host that stood aside mid-apply is the half-configured machine this host exists to // prevent, so the question is asked after an apply has finished and the answer is a clean exit — @@ -1140,6 +1161,11 @@ func runLink(ctx context.Context, opts options) error { queue.ReconcileDue() } go holdTheMachine(ctx, queue) + // And whether what it runs stays up, looked at on its own clock and said when it changes (novox/hq + // ADR 0240) — between applies, which is when a container crash-loops. + if judging != nil { + go judgeWhatRuns(aside, judging, queue, say) + } // And the core builds this host placed are judged, whichever host placed them (to-be 45 §8). go watchWhatThisHostPlaced(aside, mine.Node, queue, say) @@ -1326,6 +1352,64 @@ func rousedBySignal(ctx context.Context) link.Roused { // changed it, and then for ever. const ReconcileEvery = 5 * time.Minute +// judging is the serving host's liveness judge (novox/hq ADR 0240); nil in a one-shot command, which +// reports no health, and in a test. +var judging *liveness.Judge + +// How often a statement of health is said between reports: again every minute while anything is not +// healthy (to-be 48 §4), so a lost event is not a lost fault; and every five minutes anyway, so a +// controller that restarted knows a healthy machine's state without waiting for its next apply. +const ( + sayUnhealthyAgain = time.Minute + sayAnyway = 5 * time.Minute +) + +// judgeWhatRuns looks at every long-running resource every liveness.LookEvery and says the statement on +// the bus when a state changed, again every minute while one is not healthy, and every five minutes +// anyway. A statement the bus did not take is said at the next look. It reads; it never acts (ADR 0240 +// rule 6). +func judgeWhatRuns(ctx context.Context, j *liveness.Judge, queue *link.Queue, say link.Announce) { + ticker := time.NewTicker(liveness.LookEvery) + defer ticker.Stop() + var lastSaid time.Time + owed := true + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + st, changed := j.Look(ctx) + owed = owed || changed + since := time.Since(lastSaid) + if !owed && !(!st.Healthy() && since >= sayUnhealthyAgain) && since < sayAnyway { + continue + } + if changed { + for _, r := range st.Resources { + if r.State == liveness.Unhealthy { + say(fmt.Sprintf("%s (%s %s) is unhealthy: %s, %d restart(s) counted", r.ID, r.Kind, r.Target, + r.Reason, r.Restarts)) + } + } + } + if queue.SayHealth(ctx, *healthAsReported(st)) { + lastSaid, owed = time.Now(), false + } + } +} + +// healthAsReported is a statement as the report and the event carry it. +func healthAsReported(st liveness.Statement) *link.Health { + h := &link.Health{Contract: link.LivenessContract, At: st.At.UTC(), Resources: []link.ResourceHealth{}} + for _, r := range st.Resources { + h.Resources = append(h.Resources, link.ResourceHealth{Module: r.Module, Resource: r.ID, Kind: r.Kind, + Target: r.Target, State: r.State, Reason: r.Reason, Since: r.Since.UTC(), Streak: r.Streak, + Restarts: r.Restarts}) + } + return h +} + // holdTheMachine asks for a reconcile every ReconcileEvery. **It asks; it does not apply** (novox/hq // to-be 45 §6): the queue's worker does, when its turn comes, and a reconcile due while a delivery is // waiting is that delivery's apply. @@ -1541,17 +1625,24 @@ func applyAndKeepHeld(ctx context.Context, opts options, raw []byte, signed *sto // is armed, one whose image, environment or cadence changed is re-armed, and one no longer // declared is forgotten — and after a host restart the first apply rebuilds them all. A nil // scheduler is the one-shot CLI path, which exits rather than staying up to fire anything. + held := map[string]bool{} + for _, h := range updated.Held { + held[h.ID] = true + } if sched != nil { - held := map[string]bool{} - for _, h := range updated.Held { - held[h.ID] = true - } sched.Sync(declared, held) } report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw), Order: order, Host: runningVersion(), Rollbacks: standingRollbacks(opts.state), Witness: link.WitnessContract, Profile: profileAsReported(profile.Detect(ctx, profile.Default(nil), opts.timeout))} + // What runs here for a module, and how each is (novox/hq ADR 0240): judged from what was just + // applied, so every resource this apply started says `starting` and the gate waits for it. + if j := judging; j != nil { + j.Set(liveness.LongRunning(declared, held)) + st, _ := j.Look(ctx) + report.Health = healthAsReported(st) + } // Which of this machine's links face outside, for the filter the mesh writes around them // (novox/hq ADR 0140). Reported whatever the node's mode: a converged node's filter needs it, // and an adopted one becomes converged without a further round trip. A machine that cannot read diff --git a/internal/link/bus.go b/internal/link/bus.go index 540ff89..6266ed7 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -46,6 +46,27 @@ type OverNATS struct { func ReportSubject(node string) string { return "mesh.control." + node + ".report" } func AliveSubject(node string) string { return "mesh.control." + node + ".alive" } +// HealthSubject is where this node says its long-running resources' health between reports (novox/hq +// ADR 0240): its own, inside the grant every host already has, and on core NATS like the heartbeat — a +// statement lost is said again within a minute while anything is not healthy. +func HealthSubject(node string) string { return "mesh.control." + node + ".health" } + +// HealthBus is a link that can say a health statement. Separate from Bus so a link that cannot is still +// a link: the statement then waits for the next report, which carries it too. +type HealthBus interface { + Health(ctx context.Context, node string, body []byte) error +} + +// Health says a statement, core and flushed, as a heartbeat is. +func (b OverNATS) Health(ctx context.Context, node string, body []byte) error { + if err := b.Conn.Publish(HealthSubject(node), body); err != nil { + return err + } + flush, cancel := context.WithTimeout(ctx, 2*time.Second) + defer cancel() + return b.Conn.FlushWithContext(flush) +} + func (b OverNATS) Report(ctx context.Context, node string, body []byte) error { // Into the CONTROL stream and awaited: this is the message the store-window guarantee is // about (novox/hq ADR 0083). The controller naks with a delay while its store is away and diff --git a/internal/link/health_test.go b/internal/link/health_test.go new file mode 100644 index 0000000..38e999a --- /dev/null +++ b/internal/link/health_test.go @@ -0,0 +1,52 @@ +package link + +import ( + "context" + "encoding/json" + "testing" + "time" +) + +// A health statement is said on the link open now, as the machine's own word on its own subject, and +// not said — for the caller to say again — while no link is open (novox/hq ADR 0240, to-be 48 §4). + +// sayingLink is a link that can say health. +type sayingLink struct { + *quietLink + said [][]byte +} + +func (l *sayingLink) Health(_ context.Context, node string, body []byte) error { + l.said = append(l.said, body) + return nil +} + +func TestAHealthStatementIsSaidOnTheLinkOpenNow(t *testing.T) { + q := &Queue{Membership: Membership{Node: "anchor"}} + h := Health{Contract: LivenessContract, At: time.Now().UTC(), + Resources: []ResourceHealth{{Module: "letta", Resource: "letta.server", State: StateUnhealthy, Reason: ReasonRestarting}}} + if q.SayHealth(t.Context(), h) { + t.Fatal("said with no link open") + } + l := &sayingLink{quietLink: newQuietLink()} + q.attach(t.Context(), l) + if !q.SayHealth(t.Context(), h) || len(l.said) != 1 { + t.Fatalf("not said on the open link: %d", len(l.said)) + } + var got HealthSaid + if err := json.Unmarshal(l.said[0], &got); err != nil { + t.Fatal(err) + } + if got.Node != "anchor" || len(got.Health.Resources) != 1 || got.Health.Resources[0].State != StateUnhealthy { + t.Fatalf("said %+v", got) + } + // A link that cannot say it is not an error: the next report carries the statement. + q.detach(l) + q.attach(t.Context(), newQuietLink()) + if q.SayHealth(t.Context(), h) { + t.Fatal("said on a link that cannot") + } + if HealthSubject("anchor") != "mesh.control.anchor.health" { + t.Fatalf("the subject %s is not inside the host's own grant", HealthSubject("anchor")) + } +} diff --git a/internal/link/hearing_nats.go b/internal/link/hearing_nats.go index 3207772..32e8623 100644 --- a/internal/link/hearing_nats.go +++ b/internal/link/hearing_nats.go @@ -216,6 +216,10 @@ func (l *natsLink) Alive(ctx context.Context, node string, body []byte) error { return OverNATS{Conn: l.conn, JS: l.js}.Alive(ctx, node, body) } +func (l *natsLink) Health(ctx context.Context, node string, body []byte) error { + return OverNATS{Conn: l.conn, JS: l.js}.Health(ctx, node, body) +} + // natsDeclaration is one declaration off the NODES stream. type natsDeclaration struct{ msg *nats.Msg } diff --git a/internal/link/messages.go b/internal/link/messages.go index 9242756..cfd0eaf 100644 --- a/internal/link/messages.go +++ b/internal/link/messages.go @@ -208,6 +208,74 @@ type Report struct { // older than it refuses — so the controller sends either only to a machine whose report carries // it, as it sends `epoch` only where `report_sequence` was seen (ADR 0229). Witness int `json:"witness,omitempty"` + + // Health is this node-engine's word on every long-running resource it runs for a module (novox/hq + // ADR 0240, to-be 48 §4): its state, since when, its failing streak and the restarts it counted. + // **Its presence says this engine judges**: a health with no resources is a machine that runs nothing + // long-lived, and no health at all is an engine older than the judging — which the controller reads + // as "not known", never as healthy and never as a reason to raise anything. Fresh on every report + // this engine makes, and said between reports on HealthSubject when it changes. + Health *Health `json:"health,omitempty"` +} + +// LivenessContract is the version of the health statement this host keeps (ADR 0240 Phase A: liveness, +// judged with no declaration). +const LivenessContract = 1 + +// Health is one statement of every long-running resource's state on this machine (to-be 48 §4). Said in +// every report, as an event on each change, and again every minute while anything is not healthy — so a +// lost statement is not a lost fault. +type Health struct { + // Contract is LivenessContract. + Contract int `json:"contract"` + // At is when the engine looked, on this machine's clock: the order of its statements. A controller + // keeps the newest it heard and refuses an older one arriving late. + At time.Time `json:"at"` + // Resources are every long-running resource of a module this machine runs, by module and id. + Resources []ResourceHealth `json:"resources"` +} + +// The states a long-running resource is said in (ADR 0240 §4). +const ( + StateHealthy = "healthy" + StateUnhealthy = "unhealthy" + // StateStarting is inside its grace after a start: not yet a verdict either way. + StateStarting = "starting" + // StateHeld is held still by an open maintenance window (ADR 0189): neither alive nor dead. + StateHeld = "held" + // StateUnknown is a resource whose state could not be read. + StateUnknown = "unknown" +) + +// The reasons an unhealthy resource is said with, when it is liveness that judged it. +const ( + ReasonDown = "down" + ReasonRestarting = "restarting" +) + +// ResourceHealth is one long-running resource's state (to-be 48 §4). +type ResourceHealth struct { + // Module owns the resource; Resource is its id in the declaration. + Module string `json:"module"` + Resource string `json:"resource"` + // Kind is container, service or process; Target the container's name or the unit. + Kind string `json:"kind"` + Target string `json:"target"` + State string `json:"state"` + // Reason is why it is not healthy: down, restarting, or what could not be read. + Reason string `json:"reason,omitempty"` + Since time.Time `json:"since"` + // Streak is how many looks in a row have found it unhealthy. + Streak int `json:"streak,omitempty"` + // Restarts is how many restarts the engine counted after a grace, kept across recreates and across + // its own restarts: the runtime's own count is lost on every recreate. + Restarts int `json:"restarts,omitempty"` +} + +// HealthSaid is the health event: a machine's statement between its reports, on HealthSubject. +type HealthSaid struct { + Node string `json:"node"` + Health Health `json:"health"` } // WitnessContract is the version of the witness contract this host keeps: the fields above, the diff --git a/internal/link/messages_test.go b/internal/link/messages_test.go index a1b242d..d3aa3a7 100644 --- a/internal/link/messages_test.go +++ b/internal/link/messages_test.go @@ -134,6 +134,13 @@ func TestTheWireFormatIsExactlyTheseFieldNames(t *testing.T) { []string{"node", "rollbacks", "witness"}}, {Rollback{Component: ComponentNodeTools, From: "sha256:b", To: "sha256:a", Outcome: RolledBack, Why: "w"}, []string{"component", "from", "to", "outcome", "why", "at"}}, + // novox/hq ADR 0240: every long-running resource's health, in every report and in its own event. + {Report{Node: "n", Health: &Health{Contract: LivenessContract}}, []string{"node", "health"}}, + {Health{Contract: LivenessContract, Resources: []ResourceHealth{}}, []string{"contract", "at", "resources"}}, + {ResourceHealth{Module: "m", Resource: "m.r", Kind: "container", Target: "t", State: StateUnhealthy, + Reason: ReasonRestarting, Streak: 2, Restarts: 3}, + []string{"module", "resource", "kind", "target", "state", "reason", "since", "streak", "restarts"}}, + {HealthSaid{Node: "n"}, []string{"node", "health"}}, } { raw, err := json.Marshal(c.value) if err != nil { diff --git a/internal/link/queue.go b/internal/link/queue.go index e50899e..3efa30d 100644 --- a/internal/link/queue.go +++ b/internal/link/queue.go @@ -477,3 +477,37 @@ func declaredIn(body []byte) string { sum := sha256.Sum256(signed.Declaration) return hex.EncodeToString(sum[:]) } + +// SayHealth says a health statement on the link open now (novox/hq ADR 0240, to-be 48 §4), and whether +// the bus took it. **Not through the worker**: a statement is a word about the machine as it is, not an +// account of an apply, and it must not wait behind one — a container crash-looping while a long apply +// runs is said while it runs. False with no link open, or one that cannot say it; the caller says it +// again on its next look, and the next report carries it whatever happens. +func (q *Queue) SayHealth(ctx context.Context, h Health) bool { + q.init() + q.mu.Lock() + bus := q.bus + q.mu.Unlock() + sayer, ok := bus.(HealthBus) + if bus == nil || !ok { + return false + } + body, err := json.Marshal(HealthSaid{Node: q.Membership.Node, Health: h}) + if err != nil { + return false + } + saying, cancel := context.WithTimeout(ctx, q.timeout()) + defer cancel() + if err := sayer.Health(saying, q.Membership.Node, body); err != nil { + q.Say("could not say this machine's health: " + err.Error()) + return false + } + return true +} + +func (q *Queue) timeout() time.Duration { + if q.Timeout <= 0 { + return 10 * time.Second + } + return q.Timeout +} diff --git a/internal/liveness/liveness.go b/internal/liveness/liveness.go new file mode 100644 index 0000000..86f773f --- /dev/null +++ b/internal/liveness/liveness.go @@ -0,0 +1,508 @@ +// Package liveness is the node-engine judging whether what a module runs stays up (novox/hq ADR 0240 +// rule 1, to-be 48 §1 and §4, Phase A). +// +// **A module's container that crash-looped behind every passing check** is what this exists for: the +// agent server restarted about a hundred times while the mesh read it applied and its tools served, and a +// person found it reading its log for something else (research 032). The release gate judged a module by +// what the mesh saw from outside; nothing looked at what the module runs. +// +// Every long-running resource — a container that stays up, a process the mesh runs that stays up, a +// service stated `running` — is judged on every look, with no declaration: +// +// - **alive** when it is running and has not restarted twice within the settle window after its grace; +// - **unhealthy** when it is not running after its grace (`down`), or restarted twice in the window +// (`restarting`); +// - **starting** inside its grace after a start — a new build, a recreate, a restart the engine or a +// person made — in which a restart is not counted: churn that stops is not a crash loop (issue 058); +// - **held** while an open maintenance window holds it still (ADR 0189), and judged from a fresh grace +// when the window ends; +// - **unknown** when nothing could be read. +// +// **The engine counts restarts itself and keeps the count**, in a file beside the node's state, across a +// recreate of the container and across its own restarts: the runtime's count is lost on every recreate +// and its event history is a minute long on a busy machine (research 032 §2). Neither is read for a +// verdict. +// +// **It reads; it never acts** (ADR 0240 rule 6). One read of every container's state per look, one of +// every unit's per service manager — tens of milliseconds on the busiest machine — and never an +// execution inside a container. Nothing here restarts, recreates or stops anything. +package liveness + +import ( + "context" + "encoding/json" + "fmt" + "os" + "path/filepath" + "sort" + "strings" + "sync" + "time" + + "github.com/novox/mesh-host/internal/declaration" +) + +// The bounds of ADR 0240 rule 1 and to-be 48 §1. Grace is the default a declaration will override in +// Phase B; the settle window is the gate's bound. +const ( + DefaultGrace = 60 * time.Second + SettleWindow = 10 * time.Minute + // LookEvery is how often the engine looks. A crash loop is said within the grace, two restarts and a + // look — well inside the gate's bound — and a look is one read of the runtime. + LookEvery = 15 * time.Second +) + +// The kinds of long-running resource. +const ( + KindContainer = "container" + KindService = "service" + KindProcess = "process" +) + +// The states, as the report says them (mesh-host internal/link and the controller agree on the words). +const ( + Healthy = "healthy" + Unhealthy = "unhealthy" + Starting = "starting" + Held = "held" + Unknown = "unknown" + + ReasonDown = "down" + ReasonRestarting = "restarting" +) + +// Resource is one long-running resource of a module, as the judge needs it. +type Resource struct { + Module string `json:"module"` + ID string `json:"id"` + Kind string `json:"kind"` + // Target is the container's name, or the unit. + Target string `json:"target"` + // Scope and User are a service's manager: "user" and the account for a unit in an account's own. + Scope string `json:"scope,omitempty"` + User string `json:"user,omitempty"` +} + +// LongRunning is every long-running resource a declaration asks this machine to run for a module: its +// containers that stay up, its services stated running and its processes that stay up. A step, anything +// on a schedule, a service whose lifecycle is the machine's, what the mesh declares in its own right (no +// module) and what an adopted machine holds as it was found are not judged here. +func LongRunning(d *declaration.Declaration, held map[string]bool) []Resource { + var out []Resource + if d == nil { + return nil + } + for _, r := range d.Resources { + module, ok := ModuleOf(r.Identity()) + if !ok || held[r.Identity()] { + continue + } + switch v := r.(type) { + case *declaration.Container: + if v.RunOnce || v.Schedule != "" { + continue + } + out = append(out, Resource{Module: module, ID: v.ID, Kind: KindContainer, Target: v.Name}) + case *declaration.Service: + if v.State != "running" { + continue + } + res := Resource{Module: module, ID: v.ID, Kind: KindService, Target: v.Unit} + if v.UserScoped() { + res.Scope, res.User = declaration.ScopeUser, v.User + } + out = append(out, res) + case *declaration.Process: + if v.RunOnce || v.Schedule != "" { + continue + } + out = append(out, Resource{Module: module, ID: v.ID, Kind: KindProcess, Target: v.Name + ".service"}) + } + } + return out +} + +// ModuleOf is the module a resource's id belongs to: everything before its last dot — a module's name may +// carry a dot, its resources' own ids never do (the apply's moduleOf, for the same reason). False for +// what the mesh declares in its own right. +func ModuleOf(id string) (string, bool) { + if strings.HasPrefix(id, declaration.AdoptionPrefix) { + return "", false + } + at := strings.LastIndex(id, ".") + if at <= 0 { + return "", false + } + return id[:at], true +} + +// Observed is one resource as one read found it. +type Observed struct { + // Found is false for a container that is not there and a unit the manager does not know. + Found bool + // Identity is what changes when the thing is made again: the container's id, the unit's invocation. + Identity string + Running bool + // Restarting is the runtime or the manager restarting it after it exited. + Restarting bool + // Restarts is the runtime's own count of the restarts it made — read only to see it move. + Restarts int64 + // Started is when the runtime says the current run started, as it says it: a change with no restart + // counted is a restart somebody made, which is a new start. + Started string +} + +// Runtime is what one look reads: every container's state in one read, and every unit's per manager. +// An error is "could not be read" — said as unknown, never as down. +type Runtime interface { + Containers(ctx context.Context, names []string) (map[string]Observed, error) + Units(ctx context.Context, scope, user string, units []string) (map[string]Observed, error) +} + +// State is one resource's state as the judge says it. +type State struct { + Resource + State string `json:"state"` + Reason string `json:"reason,omitempty"` + Since time.Time `json:"since"` + Streak int `json:"streak,omitempty"` + Restarts int `json:"restarts,omitempty"` +} + +// Statement is one look at every long-running resource: when, and each resource's state. +type Statement struct { + At time.Time + Resources []State +} + +// Healthy says every resource in it is healthy. +func (s Statement) Healthy() bool { + for _, r := range s.Resources { + if r.State != Healthy { + return false + } + } + return true +} + +// kept is what the judge keeps about one resource, on disk, across its own restarts. +type kept struct { + Resource + // Seen says a read has found it at least once; Identity, RuntimeRestarts and RuntimeStarted are what + // the last read found, to see a restart or a recreate by. + Seen bool `json:"seen,omitempty"` + Identity string `json:"identity,omitempty"` + RuntimeRestarts int64 `json:"runtime-restarts,omitempty"` + RuntimeStarted string `json:"runtime-started,omitempty"` + // Started is when the current start began: its grace is counted from here. + Started time.Time `json:"started"` + // Counted is every restart counted after a grace, kept across recreates; Recent those inside the + // settle window since the current start. + Counted int `json:"counted,omitempty"` + Recent []time.Time `json:"recent,omitempty"` + WasHeld bool `json:"held,omitempty"` + + State string `json:"state,omitempty"` + Reason string `json:"reason,omitempty"` + Since time.Time `json:"since"` + Streak int `json:"streak,omitempty"` +} + +// file is the judge's file beside the node's state. +type file struct { + Resources []Resource `json:"resources"` + Kept map[string]*kept `json:"kept"` +} + +// FileName is the judge's file, beside the node's state. +const FileName = "liveness.json" + +// Judge is the one judge of liveness on a machine. Safe for the apply and the looking loop at once. +type Judge struct { + path string + runtime Runtime + // Now, Grace and Settle are the clock and the bounds; replaced in tests. + Now func() time.Time + Grace time.Duration + Settle time.Duration + // HeldNow is the containers an open maintenance window holds still, by runtime name. Nil holds none. + HeldNow func(now time.Time) map[string]bool + + mu sync.Mutex + f file + dirty bool + // said is the statement said last, to know a change by. + said map[string]string +} + +// Open is the judge whose file is at path, reading what it kept. A file that cannot be read is said and +// started afresh: a count lost is a crash loop judged from now, never a machine left unjudged. +func Open(path string, rt Runtime) (*Judge, error) { + j := &Judge{path: path, runtime: rt, Now: time.Now, Grace: DefaultGrace, Settle: SettleWindow, + f: file{Kept: map[string]*kept{}}, said: map[string]string{}} + raw, err := os.ReadFile(path) + switch { + case os.IsNotExist(err): + return j, nil + case err != nil: + return j, fmt.Errorf("the liveness kept at %s cannot be read, so restarts are counted from now: %w", path, err) + } + var f file + if err := json.Unmarshal(raw, &f); err != nil { + return j, fmt.Errorf("the liveness kept at %s cannot be read, so restarts are counted from now: %w", path, err) + } + if f.Kept == nil { + f.Kept = map[string]*kept{} + } + j.f = f + return j, nil +} + +// Set is what this machine runs now, from the declaration the apply just applied. A resource no longer +// declared is forgotten; one newly declared is judged from its first look. +func (j *Judge) Set(resources []Resource) { + j.mu.Lock() + defer j.mu.Unlock() + sorted := append([]Resource(nil), resources...) + sort.Slice(sorted, func(a, b int) bool { + if sorted[a].Module != sorted[b].Module { + return sorted[a].Module < sorted[b].Module + } + return sorted[a].ID < sorted[b].ID + }) + declared := map[string]bool{} + for _, r := range sorted { + declared[r.ID] = true + if k, ok := j.f.Kept[r.ID]; ok && (k.Kind != r.Kind || k.Target != r.Target || k.Scope != r.Scope || k.User != r.User) { + // The same id now names another thing: judged as a new one, its count kept. + counted := k.Counted + j.f.Kept[r.ID] = &kept{Resource: r, Counted: counted} + } + } + for id := range j.f.Kept { + if !declared[id] { + delete(j.f.Kept, id) + } + } + j.f.Resources = sorted + j.dirty = true +} + +// Look reads every long-running resource once, judges each, keeps what it counted, and answers the +// statement and whether any resource's state or reason changed since the last look. +func (j *Judge) Look(ctx context.Context) (Statement, bool) { + j.mu.Lock() + defer j.mu.Unlock() + now := j.Now() + var held map[string]bool + if j.HeldNow != nil { + held = j.HeldNow(now) + } + observed, unread := j.read(ctx) + + st := Statement{At: now} + changed := false + seen := map[string]bool{} + for _, r := range j.f.Resources { + k := j.f.Kept[r.ID] + if k == nil { + k = &kept{Resource: r} + j.f.Kept[r.ID] = k + } + k.Resource = r + why, blind := unread[groupOf(r)] + o := observed[keyOf(r)] + j.judge(k, o, now, r.Kind == KindContainer && held[r.Target], blind, why) + seen[r.ID] = true + st.Resources = append(st.Resources, State{Resource: r, State: k.State, Reason: k.Reason, Since: k.Since, + Streak: k.Streak, Restarts: k.Counted}) + word := k.State + "/" + k.Reason + if j.said[r.ID] != word { + changed = true + j.said[r.ID] = word + } + } + for id := range j.said { + if !seen[id] { + delete(j.said, id) + changed = true + } + } + j.save() + return st, changed +} + +// judge is one resource's verdict on one look. +func (j *Judge) judge(k *kept, o Observed, now time.Time, held, blind bool, why string) { + set := func(state, reason string) { + if k.State != state || k.Reason != reason || k.Since.IsZero() { + k.State, k.Reason, k.Since = state, reason, now + } + if state == Unhealthy { + k.Streak++ + } else { + k.Streak = 0 + } + j.dirty = true + } + fresh := func(o Observed, started time.Time) { + k.Seen, k.Identity, k.RuntimeRestarts, k.RuntimeStarted = o.Found, o.Identity, o.Restarts, o.Started + k.Started, k.Recent = started, nil + } + + // **Held is neither alive nor dead**, and the window ending is a start: judged from a fresh grace. + if held { + k.WasHeld = true + set(Held, "") + return + } + if k.WasHeld { + k.WasHeld = false + fresh(o, now) + } + if blind { + set(Unknown, why) + return + } + + switch { + case !k.Seen: + // First sight: of a resource just applied, or of one running before this engine judged. Grace is + // counted from when the runtime says it started, where it says, so an engine restarted beside + // a long-running container does not call it starting for a minute. + started := now + if t, err := time.Parse(time.RFC3339Nano, o.Started); err == nil && t.Before(now) && t.Year() > 1 { + started = t + } + switch { + case o.Found: + fresh(o, started) + case k.Started.IsZero(): + // Not there at its first look: its grace runs from now, and it is down after it. + k.Started = now + } + case o.Found && o.Identity != "" && o.Identity != k.Identity && o.Restarts <= k.RuntimeRestarts, + o.Found && o.Identity == k.Identity && o.Restarts == k.RuntimeRestarts && o.Started != k.RuntimeStarted: + // Made again — recreated by an apply, or restarted by somebody — and not by the runtime's policy: + // a new start, judged from a fresh grace. The count is kept. + fresh(o, now) + case o.Found && o.Restarts > k.RuntimeRestarts: + // The runtime restarted it after it exited. Counted only after its grace. + delta := int(o.Restarts - k.RuntimeRestarts) + k.Identity, k.RuntimeRestarts, k.RuntimeStarted = o.Identity, o.Restarts, o.Started + if !now.Before(k.Started.Add(j.Grace)) { + k.Counted += delta + for i := 0; i < delta; i++ { + k.Recent = append(k.Recent, now) + } + } + j.dirty = true + } + // Restarts older than the settle window no longer say anything. + recent := k.Recent[:0] + for _, t := range k.Recent { + if now.Sub(t) <= j.Settle { + recent = append(recent, t) + } + } + k.Recent = recent + + inGrace := now.Before(k.Started.Add(j.Grace)) + switch { + case len(k.Recent) >= 2: + set(Unhealthy, ReasonRestarting) + case inGrace: + set(Starting, "") + case !o.Found || !o.Running: + if o.Restarting { + set(Unhealthy, ReasonRestarting) + } else { + set(Unhealthy, ReasonDown) + } + default: + set(Healthy, "") + } +} + +// read is one read of everything: every container at once, every unit per manager. unread names each +// group that could not be read, with why. +func (j *Judge) read(ctx context.Context) (map[string]Observed, map[string]string) { + observed := map[string]Observed{} + unread := map[string]string{} + groups := map[string][]Resource{} + var order []string + for _, r := range j.f.Resources { + g := groupOf(r) + if _, ok := groups[g]; !ok { + order = append(order, g) + } + groups[g] = append(groups[g], r) + } + for _, g := range order { + rs := groups[g] + targets := make([]string, 0, len(rs)) + for _, r := range rs { + targets = append(targets, r.Target) + } + var found map[string]Observed + var err error + if rs[0].Kind == KindContainer { + found, err = j.runtime.Containers(ctx, targets) + } else { + found, err = j.runtime.Units(ctx, rs[0].Scope, rs[0].User, targets) + } + if err != nil { + unread[g] = firstLine(err.Error()) + continue + } + for _, r := range rs { + observed[keyOf(r)] = found[r.Target] + } + } + return observed, unread +} + +// groupOf is the one read a resource is in: the containers, or one service manager's units. +func groupOf(r Resource) string { + if r.Kind == KindContainer { + return KindContainer + } + return "units:" + r.Scope + ":" + r.User +} + +func keyOf(r Resource) string { return groupOf(r) + "/" + r.Target } + +// save keeps what was judged, when anything moved. Written whole and renamed, so a reader never sees +// half of it; a failure is said by the next Open, which counts from then. +func (j *Judge) save() { + if !j.dirty || j.path == "" { + return + } + raw, err := json.Marshal(j.f) + if err != nil { + return + } + if err := os.MkdirAll(filepath.Dir(j.path), 0o700); err != nil { + return + } + tmp, err := os.CreateTemp(filepath.Dir(j.path), ".liveness-*.json") + if err != nil { + return + } + defer os.Remove(tmp.Name()) + if _, err := tmp.Write(raw); err != nil { + tmp.Close() + return + } + if err := tmp.Close(); err != nil { + return + } + if os.Rename(tmp.Name(), j.path) == nil { + j.dirty = false + } +} + +func firstLine(s string) string { + line, _, _ := strings.Cut(strings.TrimSpace(s), "\n") + return line +} diff --git a/internal/liveness/liveness_test.go b/internal/liveness/liveness_test.go new file mode 100644 index 0000000..b304706 --- /dev/null +++ b/internal/liveness/liveness_test.go @@ -0,0 +1,332 @@ +package liveness + +import ( + "context" + "errors" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/novox/mesh-host/internal/declaration" +) + +// The judge, over a fake runtime and service manager (novox/hq ADR 0240, "how it is checked", rule 1): +// a container recreated keeps its counted restarts, and so does the engine restarted; a restart inside +// grace is not counted; two restarts within the settle window after grace make it unhealthy; a resource +// under a maintenance step is held. + +// fakeRuntime is what a look reads, set by the test. +type fakeRuntime struct { + containers map[string]Observed + units map[string]Observed + err error +} + +func (f *fakeRuntime) Containers(_ context.Context, names []string) (map[string]Observed, error) { + if f.err != nil { + return nil, f.err + } + out := map[string]Observed{} + for _, n := range names { + if o, ok := f.containers[n]; ok { + out[n] = o + } + } + return out, nil +} + +func (f *fakeRuntime) Units(_ context.Context, _, _ string, units []string) (map[string]Observed, error) { + out := map[string]Observed{} + for _, u := range units { + out[u] = f.units[u] + } + return out, nil +} + +// clock is a time the test moves. +type clock struct{ now time.Time } + +func (c *clock) Now() time.Time { return c.now } + +var t0 = time.Date(2026, 10, 7, 12, 0, 0, 0, time.UTC) + +func aJudge(t *testing.T, rt Runtime, c *clock, path string) *Judge { + t.Helper() + j, err := Open(path, rt) + if err != nil { + t.Fatal(err) + } + j.Now = c.Now + return j +} + +var server = Resource{Module: "letta", ID: "letta.server", Kind: KindContainer, Target: "letta-server"} + +func running(id string, restarts int64, started string) Observed { + return Observed{Found: true, Identity: id, Running: true, Restarts: restarts, Started: started} +} + +func stateOf(t *testing.T, st Statement, id string) State { + t.Helper() + for _, r := range st.Resources { + if r.ID == id { + return r + } + } + t.Fatalf("%s is not in the statement: %+v", id, st) + return State{} +} + +func TestARestartInsideGraceIsNotCountedAndTwoAfterItAreUnhealthy(t *testing.T) { + c := &clock{now: t0} + rt := &fakeRuntime{containers: map[string]Observed{"letta-server": running("c1", 0, "a")}} + j := aJudge(t, rt, c, filepath.Join(t.TempDir(), FileName)) + j.Set([]Resource{server}) + + st, changed := j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Starting || !changed { + t.Fatalf("a fresh start is starting: %+v", s) + } + + // Churn that stops (issue 058): restarted inside its grace, then up. + c.now = t0.Add(20 * time.Second) + rt.containers["letta-server"] = running("c1", 2, "b") + j.Look(t.Context()) + c.now = t0.Add(70 * time.Second) + st, _ = j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Healthy || s.Restarts != 0 { + t.Fatalf("restarts inside grace counted, or not healthy after it: %+v", s) + } + + // One restart after grace is not yet a crash loop. + c.now = t0.Add(2 * time.Minute) + rt.containers["letta-server"] = running("c1", 3, "c") + st, _ = j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Healthy || s.Restarts != 1 { + t.Fatalf("one restart after grace: %+v", s) + } + + // A second inside the settle window is. + c.now = t0.Add(3 * time.Minute) + rt.containers["letta-server"] = Observed{Found: true, Identity: "c1", Restarting: true, Restarts: 4, Started: "d"} + st, changed = j.Look(t.Context()) + s := stateOf(t, st, "letta.server") + if s.State != Unhealthy || s.Reason != ReasonRestarting || s.Restarts != 2 || s.Streak != 1 || !changed { + t.Fatalf("two restarts within the settle window after grace: %+v", s) + } + c.now = t0.Add(3*time.Minute + 15*time.Second) + st, changed = j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.Streak != 2 || changed || !s.Since.Equal(t0.Add(3*time.Minute)) { + t.Fatalf("the streak and since of a state that holds: %+v (changed %v)", s, changed) + } + + // Past the settle window with no restart, it is alive again. + c.now = t0.Add(14 * time.Minute) + rt.containers["letta-server"] = running("c1", 4, "d") + st, _ = j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Healthy || s.Restarts != 2 { + t.Fatalf("a crash loop that stopped ten minutes ago: %+v", s) + } +} + +func TestARecreatedContainerKeepsItsCountedRestartsAndStartsAgain(t *testing.T) { + c := &clock{now: t0} + rt := &fakeRuntime{containers: map[string]Observed{"letta-server": running("c1", 0, "a")}} + j := aJudge(t, rt, c, filepath.Join(t.TempDir(), FileName)) + j.Set([]Resource{server}) + j.Look(t.Context()) + for i, at := range []time.Duration{2 * time.Minute, 3 * time.Minute} { + c.now = t0.Add(at) + rt.containers["letta-server"] = running("c1", int64(i+1), "x"+at.String()) + j.Look(t.Context()) + } + + // The apply recreates it: the runtime's count is back to nothing, the engine's is not. + c.now = t0.Add(4 * time.Minute) + rt.containers["letta-server"] = running("c2", 0, "new") + st, _ := j.Look(t.Context()) + s := stateOf(t, st, "letta.server") + if s.State != Starting || s.Restarts != 2 { + t.Fatalf("a recreate is a new start, with its count kept: %+v", s) + } + c.now = t0.Add(6 * time.Minute) + st, _ = j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Healthy || s.Restarts != 2 { + t.Fatalf("the new build, up after its grace, is healthy — the old build's restarts are not its: %+v", s) + } +} + +func TestTheEngineRestartedKeepsWhatItCounted(t *testing.T) { + path := filepath.Join(t.TempDir(), FileName) + c := &clock{now: t0} + rt := &fakeRuntime{containers: map[string]Observed{"letta-server": running("c1", 0, "a")}} + j := aJudge(t, rt, c, path) + j.Set([]Resource{server}) + j.Look(t.Context()) + c.now = t0.Add(2 * time.Minute) + rt.containers["letta-server"] = running("c1", 1, "b") + j.Look(t.Context()) + + // Another engine — this one restarted, or its successor — reads the file and goes on. + again := aJudge(t, rt, c, path) + c.now = t0.Add(2*time.Minute + 30*time.Second) + rt.containers["letta-server"] = running("c1", 2, "c") + st, _ := again.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Unhealthy || s.Restarts != 2 { + t.Fatalf("the restarted engine forgot what it counted: %+v", s) + } +} + +func TestAContainerHeldByAWindowIsHeldAndJudgedAfreshAfter(t *testing.T) { + c := &clock{now: t0} + rt := &fakeRuntime{containers: map[string]Observed{"letta-server": running("c1", 0, "a")}} + j := aJudge(t, rt, c, filepath.Join(t.TempDir(), FileName)) + windowOpen := true + j.HeldNow = func(time.Time) map[string]bool { return map[string]bool{"letta-server": windowOpen} } + j.Set([]Resource{server}) + c.now = t0.Add(5 * time.Minute) + rt.containers["letta-server"] = Observed{Found: true, Identity: "c1", Restarts: 0, Started: "a"} + st, _ := j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Held { + t.Fatalf("a container a window holds still is held, neither alive nor dead: %+v", s) + } + windowOpen = false + rt.containers["letta-server"] = running("c1", 0, "after") + st, _ = j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Starting { + t.Fatalf("the window ended is a start, judged from a fresh grace: %+v", s) + } +} + +func TestDownAfterGraceAndUnknownWhenNothingCouldBeRead(t *testing.T) { + c := &clock{now: t0} + rt := &fakeRuntime{containers: map[string]Observed{}} + j := aJudge(t, rt, c, filepath.Join(t.TempDir(), FileName)) + j.Set([]Resource{server}) + if s := stateOf(t, func() Statement { st, _ := j.Look(t.Context()); return st }(), "letta.server"); s.State != Starting { + t.Fatalf("not there at its first look is still in its grace: %+v", s) + } + c.now = t0.Add(2 * time.Minute) + st, _ := j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Unhealthy || s.Reason != ReasonDown { + t.Fatalf("not there after its grace is down: %+v", s) + } + rt.err = errors.New("cannot connect to the runtime") + st, _ = j.Look(t.Context()) + if s := stateOf(t, st, "letta.server"); s.State != Unknown || !strings.Contains(s.Reason, "cannot connect") { + t.Fatalf("a runtime that does not answer is unknown, never down: %+v", s) + } +} + +func TestAUnitRestartedByItsManagerIsCountedAndOneStartedAgainIsNot(t *testing.T) { + unit := Resource{Module: "mqtt", ID: "mqtt.broker", Kind: KindService, Target: "mosquitto.service"} + c := &clock{now: t0} + rt := &fakeRuntime{units: map[string]Observed{"mosquitto.service": running("i1", 0, "")}} + j := aJudge(t, rt, c, filepath.Join(t.TempDir(), FileName)) + j.Set([]Resource{unit}) + j.Look(t.Context()) + c.now = t0.Add(2 * time.Minute) + rt.units["mosquitto.service"] = running("i2", 1, "") // the manager's Restart= + j.Look(t.Context()) + c.now = t0.Add(3 * time.Minute) + rt.units["mosquitto.service"] = running("i3", 0, "") // restarted by the apply: NRestarts reset + st, _ := j.Look(t.Context()) + if s := stateOf(t, st, "mqtt.broker"); s.State != Starting || s.Restarts != 1 { + t.Fatalf("a restart somebody made is a new start, the manager's is counted: %+v", s) + } + c.now = t0.Add(5 * time.Minute) + rt.units["mosquitto.service"] = Observed{Found: true, Identity: "", Running: false} + st, _ = j.Look(t.Context()) + if s := stateOf(t, st, "mqtt.broker"); s.State != Unhealthy || s.Reason != ReasonDown { + t.Fatalf("an inactive unit stated running is down: %+v", s) + } +} + +// The long-running resources of a declaration: what stays up, by module; never a step, a schedule, a +// service whose lifecycle is the machine's, the mesh's own resources, or what an adopted machine holds. +func TestWhatIsLongRunning(t *testing.T) { + const image = "docker.io/library/postgres@sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + d, err := declaration.ParseTrusted([]byte(`{"declaration":1,"resources":[ + {"id":"store","type":"container","name":"mesh-store","image":"` + image + `"}, + {"id":"letta.server","type":"container","name":"letta-server","image":"` + image + `"}, + {"id":"letta.migrate","type":"container","name":"letta-migrate","image":"` + image + `","run-once":true}, + {"id":"letta.sweep","type":"container","name":"letta-sweep","image":"` + image + `","schedule":"0 3 * * *"}, + {"id":"novox.be.web","type":"container","name":"novox-web","image":"` + image + `"}, + {"id":"mqtt.broker","type":"service","unit":"mosquitto.service","state":"running"}, + {"id":"uplink.nm","type":"service","unit":"NetworkManager.service","restart-on":["mqtt.broker"]}, + {"id":"found.thing","type":"container","name":"found","image":"` + image + `"} + ]}`)) + if err != nil { + t.Fatal(err) + } + got := LongRunning(d, map[string]bool{"found.thing": true}) + var ids []string + for _, r := range got { + ids = append(ids, r.Module+"/"+r.ID) + } + want := "letta/letta.server novox.be/novox.be.web mqtt/mqtt.broker" + if strings.Join(ids, " ") != want { + t.Fatalf("long-running: %v, want %s", ids, want) + } +} + +// **Nothing is restarted for being unhealthy** (ADR 0240 rule 6): a crash-looping container is judged +// unhealthy, and everything the judge asked the runtime and the manager was a read. +func TestAnUnhealthyContainerIsNeitherRestartedRecreatedNorStopped(t *testing.T) { + var asked []string + restarts := 0 + run := func(_ context.Context, name string, args ...string) (string, error) { + asked = append(asked, name+" "+strings.Join(args, " ")) + switch { + case name == "docker" && len(args) > 0 && args[0] == "version": + return "27.0.0\n", nil + case name == "docker": + restarts++ + return "/letta-server\tc1\trestarting\t" + string(rune('0'+restarts)) + "\t2026-10-07T12:00:00Z\n", nil + case name == "systemctl": + return "Id=mosquitto.service\nLoadState=loaded\nActiveState=failed\nSubState=failed\nNRestarts=5\nInvocationID=\n", nil + } + return "", errors.New("not a command the judge may run") + } + c := &clock{now: t0} + j := aJudge(t, &Exec{Run: run}, c, filepath.Join(t.TempDir(), FileName)) + j.Set([]Resource{server, {Module: "mqtt", ID: "mqtt.broker", Kind: KindService, Target: "mosquitto.service"}}) + var st Statement + for i := 0; i < 6; i++ { + c.now = t0.Add(time.Duration(i) * time.Minute) + st, _ = j.Look(t.Context()) + } + if s := stateOf(t, st, "letta.server"); s.State != Unhealthy { + t.Fatalf("the crash loop was not judged unhealthy: %+v", s) + } + for _, a := range asked { + read := strings.HasPrefix(a, "docker version ") || strings.HasPrefix(a, "docker container inspect ") || + strings.HasPrefix(a, "systemctl show ") + if !read { + t.Errorf("the judge asked something that is not a read: %q", a) + } + } +} + +// One read of every container per look, never one per container (ADR 0240: the engine stays cheap). +func TestOneLookIsOneReadOfEveryContainer(t *testing.T) { + inspects := 0 + run := func(_ context.Context, name string, args ...string) (string, error) { + if args[0] == "container" { + inspects++ + return "/a\tc1\trunning\t0\tx\n/b\tc2\trunning\t0\tx\n", errors.New("docker exited 1: Error: No such container: c") + } + return "27\n", nil + } + j := aJudge(t, &Exec{Run: run}, &clock{now: t0}, "") + j.Set([]Resource{{Module: "m", ID: "m.a", Kind: KindContainer, Target: "a"}, + {Module: "m", ID: "m.b", Kind: KindContainer, Target: "b"}, {Module: "m", ID: "m.c", Kind: KindContainer, Target: "c"}}) + st, _ := j.Look(t.Context()) + if inspects != 1 { + t.Fatalf("%d inspects for one look", inspects) + } + if s := stateOf(t, st, "m.b"); s.State != Starting { + t.Fatalf("a container the inspect found, beside one it did not: %+v", s) + } +} diff --git a/internal/liveness/replay_test.go b/internal/liveness/replay_test.go new file mode 100644 index 0000000..1805bf5 --- /dev/null +++ b/internal/liveness/replay_test.go @@ -0,0 +1,91 @@ +package liveness + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "os" + "os/exec" + "strings" + "testing" + "time" + + "github.com/novox/mesh-host/internal/link" +) + +// **The crash loop, on a real runtime** — the engine's half of mesh-lab's replay R-crashloop (novox/hq +// ADR 0240, "how it is checked", rule 1: "a mesh-lab replay of the crash loop (a container whose program +// exits at start) fails its gate within the bound"). +// +// A container whose program exits at once, started as the apply starts every container — detached, +// restarted by the runtime unless stopped, labelled with its resource — is judged by this engine through +// the runtime's own command line, and said unhealthy, restarting, inside the gate's bound. What it said is +// written to MESH_REPLAY_STATEMENT, for the controller's half to judge its gate from. +// +// Run by the lab (`go run ./replays/cmd/prove R-crashloop`), which starts the container and names it in +// MESH_REPLAY_CONTAINER; or by hand with MESH_REPLAY_RUNTIME=1, starting its own. Skipped otherwise: the +// suite raises no container unasked. The grace is shortened to seconds, which changes the clock and not +// the rule. +func TestReplayCrashLoopIsSaidUnhealthy(t *testing.T) { + name := os.Getenv("MESH_REPLAY_CONTAINER") + if name == "" && os.Getenv("MESH_REPLAY_RUNTIME") != "1" { + t.Skip("no MESH_REPLAY_CONTAINER, and MESH_REPLAY_RUNTIME is not 1: this replay raises a container") + } + run := func(ctx context.Context, cmd string, args ...string) (string, error) { + c := exec.CommandContext(ctx, cmd, args...) + out, err := c.Output() + var exit *exec.ExitError + if errors.As(err, &exit) { + return string(out), fmt.Errorf("%s exited %d: %s", cmd, exit.ExitCode(), strings.TrimSpace(string(exit.Stderr))) + } + return string(out), err + } + if _, err := run(t.Context(), "docker", "version", "--format", "{{.Server.Version}}"); err != nil { + t.Skipf("no container runtime answers here: %v", err) + } + if name == "" { + name = fmt.Sprintf("mesh-replay-crashloop-%d", time.Now().UnixNano()) + if out, err := run(t.Context(), "docker", "run", "--detach", "--name", name, "--restart", "unless-stopped", + "--label", "mesh-host.id=app.server", "--label", "mesh.replay=1", "alpine:3.20", "sh", "-c", "exit 3"); err != nil { + t.Fatalf("starting the crash loop: %v %s", err, out) + } + t.Cleanup(func() { _, _ = run(context.Background(), "docker", "rm", "-f", name) }) + } + + j, err := Open("", &Exec{Run: run}) + if err != nil { + t.Fatal(err) + } + j.Grace = 3 * time.Second + j.Set([]Resource{{Module: "app", ID: "app.server", Kind: KindContainer, Target: name}}) + began := time.Now() + deadline := began.Add(2 * time.Minute) + var st Statement + for time.Now().Before(deadline) { + st, _ = j.Look(t.Context()) + if r := st.Resources[0]; r.State == Unhealthy && r.Restarts >= 2 { + break + } + time.Sleep(time.Second) + } + s := st.Resources[0] + if s.State != Unhealthy || s.Restarts < 2 { + t.Fatalf("a container whose program exits at start was not said unhealthy, with its restarts counted, in two minutes: %+v", s) + } + t.Logf("said %s (%s), %d restarts counted, %s after the judging began", s.State, s.Reason, s.Restarts, + time.Since(began).Round(time.Second)) + + if out := os.Getenv("MESH_REPLAY_STATEMENT"); out != "" { + h := link.Health{Contract: link.LivenessContract, At: st.At.UTC(), Resources: []link.ResourceHealth{{ + Module: s.Module, Resource: s.ID, Kind: s.Kind, Target: s.Target, State: s.State, Reason: s.Reason, + Since: s.Since.UTC(), Streak: s.Streak, Restarts: s.Restarts}}} + raw, err := json.Marshal(h) + if err != nil { + t.Fatal(err) + } + if err := os.WriteFile(out, raw, 0o644); err != nil { + t.Fatal(err) + } + } +} diff --git a/internal/liveness/runtime.go b/internal/liveness/runtime.go new file mode 100644 index 0000000..376f352 --- /dev/null +++ b/internal/liveness/runtime.go @@ -0,0 +1,134 @@ +package liveness + +import ( + "context" + "errors" + "fmt" + "strconv" + "strings" + "sync" +) + +// Runner runs a command and answers what it printed — the apply's own (apply.ExecRunner), so a look +// reads the machine exactly as an apply does. +type Runner func(ctx context.Context, name string, args ...string) (string, error) + +// Exec reads the container runtime and the service manager through their command lines: one inspect of +// every container, one show of every unit per manager. **Reads only** — the whole of what this file may +// ask either of them is `container inspect`, `version` and `show` (ADR 0240 rule 6, and a test holds it). +type Exec struct { + Run Runner + + mu sync.Mutex + // cri is the runtime that answered, once one has. + cri string +} + +// runtimes are the container runtimes the apply knows, asked for their version in the form each +// understands (apply.containerRuntimes): the first that answers is the machine's. +var runtimes = []struct { + command string + probe []string +}{ + {"docker", []string{"version", "--format", "{{.Server.Version}}"}}, + {"podman", []string{"version", "--format", "{{.Version}}"}}, +} + +func (e *Exec) runtime(ctx context.Context) (string, error) { + e.mu.Lock() + defer e.mu.Unlock() + if e.cri != "" { + return e.cri, nil + } + var tried []string + for _, rt := range runtimes { + if _, err := e.Run(ctx, rt.command, rt.probe...); err == nil { + e.cri = rt.command + return rt.command, nil + } + tried = append(tried, rt.command) + } + return "", fmt.Errorf("no container runtime answers on this machine (tried %s)", strings.Join(tried, ", ")) +} + +// inspectFormat is what one inspect says of each container: its name, id, status, the runtime's restart +// count and when its current run started. +const inspectFormat = "{{.Name}}\t{{.Id}}\t{{.State.Status}}\t{{.RestartCount}}\t{{.State.StartedAt}}" + +// Containers is one inspect of every container named. A name the runtime does not have is not found; +// a runtime that does not answer is an error — unknown, never down. +func (e *Exec) Containers(ctx context.Context, names []string) (map[string]Observed, error) { + out := map[string]Observed{} + if len(names) == 0 { + return out, nil + } + cri, err := e.runtime(ctx) + if err != nil { + return nil, err + } + args := append([]string{"container", "inspect", "--format", inspectFormat}, names...) + said, err := e.Run(ctx, cri, args...) + if err != nil && !missing(err) { + // Asked again next look, in case another runtime is what answers now. + e.mu.Lock() + e.cri = "" + e.mu.Unlock() + return nil, fmt.Errorf("the container runtime could not be read: %w", err) + } + for _, line := range strings.Split(strings.TrimSpace(said), "\n") { + parts := strings.Split(line, "\t") + if len(parts) < 5 { + continue + } + name := strings.TrimPrefix(strings.TrimSpace(parts[0]), "/") + restarts, _ := strconv.ParseInt(strings.TrimSpace(parts[3]), 10, 64) + status := strings.TrimSpace(parts[2]) + out[name] = Observed{Found: true, Identity: strings.TrimSpace(parts[1]), Running: status == "running", + Restarting: status == "restarting", Restarts: restarts, Started: strings.TrimSpace(parts[4])} + } + return out, nil +} + +// missing is an inspect that failed only because a name is not there: what it printed for the others is +// still the answer. +func missing(err error) bool { + s := strings.ToLower(err.Error()) + return strings.Contains(s, "no such container") || strings.Contains(s, "no such object") +} + +// unitProperties are what one show says of each unit. +const unitProperties = "Id,LoadState,ActiveState,SubState,NRestarts,InvocationID" + +// Units is one show of every unit named, in the manager it is in: the machine's, or an account's own. +func (e *Exec) Units(ctx context.Context, scope, user string, units []string) (map[string]Observed, error) { + out := map[string]Observed{} + if len(units) == 0 { + return out, nil + } + args := []string{"show", "--property=" + unitProperties, "--"} + if scope == "user" && user != "" { + args = append([]string{"--user", "--machine=" + user + "@"}, args...) + } + said, err := e.Run(ctx, "systemctl", append(args, units...)...) + if err != nil { + return nil, fmt.Errorf("the service manager could not be read: %w", err) + } + blocks := strings.Split(strings.TrimSpace(said), "\n\n") + if len(blocks) != len(units) { + return nil, errors.New("the service manager answered for a different number of units than were asked") + } + for i, block := range blocks { + props := map[string]string{} + for _, line := range strings.Split(block, "\n") { + if k, v, ok := strings.Cut(line, "="); ok { + props[strings.TrimSpace(k)] = strings.TrimSpace(v) + } + } + restarts, _ := strconv.ParseInt(props["NRestarts"], 10, 64) + out[units[i]] = Observed{Found: props["LoadState"] != "not-found", Identity: props["InvocationID"], + Running: props["ActiveState"] == "active", + Restarting: props["ActiveState"] == "activating" && props["SubState"] == "auto-restart", + Restarts: restarts} + } + return out, nil +}