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 +}