Merge pull request 'Judge whether what a module runs stays up, and say it (hq ADR 0240, to-be 48 Phase A)' (#47) from feat/a-module-says-how-it-is-healthy into main
This commit was merged in pull request #47.
This commit is contained in:
+95
-4
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"))
|
||||
}
|
||||
}
|
||||
@@ -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 }
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user