diff --git a/cmd/mesh-controller/gate.go b/cmd/mesh-controller/gate.go index 30096a6..55f3098 100644 --- a/cmd/mesh-controller/gate.go +++ b/cmd/mesh-controller/gate.go @@ -105,6 +105,10 @@ type gateFacts struct { // holder is who holds the controller lease, for judging the controller. holder *lease.Holder holderErr error + // health is each machine's newest health statement (ADR 0240); a machine absent never stated one. + // healthErr is why they could not be read. + health map[string]inventory.NodeHealth + healthErr error } // gatherGateFacts reads what a judging needs, from the store, the bus and this controller's memory. A @@ -127,6 +131,9 @@ var gatherGateFacts = func(ctx context.Context, open *stores, component string) for _, n := range nodes { f.engines[n.Name] = n.HostVersion } + // What each machine says of its long-running resources (ADR 0240): unreadable is said, never read as + // healthy. + f.health, f.healthErr = inv.Healths(ctx) if d := doctorFrom; d != nil { if d.keeper != nil { f.judged = true @@ -193,7 +200,9 @@ func judgeHealth(module, component string, m catalogue.Manifest, machine string, if !onIt { continue } - if c.Subject.Scope == conditions.ScopeMachine || slices.Contains(strings.Split(c.Subject.ID, "."), module) { + // A module's own health condition names it whole, its name's dots and all (ADR 0240). + ownHealth := c.Subject.Scope == conditions.ScopeModule && c.Subject.ID == module+"."+machine + if c.Subject.Scope == conditions.ScopeMachine || ownHealth || slices.Contains(strings.Split(c.Subject.ID, "."), module) { return healthNotYet, fmt.Sprintf("raised since it was sent: %s — %s", c.Key, c.Summary) } } @@ -241,6 +250,11 @@ func judgeHealth(module, component string, m catalogue.Manifest, machine string, return healthNotYet, fmt.Sprintf("the node tools on %s do not serve %s's tools", machine, module) } } + // **And what it runs is stated healthy** (ADR 0240 §4): every long-running resource of it on that + // machine, in a statement heard since the send. A resource still starting makes the judging wait. + if h, why := moduleHealthWord(module, machine, since, f); h != healthGood { + return h, why + } } return healthGood, "" } diff --git a/cmd/mesh-controller/gate_test.go b/cmd/mesh-controller/gate_test.go index 0ad1bd4..e8d55ce 100644 --- a/cmd/mesh-controller/gate_test.go +++ b/cmd/mesh-controller/gate_test.go @@ -346,6 +346,48 @@ func TestTheHealthDefinitions(t *testing.T) { !strings.Contains(why, "first run") { t.Errorf("a controller not ready is %v (%s)", h, why) } + // What the module runs (novox/hq ADR 0240 §4): a judging passes only when every long-running + // resource of it on that machine is stated healthy since the send; starting and unhealthy wait. + f = applied("anchor") + stated := func(state, reason string, heard time.Time) { + f.health = map[string]inventory.NodeHealth{"anchor": {Node: "anchor", Contract: 1, SaidAt: heard, HeardAt: heard, + Resources: []inventory.ResourceHealth{ + {Module: "app", Resource: "app.server", Kind: "container", Target: "app-server", State: state, Reason: reason}, + {Module: "other", Resource: "other.server", Kind: "container", State: link.StateUnhealthy, Reason: "down"}}}} + } + plain := catalogue.Manifest{Module: "app"} + if h, why := judgeHealth("app", "", plain, "anchor", since, f); h != healthGood { + t.Errorf("a machine whose engine states no health is judged as before, not %v (%s)", h, why) + } + stated(link.StateStarting, "", now) + if h, why := judgeHealth("app", "", plain, "anchor", since, f); h != healthNotYet || !strings.Contains(why, "starting") { + t.Errorf("a resource still starting is not yet a pass: %v (%s)", h, why) + } + stated(link.StateUnhealthy, "restarting", now) + if h, why := judgeHealth("app", "", plain, "anchor", since, f); h != healthNotYet || !strings.Contains(why, "restarting") { + t.Errorf("a resource unhealthy fails the judging: %v (%s)", h, why) + } + stated(link.StateHealthy, "", since.Add(-time.Minute)) + if h, why := judgeHealth("app", "", plain, "anchor", since, f); h != healthNotYet { + t.Errorf("a statement from before the send says nothing of the new build: %v (%s)", h, why) + } + stated(link.StateHealthy, "", now) + if h, why := judgeHealth("app", "", plain, "anchor", since, f); h != healthGood { + t.Errorf("every resource of it healthy since the send — another module's state is not its — is %v (%s)", h, why) + } + f.healthErr = errors.New("the store is away") + if h, _ := judgeHealth("app", "", plain, "anchor", since, f); h != healthNotYet { + t.Errorf("health that cannot be read is never read as healthy: %v", h) + } + // And the condition it raises after the send holds it, as every condition about it does. + f.healthErr = nil + f.judged = true + f.open = []conditions.Condition{{Key: "module.app.anchor.unhealthy", Kind: kindModuleUnhealthy, Subject: conditions.Subject{ + Scope: conditions.ScopeModule, ID: "app.anchor", Machine: "anchor"}, Raised: now, Summary: "app on anchor is not healthy"}} + if h, why := judgeHealth("app", "", plain, "anchor", since, f); h != healthNotYet || !strings.Contains(why, "module.app.anchor.unhealthy") { + t.Errorf("a module held on its own unhealthy condition is %v (%s)", h, why) + } + // A machine that refused what it was sent: broken. f = applied("anchor") r := f.reports["anchor"] diff --git a/cmd/mesh-controller/module_health.go b/cmd/mesh-controller/module_health.go new file mode 100644 index 0000000..5529be9 --- /dev/null +++ b/cmd/mesh-controller/module_health.go @@ -0,0 +1,249 @@ +package main + +import ( + "context" + "fmt" + "sort" + "strings" + "sync/atomic" + "time" + + "github.com/novox/mesh-controller/internal/conditions" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// A module says how it is healthy, and the node-engine judges it (novox/hq ADR 0240, to-be 48 §4 and §5, +// Phase A). +// +// **The node-engine owns every verdict; the controller keeps the last word and raises the condition.** A +// machine states, in every report and as an event between reports, the state of every long-running +// resource it runs for a module. The controller keeps the newest statement per machine (node_health), and +// raises `module...unhealthy` when two statements in a row say a resource of the module +// is unhealthy — one statement is listed as unconfirmed, as the self-check does a finding one look can be +// wrong about (to-be 45 §4, issue 277) — and clears it on the first that does not. The release gate reads +// the stated health: a judging passes a module only when every long-running resource of it on that +// machine is stated healthy, so a resource still starting is not yet a pass. +// +// **An engine older than the judging states nothing**, and its machine's health is not known: never +// healthy, never a reason to raise anything, and the gate judges it as it did before. + +// The condition a module's health raises. +const ( + kindModuleUnhealthy = "module-unhealthy" + // sourceHealth is what raised it: the machine's own statement. + sourceHealth = "health" + // moduleUnhealthyUrgentAfter is how long it stands before it is urgent (to-be 48 §4). + moduleUnhealthyUrgentAfter = 4 * time.Hour + // moduleUnhealthyAfter is how many statements in a row raise it. + moduleUnhealthyAfter = 2 +) + +// healthRefused counts the statements refused as older than the one kept, for the log and a test. +var healthRefused atomic.Int64 + +// moduleHealth keeps what the machines state, for the link (link.Healths). +type moduleHealth struct { + inv *inventory.Inventory + keeper func() *conditions.Keeper +} + +func (m moduleHealth) Stated(ctx context.Context, node string, h link.Health) error { + return stateHealth(ctx, m.inv, m.keeper(), node, h, time.Now()) +} + +// stateHealth keeps one machine's statement and raises or clears its modules' conditions from it. An +// older statement than the one kept is refused, by when the engine looked. +func stateHealth(ctx context.Context, inv *inventory.Inventory, k *conditions.Keeper, node string, h link.Health, + now time.Time) error { + if h.Contract == 0 { + return nil + } + prev, had, err := inv.HealthOf(ctx, node) + if err != nil { + return err + } + if had && h.At.Before(prev.SaidAt) { + healthRefused.Add(1) + return nil + } + unhealthy := map[string][]inventory.ResourceHealth{} + resources := make([]inventory.ResourceHealth, 0, len(h.Resources)) + for _, r := range h.Resources { + kept := inventory.ResourceHealth{Module: r.Module, Resource: r.Resource, Kind: r.Kind, Target: r.Target, + State: r.State, Reason: r.Reason, Since: r.Since, Streak: r.Streak, Restarts: r.Restarts} + resources = append(resources, kept) + if r.State == link.StateUnhealthy && r.Module != "" { + unhealthy[r.Module] = append(unhealthy[r.Module], kept) + } + } + streaks := map[string]int{} + for module := range unhealthy { + streaks[module] = prev.Streaks[module] + 1 + } + stored, err := inv.RecordHealth(ctx, inventory.NodeHealth{Node: node, Contract: h.Contract, SaidAt: h.At, + HeardAt: now, Resources: resources, Streaks: streaks}) + if err != nil || !stored { + if err == nil { + healthRefused.Add(1) + } + return err + } + if k == nil { + return nil + } + return judgeModuleHealth(ctx, k, node, unhealthy, streaks, now) +} + +// judgeModuleHealth raises a module's condition on a machine on the second statement in a row that says a +// resource of it is unhealthy — or on the first while it is already open — and clears every one this +// statement no longer says. +func judgeModuleHealth(ctx context.Context, k *conditions.Keeper, node string, + unhealthy map[string][]inventory.ResourceHealth, streaks map[string]int, now time.Time) error { + open, err := k.Open(ctx) + if err != nil { + return err + } + standing := map[string]conditions.Condition{} + for _, c := range open { + if c.Kind == kindModuleUnhealthy && c.Subject.Machine == node { + standing[c.Key] = c + } + } + var problems []string + modules := make([]string, 0, len(unhealthy)) + for m := range unhealthy { + modules = append(modules, m) + } + sort.Strings(modules) + seen := map[string]bool{} + for _, m := range modules { + o := moduleUnhealthyObservation(m, node, unhealthy[m]) + seen[o.Key()] = true + c, isOpen := standing[o.Key()] + if streaks[m] < moduleUnhealthyAfter && !isOpen { + continue // unconfirmed: one statement can be wrong; `node show` lists it + } + if isOpen && now.Sub(c.Raised) >= moduleUnhealthyUrgentAfter { + o.Severity = conditions.Urgent + } + if _, err := k.Observe(ctx, o); err != nil { + problems = append(problems, err.Error()) + } + } + for key, c := range standing { + if seen[key] { + continue + } + module := strings.TrimSuffix(c.Subject.ID, "."+node) + if _, err := k.Clear(ctx, key, fmt.Sprintf("%s says no resource of %s is unhealthy", node, module)); err != nil { + problems = append(problems, err.Error()) + } + } + if len(problems) > 0 { + return fmt.Errorf("%s", strings.Join(problems, "; ")) + } + return nil +} + +// moduleUnhealthyObservation is a module unhealthy on a machine, in words: the summary names the module, +// the machine and what is wrong with each resource; the detail — targets, streaks, since — is evidence. +func moduleUnhealthyObservation(module, node string, rs []inventory.ResourceHealth) conditions.Observation { + var words, said []string + for _, r := range rs { + words = append(words, fmt.Sprintf("its %s %s %s", r.Kind, r.Resource, reasonWords(r))) + said = append(said, fmt.Sprintf("%s (%s %s): %s, %d look(s) in a row, %d restart(s) counted, since %s", + r.Resource, r.Kind, r.Target, orNotSaid(r.Reason), r.Streak, r.Restarts, + r.Since.UTC().Format("2006-01-02 15:04:05 MST"))) + } + return conditions.Observation{Scope: conditions.ScopeModule, ID: module + "." + node, Token: "unhealthy", + Kind: kindModuleUnhealthy, Machine: node, Severity: conditions.Warning, Source: sourceHealth, + Summary: fmt.Sprintf("%s on %s is not healthy: %s", module, node, strings.Join(words, "; ")), + Said: strings.Join(said, "; ")} +} + +// reasonWords is why a resource is unhealthy, as a person reads it. +func reasonWords(r inventory.ResourceHealth) string { + switch r.Reason { + case "restarting": + return fmt.Sprintf("keeps restarting (%d restart(s) counted)", r.Restarts) + case "down": + return "is not running" + case "": + return "is unhealthy" + } + return "is unhealthy: " + r.Reason +} + +func orNotSaid(s string) string { + if s == "" { + return "no reason said" + } + return s +} + +// moduleHealthWord is the gate's reading of a module's stated health on a machine (ADR 0240 §4, ADR 0236 +// §2 as amended): good when every long-running resource of it is stated healthy in a statement heard since +// the send; not yet otherwise, saying which. A machine that never stated health is judged as before. +func moduleHealthWord(module, machine string, since time.Time, f gateFacts) (health, string) { + if f.healthErr != nil { + return healthNotYet, "what " + machine + " says of its resources' health cannot be read: " + firstLine(f.healthErr.Error()) + } + h, states := f.health[machine] + if !states { + return healthGood, "" + } + if h.HeardAt.Before(since) { + return healthNotYet, fmt.Sprintf("%s has not said how what %s runs is since it was sent", machine, module) + } + for _, r := range h.Resources { + if r.Module != module { + continue + } + switch r.State { + case link.StateHealthy: + case link.StateStarting: + return healthNotYet, fmt.Sprintf("its %s %s on %s is still starting", r.Kind, r.Resource, machine) + case link.StateUnhealthy: + return healthNotYet, fmt.Sprintf("its %s %s on %s %s", r.Kind, r.Resource, machine, reasonWords(r)) + default: + return healthNotYet, fmt.Sprintf("its %s %s on %s is %s%s", r.Kind, r.Resource, machine, r.State, + reasonAfter(r.Reason)) + } + } + return healthGood, "" +} + +func reasonAfter(s string) string { + if s == "" { + return "" + } + return ": " + s +} + +// healthLines is what `node show` says of a machine's long-running resources: each with its state and +// since when, an unhealthy one said once marked unconfirmed. +func healthLines(h inventory.NodeHealth, had bool, now time.Time) []string { + if !had { + return []string{" its node-engine does not say how what it runs is — it is older than the judging (ADR 0240)"} + } + if len(h.Resources) == 0 { + return []string{fmt.Sprintf(" it runs nothing long-lived for a module (said %s ago)", roughly(now.Sub(h.HeardAt)))} + } + out := []string{fmt.Sprintf(" what it runs, as it said %s ago:", roughly(now.Sub(h.HeardAt)))} + for _, r := range h.Resources { + line := fmt.Sprintf(" %-10s %-34s %s %s, since %s", r.State, r.Resource, r.Kind, r.Target, + r.Since.Local().Format("2006-01-02 15:04")) + if r.Reason != "" { + line += " — " + r.Reason + } + if r.Restarts > 0 { + line += fmt.Sprintf(", %d restart(s) counted", r.Restarts) + } + if r.State == link.StateUnhealthy && h.Streaks[r.Module] < moduleUnhealthyAfter { + line += " (unconfirmed: said once)" + } + out = append(out, line) + } + return out +} diff --git a/cmd/mesh-controller/module_health_test.go b/cmd/mesh-controller/module_health_test.go new file mode 100644 index 0000000..7eb4c03 --- /dev/null +++ b/cmd/mesh-controller/module_health_test.go @@ -0,0 +1,163 @@ +package main + +import ( + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/conditions" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// A module's stated health, as the controller keeps it and raises from it (novox/hq ADR 0240, "how it is +// checked", rule 4): one unhealthy statement raises nothing and is listed unconfirmed; two raise; a +// healthy one clears; a condition from before the send clears at the new build's start; an older +// statement is refused; and an engine that states nothing raises nothing. + +var h0 = time.Date(2026, 10, 7, 12, 0, 0, 0, time.UTC) + +func aStatement(at time.Time, states ...string) link.Health { + h := link.Health{Contract: link.LivenessContract, At: at} + for i, s := range states { + r := link.ResourceHealth{Module: "letta", Resource: "letta.server", Kind: "container", Target: "letta-server", + State: s, Since: at} + if i > 0 { + r.Module, r.Resource, r.Target = "mqtt", "mqtt.broker", "mosquitto.service" + } + if s == link.StateUnhealthy { + r.Reason, r.Restarts, r.Streak = "restarting", 4, 2 + } + h.Resources = append(h.Resources, r) + } + return h +} + +func TestTwoUnhealthyStatementsRaiseTheModulesConditionAndAHealthyOneClearsIt(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + inv := open.inventory + k := conditionsFrom + const key = "module.letta.anchor.unhealthy" + openKeys := func() []string { + t.Helper() + list, err := k.Open(ctx) + if err != nil { + t.Fatal(err) + } + var keys []string + for _, c := range list { + keys = append(keys, c.Key) + } + return keys + } + + if err := stateHealth(ctx, inv, k, "anchor", aStatement(h0, link.StateHealthy, link.StateHealthy), h0); err != nil { + t.Fatal(err) + } + // One statement: nothing raised, and `node show` lists it as unconfirmed. + if err := stateHealth(ctx, inv, k, "anchor", aStatement(h0.Add(time.Minute), link.StateUnhealthy, link.StateHealthy), + h0.Add(time.Minute)); err != nil { + t.Fatal(err) + } + if keys := openKeys(); len(keys) != 0 { + t.Fatalf("one statement raised %v", keys) + } + kept, had, err := inv.HealthOf(ctx, "anchor") + if err != nil || !had { + t.Fatalf("the statement was not kept: %v %v", had, err) + } + if lines := strings.Join(healthLines(kept, had, h0.Add(time.Minute)), "\n"); !strings.Contains(lines, "unconfirmed") || + !strings.Contains(lines, "letta.server") { + t.Fatalf("node show does not list the first statement as unconfirmed:\n%s", lines) + } + + // The second in a row raises it — the module's own, never the other module's on the machine. + if err := stateHealth(ctx, inv, k, "anchor", aStatement(h0.Add(2*time.Minute), link.StateUnhealthy, link.StateHealthy), + h0.Add(2*time.Minute)); err != nil { + t.Fatal(err) + } + if keys := openKeys(); len(keys) != 1 || keys[0] != key { + t.Fatalf("two statements raised %v, not %s", keys, key) + } + c, _, _ := k.Get(ctx, key) + if c.Severity != conditions.Warning || c.Resolver != conditions.ResolverSelf || c.Subject.Machine != "anchor" || + !strings.Contains(c.Summary, "letta on anchor") || !strings.Contains(c.Summary, "keeps restarting") || + !strings.Contains(c.Evidence[0].Said, "letta-server") { + t.Fatalf("the condition does not say it in words with its evidence: %+v", c) + } + + // Standing four hours, it is urgent. + if err := stateHealth(ctx, inv, k, "anchor", aStatement(h0.Add(5*time.Hour), link.StateUnhealthy, link.StateHealthy), + time.Now().Add(5*time.Hour)); err != nil { + t.Fatal(err) + } + if c, _, _ := k.Get(ctx, key); c.Severity != conditions.Urgent { + t.Fatalf("unhealthy for four hours is still %s", c.Severity) + } + + // A new build's start — every start begins in `starting` — clears it: what follows is the new build's. + if err := stateHealth(ctx, inv, k, "anchor", aStatement(h0.Add(6*time.Hour), link.StateStarting, link.StateHealthy), + h0.Add(6*time.Hour)); err != nil { + t.Fatal(err) + } + if keys := openKeys(); len(keys) != 0 { + t.Fatalf("the first statement that says no resource is unhealthy did not clear it: %v", keys) + } + // And the streak starts again: one unhealthy statement after it raises nothing. + if err := stateHealth(ctx, inv, k, "anchor", aStatement(h0.Add(7*time.Hour), link.StateUnhealthy), h0.Add(7*time.Hour)); err != nil { + t.Fatal(err) + } + if keys := openKeys(); len(keys) != 0 { + t.Fatalf("one statement after a clearing raised %v", keys) + } +} + +func TestAnOlderHealthStatementIsRefused(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + inv := open.inventory + before := healthRefused.Load() + if err := stateHealth(ctx, inv, nil, "anchor", aStatement(h0.Add(time.Minute), link.StateHealthy), h0); err != nil { + t.Fatal(err) + } + // An event said before the report that overtook it arrives late. + if err := stateHealth(ctx, inv, nil, "anchor", aStatement(h0, link.StateUnhealthy), h0); err != nil { + t.Fatal(err) + } + kept, _, err := inv.HealthOf(ctx, "anchor") + if err != nil { + t.Fatal(err) + } + if !kept.SaidAt.Equal(h0.Add(time.Minute)) || kept.Resources[0].State != link.StateHealthy { + t.Fatalf("the older statement replaced the newer: %+v", kept) + } + if healthRefused.Load() != before+1 { + t.Fatalf("the refusal was not counted") + } +} + +// **An engine older than the judging states nothing**: its reports raise nothing, keep nothing, and the +// gate judges its machine as before — never healthy for having said nothing, never unhealthy. +func TestAReportWithNoHealthRaisesAndKeepsNothing(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + l := nudgingListener{Enrolment: link.Enrolment{Inventory: open.inventory}} + if _, err := l.Heard(ctx, link.Report{Node: "anchor", Declared: "d1", Applied: []string{"letta.server"}}); err != nil { + t.Fatal(err) + } + if _, had, err := open.inventory.HealthOf(ctx, "anchor"); err != nil || had { + t.Fatalf("health kept for an engine that said none: %v %v", had, err) + } + if lines := healthLines(inventory.NodeHealth{}, false, time.Now()); !strings.Contains(lines[0], "older than the judging") { + t.Fatalf("node show: %v", lines) + } + // And one that does, through the report, is kept. + h := aStatement(time.Now().UTC(), link.StateHealthy) + if _, err := l.Heard(ctx, link.Report{Node: "anchor", Declared: "d1", Health: &h}); err != nil { + t.Fatal(err) + } + if kept, had, err := open.inventory.HealthOf(ctx, "anchor"); err != nil || !had || len(kept.Resources) != 1 { + t.Fatalf("the report's health was not kept: %+v %v %v", kept, had, err) + } +} diff --git a/cmd/mesh-controller/nodes.go b/cmd/mesh-controller/nodes.go index f2bc0b9..701b010 100644 --- a/cmd/mesh-controller/nodes.go +++ b/cmd/mesh-controller/nodes.go @@ -537,6 +537,17 @@ func showNode(ctx context.Context, inv *inventory.Inventory, name string) error } } + // What it runs for its modules, and how each is (novox/hq ADR 0240): "is it working" answered without + // a terminal on the machine. + if h, had, err := inv.HealthOf(ctx, name); err != nil { + fmt.Printf("\n what it says of what it runs could NOT be read: %v\n", err) + } else { + fmt.Println() + for _, line := range healthLines(h, had, time.Now()) { + fmt.Println(line) + } + } + held, err := inv.Profile(ctx, name) if err != nil { return err diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index fdc34ca..e6e78e9 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -182,6 +182,9 @@ func serve(ctx context.Context) (err error) { if err := server.Watches(standings{keeper: func() *conditions.Keeper { return conditionsFrom }}); err != nil { return err } + // And what each machine says of what it runs between its reports, kept and raised from (novox/hq ADR + // 0240): the gate, `node show` and the conditions read it. + server.Hears(moduleHealth{inv: inv, keeper: func() *conditions.Keeper { return conditionsFrom }}) // And what they say about consumers the mesh stopped asking for: one waiting for a person is an // urgent condition, and an act asked of a provider some other way is recorded by hand (ADR 0230). if err := server.KeepsRetirements(retirements{keeper: func() *conditions.Keeper { return conditionsFrom }, diff --git a/cmd/mesh-controller/replays_test.go b/cmd/mesh-controller/replays_test.go index 636b0c2..dc0d0d1 100644 --- a/cmd/mesh-controller/replays_test.go +++ b/cmd/mesh-controller/replays_test.go @@ -1,9 +1,17 @@ package main import ( + "context" + "encoding/json" + "fmt" + "os" + "slices" "testing" + "time" "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" ) // The replays of the controller's incidents (novox/hq to-be 45 §9, M9): each a scripted replay of what @@ -111,3 +119,170 @@ func TestReplay273AConsumerBesideItsStoreStaysBoundToIt(t *testing.T) { t.Fatalf("the resolver was bound to %q; its seat is held on anchor (issue 258)", network) } } + +// **R-crashloop — a module that crash-loops after it applied fails its gate on its first machine** (novox/hq +// ADR 0240; research 032 §6). On the home server the agent server restarted about a hundred times while the +// mesh read it applied, its tools served and no condition raised — the gate judged a module by what the +// mesh saw from outside, and nothing looked at what it ran. The outcome asserted: a build whose container +// exits at start never passes its gate on the first machine, is put back there at the bound, and never +// reaches the second. +// +// The machine is heard as a node-engine says it: its report of the apply, then the same account said again +// later, each carrying what the engine states of what it runs — `starting` as the apply ends, then the +// crash loop. What it states of the crash loop is MESH_REPLAY_STATEMENT when the lab's replay ran the +// engine against a real container (mesh-lab replays, R-crashloop), and otherwise what the engine said of +// one, kept below. Heard as bytes, so an older controller reads them as it reads any report — this file is +// written only with what the controller had before the judging, for the prover to lay over that commit. +func TestReplayCrashLoopFailsItsGateOnTheFirstMachine(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + inv := open.inventory + + crashLoop := []byte(`{"contract":1,"at":"2026-10-07T00:00:00Z","resources":[{"module":"app","resource":"app.server",` + + `"kind":"container","target":"app-server","state":"unhealthy","reason":"restarting","since":"2026-10-07T00:00:00Z",` + + `"streak":5,"restarts":2}]}`) + if path := os.Getenv("MESH_REPLAY_STATEMENT"); path != "" { + raw, err := os.ReadFile(path) + if err != nil { + t.Fatalf("the engine's statement of the crash loop: %v", err) + } + crashLoop = raw + } + // `null` is an engine that states nothing — older than the judging — and its reports carry no health. + var stated map[string]any + if err := json.Unmarshal(crashLoop, &stated); err != nil { + t.Fatal(err) + } + + for _, b := range []inventory.Build{ + {ID: "build-1", Module: "app", Commit: "c1", Repository: "novox/mesh-catalog", Path: "modules/app", + Asked: time.Now().Add(-2 * time.Hour), At: time.Now().Add(-2 * time.Hour)}, + {ID: "build-2", Module: "app", Commit: "c2", Repository: "novox/mesh-catalog", Path: "modules/app", + Asked: time.Now().Add(-time.Minute), At: time.Now().Add(-time.Minute)}, + } { + manifest, _ := json.Marshal(catalogue.Manifest{Module: "app", Version: b.Commit}) + b.Manifest = manifest + if err := inv.RecordBuild(ctx, b); err != nil { + t.Fatal(err) + } + } + registerAt := func(commit string, asked time.Time) { + if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "app", Version: commit}, + inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/app", BuiltFrom: commit, + Head: commit, Asked: asked}); err != nil { + t.Fatal(err) + } + } + registerAt("c1", time.Now().Add(-2*time.Hour)) + for _, n := range []string{"anchor", "laptop"} { + if _, err := inv.Assign(ctx, n, "app"); err != nil { + t.Fatal(err) + } + if err := inv.RecordSent(ctx, nodeID(t, open, n), "d-"+n+"-c1", map[string]string{"app": "c1"}); err != nil { + t.Fatal(err) + } + } + registerAt("c2", time.Now().Add(-time.Minute)) + + // The machine: what it is sent is applied, and its report says what runs is starting. + listener := nudgingListener{Enrolment: link.Enrolment{Inventory: inv}} + sequence := int64(0) + say := func(node, digest, state string) { + t.Helper() + sequence++ + health := map[string]any{} + for k, v := range stated { + health[k] = v + } + health["at"] = time.Now().UTC().Format(time.RFC3339Nano) + if state != "" && stated != nil { + resources := []any{} + for _, r := range stated["resources"].([]any) { + kept := map[string]any{} + for k, v := range r.(map[string]any) { + kept[k] = v + } + kept["state"], kept["reason"] = state, "" + resources = append(resources, kept) + } + health["resources"] = resources + } + said := map[string]any{"node": node, "applied": []string{"app.server"}, "declared": digest, + "report_sequence": sequence} + if stated != nil { + said["health"] = health + } + body, _ := json.Marshal(said) + var report link.Report + if err := json.Unmarshal(body, &report); err != nil { + t.Fatal(err) + } + if _, err := listener.Heard(ctx, report); err != nil { + t.Fatal(err) + } + } + var sent [][]string + last := map[string]string{} + n := 0 + wasSend := sendRollout + sendRollout = func(ctx context.Context, open *stores, names []string) ([]string, error) { + sent = append(sent, append([]string(nil), names...)) + current, err := open.inventory.CurrentBuilds(ctx) + if err != nil { + return nil, err + } + for _, node := range names { + n++ + digest := fmt.Sprintf("d-%s-%d", node, n) + if err := open.inventory.RecordSent(ctx, nodeID(t, open, node), digest, map[string]string{"app": current["app"].Commit}); err != nil { + return nil, err + } + last[node] = digest + say(node, digest, "starting") + } + return names, nil + } + t.Cleanup(func() { sendRollout = wasSend }) + wasSettle, wasEvery, wasBound := gateSettle, gateEvery, gateBound + gateSettle, gateEvery = 0, 0 + t.Cleanup(func() { gateSettle, gateEvery, gateBound = wasSettle, wasEvery, wasBound }) + + built := time.Now().UTC() + plan := inventory.Plan{ID: "plan-crashloop", Repository: "novox/mesh-catalog", Branch: "main", Commit: "c2", + Created: built, State: inventory.PlanBuilding, Tiers: [][]string{{"app"}}, + Modules: map[string]*inventory.PlanModule{"app": {State: "built", BuiltAt: &built, Commit: "c2", Build: "build-2"}}} + if err := inv.SavePlan(ctx, &plan); err != nil { + t.Fatal(err) + } + + advancePlans(ctx, open) // the first machine is sent the new build, and starts it + if len(sent) == 0 || !slices.Equal(sent[0], []string{"anchor"}) { + t.Fatalf("sent %v, not the first machine first", sent) + } + advancePlans(ctx, open) // judged while it starts + // Its container exits at start, and the runtime restarts it: the machine says so, again and again. + for i := 0; i < 6 && len(sent) == 1; i++ { + say("anchor", last["anchor"], "") + advancePlans(ctx, open) + } + gateBound = -time.Second // and the bound passes + advancePlans(ctx, open) + + p, err := inv.PlanByID(ctx, "plan-crashloop") + if err != nil { + t.Fatal(err) + } + gate := p.Modules["app"].Gate + for _, names := range sent { + if slices.Contains(names, "laptop") { + t.Fatalf("the crash loop passed its gate on anchor and was sent to laptop: sent %v, the gate %+v", sent, gate) + } + } + if gate == nil || gate.Verdict != inventory.GateFailed || p.State != inventory.PlanFailed { + t.Fatalf("a crash-looping build did not fail its gate on the first machine: the plan is %s (%s), its gate %+v", + p.State, p.Note, gate) + } + if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c1" { + t.Fatalf("the module is registered at %s, not put back to c1", current["app"].Commit) + } +} diff --git a/cmd/mesh-controller/status_summary.go b/cmd/mesh-controller/status_summary.go index 053f8b6..7974abd 100644 --- a/cmd/mesh-controller/status_summary.go +++ b/cmd/mesh-controller/status_summary.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "fmt" + "os" "sync" "time" @@ -204,6 +205,15 @@ func (l nudgingListener) Heard(ctx context.Context, report link.Report) (bool, e if news { l.summary.nudge() } + // What it says of its long-running resources (novox/hq ADR 0240): kept, and its modules' conditions + // raised or cleared from it. Only from an account of the machine the store took. + if err == nil && report.Health != nil && report.Superseded == "" && report.Rekey == nil && report.Node != "" && + l.Enrolment.Inventory != nil { + if herr := stateHealth(ctx, l.Enrolment.Inventory, conditionsFrom, report.Node, *report.Health, now); herr != nil { + fmt.Fprintf(os.Stderr, "mesh-controller: could not keep what %s says of its resources' health: %v\n", + report.Node, herr) + } + } if err == nil && l.open != nil && startedWell(report) { // Off the report's path: replacing a given value sends the machine, and a report waits for // nothing it caused (novox/hq ADR 0228). diff --git a/internal/broker/writers.go b/internal/broker/writers.go index 589f44b..9c18ed1 100644 --- a/internal/broker/writers.go +++ b/internal/broker/writers.go @@ -86,7 +86,9 @@ var WritersTable = []WriterRow{ Others: "read", Subjects: []string{"mesh.node.*.declare"}, Writes: isController}, {State: "a machine's applied state and its report", Writer: "the node-engine's apply queue", KeptIn: "the machine; the report on the bus", Others: "the reconcile and a delivery enqueue, never apply", - Subjects: []string{"mesh.control.*.report"}, Writes: ownMachine}, + // And its health statement between reports (novox/hq ADR 0240): the same writer stating the same + // machine, inside the grant it already had (`mesh.control..>`). + Subjects: []string{"mesh.control.*.report", "mesh.control.*.health"}, Writes: ownMachine}, {State: "the controller lease", Writer: "the controller instance holding it", KeptIn: "key-value " + LeaseBucket, Others: "a candidate waits", Subjects: kvOf(LeaseBucket), Writes: isController}, {State: "plans and their tiers", Writer: "controller (lease holder), compare-and-set on the plan's revision", diff --git a/internal/conditions/condition.go b/internal/conditions/condition.go index 11b6659..500e793 100644 --- a/internal/conditions/condition.go +++ b/internal/conditions/condition.go @@ -46,11 +46,14 @@ const ( ScopeMesh = "mesh" // ScopeDelivery is a delivery mesh-delivery owns, by its id (novox/hq ADR 0239). ScopeDelivery = "delivery" + // ScopeModule is a module on a machine, by `.`: what it runs is not healthy there + // (novox/hq ADR 0240). + ScopeModule = "module" ) // Scopes is every scope, in the order a person reads them. var Scopes = []string{ScopeMachine, ScopePlan, ScopeCall, ScopeBuild, ScopeMerge, ScopeProvider, - ScopeSeat, ScopeBus, ScopeCore, ScopeProbe, ScopeMesh, ScopeDelivery} + ScopeSeat, ScopeBus, ScopeCore, ScopeProbe, ScopeMesh, ScopeDelivery, ScopeModule} // Who resolves a condition. const ( diff --git a/internal/inventory/health.go b/internal/inventory/health.go new file mode 100644 index 0000000..6d25898 --- /dev/null +++ b/internal/inventory/health.go @@ -0,0 +1,120 @@ +package inventory + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" +) + +// What each machine says of its long-running resources (novox/hq ADR 0240, to-be 48 §4): the newest +// statement per machine, and per module how many statements in a row said a resource of it was unhealthy. + +// ResourceHealth is one long-running resource's state as the node-engine said it. +type ResourceHealth struct { + Module string `json:"module"` + Resource string `json:"resource"` + Kind string `json:"kind"` + Target string `json:"target"` + State string `json:"state"` + Reason string `json:"reason,omitempty"` + Since time.Time `json:"since"` + Streak int `json:"streak,omitempty"` + Restarts int `json:"restarts,omitempty"` +} + +// NodeHealth is a machine's newest statement, as kept. +type NodeHealth struct { + Node string + Contract int + // SaidAt is when the engine looked (the machine's clock); HeardAt when this controller heard it. + SaidAt time.Time + HeardAt time.Time + Resources []ResourceHealth + // Streaks is, per module, how many statements in a row said a resource of it was unhealthy. + Streaks map[string]int +} + +// HealthOf is a machine's newest statement; false when its node-engine has never stated one. +func (i *Inventory) HealthOf(ctx context.Context, nodeName string) (NodeHealth, bool, error) { + all, err := i.healths(ctx, nodeName) + if err != nil { + return NodeHealth{}, false, err + } + h, ok := all[nodeName] + return h, ok, nil +} + +// Healths is every machine's newest statement, by machine. A machine absent never stated one: its +// node-engine is older than the judging, and its health is not known. +func (i *Inventory) Healths(ctx context.Context) (map[string]NodeHealth, error) { + return i.healths(ctx, "") +} + +func (i *Inventory) healths(ctx context.Context, only string) (map[string]NodeHealth, error) { + rows, err := i.store.Pool().Query(ctx, + `select n.name, h.contract, h.said_at, h.heard_at, h.resources, h.streaks + from node_health h join node n on n.id = h.node + where $1 = '' or n.name = $1`, only) + if err != nil { + return nil, err + } + defer rows.Close() + out := map[string]NodeHealth{} + for rows.Next() { + var h NodeHealth + var resources, streaks []byte + if err := rows.Scan(&h.Node, &h.Contract, &h.SaidAt, &h.HeardAt, &resources, &streaks); err != nil { + return nil, err + } + if err := json.Unmarshal(resources, &h.Resources); err != nil { + return nil, fmt.Errorf("%s's health cannot be read: %w", h.Node, err) + } + if err := json.Unmarshal(streaks, &h.Streaks); err != nil { + return nil, fmt.Errorf("%s's health cannot be read: %w", h.Node, err) + } + out[h.Node] = h + } + return out, rows.Err() +} + +// RecordHealth keeps a machine's statement in place of the one kept — unless the one kept is newer, by +// when the engine looked: then nothing is written, and false says it was refused as older. +func (i *Inventory) RecordHealth(ctx context.Context, h NodeHealth) (bool, error) { + if h.Resources == nil { + h.Resources = []ResourceHealth{} + } + if h.Streaks == nil { + h.Streaks = map[string]int{} + } + resources, err := json.Marshal(h.Resources) + if err != nil { + return false, err + } + streaks, err := json.Marshal(h.Streaks) + if err != nil { + return false, err + } + heard := h.HeardAt + if heard.IsZero() { + heard = time.Now() + } + var node string + err = i.store.Pool().QueryRow(ctx, + `insert into node_health (node, contract, said_at, heard_at, resources, streaks) + select id, $2, $3, $4, $5, $6 from node where name = $1 + on conflict (node) do update set contract = excluded.contract, said_at = excluded.said_at, + heard_at = excluded.heard_at, resources = excluded.resources, streaks = excluded.streaks + where node_health.said_at <= excluded.said_at + returning node`, h.Node, h.Contract, h.SaidAt, heard, resources, streaks).Scan(&node) + if errors.Is(err, pgx.ErrNoRows) { + if _, nerr := i.NodeByName(ctx, h.Node); nerr != nil { + return false, nerr + } + return false, nil + } + return err == nil, err +} diff --git a/internal/inventory/migrations/0076-a-module-says-how-it-is-healthy.sql b/internal/inventory/migrations/0076-a-module-says-how-it-is-healthy.sql new file mode 100644 index 0000000..9d04b7d --- /dev/null +++ b/internal/inventory/migrations/0076-a-module-says-how-it-is-healthy.sql @@ -0,0 +1,26 @@ +-- A module says how it is healthy, and the node-engine judges it (novox/hq ADR 0240, to-be 48 Phase A). +-- +-- Every machine's node-engine states the health of every long-running resource it runs for a module — a +-- container that stays up, a process that stays up, a service stated running — in every report and as an +-- event between reports. The controller keeps the newest statement per machine, here, so the release +-- gate, `node show` and a controller started again all read the same word; and with it, per module, how +-- many statements in a row said a resource of it was unhealthy — `module...unhealthy` is +-- raised on the second (ADR 0240 §4). +-- +-- One row per machine, replaced, as node_report is: the question is the machine's state now. A machine +-- whose node-engine is older than the judging has no row, and its health is not known — never healthy, +-- never unhealthy. +create table node_health ( + node uuid primary key references node(id) on delete cascade, + -- The statement's version (the engine's liveness contract). + contract int not null, + -- When the node-engine looked, on the machine's clock: the order of its statements. An older one + -- than this is refused. + said_at timestamptz not null, + -- When this controller heard it, on its own: what "since the send" is judged by. + heard_at timestamptz not null default now(), + -- [{module, resource, kind, target, state, reason, since, streak, restarts}], as the engine said them. + resources jsonb not null default '[]', + -- {module: statements in a row saying a resource of it is unhealthy}. + streaks jsonb not null default '{}' +); diff --git a/internal/link/bus.go b/internal/link/bus.go index e03c8e8..e2e227d 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -75,8 +75,16 @@ const ( // ToolsAliveSubjects is every machine's node tools saying they are there (novox/hq to-be 45 §3, // S11): core NATS like the host's, for the same reason. ToolsAliveSubjects = "mesh.control.*.tools-alive" + + // HealthSubjects is every machine's health statement between its reports (novox/hq ADR 0240): core + // NATS like the heartbeat, because a statement lost is said again within a minute while anything is + // not healthy, and the next report carries it whatever happens. + HealthSubjects = "mesh.control.*.health" ) +// HealthSubject is one machine's health statement. +func HealthSubject(node string) string { return "mesh.control." + node + ".health" } + // ReportSubject is where one node says what it did. On the CONTROL stream, because it is the // message the store-window guarantee is about (ADR 0083). func ReportSubject(node string) string { return "mesh.control." + node + ".report" } diff --git a/internal/link/contracts.go b/internal/link/contracts.go index f5e759a..b3acb7c 100644 --- a/internal/link/contracts.go +++ b/internal/link/contracts.go @@ -35,6 +35,10 @@ var Contracts = map[string]Contract{ KindHeartbeat: {Unordered: "a word that the machine is there: the newest heard is the newest said, and one " + "lost is the next one"}, KindToolsHeartbeat: {Unordered: "a word that the node tools are there, as a machine's heartbeat"}, + KindHealth: {Ordered: "by the time the node-engine looked (`at`, on the machine's clock): a statement older " + + "than the one kept for that machine is refused, so an event arriving after a newer report does not undo it " + + "(novox/hq ADR 0240)", + Tests: []string{"TestAnOlderHealthStatementIsRefused"}}, KindEnrolment: {Unordered: "a request answered once, under a token spent once: a second presentation is " + "refused by the token, not by an order (issue 083)", Tests: []string{"TestAnEnrolmentMetByAHeldTokenIsAskedToTryAgain"}}, diff --git a/internal/link/contracts_test.go b/internal/link/contracts_test.go index c2995fd..20a5a41 100644 --- a/internal/link/contracts_test.go +++ b/internal/link/contracts_test.go @@ -16,7 +16,7 @@ import ( func TestEveryConsumedKindHasAContract(t *testing.T) { // What the controller can be handed, from the subjects it is granted and the ones it derives. subjects := []string{EnrolSubject, BuiltSubject, ReportSubject("anchor"), AliveSubject("anchor"), - ToolsAliveSubject("anchor"), "mesh.mod.postgres.event.provisioner.failing"} + ToolsAliveSubject("anchor"), HealthSubject("anchor"), "mesh.mod.postgres.event.provisioner.failing"} subjects = append(subjects, broker.ControllerFollows...) kinds := map[string]bool{} for _, s := range subjects { diff --git a/internal/link/health.go b/internal/link/health.go new file mode 100644 index 0000000..7841386 --- /dev/null +++ b/internal/link/health.go @@ -0,0 +1,56 @@ +package link + +import ( + "context" + "encoding/json" + "strings" +) + +// A machine's health statement (novox/hq ADR 0240, to-be 48 §4). +// +// **The node-engine owns every verdict; the controller keeps the last word per machine.** A statement +// arrives in every report and as an event between reports — on each change, and again every minute while +// a resource is not healthy. What the controller does with it — keep it, refuse an older one, raise +// `module...unhealthy` on the second statement that says so — is the Healths it is given. + +// Healths keeps what machines state of their long-running resources. +type Healths interface { + // Stated keeps one machine's statement, from a report or the event. An older statement than the one + // kept is refused there, by its time on the machine. + Stated(ctx context.Context, node string, h Health) error +} + +// Hears says where the machines' health statements are kept. Their subscription is the heartbeats': +// core, and always made, so nothing is asked of the bus here. +func (s *Server) Hears(h Healths) { s.healths = h } + +// healthSaid acts on one health event. Core NATS, so there is nothing to hold: a statement that could not +// be kept is said in the log, and the next one — a minute away while anything is not healthy — is kept. +func (s *Server) healthSaid(ctx context.Context, m Control) { + defer func() { _ = m.Took() }() + var said HealthSaid + if err := json.Unmarshal(m.Body(), &said); err != nil || said.Health.Contract == 0 { + return + } + // The machine is the one in the subject the bus let it publish on, never the body's. + node, ok := nodeOfHealth(m.Subject()) + if !ok || s.healths == nil { + return + } + if err := s.healths.Stated(ctx, node, said.Health); err != nil { + s.log.Printf("could not keep what %s says of its resources' health: %v", node, err) + } +} + +// nodeOfHealth is the machine a health statement names in its subject. +func nodeOfHealth(subject string) (string, bool) { + rest, ok := strings.CutPrefix(subject, "mesh.control.") + if !ok { + return "", false + } + node, kind, ok := strings.Cut(rest, ".") + if !ok || kind != "health" || node == "" || strings.Contains(node, ".") { + return "", false + } + return node, true +} diff --git a/internal/link/health_test.go b/internal/link/health_test.go new file mode 100644 index 0000000..5a58063 --- /dev/null +++ b/internal/link/health_test.go @@ -0,0 +1,37 @@ +package link + +import ( + "context" + "testing" + "time" +) + +// A health event is kept as the machine its subject names says it, whatever its body claims, and taken +// (novox/hq ADR 0240): core NATS, nothing to hold. + +type keptHealths struct{ by map[string]Health } + +func (k *keptHealths) Stated(_ context.Context, node string, h Health) error { + k.by[node] = h + return nil +} + +func TestAHealthEventIsKeptAsTheMachineInItsSubject(t *testing.T) { + s, in := serving() + kept := &keptHealths{by: map[string]Health{}} + s.Hears(kept) + to := &settled{} + m := in.sends(t, to, KindHealth, HealthSaid{Node: "laptop", Health: Health{Contract: LivenessContract, At: time.Now(), + Resources: []ResourceHealth{{Module: "letta", Resource: "letta.server", State: StateUnhealthy}}}}).(*fakeControl) + m.subject = HealthSubject("anchor") + s.act(t.Context(), m) + if !to.acked { + t.Fatal("a health event was not taken") + } + if _, lied := kept.by["laptop"]; lied || len(kept.by["anchor"].Resources) != 1 { + t.Fatalf("kept %+v: the machine is the subject's, never the body's", kept.by) + } + if kind, ok := kindOfSubject(HealthSubject("anchor")); !ok || kind != KindHealth { + t.Fatalf("%s is not read as a health statement", HealthSubject("anchor")) + } +} diff --git a/internal/link/protocol.go b/internal/link/protocol.go index e289ada..c56e6c6 100644 --- a/internal/link/protocol.go +++ b/internal/link/protocol.go @@ -262,6 +262,13 @@ type Report struct { // `witness` and `not-reversible` are sent only to a machine whose report carries it. Witness int `json:"witness,omitempty"` + // Health is the machine's word on every long-running resource it runs for a module (novox/hq ADR + // 0240, to-be 48 §4; mesh-host's internal/liveness): its state, since when, its failing streak and + // the restarts its node-engine counted. **Absent from an engine older than the judging**, which is + // read as "not known" — never as healthy, never as a reason to raise anything; present with no + // resources from a machine that runs nothing long-lived. + Health *Health `json:"health,omitempty"` + // Rekey is a node taking a found tunnel's key as its overlay key after enrolment (novox/hq // ADR 0105). A report carrying one is not an account of the machine: it moves the node's // overlay key and tunnel and nothing else. @@ -404,3 +411,45 @@ func EnrolProof(secret string, public []byte, overlay, sealing, serving string) return []byte("novox-mesh-enrol\x00" + secret + "\x00" + base64.StdEncoding.EncodeToString(public) + "\x00" + overlay + "\x00" + sealing + "\x00" + serving) } + +// LivenessContract is the version of the health statement this controller reads (ADR 0240 Phase A). +const LivenessContract = 1 + +// Health is one statement of a machine's long-running resources (to-be 48 §4): in every report, as the +// event HealthSubject between reports on each change, and again every minute while one is not healthy. +// The node-engine's own (mesh-host internal/link Health); a test on each side holds the field names. +type Health struct { + Contract int `json:"contract"` + // At is when the engine looked, on the machine's clock: the order of its statements. + At time.Time `json:"at"` + Resources []ResourceHealth `json:"resources"` +} + +// The states a resource is said in (ADR 0240 §4). +const ( + StateHealthy = "healthy" + StateUnhealthy = "unhealthy" + StateStarting = "starting" + StateHeld = "held" + StateUnknown = "unknown" +) + +// ResourceHealth is one long-running resource's state. +type ResourceHealth struct { + Module string `json:"module"` + Resource string `json:"resource"` + Kind string `json:"kind"` + Target string `json:"target"` + State string `json:"state"` + Reason string `json:"reason,omitempty"` + Since time.Time `json:"since"` + Streak int `json:"streak,omitempty"` + Restarts int `json:"restarts,omitempty"` +} + +// HealthSaid is the health event's body: the machine and its statement. The machine is read from the +// subject the bus let it publish on, never from here. +type HealthSaid struct { + Node string `json:"node"` + Health Health `json:"health"` +} diff --git a/internal/link/protocol_test.go b/internal/link/protocol_test.go index 53ef4ce..7c1e89c 100644 --- a/internal/link/protocol_test.go +++ b/internal/link/protocol_test.go @@ -25,6 +25,13 @@ func TestTheWireFormatIsExactlyTheseFieldNames(t *testing.T) { []string{"protocol", "address", "port", "by", "published", "container-port"}}, {EnrolRequest{Node: "n", Secret: "s", PublicKey: []byte("k")}, []string{"node", "secret", "public_key"}}, + // 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: "restarting", 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/receive.go b/internal/link/receive.go index 7222b38..3e4cbdb 100644 --- a/internal/link/receive.go +++ b/internal/link/receive.go @@ -30,8 +30,10 @@ const ( KindHeartbeat = "heartbeat" // KindToolsHeartbeat is a machine's node tools saying they are there (novox/hq to-be 45 S11). KindToolsHeartbeat = "tools-heartbeat" - KindBuilt = "built" - KindModuleMoved = "module-moved" + // KindHealth is a machine's statement of its long-running resources' health (novox/hq ADR 0240). + KindHealth = "health" + KindBuilt = "built" + KindModuleMoved = "module-moved" // KindSourceMoved is the forge announcing a merge: a source moved, and what it produces is // built without anybody telling the mesh (novox/hq 04-ISSUES/131). KindSourceMoved = "source-moved" diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index 3200d81..f3d9d5e 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -108,6 +108,13 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con return fmt.Errorf("subscribing to the node tools' heartbeats: %w", err) } defer func() { _ = toolsAlive.Unsubscribe() }() + // And what each machine says of its long-running resources between reports (novox/hq ADR 0240): core + // too, and said again while it matters. + health, err := conn.ChanSubscribe(HealthSubjects, beats) + if err != nil { + return fmt.Errorf("subscribing to the machines' health: %w", err) + } + defer func() { _ = health.Unsubscribe() }() // The events the controller follows, when something is listening for them. var events chan *nats.Msg @@ -239,6 +246,8 @@ func kindOfSubject(subject string) (string, bool) { return KindHeartbeat, true case "tools-alive": return KindToolsHeartbeat, true + case "health": + return KindHealth, true } } switch subject { diff --git a/internal/link/serve.go b/internal/link/serve.go index 4a8bac8..f07627e 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -94,6 +94,8 @@ type Server struct { retirements Retirements // checker asks for a pull request's merge check (novox/hq to-be 45 §9). checker Checker + // healths keeps what machines say of their long-running resources (novox/hq ADR 0240). + healths Healths log *log.Logger // giveUp is how long one message is held for the store; zero means GiveUpAfter. @@ -194,6 +196,8 @@ func (s *Server) act(ctx context.Context, m Control) { s.heartbeat(m) case KindToolsHeartbeat: s.toolsHeartbeat(m) + case KindHealth: + s.healthSaid(ctx, m) case KindBuilt: s.wasBuilt(ctx, m) case KindModuleMoved: