Merge pull request 'Heal what is known, under a brake, and say every repair (hq to-be 45 Phase 3)' (#85) from feat/a-core-that-cannot-fail-silently-phase-3 into main
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on

This commit was merged in pull request #85.
This commit is contained in:
2026-10-06 12:29:15 +00:00
31 changed files with 2154 additions and 21 deletions
+4
View File
@@ -691,6 +691,10 @@ type answers struct {
// where the log is not on hand; handActsUnread why it could not be read when it could not.
handActs *int
handActsUnread string
// heals is what the healers did in the last seven days (novox/hq to-be 45 §7): acts, and conditions
// handed to the operator; healsUnread why it could not be read.
heals *healsCount
healsUnread string
}
// heldBy is every artifact this mesh has built, for a build that may need one as its base.
+9 -1
View File
@@ -250,7 +250,7 @@ func showCondition(ctx context.Context, args []string) error {
if len(c.Tried) > 0 {
fmt.Println("\n tried:")
for _, t := range c.Tried {
fmt.Printf(" %s %s: %s\n", t.At.Local().Format("2006-01-02 15:04"), t.What, t.Outcome)
fmt.Printf(" %s %s — %s: %s\n", t.At.Local().Format("2006-01-02 15:04"), orHealer(t.By), t.What, t.Outcome)
}
}
fmt.Println("\n evidence, newest first:")
@@ -431,3 +431,11 @@ func printConditions(list []conditions.Condition, unread string, now time.Time)
fmt.Printf("\n `conditions show <key>` says more; each clears when observation says it is resolved, " +
"never by hand — `conditions silence <key> --for <d> --why <text>` stops its messages\n\n")
}
// orHealer is who tried, as an attempt names it.
func orHealer(by string) string {
if by == "" {
return "a healer"
}
return by
}
+10 -3
View File
@@ -57,8 +57,11 @@ type probe struct {
Asserts string
From string
// Kind is the condition kind raised when the invariant does not hold.
Kind string
Phase int
Kind string
// Raises are the other kinds its findings carry, each its own (a consumer missing, one far behind):
// what a healer may be registered against (healers.go).
Raises []string
Phase int
// Deferred says why it is not run yet; empty for one that is.
Deferred string
// Asks are the seat verbs it calls. A probe may call no other (askSeatTool refuses), and the
@@ -83,7 +86,8 @@ var probeRegistry = []probe{
"holds no other epoch open; no message from a stale epoch refused in the last interval",
From: "issue 204", Kind: "lease-split", Phase: 2, run: probeLease},
{ID: "D6", Asserts: "every durable consumer the mesh expects exists with its definition, and is near its " +
"stream's head", From: "issues 248, 266", Kind: "consumer-wrong", Phase: 1, run: probeConsumers},
"stream's head", From: "issues 248, 266", Kind: "consumer-wrong", Raises: []string{"consumer-lost", kindConsumerBehind},
Phase: 1, run: probeConsumers},
{ID: "D7", Asserts: "every stream the controller defines exists with its definition, and its own buckets",
From: "issue 208", Kind: "stream-wrong", Phase: 1, run: probeStreams},
{ID: "D8", Asserts: "no address the mesh owns — a machine's private address or its endpoint — is in a ban list",
@@ -396,6 +400,9 @@ func probesAnswer() map[string]any {
var out []map[string]any
for _, p := range probeRegistry {
row := map[string]any{"id": p.ID, "asserts": p.Asserts, "from": p.From, "kind": p.Kind, "phase": p.Phase}
if len(p.Raises) > 0 {
row["raises"] = p.Raises
}
if p.Deferred != "" {
row["deferred"] = p.Deferred
}
+865
View File
@@ -0,0 +1,865 @@
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"slices"
"sort"
"strings"
"sync"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// The healers (novox/hq to-be 45 §7, ADR 0227 rule 7, Phase 3).
//
// **A known failure heals itself, under a brake, and every repair is said.** Research 031 counted the
// repairs people made by hand in six days: a push to unstick a plan waiting on a report (four times),
// a controller restarted to make an object again (twice), a plan closed (twice), a consumer re-made from
// now (once). Each was the ordinary path, taken again by a person who noticed. A healer is that act,
// registered against the one condition kind it answers, so the mesh takes it itself:
//
// - **its repair is the ordinary path again** — the send a push makes, the plan's own close, the
// assertion every send makes, the consumer reset the verb makes — never a withdrawal, a deletion of
// data or a recreation of it;
// - **its budget** is how many acts it may take against one thing in a window, and **its settle** how
// long after an act it leaves the observation to say whether it worked;
// - **success is never the healer's to say.** The condition clears when the watchdog or the probe that
// raised it no longer sees it. A healer marking its own work done would be the second opinion of a
// fact that rule 1 forbids;
// - **a spent budget stops it**: the condition is the operator's, urgent, with every attempt in
// `tried`, and no healer touches it again until observation clears it;
// - **every act is said**: begun in the store before it is made (so a controller dying mid-act has
// still spent it), kept in the condition's `tried` as `healer Hn`, and said on the bus as the
// controller seat's `healer-acted`. A heal is not a hand act and is never in the hand-act log —
// which is how S15 can tell a repair the mesh made from one a person had to.
//
// **And the mesh-wide brake**: more than healBrakeLimit acts in an hour, all healers together, and every
// healer stops — an urgent condition says so — until an hour has passed with none. A healer looping is
// then at most a dozen acts, said, and never the incident itself.
//
// Only the controller holding the lease heals, under its epoch: a controller serving without the lease
// (S12) heals nothing, because a repair is the one act that can always wait for the lease.
// The mesh-wide brake (to-be 45 §7): how many acts all healers together may take in an hour.
const (
healBrakeLimit = 12
healBrakeWindow = time.Hour
)
// healEvery is how often the healers look at what is open.
var healEvery = 30 * time.Second
// What raises the healers' own conditions, and their kinds.
const (
sourceHealers = "healers"
kindHealersBraked = "healers-braked"
kindConsumerBehind = "consumer-behind"
kindHealersBlind = "probe-failed"
)
// healerRow is one row of the registry.
type healerRow struct {
ID string
// Kinds are the condition kinds it answers: from the signals table, the probe registry, or an event.
Kinds []string
// Condition, Repair, Then, From are the row in words, as to-be 45 §7 and `healers` say it.
Condition string
Repair string
Then string
From string
// Budget acts against one budget key within Window; Settle after an act before the next, or the
// escalation, so the observation has had its turn to say whether it worked.
Budget int
Window time.Duration
Settle time.Duration
// ActsIn is where the repair runs: the controller, or a module that repairs its own (H5).
ActsIn string
// Event is what each act is said as.
Event string
// applies says whether this condition is one the healer may act on, and what its budget is
// counted against; why when it is not. Reads, never acts. Nil where the repair is a module's.
applies func(ctx context.Context, h *healing, c conditions.Condition) (budget string, ok bool, why string, err error)
// repair is the act: what it did in words, what came of it, or an error when it could not be done.
repair func(ctx context.Context, h *healing, c conditions.Condition) (act, said string, err error)
}
// actsInController is a healer whose repair the controller makes.
const actsInController = "controller"
// healerRegistry is to-be 45 §7's table, compiled in. **The registry is the design's live form**: a
// test generated from it holds every row to a condition kind the mesh raises, a budget, a brake and its
// event (healers_test.go).
var healerRegistry = []healerRow{
{ID: "H1", Kinds: []string{"sent-not-reported"},
Condition: "a machine was sent a declaration and has not reported it (S2)",
Repair: "ask the machine's node-engine to say again what it last applied (the `report` verb); if what " +
"comes back is not the declaration it was sent, or nothing comes, send it its current declaration " +
"again — what a push of that one machine does, never moving a build a policy or a plan holds back",
Then: "resolver operator, urgent", From: "a push by hand to unstick a plan waiting on a report (230, 257, 264, 267)",
Budget: 2, Window: 6 * time.Hour, Settle: 3 * time.Minute, ActsIn: actsInController, Event: link.KeyHealerActed,
applies: appliesToAMachine, repair: repairReport},
{ID: "H2", Kinds: []string{"stalled"},
Condition: "a plan stalled on a wait that is superseded — a newer plan of its repository and branch " +
"exists — or already finished — every module of every tier built or failed, and every one that " +
"rolls out sent (S3)",
Repair: "close the plan with its note, as `plans close` does: superseded, naming the newer plan, or " +
"done; what it asked still builds and registers",
Then: "resolver operator, urgent", From: "a stuck plan closed by hand (214, 254)",
Budget: 1, Window: 24 * time.Hour, Settle: 2 * time.Minute, ActsIn: actsInController, Event: link.KeyHealerActed,
applies: appliesToAStalePlan, repair: repairPlan},
{ID: "H3", Kinds: []string{"holder-silent", "consumer-lost"},
Condition: "a seat's holder that does not answer (D3), or a durable consumer the mesh expects and the " +
"bus does not hold (D6, S9)",
Repair: "assert the bus's streams, consumers and seat workers again — the assertion every send makes " +
"(issue 208)",
Then: "resolver operator, urgent", From: "a controller restarted to make a missing object again (208, 248)",
Budget: 1, Window: time.Hour, Settle: 6 * time.Minute, ActsIn: actsInController, Event: link.KeyHealerActed,
applies: appliesToAnObject, repair: repairObjects},
{ID: "H4", Kinds: []string{kindConsumerBehind},
Condition: "a durable consumer far behind its stream's head (D6) that the stream table marks resettable",
Repair: "re-make the consumer to deliver from now — `broker consumer-reset` (issue 248); what it drops " +
"is caught up where the table says",
Then: "resolver operator, urgent", From: "a consumer re-made from now by hand (248)",
Budget: 1, Window: 24 * time.Hour, Settle: 6 * time.Minute, ActsIn: actsInController, Event: link.KeyHealerActed,
applies: appliesToAResettableConsumer, repair: repairConsumer},
{ID: "H5", Kinds: []string{kindProviderFailing},
Condition: "the identity provider's administrator refusing the mesh's secret (provider-failing, " +
"credentials-rejected)",
Repair: "the provider repairs it itself through the server's own bootstrap command, checks again and " +
"says what it did (ADR 0224 §5); the controller keeps its word as the condition",
Then: "announced failing by the provider, braked from ten minutes doubling to six hours (ADR 0224 §5)",
From: "the identity provider's admin reset through its bootstrap command (179)",
// Its budget and brake are the module's own: ten minutes, doubling, to six hours.
Budget: 1, Window: 10 * time.Minute, Settle: 10 * time.Minute,
ActsIn: "the provider module (keycloak), ADR 0224 §5", Event: "provisioner.failing, provisioner.recovered"},
}
// healerFor is the healer registered for a kind, and false when none acts in the controller.
func healerFor(kind string) (healerRow, bool) {
for _, r := range healerRegistry {
if r.repair != nil && slices.Contains(r.Kinds, kind) {
return r, true
}
}
return healerRow{}, false
}
// healerNamedFor is any healer registered for a kind, wherever it acts: what S15 says beside a cause.
func healerNamedFor(kind string) string {
for _, r := range healerRegistry {
if slices.Contains(r.Kinds, kind) {
return r.ID
}
}
return ""
}
// healing is the healers' runner and what their repairs reach.
type healing struct {
open *stores
keeper *conditions.Keeper
teller conditions.Teller
// js is the bus, for the repairs that assert or reset its objects; nil where there is none.
js *broker.JetStream
// epoch is the lease's gate; acting says this controller is the one acting, not one standing by.
epoch func(ctx context.Context) (uint64, error)
acting func() bool
now func() time.Time
say func(format string, args ...any)
// The acts, as seams a test replaces: asking a machine to report, sending it again, asserting the
// bus's objects, resetting a consumer.
askReport func(ctx context.Context, node string) error
sendAgain func(ctx context.Context, node string) error
assertObjects func(ctx context.Context) error
resetConsumer func(stream, name string) (string, error)
// reportWait is how long H1 waits for the machine's report after asking.
reportWait time.Duration
mu sync.Mutex
// declined is why each condition was last passed over, so it is said once and `healers` can show it.
declined map[string]string
paused string
}
// newHealing is the serving controller's runner, its acts the real ones.
func newHealing(open *stores, keeper *conditions.Keeper, teller conditions.Teller, js *broker.JetStream) *healing {
h := &healing{open: open, keeper: keeper, teller: teller, js: js, acting: link.Holding,
epoch: func(ctx context.Context) (uint64, error) { return theLease.epoch(ctx) },
now: time.Now, say: func(format string, args ...any) { fmt.Printf(format+"\n", args...) },
reportWait: 45 * time.Second, declined: map[string]string{}}
h.askReport = func(ctx context.Context, node string) error {
if h.js == nil {
return errors.New("this controller is not on the bus")
}
return askToReport(ctx, h.js.Conn(), node)
}
h.sendAgain = func(ctx context.Context, node string) error {
// **Never what a policy or a plan holds back** (ADR 0221): a send by the mesh itself is a push
// that did not name the machine — a person's word sends a held build, a healer's does not.
held, err := heldMachines(ctx, open, []string{node})
if err != nil {
return err
}
if why := held[node]; len(why) > 0 {
return fmt.Errorf("not sent again: it would move what a policy or a plan holds back (ADR 0221) — %s; "+
"`push %s` sends it, on a person's word", strings.Join(why, "; "), node)
}
return sendTo(ctx, open, []string{node})
}
h.assertObjects = func(ctx context.Context) error {
if h.js == nil {
return errors.New("this controller is not on the bus")
}
return assertOnSend(ctx, open.inventory, h.js, " ")
}
h.resetConsumer = func(stream, name string) (string, error) {
if h.js == nil {
return "", errors.New("this controller is not on the bus")
}
before, after, err := h.js.ResetConsumer(stream, name)
if err != nil {
return "", err
}
return fmt.Sprintf("it was %d behind with %d unacknowledged; it delivers from now, %d pending", before.Pending,
before.AckPending, after.Pending), nil
}
return h
}
// askToReport asks one machine's node-engine to say again what it last applied (to-be 45 §6). On core
// NATS, fired and flushed: the answer is the machine's ordinary report, read from the store.
func askToReport(ctx context.Context, conn *nats.Conn, node string) error {
body, err := json.Marshal(map[string]any{"asked": time.Now().UTC(), "by": "the controller's healer H1"})
if err != nil {
return err
}
if err := conn.Publish(broker.AskReportSubject(node), body); err != nil {
return err
}
flushing, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
return conn.FlushWithContext(flushing)
}
// keep runs the healers until ctx ends.
func (h *healing) keep(ctx context.Context) {
tick := time.NewTicker(healEvery)
defer tick.Stop()
for {
h.tick(ctx)
select {
case <-ctx.Done():
return
case <-tick.C:
}
}
}
// tick is one look at what is open, and every act it calls for.
func (h *healing) tick(ctx context.Context) {
if h.acting != nil && !h.acting() {
return
}
epoch, err := h.epoch(ctx)
switch {
case err != nil:
h.pause(fmt.Sprintf("this controller may not act: %v", err))
return
case epoch == 0:
h.pause("this controller serves without the lease (S12): a repair waits for it")
return
}
inv := h.open.inventory
now := h.now()
heals, err := inv.HealsSince(ctx, now.Add(-24*time.Hour))
if err != nil {
// Blind: nothing is done, and that is said — never read as "no heals, budgets whole".
h.reconcile(ctx, []conditions.Observation{{Scope: conditions.ScopeProbe, ID: "healers", Token: "failed",
Kind: kindHealersBlind, Severity: conditions.Warning,
Summary: "the healers cannot read what they did, so they count no budget and act on nothing",
Said: firstLine(err.Error())}})
return
}
open, err := h.keeper.Open(ctx)
if err != nil {
h.say("the healers cannot read the open conditions, and act on nothing: %v", err)
return
}
acts := actsWithin(heals, now.Add(-healBrakeWindow))
if braked, said := brakeHolds(acts, open, now); braked {
h.reconcile(ctx, []conditions.Observation{brakeObservation(acts, said)})
h.pause("the mesh-wide brake holds: " + said)
return
}
h.reconcile(ctx, nil)
h.resume()
for _, c := range open {
row, ok := healerFor(c.Kind)
if !ok || c.Escalated() {
continue
}
budget, applies, why, err := row.applies(ctx, h, c)
if err != nil {
h.decline(c.Key, row.ID, "could not tell whether it applies: "+err.Error())
continue
}
if !applies {
h.decline(c.Key, row.ID, why)
continue
}
h.forget(c.Key)
spent := spentOn(heals, row, budget, now)
if n := len(spent); n > 0 && now.Sub(spent[n-1].At) < row.Settle {
continue // its last act is still the observation's to judge
}
if len(spent) >= row.Budget {
h.escalate(ctx, row, c, budget, spent)
continue
}
if len(acts) >= healBrakeLimit {
return // the brake takes hold on the next look, said there
}
if h.act(ctx, row, c, budget, len(spent)) {
acts = append(acts, inventory.Heal{At: now})
}
}
}
// act is one healer's act on one condition: begun in the store, made, finished, kept in the
// condition's tried and said. False when it could not even be begun.
func (h *healing) act(ctx context.Context, row healerRow, c conditions.Condition, budget string, spent int) bool {
inv := h.open.inventory
begun, err := inv.BeginHeal(ctx, inventory.Heal{Healer: row.ID, ConditionKey: c.Key, Kind: c.Kind,
BudgetKey: budget, Act: row.Repair, At: h.now()})
if err != nil {
h.say("healer %s did not act on %s: %v", row.ID, c.Key, err)
return false
}
act, said, err := row.repair(ctx, h, c)
outcome := inventory.HealActed
if err != nil {
outcome, said = inventory.HealFailed, err.Error()
}
if act == "" {
act = row.Repair
}
finishing, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
defer cancel()
if ferr := inv.FinishHeal(finishing, begun.ID, outcome, act+" — "+said); ferr != nil {
h.say("healer %s acted on %s and could not record what came of it: %v", row.ID, c.Key, ferr)
}
budgetWords := fmt.Sprintf("act %d of %d within %s", spent+1, row.Budget, row.Window)
attempt := conditions.Attempt{What: act, Outcome: outcome + ": " + said + " (" + budgetWords + ")",
By: "healer " + row.ID}
if _, _, terr := h.keeper.Tried(finishing, c.Key, attempt, conditions.ResolverHealer(row.ID)); terr != nil {
h.say("healer %s acted on %s and could not keep it in the condition: %v", row.ID, c.Key, terr)
}
h.tell(finishing, healerActed{Healer: row.ID, Condition: c.Key, Kind: c.Kind, Act: act, Outcome: outcome,
Said: said, Budget: budgetWords, Epoch: begun.Epoch})
h.say("healer %s %s %s: %s — %s (%s)", row.ID, outcome, c.Key, act, said, budgetWords)
return true
}
// escalate is a spent budget: the condition is the operator's now, urgent, with what was tried. Said
// once: an escalated condition is passed over by every healer until observation clears it.
func (h *healing) escalate(ctx context.Context, row healerRow, c conditions.Condition, budget string, spent []inventory.Heal) {
var tried []string
for _, s := range spent {
tried = append(tried, fmt.Sprintf("%s at %s (%s)", s.Outcome, s.At.UTC().Format("15:04 MST"), firstLine(s.Said)))
}
said := fmt.Sprintf("its budget of %d act(s) within %s is spent and %s is still open: %s", row.Budget, row.Window,
c.Key, strings.Join(tried, "; "))
if _, err := h.open.inventory.BeginHeal(ctx, inventory.Heal{Healer: row.ID, ConditionKey: c.Key, Kind: c.Kind,
BudgetKey: budget, Act: "handed the condition to the operator", Outcome: inventory.HealEscalated, Said: said,
At: h.now()}); err != nil {
h.say("healer %s could not record handing %s to the operator, and does not: %v", row.ID, c.Key, err)
return
}
attempt := conditions.Attempt{What: "handed to the operator: nothing in the mesh will repair it now",
Outcome: inventory.HealEscalated + ": " + said, By: "healer " + row.ID}
if _, _, err := h.keeper.Escalate(ctx, c.Key, attempt); err != nil {
h.say("healer %s could not hand %s to the operator: %v", row.ID, c.Key, err)
}
h.tell(ctx, healerActed{Healer: row.ID, Condition: c.Key, Kind: c.Kind, Act: "handed to the operator",
Outcome: inventory.HealEscalated, Said: said, Budget: fmt.Sprintf("%d of %d within %s", len(spent), row.Budget, row.Window)})
h.say("healer %s handed %s to the operator: %s", row.ID, c.Key, said)
}
// healerActed is the body of `healer-acted` (to-be 45 §7): the healer, the condition, the act and its
// outcome, and where the budget stands. **A contract**, like a condition's events: the operator-channel's
// holder may say it.
type healerActed struct {
Event string `json:"event"`
At time.Time `json:"at"`
Healer string `json:"healer"`
By string `json:"by"`
Condition string `json:"condition"`
Kind string `json:"kind"`
Act string `json:"act"`
// Outcome is acted, failed or escalated. Acted says the act was made, not that it worked: the
// condition clearing says that.
Outcome string `json:"outcome"`
Said string `json:"said"`
Budget string `json:"budget"`
Epoch uint64 `json:"epoch,omitempty"`
Show string `json:"show"`
}
// tell says a heal on the bus. Not said is said in the log: the act and the condition's tried still
// hold it.
func (h *healing) tell(ctx context.Context, e healerActed) {
e.Event, e.At, e.By = link.KeyHealerActed, h.now().UTC(), "healer "+e.Healer
e.Show = "mesh-controller.conditions key=" + e.Condition
if h.teller == nil {
return
}
body, err := json.Marshal(e)
if err != nil {
return
}
saying, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
if err := h.teller.PublishSeatEvent(saying, conditions.Seat, link.KeyHealerActed, body); err != nil {
h.say("healer %s %s %s, and that could NOT be said on the bus: %v", e.Healer, e.Outcome, e.Condition, err)
}
}
// reconcile keeps the healers' own conditions: the brake, and their being blind.
func (h *healing) reconcile(ctx context.Context, observed []conditions.Observation) {
if err := h.keeper.Reconcile(ctx, sourceHealers, observed); err != nil {
h.say("the healers' own conditions could not be kept: %v", err)
}
}
func (h *healing) pause(why string) {
h.mu.Lock()
defer h.mu.Unlock()
if h.paused != why {
h.say("the healers act on nothing: %s", why)
}
h.paused = why
}
func (h *healing) resume() {
h.mu.Lock()
defer h.mu.Unlock()
if h.paused != "" {
h.say("the healers act again")
}
h.paused = ""
}
// decline remembers why a healer passed a condition over, said once.
func (h *healing) decline(key, healer, why string) {
h.mu.Lock()
defer h.mu.Unlock()
if h.declined[key] != why {
h.say("healer %s leaves %s alone: %s", healer, key, why)
}
h.declined[key] = why
}
func (h *healing) forget(key string) {
h.mu.Lock()
defer h.mu.Unlock()
delete(h.declined, key)
}
// actsWithin are the heals since a moment that were acts — begun, made or failed; an escalation is a
// saying, not an act, and the brake does not count it.
func actsWithin(heals []inventory.Heal, since time.Time) []inventory.Heal {
var out []inventory.Heal
for _, h := range heals {
if !h.At.Before(since) && h.Outcome != inventory.HealEscalated {
out = append(out, h)
}
}
return out
}
// spentOn is one healer's acts against one budget key within its window, oldest first.
func spentOn(heals []inventory.Heal, row healerRow, budget string, now time.Time) []inventory.Heal {
var out []inventory.Heal
for _, h := range actsWithin(heals, now.Add(-row.Window)) {
if h.Healer == row.ID && h.BudgetKey == budget {
out = append(out, h)
}
}
return out
}
// brakeHolds says whether the mesh-wide brake holds: the limit reached within the hour, or the brake
// already said and an act within the hour still — it lets go an hour after the last act, not the first.
func brakeHolds(acts []inventory.Heal, open []conditions.Condition, now time.Time) (bool, string) {
said := fmt.Sprintf("%d act(s) by the healers in the last %s, the limit %d", len(acts), healBrakeWindow, healBrakeLimit)
if len(acts) >= healBrakeLimit {
return true, said
}
held := slices.ContainsFunc(open, func(c conditions.Condition) bool { return c.Kind == kindHealersBraked })
if held && len(acts) > 0 {
last := acts[len(acts)-1].At
return true, fmt.Sprintf("%s; held until an hour after the last, at %s", said,
last.Add(healBrakeWindow).UTC().Format("15:04 MST"))
}
return false, said
}
// brakeObservation is the brake, as the urgent condition that says it.
func brakeObservation(acts []inventory.Heal, said string) conditions.Observation {
counts := map[string]int{}
for _, a := range acts {
if a.Healer != "" {
counts[a.Healer+" on "+a.ConditionKey]++
}
}
var most []string
for what, n := range counts {
most = append(most, fmt.Sprintf("%s ×%d", what, n))
}
sort.Strings(most)
return conditions.Observation{Scope: conditions.ScopeMesh, ID: "healers", Token: "braked", Kind: kindHealersBraked,
Severity: conditions.Urgent,
Summary: "every healer has stopped: the healers acted more often in an hour than the mesh allows, which is a " +
"healer looping, not repairing — " + said,
Said: strings.Join(most, ", ")}
}
// --- H1 ----------------------------------------------------------------------------------------------
// appliesToAMachine is H1's: a machine named, its budget counted against the condition.
func appliesToAMachine(_ context.Context, _ *healing, c conditions.Condition) (string, bool, string, error) {
if c.Subject.Machine == "" {
return "", false, "it names no machine", nil
}
return c.Key, true, "", nil
}
// repairReport is H1: ask the machine to report; if what comes back is not what it was sent, or nothing
// comes, send it again.
func repairReport(ctx context.Context, h *healing, c conditions.Condition) (string, string, error) {
node := c.Subject.Machine
inv := h.open.inventory
// The store's clock, as the report's time is: not the runner's, which a test moves by hand.
asked := time.Now()
askErr := h.askReport(ctx, node)
current, heard := false, false
if askErr == nil {
deadline := time.Now().Add(h.reportWait)
for {
r, found, err := lastReportOf(ctx, inv, node)
if err != nil {
return "", "", err
}
if found && r.At != nil && r.At.After(asked) {
heard, current = true, r.Current
if current {
break
}
}
if time.Now().After(deadline) {
break
}
select {
case <-ctx.Done():
return "", "", ctx.Err()
case <-time.After(min(2*time.Second, h.reportWait/4+time.Millisecond)):
}
}
}
if current {
return "asked " + node + "'s node-engine to say again what it last applied",
"it reported the declaration it was sent: its report had not reached the mesh", nil
}
why := "it said nothing within " + h.reportWait.String() + " — its node-engine may be older than the report verb"
switch {
case askErr != nil:
why = "it could not be asked: " + askErr.Error()
case heard:
why = "it reported a declaration other than the one it was sent"
}
if err := h.sendAgain(ctx, node); err != nil {
return "asked " + node + " to report, then sent it its current declaration again", why + "; the send failed",
err
}
return "asked " + node + " to report, then sent it its current declaration again", why + "; sent again", nil
}
// lastReportOf is one machine's last report beside its last send.
func lastReportOf(ctx context.Context, inv *inventory.Inventory, node string) (inventory.Reported, bool, error) {
reports, err := inv.LastReports(ctx)
if err != nil {
return inventory.Reported{}, false, err
}
for _, r := range reports {
if r.Node == node {
return r, true, nil
}
}
return inventory.Reported{}, false, nil
}
// --- H2 ----------------------------------------------------------------------------------------------
// planStale is why a plan's wait is superseded or finished, and the state closing it leaves it in; empty
// when it is neither, which is not H2's to repair.
func planStale(ctx context.Context, inv *inventory.Inventory, p inventory.Plan) (state, why string, err error) {
recent, err := inv.RecentPlans(ctx, 50)
if err != nil {
return "", "", err
}
for _, newer := range recent {
if newer.ID == p.ID || !newer.Created.After(p.Created) || !repositoryMatches(newer.Repository, p.Repository) ||
newer.Branch != p.Branch || newer.State == inventory.PlanSuperseded {
continue
}
return inventory.PlanSuperseded, fmt.Sprintf("superseded by %s (%s %s), a newer merge of the same repository "+
"and branch", newer.ID, newer.Repository, short(newer.Commit)), nil
}
for _, tier := range p.Tiers {
for _, m := range tier {
s := p.Modules[m]
if s == nil || (s.State != "built" && s.State != "failed") {
return "", "", nil
}
if s.State == "built" && s.SentAt == nil {
u, err := inv.UpgradeOf(ctx, m)
if err != nil {
return "", "", err
}
if u.RollOut {
return "", "", nil
}
}
}
}
return inventory.PlanDone, "finished: every module of every tier is built or failed, and every one that rolls " +
"out was sent — nothing is left to wait on", nil
}
// appliesToAStalePlan is H2's: the plan is open, and its wait is superseded or finished.
func appliesToAStalePlan(ctx context.Context, h *healing, c conditions.Condition) (string, bool, string, error) {
p, err := h.open.inventory.PlanByID(ctx, c.Subject.ID)
if err != nil {
return "", false, "", err
}
if !p.Open() {
return "", false, "the plan is " + p.State + " already: its condition clears on the next look", nil
}
state, _, err := planStale(ctx, h.open.inventory, p)
if err != nil {
return "", false, "", err
}
if state == "" {
return "", false, "its wait is neither superseded nor finished: it waits on something still to come, " +
"which its own signal says", nil
}
return c.Key, true, "", nil
}
// repairPlan is H2: the plan closed with its note, under the plans' hold, by compare-and-set.
func repairPlan(ctx context.Context, h *healing, c conditions.Condition) (string, string, error) {
inv := h.open.inventory
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return "", "", err
}
defer release()
p, err := inv.PlanByID(ctx, c.Subject.ID)
if err != nil {
return "", "", err
}
if !p.Open() {
return "closed nothing", "the plan is " + p.State + " already", nil
}
state, why, err := planStale(ctx, inv, p)
if err != nil {
return "", "", err
}
if state == "" {
return "closed nothing", "its wait is no longer superseded or finished", nil
}
p.State = state
p.Note = fmt.Sprintf("closed by healer H2 at tier %d: %s", p.Tier, why)
if state != inventory.PlanDone {
sayUnsent(&p, func(m string) bool {
u, err := inv.UpgradeOf(ctx, m)
return err == nil && u.RollOut
})
}
if err := inv.SavePlan(ctx, &p); err != nil {
return "", "", err
}
return "closed the plan " + p.ID + " (" + state + ")", why, nil
}
// --- H3 ----------------------------------------------------------------------------------------------
// appliesToAnObject is H3's: every holder or consumer the probe or the bus names, its budget per object.
func appliesToAnObject(_ context.Context, h *healing, c conditions.Condition) (string, bool, string, error) {
if h.assertObjects == nil {
return "", false, "this controller is not on the bus", nil
}
return c.Key, true, "", nil
}
// repairObjects is H3: the send's own assertion of every stream, consumer and seat worker.
func repairObjects(ctx context.Context, h *healing, _ conditions.Condition) (string, string, error) {
if err := h.assertObjects(ctx); err != nil {
return "", "", err
}
return "asserted the bus's streams, consumers and seat workers again",
"every object the mesh defines is asserted; the probe that raised it says on its next run whether it is there", nil
}
// --- H4 ----------------------------------------------------------------------------------------------
// resettable is why a consumer may be reset by the mesh itself, from the stream table; empty when not.
func resettable(stream, name string) string {
for _, c := range broker.MeshConsumers() {
if c.Stream == stream && c.Name == name {
return c.Resettable
}
}
return ""
}
// consumerOf is the stream and consumer a bus condition names: `<stream>.<consumer>`.
func consumerOf(c conditions.Condition) (string, string, bool) {
stream, name, found := strings.Cut(c.Subject.ID, ".")
return stream, name, found && stream != "" && name != ""
}
// appliesToAResettableConsumer is H4's: only a consumer the stream table marks resettable.
func appliesToAResettableConsumer(_ context.Context, _ *healing, c conditions.Condition) (string, bool, string, error) {
stream, name, ok := consumerOf(c)
if !ok {
return "", false, "it names no consumer", nil
}
if resettable(stream, name) == "" {
return "", false, fmt.Sprintf("%s on %s is not marked resettable in the stream table: what a reset drops, "+
"nothing would catch up", name, stream), nil
}
return c.Key, true, "", nil
}
// repairConsumer is H4: `broker consumer-reset`, made by the mesh.
func repairConsumer(_ context.Context, h *healing, c conditions.Condition) (string, string, error) {
stream, name, _ := consumerOf(c)
said, err := h.resetConsumer(stream, name)
if err != nil {
return "", "", err
}
return fmt.Sprintf("re-made %s on %s to deliver from now", name, stream),
said + "; " + resettable(stream, name), nil
}
// --- the verb ----------------------------------------------------------------------------------------
// healersCommand is `healers`: the registry, what the healers did lately, and the brake.
func healersCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("healers", flag.ContinueOnError)
daysFlag := set.Int("days", 7, "how many days of acts back")
jsonFlag := set.Bool("json", false, "as data")
if rest, err := parseAround(set, args); err != nil {
return err
} else if len(rest) > 0 {
return errors.New("healers [--days N] [--json]")
}
days, asJSON := *daysFlag, *jsonFlag
if days <= 0 {
return fmt.Errorf("healers --days takes a number of days, not %d", days)
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
now := time.Now()
heals, err := open.inventory.HealsSince(ctx, now.Add(-time.Duration(days)*24*time.Hour))
if err != nil {
return fmt.Errorf("what the healers did cannot be read: %w", err)
}
conds, condErr := openConditions(ctx)
braked, said := brakeHolds(actsWithin(heals, now.Add(-healBrakeWindow)), conds, now)
brake := map[string]any{"holds": braked, "said": said, "limit": healBrakeLimit, "window": healBrakeWindow.String()}
if condErr != nil {
brake["conditions"] = "the open conditions could not be read, so whether the brake was said is not known: " +
condErr.Error()
}
var rows []map[string]any
for _, r := range healerRegistry {
rows = append(rows, map[string]any{"id": r.ID, "kinds": r.Kinds, "condition": r.Condition, "repair": r.Repair,
"budget": fmt.Sprintf("%d act(s) within %s, %s apart at least", r.Budget, r.Window, r.Settle),
"then": r.Then, "acts-in": r.ActsIn, "event": r.Event, "from": r.From})
}
if heals == nil {
heals = []inventory.Heal{}
}
if asJSON {
return printJSON(map[string]any{"healers": rows, "heals": heals, "brake": brake, "days": days,
"note": "a heal is never a hand act; whether it repaired anything is the condition's clearing to say"})
}
for _, r := range healerRegistry {
fmt.Printf("%s %s\n → %s\n budget %d within %s; then %s (in %s)\n", r.ID, r.Condition, r.Repair, r.Budget,
r.Window, r.Then, r.ActsIn)
}
fmt.Printf("\nthe brake: %s — %s\n\n", map[bool]string{true: "HOLDS, every healer has stopped", false: "off"}[braked], said)
if len(heals) == 0 {
fmt.Printf("no healer acted in the last %d day(s)\n", days)
return nil
}
for i := len(heals) - 1; i >= 0; i-- {
x := heals[i]
fmt.Printf("%s %s %-9s %s\n %s\n", x.At.Local().Format("2006-01-02 15:04"), x.Healer, x.Outcome, x.ConditionKey,
firstLine(x.Said))
}
return nil
}
// forgettingOldHeals removes heals past their keeping once a day, by the controller acting only.
func forgettingOldHeals(ctx context.Context, inv *inventory.Inventory) {
for {
if link.Holding() {
if n, err := inv.ForgetOldHeals(ctx); err != nil {
fmt.Printf("heals older than %s could not be removed: %v\n", inventory.HealsKeptFor, err)
} else if n > 0 {
fmt.Printf("removed %d heal(s) older than %s\n", n, inventory.HealsKeptFor)
}
}
select {
case <-ctx.Done():
return
case <-time.After(24 * time.Hour):
}
}
}
// healsCount is what the healers did, counted for `status`.
type healsCount struct {
Acts int `json:"acts"`
Escalated int `json:"escalated"`
}
// countHeals counts acts and escalations.
func countHeals(heals []inventory.Heal) *healsCount {
out := &healsCount{}
for _, h := range heals {
if h.Outcome == inventory.HealEscalated {
out.Escalated++
} else {
out.Acts++
}
}
return out
}
+647
View File
@@ -0,0 +1,647 @@
package main
import (
"context"
"encoding/json"
"errors"
"os"
"slices"
"strconv"
"strings"
"sync"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// The test generated from the healer registry (novox/hq to-be 45 §7, ADR 0227 rule 7 "how it is
// checked"): **every row is walked.** Its kinds are ones the mesh raises — a row of the signals table, a
// probe of the self-check, or a provider's event — it has a budget, a window and a settle inside it,
// what happens when the budget is spent, and the event each act is said as; a healer the controller runs
// has its applies and its repair, and an induced failure below that sees it act, say so and brake. A
// row added without one fails, so the registry cannot grow a healer nobody has seen act.
// raisedKinds is every condition kind the mesh raises, and what raises it.
func raisedKinds() map[string]string {
out := map[string]string{kindProviderFailing: "the provisioner.failing event (ADR 0224)"}
for _, r := range signalsTable {
for _, k := range kindsOf(r) {
out[k] = r.Row
}
}
for _, p := range probeRegistry {
out[p.Kind] = p.ID
for _, k := range p.Raises {
out[k] = p.ID
}
}
return out
}
// inducedFailures are the healers seen acting in this file, by id: a row without one fails.
var inducedFailures = map[string]string{
"H1": "TestH1AsksAMachineToReportAndSendsItAgain",
"H2": "TestH2ClosesAPlanAnotherHasTakenOver",
"H3": "TestNatsH3AssertsAMissingConsumerAgainAndBrakesAfterItsBudget",
"H4": "TestNatsH4ResetsTheControllersEventsConsumerAndItStillDelivers",
}
func TestEveryHealerAnswersAKindTheMeshRaisesWithABudgetABrakeAndItsEvent(t *testing.T) {
kinds := raisedKinds()
source, err := os.ReadFile("healers_test.go")
if err != nil {
t.Fatal(err)
}
seen := map[string]bool{}
answered := map[string]string{}
for _, r := range healerRegistry {
t.Run(r.ID, func(t *testing.T) {
if seen[r.ID] {
t.Fatalf("%s is in the registry twice", r.ID)
}
seen[r.ID] = true
if len(r.Kinds) == 0 || r.Condition == "" || r.Repair == "" || r.Then == "" || r.From == "" || r.ActsIn == "" {
t.Fatalf("%s does not say what it answers, what it does, what then, and which hand act it replaces: %+v", r.ID, r)
}
for _, k := range r.Kinds {
if _, ok := kinds[k]; !ok {
t.Errorf("%s answers %q, which nothing in the mesh raises", r.ID, k)
}
if other, twice := answered[k]; twice {
t.Errorf("%q is answered by %s and %s: one healer per kind", k, other, r.ID)
}
answered[k] = r.ID
}
if r.Budget <= 0 || r.Window <= 0 || r.Settle <= 0 || r.Settle > r.Window {
t.Errorf("%s has no budget it can spend: %d within %s, settled after %s", r.ID, r.Budget, r.Window, r.Settle)
}
if r.Event == "" {
t.Errorf("%s says nothing when it acts", r.ID)
}
if r.ActsIn != actsInController {
if r.repair != nil || r.applies != nil {
t.Errorf("%s acts in %s and the controller would act for it too", r.ID, r.ActsIn)
}
return
}
if r.repair == nil || r.applies == nil {
t.Fatalf("%s acts in the controller and has no repair or no applies", r.ID)
}
if r.Event != link.KeyHealerActed || !slices.Contains(broker.ControllerStates, r.Event) {
t.Errorf("%s is said as %q, which the controller's grant does not permit", r.ID, r.Event)
}
if inducedFailures[r.ID] == "" {
t.Errorf("%s has no induced failure: a healer nobody has seen act", r.ID)
} else if !strings.Contains(string(source), "func "+inducedFailures[r.ID]+"(t *testing.T)") {
t.Errorf("%s's induced failure %s is not a test in this file", r.ID, inducedFailures[r.ID])
}
})
}
for id := range inducedFailures {
if !seen[id] {
t.Errorf("an induced failure for %s, which the registry does not have", id)
}
}
if healBrakeLimit <= 0 || healBrakeWindow <= 0 {
t.Error("the mesh-wide brake holds nothing")
}
}
// heard keeps the healer-acted events said.
type heard struct {
mu sync.Mutex
acts []healerActed
}
func (h *heard) PublishSeatEvent(_ context.Context, seat, event string, body []byte) error {
if seat != conditions.Seat || event != link.KeyHealerActed {
return errors.New("said under the wrong seat or name: " + seat + " " + event)
}
var e healerActed
if err := json.Unmarshal(body, &e); err != nil {
return err
}
h.mu.Lock()
defer h.mu.Unlock()
h.acts = append(h.acts, e)
return nil
}
func (h *heard) said() []healerActed {
h.mu.Lock()
defer h.mu.Unlock()
return append([]healerActed(nil), h.acts...)
}
// testClock is a moment a test moves by hand.
type testClock struct {
mu sync.Mutex
at time.Time
}
func (c *testClock) now() time.Time {
c.mu.Lock()
defer c.mu.Unlock()
return c.at
}
func (c *testClock) pass(d time.Duration) {
c.mu.Lock()
defer c.mu.Unlock()
c.at = c.at.Add(d)
}
// healingOn is a runner over a mesh's stores and its condition store, under epoch 57, every act a
// fake that fails the test unless the test gives it.
func healingOn(t *testing.T, open *stores) (*healing, *heard, *testClock) {
t.Helper()
if conditionsFrom == nil {
t.Fatal("the mesh has no condition store")
}
told, clock := &heard{}, &testClock{at: time.Now()}
// The store's gate, as the serving controller's is the lease's (stores.go).
open.inventory.ActsUnder(func(context.Context) (uint64, error) { return 57, nil })
h := &healing{open: open, keeper: conditionsFrom, teller: told,
epoch: func(context.Context) (uint64, error) { return 57, nil }, now: clock.now,
say: func(f string, a ...any) { t.Logf(f, a...) }, reportWait: 300 * time.Millisecond,
declined: map[string]string{}}
h.askReport = func(context.Context, string) error { t.Error("asked a machine to report"); return nil }
h.sendAgain = func(context.Context, string) error { t.Error("sent a machine again"); return nil }
h.assertObjects = func(context.Context) error { t.Error("asserted the bus's objects"); return nil }
h.resetConsumer = func(string, string) (string, error) { t.Error("reset a consumer"); return "", nil }
return h, told, clock
}
// sentNotReported raises S2 for a machine, as the watchdog does.
func sentNotReported(t *testing.T, node string) string {
t.Helper()
o := conditions.Observation{Scope: conditions.ScopeMachine, ID: node, Kind: "sent-not-reported", Machine: node,
Severity: conditions.Warning, Summary: node + " was sent a declaration and has not reported it", Source: "S2"}
if _, err := conditionsFrom.Observe(t.Context(), o); err != nil {
t.Fatal(err)
}
return o.Key()
}
// **H1, the commonest hand act** (031/01 §f: a push by hand to unstick a plan waiting on a report, four
// times): the machine is asked to report; a report that names what it was sent is all it takes, and one
// that does not — or none — is a send of its current declaration again. Each act kept in `tried` as
// `healer H1`, said as healer-acted, counted; twice, then the operator's, urgent; then nothing more.
func TestH1AsksAMachineToReportAndSendsItAgain(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
h, told, clock := healingOn(t, open)
record, err := open.inventory.NodeByName(ctx, "laptop")
if err != nil {
t.Fatal(err)
}
if err := open.inventory.RecordSent(ctx, record.ID, "d2", nil); err != nil {
t.Fatal(err)
}
key := sentNotReported(t, "laptop")
// First: the report had not reached the mesh, and asking for it brings it.
asked := 0
h.askReport = func(ctx context.Context, node string) error {
asked++
if node != "laptop" {
t.Errorf("asked %s", node)
}
_, err := open.inventory.RecordDoing(ctx, record.ID, inventory.Doing{Outcome: inventory.OutcomeApplied,
At: time.Now().Add(time.Second), Declared: "d2"})
return err
}
h.tick(ctx)
c, found, err := conditionsFrom.Get(ctx, key)
if err != nil || !found {
t.Fatalf("the condition is gone after the act — a healer cleared it, not an observation: %v", err)
}
if asked != 1 || len(c.Tried) != 1 || c.Tried[0].By != "healer H1" || !strings.Contains(c.Tried[0].Outcome, "acted") ||
c.Resolver != "healer:H1" {
t.Fatalf("after the first act: asked %d, %+v", asked, c)
}
if acts := told.said(); len(acts) != 1 || acts[0].Healer != "H1" || acts[0].Outcome != inventory.HealActed ||
acts[0].Condition != key || acts[0].By != "healer H1" || acts[0].Epoch != 57 {
t.Fatalf("the act was not said as healer-acted: %+v", acts)
}
// Within its settle, nothing more: the observation has its turn.
clock.pass(time.Minute)
h.tick(ctx)
if asked != 1 {
t.Fatalf("acted again inside the settle: asked %d", asked)
}
// Second: the machine says nothing, so it is sent again.
clock.pass(3 * time.Minute)
sent := 0
h.askReport = func(context.Context, string) error { asked++; return nil }
h.sendAgain = func(_ context.Context, node string) error { sent++; return nil }
h.tick(ctx)
if asked != 2 || sent != 1 {
t.Fatalf("the second act: asked %d, sent %d", asked, sent)
}
c, _, _ = conditionsFrom.Get(ctx, key)
if len(c.Tried) != 2 || !strings.Contains(c.Tried[1].Outcome, "sent again") || !strings.Contains(c.Tried[1].Outcome, "act 2 of 2") {
t.Fatalf("tried %+v", c.Tried)
}
// The budget is spent: the operator's, urgent, said; and no healer touches it again.
clock.pass(4 * time.Minute)
h.tick(ctx)
c, _, _ = conditionsFrom.Get(ctx, key)
if !c.Escalated() || c.Severity != conditions.Urgent || len(c.Tried) != 3 ||
!strings.Contains(c.Tried[2].Outcome, "budget of 2") {
t.Fatalf("not handed to the operator: %+v", c)
}
if acts := told.said(); len(acts) != 3 || acts[2].Outcome != inventory.HealEscalated {
t.Fatalf("the escalation was not said: %+v", acts)
}
clock.pass(time.Hour)
h.tick(ctx)
if asked != 2 || sent != 1 || len(told.said()) != 3 {
t.Fatalf("a healer acted on a condition the operator holds: asked %d sent %d said %d", asked, sent, len(told.said()))
}
heals, err := open.inventory.HealsSince(ctx, time.Now().Add(-time.Hour))
if err != nil || len(heals) != 3 || heals[0].Outcome != inventory.HealActed || heals[1].Outcome != inventory.HealActed ||
heals[2].Outcome != inventory.HealEscalated || heals[0].Epoch != 57 {
t.Fatalf("the heals kept: %+v %v", heals, err)
}
// And a heal is not a hand act: what S15 counts never sees it.
for _, x := range heals {
if x.Healer == "" || strings.HasPrefix(x.Act, "hand-act") {
t.Errorf("a heal reads as a hand act: %+v", x)
}
}
}
// **Only the controller holding the lease heals** (to-be 45 §6): unleased, or standing by, nothing.
func TestNoHealerActsWithoutTheLease(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
h, told, _ := healingOn(t, open)
sentNotReported(t, "laptop")
h.epoch = func(context.Context) (uint64, error) { return 0, nil }
h.tick(ctx)
h.epoch = func(context.Context) (uint64, error) { return 0, errors.New("the lease is not held") }
h.tick(ctx)
h.epoch = func(context.Context) (uint64, error) { return 57, nil }
h.acting = func() bool { return false }
h.tick(ctx)
if len(told.said()) != 0 {
t.Fatalf("a healer acted without the lease: %+v", told.said())
}
}
// **The mesh-wide brake**: a dozen acts in an hour and every healer stops, said urgently, until an hour
// after the last.
func TestTheBrakeStopsEveryHealerAndSaysSo(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
h, told, clock := healingOn(t, open)
for i := 0; i < healBrakeLimit; i++ {
if _, err := open.inventory.BeginHeal(ctx, inventory.Heal{Healer: "H3", ConditionKey: "seat.x.anchor.silent",
Kind: "holder-silent", Act: "asserted", Outcome: inventory.HealActed,
At: clock.now().Add(-time.Duration(healBrakeLimit-i) * time.Minute)}); err != nil {
t.Fatal(err)
}
}
sentNotReported(t, "laptop")
h.tick(ctx) // the fakes fail the test if anything acts
braked, found, err := conditionsFrom.Get(ctx, "mesh.healers.braked")
if err != nil || !found || braked.Severity != conditions.Urgent || !strings.Contains(braked.Evidence[0].Said, "H3 on seat.x.anchor.silent ×12") {
t.Fatalf("the brake was not said: %+v %v", braked, err)
}
if len(told.said()) != 0 {
t.Fatalf("a healer acted under the brake: %+v", told.said())
}
// Past the hour of the first act, still held: it lets go an hour after the last.
clock.pass(30 * time.Minute)
h.tick(ctx)
if _, held, _ := conditionsFrom.Get(ctx, "mesh.healers.braked"); !held {
t.Fatal("the brake let go before an hour had passed since the last act")
}
clock.pass(31 * time.Minute)
asked := 0
h.askReport = func(context.Context, string) error { asked++; return nil }
h.sendAgain = func(context.Context, string) error { return nil }
h.tick(ctx)
if _, held, _ := conditionsFrom.Get(ctx, "mesh.healers.braked"); held {
t.Fatal("the brake held an hour after the last act")
}
if asked != 1 {
t.Fatalf("the healers did not act again after the brake let go: asked %d", asked)
}
}
// **H2 closes a plan another has taken over**, with its note — and leaves a plan alone whose wait is
// still to come: that one is its own signal's, not a healer's.
func TestH2ClosesAPlanAnotherHasTakenOver(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
h, told, _ := healingOn(t, open)
created := time.Now().UTC().Add(-time.Hour)
older := inventory.Plan{ID: "plan-old", Repository: "novox/app", Branch: "main", Commit: "0ld0ld0", Created: created,
State: inventory.PlanBuilding, Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{"a": {State: "asked"}}}
waiting := inventory.Plan{ID: "plan-other", Repository: "novox/other", Branch: "main", Commit: "07he707", Created: created,
State: inventory.PlanBuilding, Tiers: [][]string{{"b"}}, Modules: map[string]*inventory.PlanModule{"b": {State: "asked"}}}
newer := inventory.Plan{ID: "plan-new", Repository: "novox/app", Branch: "main", Commit: "new0new", Created: created.Add(time.Minute),
State: inventory.PlanDone, Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{"a": {State: "built"}}}
for _, p := range []*inventory.Plan{&older, &waiting, &newer} {
if err := open.inventory.SavePlan(ctx, p); err != nil {
t.Fatal(err)
}
}
for _, id := range []string{"plan-old", "plan-other"} {
if _, err := conditionsFrom.Observe(ctx, conditions.Observation{Scope: conditions.ScopePlan, ID: id, Kind: "stalled",
Severity: conditions.Warning, Summary: id + " is stalled", Source: "S3"}); err != nil {
t.Fatal(err)
}
}
h.tick(ctx)
closed, err := open.inventory.PlanByID(ctx, "plan-old")
if err != nil || closed.State != inventory.PlanSuperseded || !strings.Contains(closed.Note, "closed by healer H2") ||
!strings.Contains(closed.Note, "plan-new") {
t.Fatalf("the superseded plan: %+v %v", closed, err)
}
if other, _ := open.inventory.PlanByID(ctx, "plan-other"); other.State != inventory.PlanBuilding {
t.Fatalf("a plan still waiting was closed: %+v", other)
}
if c, _, _ := conditionsFrom.Get(ctx, "plan.plan-other.stalled"); len(c.Tried) != 0 || c.Resolver != conditions.ResolverSelf {
t.Fatalf("a plan H2 does not repair was touched: %+v", c)
}
if acts := told.said(); len(acts) != 1 || acts[0].Healer != "H2" || acts[0].Condition != "plan.plan-old.stalled" {
t.Fatalf("said %+v", acts)
}
// Finished: every module built, none rolled out unsent — closed as done.
done := inventory.Plan{ID: "plan-done", Repository: "novox/third", Branch: "main", Commit: "d0ned0n", Created: created,
State: inventory.PlanRolling, Tier: 0, Tiers: [][]string{{"c"}}, Modules: map[string]*inventory.PlanModule{"c": {State: "built"}}}
if err := open.inventory.SavePlan(ctx, &done); err != nil {
t.Fatal(err)
}
if state, why, err := planStale(ctx, open.inventory, done); err != nil || state != inventory.PlanDone || !strings.Contains(why, "finished") {
t.Fatalf("a finished plan reads %q %q %v", state, why, err)
}
}
// **H4 only for a consumer the stream table marks resettable**: a module's consumer far behind is said
// and left — what a reset drops, nothing would catch up for it.
func TestH4ResetsOnlyWhatTheTableMarksResettable(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
h, _, _ := healingOn(t, open)
marked := 0
for _, c := range broker.MeshConsumers() {
if c.Resettable != "" {
marked++
if c.Stream != broker.EventsStream || c.Name != broker.ControllerName {
t.Errorf("%s on %s is marked resettable: only the controller's own events consumer is", c.Name, c.Stream)
}
}
}
if marked != 1 {
t.Fatalf("%d consumers are marked resettable, want the controller's events consumer alone", marked)
}
for _, who := range []string{"EVENTS.anchor_shop", "CONTROL.controller"} {
if _, err := conditionsFrom.Observe(ctx, conditions.Observation{Scope: conditions.ScopeBus, ID: who, Token: "behind",
Kind: kindConsumerBehind, Severity: conditions.Warning, Summary: who + " is far behind", Source: "D6"}); err != nil {
t.Fatal(err)
}
}
h.tick(ctx) // the fake reset fails the test if it is called
reset := ""
h.resetConsumer = func(stream, name string) (string, error) {
reset = stream + "." + name
return "it was 1500 behind", nil
}
if _, err := conditionsFrom.Observe(ctx, conditions.Observation{Scope: conditions.ScopeBus, ID: "EVENTS.controller",
Token: "behind", Kind: kindConsumerBehind, Severity: conditions.Warning, Summary: "far behind", Source: "D6"}); err != nil {
t.Fatal(err)
}
h.tick(ctx)
if reset != "EVENTS.controller" {
t.Fatalf("reset %q", reset)
}
}
// **H3 against a real bus** (issue 208's note): a consumer the mesh expects, deleted; D6 says so; H3
// asserts the bus's objects the way a send does, and D6's next run clears it — the healer never does.
// Deleted again within the hour, the budget is spent: the operator's, urgent, and H3 stops.
func TestNatsH3AssertsAMissingConsumerAgainAndBrakesAfterItsBudget(t *testing.T) {
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
open := aMesh(t)
ctx := t.Context()
js, err := broker.Dial(url)
if err != nil {
t.Fatal(err)
}
t.Cleanup(js.Close)
for _, s := range []string{"CONTROL", "NODES", "ASSIGNMENTS", "EVENTS"} {
_ = js.Context().DeleteStream(s)
}
if _, err := assertBusObjects(ctx, open.inventory, js); err != nil {
t.Fatal(err)
}
h, told, clock := healingOn(t, open)
h.js = js
h.assertObjects = func(ctx context.Context) error { return assertOnSend(ctx, open.inventory, js, " ") }
d := &doctor{open: open, js: js}
probe := func() {
t.Helper()
found, err := probeConsumers(ctx, d)
if err != nil {
t.Fatal(err)
}
if err := conditionsFrom.Reconcile(ctx, "D6", kindedAs(found, "consumer-wrong")); err != nil {
t.Fatal(err)
}
}
const key = "bus.NODES.laptop.missing"
if err := js.Context().DeleteConsumer("NODES", "laptop"); err != nil {
t.Fatal(err)
}
probe()
if c, raised, _ := conditionsFrom.Get(ctx, key); !raised || c.Kind != "consumer-lost" {
t.Fatalf("the deleted consumer was not raised: %+v", c)
}
h.tick(ctx)
if _, err := js.Context().ConsumerInfo("NODES", "laptop"); err != nil {
t.Fatalf("H3 did not make the consumer again: %v", err)
}
c, still, _ := conditionsFrom.Get(ctx, key)
if !still || len(c.Tried) != 1 || c.Tried[0].By != "healer H3" || c.Resolver != "healer:H3" {
t.Fatalf("the act is not in the condition, or the healer cleared it: %+v", c)
}
probe()
if _, still, _ := conditionsFrom.Get(ctx, key); still {
t.Fatal("the probe's next run did not clear what H3 repaired")
}
// Kept in the history as the keeper says it, a moment later.
var cleared *conditions.Event
for wait := time.Now().Add(5 * time.Second); cleared == nil && time.Now().Before(wait); time.Sleep(20 * time.Millisecond) {
history, err := conditionsFrom.HistorySince(ctx, time.Now().Add(-time.Minute))
if err != nil {
t.Fatal(err)
}
for i := range history {
if history[i].Key == key && history[i].Change == conditions.ChangeCleared {
cleared = &history[i]
}
}
}
if cleared == nil || !strings.Contains(cleared.Why, "D6 no longer observes it") || len(cleared.Tried) != 1 {
t.Fatalf("the clearing is not the probe's, or forgets what was tried: %+v", cleared)
}
// Again within the hour: the budget is one, so after its settle the operator is told, and H3 stops.
clock.pass(10 * time.Minute)
if err := js.Context().DeleteConsumer("NODES", "laptop"); err != nil {
t.Fatal(err)
}
probe()
h.assertObjects = func(context.Context) error { t.Error("H3 acted past its budget"); return nil }
h.tick(ctx)
c, _, _ = conditionsFrom.Get(ctx, key)
if !c.Escalated() || c.Severity != conditions.Urgent || c.Count != 2 {
t.Fatalf("the spent budget was not handed to the operator: %+v", c)
}
acts := told.said()
if len(acts) != 2 || acts[0].Outcome != inventory.HealActed || acts[1].Outcome != inventory.HealEscalated {
t.Fatalf("said %+v", acts)
}
clock.pass(2 * time.Hour)
h.tick(ctx)
if len(told.said()) != 2 {
t.Fatal("a healer acted on what the operator holds")
}
}
// **H4 against a real bus** (issue 248): the controller's events consumer a long way behind — the week it
// once replayed — is re-made from now by the mesh itself, and the controller's bound subscription still
// receives what comes next: a reset that left the controller deaf would be the incident.
func TestNatsH4ResetsTheControllersEventsConsumerAndItStillDelivers(t *testing.T) {
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
open := aMesh(t)
ctx := t.Context()
js, err := broker.Dial(url)
if err != nil {
t.Fatal(err)
}
t.Cleanup(js.Close)
for _, s := range []string{"CONTROL", "NODES", "ASSIGNMENTS", "EVENTS"} {
_ = js.Context().DeleteStream(s)
}
if _, err := assertBusObjects(ctx, open.inventory, js); err != nil {
t.Fatal(err)
}
// Bound as the controller binds it (link/receive_nats.go), taking one and acknowledging none.
events := make(chan *nats.Msg, 64)
sub, err := js.Context().ChanSubscribe("", events, nats.Bind(broker.EventsStream, broker.ControllerName))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = sub.Unsubscribe() })
followed := broker.ControllerFollows[0]
for i := 0; i < consumerFarBehind+200; i++ {
if _, err := js.Context().Publish(followed, []byte(`{"n":`+strconv.Itoa(i)+`}`)); err != nil {
t.Fatal(err)
}
}
d := &doctor{open: open, js: js}
probe := func() []conditions.Observation {
t.Helper()
found, err := probeConsumers(ctx, d)
if err != nil {
t.Fatal(err)
}
if err := conditionsFrom.Reconcile(ctx, "D6", kindedAs(found, "consumer-wrong")); err != nil {
t.Fatal(err)
}
return found
}
probe()
const key = "bus.EVENTS.controller.behind"
if c, raised, _ := conditionsFrom.Get(ctx, key); !raised || c.Kind != kindConsumerBehind {
t.Fatalf("a consumer %d behind was not raised: %+v", consumerFarBehind+200, c)
}
h, told, _ := healingOn(t, open)
h.resetConsumer = func(stream, name string) (string, error) {
before, after, err := js.ResetConsumer(stream, name)
if err != nil {
return "", err
}
return "it was " + strconv.FormatUint(before.Pending, 10) + " behind; " + strconv.FormatUint(after.Pending, 10) +
" pending now", nil
}
h.tick(ctx)
if acts := told.said(); len(acts) != 1 || acts[0].Healer != "H4" || acts[0].Outcome != inventory.HealActed {
t.Fatalf("said %+v", acts)
}
if found := probe(); len(found) != 0 {
t.Fatalf("after the reset the probe still finds %+v", found)
}
if _, still, _ := conditionsFrom.Get(ctx, key); still {
t.Fatal("the probe's next run did not clear what H4 repaired")
}
// What comes next still reaches the controller.
for len(events) > 0 {
<-events
}
if _, err := js.Context().Publish(followed, []byte(`{"after":"the reset"}`)); err != nil {
t.Fatal(err)
}
deadline := time.After(10 * time.Second)
for {
select {
case m := <-events:
if strings.Contains(string(m.Data), "after") {
return
}
_ = m.Ack()
case <-deadline:
t.Fatal("the controller's subscription heard nothing after its consumer was reset: it would be deaf")
}
}
}
// **H1's question reaches the machine**, on the subject its grant lets it hear and nothing else.
func TestNatsAskToReportReachesTheMachine(t *testing.T) {
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
conn, err := nats.Connect(url)
if err != nil {
t.Fatal(err)
}
t.Cleanup(conn.Close)
asked, err := conn.SubscribeSync(broker.AskReportSubject("laptop"))
if err != nil {
t.Fatal(err)
}
if err := askToReport(t.Context(), conn, "laptop"); err != nil {
t.Fatal(err)
}
msg, err := asked.NextMsg(5 * time.Second)
if err != nil || !strings.Contains(string(msg.Data), "healer H1") {
t.Fatalf("the machine heard %v, %v", msg, err)
}
perms, err := broker.PermissionsFor(broker.Principal{Kind: broker.KindNode, Node: "laptop", PasswordHash: "x"})
if err != nil || !slices.Contains(perms.Subscribe, "mesh.node.laptop.ask.report") ||
slices.Contains(perms.Subscribe, "mesh.node.anchor.ask.report") {
t.Fatalf("a machine's grant for the question: %v %v", perms.Subscribe, err)
}
}
+4
View File
@@ -161,6 +161,9 @@ func run() error {
return conditionsCommand(ctx, args[1:])
case "doctor":
return doctorCommand(ctx, args[1:])
// What the healers did, and their brake (novox/hq to-be 45 §7).
case "healers":
return healersCommand(ctx, args[1:])
case "version":
fmt.Println(version)
return nil
@@ -253,6 +256,7 @@ func usage() {
conditions history [--days N] [--key K] every raising, change and clearing lately
doctor [run|probes|signals] [--json]
the self-check: the last verdict, a run now, the probes, the signals' ages
healers [--days N] [--json] the healers, what they did lately, and their brake (to-be 45 §7)
durations [--kind K] [--days N] [--json]
apply, heartbeat, plan-tier and build durations, per machine or module
collection [--json] kept archives held/unheld by a manifest, and what the sweep may let go
+5 -3
View File
@@ -429,10 +429,12 @@ func probeConsumers(ctx context.Context, d *doctor) ([]conditions.Observation, e
Said: differs})
}
if c.Stream != "NODES" && info.NumPending > consumerFarBehind {
// Its own kind (novox/hq to-be 45 §7): healer H4 answers a consumer far behind, and only
// this — a redefined consumer is not one a reset repairs.
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: who, Token: "behind",
Severity: conditions.Warning,
Summary: fmt.Sprintf("%s is %d message(s) behind its stream's head", consumerWords(c), info.NumPending),
Said: fmt.Sprintf("%d pending, %d handed out and not settled", info.NumPending, info.NumAckPending)})
Kind: kindConsumerBehind, Severity: conditions.Warning,
Summary: fmt.Sprintf("%s is %d message(s) behind its stream's head", consumerWords(c), info.NumPending),
Said: fmt.Sprintf("%d pending, %d handed out and not settled", info.NumPending, info.NumAckPending)})
}
}
return sortedFound(out), nil
+5
View File
@@ -84,6 +84,10 @@ type meshStatus struct {
// HandActsUnread says why when it could not be read, rather than reading as none.
HandActsThisWeek *int `json:"handActsThisWeek,omitempty"`
HandActsUnread string `json:"handActsUnread,omitempty"`
// HealsThisWeek is what the healers did in the last seven days (novox/hq to-be 45 §7); HealsUnread
// why it could not be read.
HealsThisWeek *healsCount `json:"healsThisWeek,omitempty"`
HealsUnread string `json:"healsUnread,omitempty"`
// Conditions is every open condition, urgent first and then oldest first (novox/hq to-be 45 §2):
// what is wrong, as the watchdogs, the self-check and the providers say it. Always present — an
// empty list is "none open" — unless they could not be read, which ConditionsUnread says.
@@ -230,6 +234,7 @@ func statusAsJSON(asked answers) ([]byte, error) {
}
out.Unheld = asked.unheld
out.HandActsThisWeek, out.HandActsUnread = asked.handActs, asked.handActsUnread
out.HealsThisWeek, out.HealsUnread = asked.heals, asked.healsUnread
out.Conditions, out.ConditionsUnread = asked.conditions, asked.conditionsUnread
if out.Conditions == nil {
out.Conditions = []conditions.Condition{}
+6
View File
@@ -407,6 +407,12 @@ func (a *verbArguments) commandLine() ([]string, error) {
argv = append(argv, "--days", d)
}
return argv, nil
case "healers":
argv := []string{"healers", "--json"}
if d := str("days"); d != "" {
argv = append(argv, "--days", d)
}
return argv, nil
case "durations":
argv := []string{"durations", "--json"}
if k := str("kind"); k != "" {
@@ -276,6 +276,7 @@ var accountedFlags = map[string]map[string]string{
"conditions": {"json": "set by the verb: the answer is data"},
"conditions history": {"json": "set by the verb: the answer is data"},
"conditions show": {"json": "set by the verb: the answer is data"},
"healers": {"json": "set by the verb: the answer is data"},
}
// **Every flag of the command a verb runs is in the verb's schema, or accounted for here.** Derived
+54 -3
View File
@@ -2,6 +2,7 @@ package main
import (
"fmt"
"sort"
"strings"
"time"
@@ -186,9 +187,12 @@ var signalsTable = []signalRow{
Bound: "2 days", Kind: "facts-stale", Severity: conditions.Warning, Phase: 5,
Deferred: "the facts snapshot is built in Phase 5 (to-be 45 §9): nothing exports one yet"},
{Row: "S15", Signal: "a hand act with a cause already recorded", Emitter: "hand-act log",
Trigger: "each act", Bound: "the second within 14 days", Kind: "healer-wanted",
Severity: conditions.Warning, Phase: 3,
Deferred: "Phase 3 (to-be 45 §10): `hand-acts` lists repeated causes today; the condition comes with the healers"},
Trigger: "each act", Bound: "the second within 14 days; clears when fewer than two remain within 14 days",
Kind: "healer-wanted", Severity: conditions.Warning, Phase: 3,
needs: func(f *signalFacts) error { return f.handActsErr }, watch: watchHandActs,
newest: func(f *signalFacts) time.Time {
return newestOf(f.handActs, func(a link.HandAct) time.Time { return a.At })
}},
}
// newestOf is the newest time among things.
@@ -564,6 +568,53 @@ func watchLease(f *signalFacts) []conditions.Observation {
return out
}
// handActsWithin is how far back a repeated cause counts (S15).
const handActsWithin = 14 * 24 * time.Hour
// watchHandActs is S15: a cause recorded by hand twice within a fortnight is a healer wanted, named by
// the cause. **A heal is never a hand act** (healers.go), so a cause a healer exists for and a person
// still repaired twice says the healer is not enough — its reach or its budget — and is said so.
func watchHandActs(f *signalFacts) []conditions.Observation {
recent := make([]link.HandAct, 0, len(f.handActs))
for _, a := range f.handActs {
if f.now.Sub(a.At) <= handActsWithin {
recent = append(recent, a)
}
}
repeated := link.RepeatedCauses(recent, f.now)
causes := make([]string, 0, len(repeated))
for c := range repeated {
causes = append(causes, c)
}
sort.Strings(causes)
var out []conditions.Observation
for _, cause := range causes {
var acts []string
var newest link.HandAct
for _, a := range recent {
if a.Cause != cause {
continue
}
acts = append(acts, fmt.Sprintf("%s %s by %s: %s", a.At.UTC().Format("2006-01-02 15:04"),
strings.TrimSpace(a.Verb+" "+strings.Join(a.Args, " ")), a.By, a.Why))
if a.At.After(newest.At) {
newest = a
}
}
wanted := "a healer is wanted for it"
if h := healerNamedFor(cause); h != "" {
wanted = fmt.Sprintf("healer %s answers this cause and a person still repaired it: its reach or its "+
"budget is not enough", h)
}
out = append(out, conditions.Observation{Scope: conditions.ScopeMesh, ID: "hand-acts." + cause,
Token: "healer-wanted", Kind: "healer-wanted", Severity: conditions.Warning,
Summary: fmt.Sprintf("%q was repaired by hand %d times in %d days, the last by %s: %s", cause,
repeated[cause], int(handActsWithin.Hours()/24), newest.By, wanted),
Said: strings.Join(acts, "; ")})
}
return out
}
// watchedRows are the rows a watchdog runs for.
func watchedRows() []signalRow {
var out []signalRow
+37
View File
@@ -122,6 +122,43 @@ var suppressions = map[string]suppression{
inside: func(f *signalFacts) { f.staleRefusals = refusedBy(f.now, 41, 5) },
past: func(f *signalFacts) { f.staleRefusals = refusedBy(f.now, 41, 6) },
},
// Twice by hand within a fortnight is a healer wanted; once, or the first of two a day too old, is not.
"S15": {
inside: func(f *signalFacts) {
f.handActs = []link.HandAct{actByHand(f.now.Add(-15*24*time.Hour), "consumer-behind"),
actByHand(f.now.Add(-time.Hour), "consumer-behind")}
},
past: func(f *signalFacts) {
f.handActs = []link.HandAct{actByHand(f.now.Add(-13*24*time.Hour), "consumer-behind"),
actByHand(f.now.Add(-time.Hour), "consumer-behind")}
},
},
}
// actByHand is one entry in the hand-act log, with its cause.
func actByHand(at time.Time, cause string) link.HandAct {
return link.HandAct{ID: "act-" + at.Format("150405"), At: at, By: "jochen at a shell on the laptop",
Verb: "broker consumer-reset", Args: []string{"EVENTS", "controller"}, Why: "it replayed a week", Cause: cause}
}
// **S15 names the cause and, where a healer answers it, that the healer was not enough.**
func TestARepeatedHandActNamesItsCauseAndItsHealer(t *testing.T) {
now := time.Date(2026, 10, 6, 12, 0, 0, 0, time.UTC)
f := calm(now)
f.handActs = []link.HandAct{actByHand(now.Add(-2*time.Hour), "consumer-behind"),
actByHand(now.Add(-time.Hour), "consumer-behind"), actByHand(now.Add(-time.Hour), "restarted the proxy"),
actByHand(now.Add(-time.Minute), "restarted the proxy")}
got := watchHandActs(f)
if len(got) != 2 {
t.Fatalf("%+v", got)
}
if got[0].Key() != "mesh.hand-acts.consumer-behind.healer-wanted" || !strings.Contains(got[0].Summary, "healer H4") {
t.Errorf("a cause a healer answers: %s — %s", got[0].Key(), got[0].Summary)
}
if got[1].Key() != "mesh.hand-acts.restarted_the_proxy.healer-wanted" ||
!strings.Contains(got[1].Summary, "a healer is wanted") {
t.Errorf("a cause no healer answers: %s — %s", got[1].Key(), got[1].Summary)
}
}
// standingSaid is a provider's failing word last said at a moment.
+16
View File
@@ -326,6 +326,15 @@ func printStatus(asked answers) error {
fmt.Printf("%d act(s) done by hand in the last seven days — `hand-acts` lists them, and why\n\n", *asked.handActs)
}
// And what the mesh repaired by itself (novox/hq to-be 45 §7): seen, not only done.
switch {
case asked.healsUnread != "":
fmt.Printf("what the healers did could not be read: %s\n\n", asked.healsUnread)
case asked.heals != nil && (asked.heals.Acts > 0 || asked.heals.Escalated > 0):
fmt.Printf("%d repair(s) made by the healers in the last seven days, %d handed to the operator — `healers` "+
"lists them, and each is in its condition's tried\n\n", asked.heals.Acts, asked.heals.Escalated)
}
if adopted := adoptedNodes(nodes); len(adopted) > 0 {
// Said, because nothing forces the flip: a node left adopted is visible here rather than
// read as converged (novox/hq ADR 0100). Not a fault, so it does not break "all well".
@@ -474,6 +483,13 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) {
}
}
// And what the healers did this week (novox/hq to-be 45 §7).
if heals, err := inv.HealsSince(ctx, time.Now().Add(-7*24*time.Hour)); err != nil {
out.healsUnread = err.Error()
} else {
out.heals = countHeals(heals)
}
// And which machines are not running what the mesh would send them. The same question as a
// module being behind its source, one level down: that one says the catalogue is out of date,
// this one says a machine is — and only the second has anybody's change waiting in it.
+1 -1
View File
@@ -171,7 +171,7 @@ func composeStatus(open *stores) func(context.Context) ([]byte, error) {
var readingVerbs = map[string]bool{
"tools": true, "calls": true, "status": true, "nodes": true, "node": true, "modules": true,
"seats": true, "builds": true, "plan": true, "queue": true, "durations": true, "hand-acts": true,
"doctor": true, "conditions": true,
"doctor": true, "conditions": true, "healers": true,
}
// nudgingListener is the enrolment, nudging the summary when a machine said something new.
+24
View File
@@ -78,6 +78,10 @@ type signalFacts struct {
staleRefusals []link.WriterRefusals
epochs map[int64]inventory.Epoch
// handActs are the acts done by hand within the fortnight S15 counts.
handActs []link.HandAct
handActsErr error
// lease is this controller's standing to the lease, and the epochs that ended lately (S12).
lease leaseFacts
leaseErr error
@@ -296,6 +300,7 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts {
}
f.advisories = link.Advisories.Since(now.Add(-advisoryQuiet))
f.lostConsumers, f.advisoriesErr = w.lostConsumers(ctx, f.advisories)
f.handActs, f.handActsErr = w.gatherHandActs(ctx, now)
return f
}
@@ -533,6 +538,20 @@ func (w *watchdogs) lostConsumers(ctx context.Context, heard []link.Advisory) (m
return out, nil
}
// gatherHandActs is the hand-act log's fortnight (S15), read from the bus.
func (w *watchdogs) gatherHandActs(ctx context.Context, now time.Time) ([]link.HandAct, error) {
if w.js == nil {
return nil, errors.New("this controller is not on the bus")
}
reading, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
acts, err := link.HandActs(reading, w.js.Conn(), now.Add(-handActsWithin))
if err != nil {
return nil, fmt.Errorf("the hand-act log cannot be read: %w", err)
}
return acts, nil
}
// later is the later of two moments.
func later(a, b time.Time) time.Time {
if a.After(b) {
@@ -567,6 +586,11 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li
watching, stop := context.WithCancel(ctx)
go w.keep(watching)
go d.keep(watching)
// And the healers (novox/hq to-be 45 §7): what is open that a registered healer answers, repaired
// under the lease and the brake, every act said.
healers := newHealing(open, keeper, bus, server.JetStream())
go healers.keep(watching)
go forgettingOldHeals(watching, open.inventory)
fmt.Printf("watching the mesh: %d signal(s) every %s, %d probe(s) every %s; what is wrong is kept in %s "+
"and said as %s events\n", len(watchedRows()), watchEvery, len(runnableProbes()), doctorEvery,
broker.ConditionsBucket, conditions.Seat)
+5
View File
@@ -54,6 +54,11 @@ type Consumer struct {
// that exists keeps where it is, whatever this says; only its making is decided here.
FromNow bool
Why string
// Resettable says why this consumer may be re-made to deliver from now by the mesh itself, with
// nobody asked (novox/hq to-be 45 §7, healer H4): what a reset drops, something else catches up.
// Empty for every consumer where nothing would — a module's, whose events would be lost to it.
// Not a property of the consumer on the bus: nothing here is sent to the server.
Resettable string
}
// seatStreamName is the stream holding a seat's inbound work. Named after the seat rather than
+9 -1
View File
@@ -179,6 +179,10 @@ func (p Principal) Username() string {
return ""
}
// AskReportSubject is where the mesh asks one machine's node-engine to say again what it last applied
// (novox/hq to-be 45 §6): its answer is an ordinary report, on its own report subject.
func AskReportSubject(node string) string { return "mesh.node." + node + ".ask.report" }
// inbox is a principal's own reply space. No user is ever granted a bare `_INBOX.>` (design 25
// §4): with one account, inbox privacy is the permission list or it is nothing, so each user's
// inbox is derived from its own identity and its permissions name that prefix and no other.
@@ -378,7 +382,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
// next assertion moves it, and a permission that only allowed the new shape would refuse
// every node in the mesh for exactly as long as that took.
sub = []string{"mesh.node." + p.Node + ".declare",
"_DELIVER." + p.Node, "_DELIVER." + p.Node + ".>"}
"_DELIVER." + p.Node, "_DELIVER." + p.Node + ".>",
// And the mesh asking it to say again what it last applied (novox/hq to-be 45 §6, the
// `report` verb healer H1 asks): its own machine's, on core NATS and off any stream. It
// answers through its report, the one thing it already says — no reply to anybody's inbox.
AskReportSubject(p.Node)}
case KindModule:
// 1. Its own namespace: it publishes its events there and serves its tools there. Nothing
+2
View File
@@ -26,6 +26,8 @@ func TestTheFactsTheGrantPermitsAreTheFactsTheMeshStates(t *testing.T) {
states = append(states, conditions.HeartbeatEvent)
// And a value given by hand, replaced after its module's first good start (novox/hq ADR 0228).
states = append(states, link.KeySecretReplaced)
// And every act a healer takes (novox/hq to-be 45 §7).
states = append(states, link.KeyHealerActed)
for _, event := range states {
if !slices.Contains(broker.ControllerStates, event) {
t.Errorf("the mesh states %q and its account may not publish it", event)
+11 -1
View File
@@ -210,7 +210,10 @@ var ControllerStates = []string{"applied", "refused", "built-before",
// machine that is not the control node, so the controller going quiet is itself said.
"doctor-heartbeat",
// And a value given by hand, replaced after its module's first good start (novox/hq ADR 0228).
"secret-replaced"}
"secret-replaced",
// And every act a healer takes on a condition (novox/hq to-be 45 §7, Phase 3): a repair the mesh
// made by itself is said like one a person made, never quietly.
"healer-acted"}
// BusAdvisories are what the bus server says about the mesh's own account that the controller
// reads (novox/hq to-be 45 §3, S9): a durable consumer that handed a message over as often as it
@@ -304,6 +307,13 @@ func MeshConsumers() []Consumer {
// client and come back to be acted on again.
MaxAckPending: 1,
FromNow: true,
// **The one consumer the mesh resets by itself** (healer H4, issue 248): it fell a week
// behind once and held every merge after it. What a reset drops is caught up elsewhere —
// a merge by the catch-up pass that reads the forge (issue 266), a build's outcome from the
// build records a plan settles from (issue 214), a provider's failing word by the provider
// saying it again every quarter of an hour (ADR 0224).
Resettable: "what it drops is caught up: merges by the catch-up pass (issue 266), build outcomes " +
"from the build records (issue 214), a provider's failing word said again (ADR 0224)",
Why: "the events the mesh's own controller reacts to, one at a time; after " +
"max-deliver it dead-letters, because an announcement it cannot act on will not " +
"become actionable",
+2 -2
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled
users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] }
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] }
subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] }
allow_responses: { max: 1, ttl: "1m" }
} }
@@ -34,7 +34,7 @@ accounts {
} }
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] }
subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.ask.report", "mesh.node.one.declare"] }
} }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
+4
View File
@@ -98,6 +98,10 @@ var WritersTable = []WriterRow{
Others: "read by id", Subjects: kvOf(CallsBucket), Writes: isController},
{State: "the hand-act log", Writer: "controller, through the verbs that act", KeptIn: "key-value " + HandActsBucket,
Others: "—", Subjects: kvOf(HandActsBucket), Writes: isController},
// What the healers did (novox/hq to-be 45 §7, Phase 3): the controller's alone, in its store — the
// budgets and the mesh-wide brake are counted from it, so a controller restarting cannot reset them.
{State: "the healers' acts and their brake", Writer: "controller (lease holder), each act begun before it is made",
KeptIn: "the controller's store", Others: "read through healers; each act said as healer-acted and in its condition's tried"},
{State: "stream definitions and bus permissions", Writer: "controller", KeptIn: "the bus", Others: "—",
// A stream's definition, and a durable consumer's by the API that names it so. Not every
// consumer create: a module watching its own bucket makes and deletes an ordered consumer on the
+1
View File
@@ -21,6 +21,7 @@ var designRows = []string{
"conditions",
"calls and their outcomes",
"the hand-act log",
"the healers' acts and their brake",
"stream definitions and bus permissions",
"builds and their outcomes",
"a merge announced",
+3 -1
View File
@@ -87,7 +87,9 @@ var defaultSeats = append([]Seat{
Emits: []string{"applied", "refused", "built-before",
"condition-raised", "condition-changed", "condition-cleared", "doctor-heartbeat",
// A value given by hand, replaced after its module's first good start (novox/hq ADR 0228).
"secret-replaced"},
"secret-replaced",
// Every act a healer takes (novox/hq to-be 45 §7).
"healer-acted"},
Serves: ControllerVerbs},
// The store's first verbs (novox/hq ADR 0159): the smallest set that makes the store askable,
// served by whichever module holds the seat with tools of these names.
+6
View File
@@ -252,6 +252,12 @@ var ControllerVerbs = []Verb{
"why": "with silence: why — required, and recorded in the hand-act log",
"cause": "with silence: the cause in a word (the condition's kind when absent)",
}, nil, "history")},
// What the healers did (novox/hq to-be 45 §7).
{Name: "healers", Description: "The healers (novox/hq to-be 45 §7): each registered response to one kind " +
"of condition — its repair (the ordinary path again), its budget, what happens when it is spent — every act " +
"they took lately with its outcome, and the mesh-wide brake. A heal is never a hand act; whether it " +
"repaired anything is the condition's clearing to say. Each act is also in its condition's tried.",
Input: schema(map[string]string{"days": "how many days of acts back (default 7)"}, nil)},
{Name: "doctor", Description: "The self-check (novox/hq to-be 45 §4): the last run's verdict at once — " +
"each probe of the design's live invariants passed, failed or could not run, and how long ago. With " +
"run, a run now; with probes, the registry; with signals, every row of the signals table and the age " +
+14 -1
View File
@@ -79,13 +79,21 @@ type Evidence struct {
Said string `json:"said"`
}
// Attempt is one healer's try at a condition (to-be 45 §7; written from Phase 3).
// Attempt is one healer's try at a condition (to-be 45 §7): when, what it did, what came of it, and
// which healer — so a person reading `conditions show` sees the mesh repairing itself, by whom.
type Attempt struct {
At time.Time `json:"at"`
What string `json:"what"`
Outcome string `json:"outcome"`
// By is the healer, as a person reads it: `healer H1`. Never a person: an act by hand is the
// hand-act log's, not a condition's attempt.
By string `json:"by"`
}
// KeptAttempts is how many attempts a condition keeps, newest last: a healer's budget is a few, and
// a condition reopened again and again carries its tries forward only so far.
const KeptAttempts = 10
// Silence is a person saying they know: no messages until it ends (to-be 45 §2). Recorded as a hand
// act; the condition stays open, and `status` still says it.
type Silence struct {
@@ -139,6 +147,11 @@ func (c Condition) SilencedAt(now time.Time) bool {
return c.Silenced != nil && now.Before(c.Silenced.Until)
}
// Escalated says a healer tried and its budget is spent: the operator resolves it now (to-be 45 §2,
// "budget spent ──► OPEN, resolver: operator, severity: urgent"). An observation does not lower its
// severity again: the watchdog that sees it every half minute would otherwise undo the escalation.
func (c Condition) Escalated() bool { return c.Resolver == ResolverOperator && len(c.Tried) > 0 }
// Show is the verb that shows more about a condition, as a message carries it.
func (c Condition) Show() string { return "mesh-controller.conditions key=" + c.Key }
+97 -4
View File
@@ -76,6 +76,9 @@ type clearing struct {
at time.Time
count int
silenced *Silence
// tried is what healers tried before it cleared: a reopening is the same fault, and what was
// tried on it is still what was tried.
tried []Attempt
}
// Options are what a Keeper is made with.
@@ -113,7 +116,8 @@ func NewKeeper(ctx context.Context, o Options) *Keeper {
if recent, err := k.history.Since(ctx, k.now().Add(-ReopenWithin)); err == nil {
for _, e := range recent {
if e.Change == ChangeCleared {
k.cleared[e.Key] = clearing{at: e.At, count: e.Condition.Count, silenced: e.Condition.Silenced}
k.cleared[e.Key] = clearing{at: e.At, count: e.Condition.Count, silenced: e.Condition.Silenced,
tried: e.Condition.Tried}
}
}
} else {
@@ -191,6 +195,7 @@ func (k *Keeper) Observe(ctx context.Context, o Observation) (Condition, error)
if before.silenced != nil && now.Before(before.silenced.Until) {
c.Silenced = before.silenced
}
c.Tried = before.tried
}
k.mu.Unlock()
if err := k.stamp(&c); err != nil {
@@ -216,7 +221,8 @@ func (k *Keeper) Observe(ctx context.Context, o Observation) (Condition, error)
return Condition{}, fmt.Errorf("the condition %s on the bus cannot be read: %w", key, err)
}
var changes []Event
if o.Severity != c.Severity {
// A healer's budget spent made it urgent; the watchdog seeing it again does not undo that.
if o.Severity != c.Severity && !(c.Escalated() && o.Severity != Urgent) {
changes = append(changes, Event{Change: ChangeSeverity, Was: string(c.Severity)})
c.Severity = o.Severity
}
@@ -224,7 +230,9 @@ func (k *Keeper) Observe(ctx context.Context, o Observation) (Condition, error)
changes = append(changes, Event{Change: ChangeResolver, Was: c.Resolver})
c.Resolver = r
}
c.Summary, c.Source, c.LastObserved = o.Summary, o.Source, now
// The kind as the source says it now: a source that gave the same key a kind of its own since
// (a probe's finding split out for a healer) is read by that kind from its next observation.
c.Kind, c.Summary, c.Source, c.LastObserved = o.Kind, o.Summary, o.Source, now
if o.Machine != "" {
c.Subject.Machine = o.Machine
}
@@ -282,7 +290,7 @@ func (k *Keeper) Clear(ctx context.Context, key, why string) (bool, error) {
}
now := k.now().UTC()
k.mu.Lock()
k.cleared[key] = clearing{at: now, count: c.Count, silenced: c.Silenced}
k.cleared[key] = clearing{at: now, count: c.Count, silenced: c.Silenced, tried: c.Tried}
k.mu.Unlock()
k.tell(Event{Condition: c, At: now, Change: ChangeCleared, Why: why, Cleared: &now})
return true, nil
@@ -290,6 +298,91 @@ func (k *Keeper) Clear(ctx context.Context, key, why string) (bool, error) {
return false, fmt.Errorf("the condition %s kept moving under this clearing; %d tries", key, tries)
}
// Tried records a healer's attempt on an open condition (to-be 45 §7) and makes the healer its
// resolver: said as `condition-changed` when the resolver changes, kept in `tried` either way. **It
// never clears the condition**: a repair that worked is seen by the observation that raised it, which
// clears it — a healer marking its own work done would be a second opinion of the fact. False when no
// condition is open under the key: it cleared meanwhile, and there is nothing to record against.
func (k *Keeper) Tried(ctx context.Context, key string, a Attempt, resolver string) (Condition, bool, error) {
return k.amend(ctx, key, func(c *Condition, now time.Time) []Event {
c.Tried = appendAttempt(c.Tried, a, now)
var changes []Event
if resolver != "" && resolver != c.Resolver && !c.Escalated() {
changes = append(changes, Event{Change: ChangeResolver, Was: c.Resolver, Why: a.Outcome})
c.Resolver = resolver
}
return changes
})
}
// Escalate is a healer's budget spent, or a repair it may not make (to-be 45 §2, §7): the attempt that
// says so is kept in `tried`, the resolver becomes the operator and the severity urgent, each said as
// `condition-changed`. Observation still clears it when the fault goes; nothing else does.
func (k *Keeper) Escalate(ctx context.Context, key string, a Attempt) (Condition, bool, error) {
return k.amend(ctx, key, func(c *Condition, now time.Time) []Event {
c.Tried = appendAttempt(c.Tried, a, now)
var changes []Event
if c.Severity != Urgent {
changes = append(changes, Event{Change: ChangeSeverity, Was: string(c.Severity), Why: a.Outcome})
c.Severity = Urgent
}
if c.Resolver != ResolverOperator {
changes = append(changes, Event{Change: ChangeResolver, Was: c.Resolver, Why: a.Outcome})
c.Resolver = ResolverOperator
}
return changes
})
}
// appendAttempt adds an attempt, stamped when it has no time, keeping the newest KeptAttempts.
func appendAttempt(tried []Attempt, a Attempt, now time.Time) []Attempt {
if a.At.IsZero() {
a.At = now
}
tried = append(tried, a)
if len(tried) > KeptAttempts {
tried = tried[len(tried)-KeptAttempts:]
}
return tried
}
// amend changes one open condition by compare-and-set and says what changed. False when none is open.
func (k *Keeper) amend(ctx context.Context, key string, change func(*Condition, time.Time) []Event) (Condition, bool, error) {
for i := 0; i < tries; i++ {
entry, found, err := k.store.Get(ctx, key)
if err != nil {
return Condition{}, false, fmt.Errorf("reading the condition %s: %w", key, err)
}
if !found {
return Condition{}, false, nil
}
var c Condition
if err := json.Unmarshal(entry.Value, &c); err != nil {
return Condition{}, false, fmt.Errorf("the condition %s on the bus cannot be read: %w", key, err)
}
now := k.now().UTC()
changes := change(&c, now)
if err := k.stamp(&c); err != nil {
return Condition{}, false, err
}
body, err := json.Marshal(c)
if err != nil {
return Condition{}, false, err
}
if err := k.store.Update(ctx, key, body, entry.Revision); errors.Is(err, ErrMoved) {
continue
} else if err != nil {
return Condition{}, false, fmt.Errorf("writing the condition %s: %w", key, err)
}
for _, e := range changes {
e.At, e.Condition = now, c
k.tell(e)
}
return c, true, nil
}
return Condition{}, false, fmt.Errorf("the condition %s kept moving under this write; %d tries", key, tries)
}
// Reconcile is one source's whole observation: every condition it observes is observed, and every
// condition it raised before and no longer observes is cleared — the observation says it is
// resolved. A source that could not observe must not call this: an empty observation clears all it
+105
View File
@@ -0,0 +1,105 @@
package conditions
import (
"testing"
"time"
)
// **A healer's attempt is said and kept, and never clears the condition** (to-be 45 §7): the resolver
// becomes the healer, `tried` says what it did, and only an observation clears it.
func TestAHealersAttemptIsKeptAndClearsNothing(t *testing.T) {
k, _, told, c := keeper(t)
ctx := t.Context()
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
held, open, err := k.Tried(ctx, "machine.ace.silent", Attempt{What: "asked ace to report", Outcome: "acted",
By: "healer H1"}, ResolverHealer("H1"))
if err != nil || !open {
t.Fatalf("tried: %v %v", open, err)
}
if held.Resolver != "healer:H1" || len(held.Tried) != 1 || held.Tried[0].By != "healer H1" || held.Tried[0].At.IsZero() {
t.Fatalf("after one attempt: %+v", held)
}
said := settled(t, told, 2)
if said[1].Event != EventChanged || said[1].Change != ChangeResolver || said[1].Was != ResolverSelf {
t.Errorf("the healer taking it is not said as a change of resolver: %+v", said[1])
}
// A second attempt by the same healer is kept and says nothing new.
if _, _, err := k.Tried(ctx, "machine.ace.silent", Attempt{What: "sent again", Outcome: "acted", By: "healer H1"},
ResolverHealer("H1")); err != nil {
t.Fatal(err)
}
if got, _, _ := k.Get(ctx, "machine.ace.silent"); len(got.Tried) != 2 {
t.Errorf("tried %+v", got.Tried)
}
time.Sleep(50 * time.Millisecond)
if n := len(told.Said()); n != 2 {
t.Errorf("a second attempt said %d events in all, want 2", n)
}
// Nothing open: nothing recorded, and no error.
if _, open, err := k.Tried(ctx, "machine.nobody.silent", Attempt{What: "x"}, ResolverHealer("H1")); open || err != nil {
t.Errorf("an attempt on nothing open: %v %v", open, err)
}
c.pass(time.Minute)
}
// **A spent budget is the operator's, urgent, and the watchdog seeing it again does not undo that.**
func TestAnEscalationHoldsAgainstTheNextObservation(t *testing.T) {
k, _, told, _ := keeper(t)
ctx := t.Context()
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
if _, _, err := k.Tried(ctx, "machine.ace.silent", Attempt{What: "asked", Outcome: "acted", By: "healer H1"},
ResolverHealer("H1")); err != nil {
t.Fatal(err)
}
held, _, err := k.Escalate(ctx, "machine.ace.silent", Attempt{What: "budget spent", Outcome: "escalated", By: "healer H1"})
if err != nil {
t.Fatal(err)
}
if held.Severity != Urgent || held.Resolver != ResolverOperator || !held.Escalated() {
t.Fatalf("escalated: %+v", held)
}
said := settled(t, told, 4)
if said[2].Change != ChangeSeverity || said[3].Change != ChangeResolver || said[3].Was != "healer:H1" {
t.Errorf("the escalation is not said as severity then resolver: %+v %+v", said[2], said[3])
}
again, err := k.Observe(ctx, silent("ace")) // a warning, as the watchdog says it
if err != nil {
t.Fatal(err)
}
if again.Severity != Urgent || again.Resolver != ResolverOperator {
t.Errorf("the next observation undid the escalation: %+v", again)
}
// A healer trying later does not take it back from the operator.
after, _, _ := k.Tried(ctx, "machine.ace.silent", Attempt{What: "x", By: "healer H1"}, ResolverHealer("H1"))
if after.Resolver != ResolverOperator {
t.Errorf("a later attempt took the condition back from the operator: %s", after.Resolver)
}
}
// **What was tried is carried into a reopening**: the same fault again within ten minutes is the
// same condition, and what the healers tried on it is still what they tried.
func TestAReopeningCarriesWhatWasTried(t *testing.T) {
k, _, _, c := keeper(t)
ctx := t.Context()
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
if _, _, err := k.Tried(ctx, "machine.ace.silent", Attempt{What: "asked", By: "healer H1"}, ResolverHealer("H1")); err != nil {
t.Fatal(err)
}
if _, err := k.Clear(ctx, "machine.ace.silent", "heard again"); err != nil {
t.Fatal(err)
}
c.pass(5 * time.Minute)
again, err := k.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if again.Count != 2 || len(again.Tried) != 1 || again.Tried[0].What != "asked" {
t.Errorf("reopened without what was tried: %+v", again)
}
}
+126
View File
@@ -0,0 +1,126 @@
package inventory
import (
"context"
"errors"
"fmt"
"strings"
"time"
)
// The healers' acts (novox/hq to-be 45 §7, Phase 3): see migration 0070. The controller is their
// only writer, under its lease; `healers` and `conditions` read them.
// The outcomes of a heal.
const (
// HealActing is an act begun and not yet finished: counted against the budget and the brake from
// the moment it starts, so a controller that dies mid-act has still spent it.
HealActing = "acting"
// HealActed is an act that did what it does. Whether it repaired anything is the observation's to
// say, never the healer's: the condition clears, or it does not.
HealActed = "acted"
// HealFailed is an act that could not be done.
HealFailed = "failed"
// HealEscalated is a budget spent, said: the condition was handed to the operator.
HealEscalated = "escalated"
)
// HealsKeptFor is how long a heal is kept: as long as a condition's history.
const HealsKeptFor = 90 * 24 * time.Hour
// Heal is one act of a healer.
type Heal struct {
ID int64 `json:"id"`
Healer string `json:"healer"`
ConditionKey string `json:"condition"`
Kind string `json:"kind"`
BudgetKey string `json:"budget-key"`
Act string `json:"act"`
Outcome string `json:"outcome"`
Said string `json:"said,omitempty"`
At time.Time `json:"at"`
Finished *time.Time `json:"finished,omitempty"`
Epoch uint64 `json:"epoch,omitempty"`
}
// BeginHeal writes an act about to be taken, under the lease: the act is counted from here. Refused,
// and nothing written, by a process that may not act.
func (i *Inventory) BeginHeal(ctx context.Context, h Heal) (Heal, error) {
epoch, err := i.actingEpoch(ctx)
if err != nil {
return h, fmt.Errorf("the heal %s of %s is not begun: %w", h.Healer, h.ConditionKey, err)
}
if strings.TrimSpace(h.Healer) == "" || strings.TrimSpace(h.ConditionKey) == "" || strings.TrimSpace(h.Act) == "" {
return h, errors.New("a heal names its healer, its condition and its act")
}
if h.BudgetKey == "" {
h.BudgetKey = h.ConditionKey
}
if h.Outcome == "" {
h.Outcome = HealActing
}
// At the runner's moment when it gives one, so its budgets and the brake are counted on one clock.
var at *time.Time
if !h.At.IsZero() {
at = &h.At
}
err = i.store.Pool().QueryRow(ctx,
`insert into heal (healer, condition_key, kind, budget_key, act, outcome, said, epoch, at)
values ($1, $2, $3, $4, $5, $6, $7, $8, coalesce($9, now())) returning id, at`,
h.Healer, h.ConditionKey, h.Kind, h.BudgetKey, h.Act, h.Outcome, h.Said, epoch, at).Scan(&h.ID, &h.At)
if err != nil {
return h, err
}
if epoch != nil {
h.Epoch = uint64(*epoch)
}
return h, nil
}
// FinishHeal says what came of an act begun. Written whatever the lease says by now: what was done was
// done, and its record must not depend on the doer still holding the lease a moment later.
func (i *Inventory) FinishHeal(ctx context.Context, id int64, outcome, said string) error {
tag, err := i.store.Pool().Exec(ctx,
`update heal set outcome = $2, said = $3, finished = now() where id = $1`, id, outcome, said)
if err != nil {
return err
}
if tag.RowsAffected() != 1 {
return fmt.Errorf("no heal %d to finish", id)
}
return nil
}
// HealsSince is every heal from a moment, oldest first.
func (i *Inventory) HealsSince(ctx context.Context, since time.Time) ([]Heal, error) {
rows, err := i.store.Pool().Query(ctx,
`select id, healer, condition_key, kind, budget_key, act, outcome, said, at, finished, epoch
from heal where at >= $1 order by at, id`, since)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Heal
for rows.Next() {
var h Heal
var epoch *int64
if err := rows.Scan(&h.ID, &h.Healer, &h.ConditionKey, &h.Kind, &h.BudgetKey, &h.Act, &h.Outcome, &h.Said,
&h.At, &h.Finished, &epoch); err != nil {
return nil, err
}
if epoch != nil {
h.Epoch = uint64(*epoch)
}
out = append(out, h)
}
return out, rows.Err()
}
// ForgetOldHeals removes what is older than HealsKeptFor, and says how many.
func (i *Inventory) ForgetOldHeals(ctx context.Context) (int64, error) {
tag, err := i.store.Pool().Exec(ctx, `delete from heal where at < $1`, time.Now().Add(-HealsKeptFor))
if err != nil {
return 0, err
}
return tag.RowsAffected(), nil
}
+43
View File
@@ -0,0 +1,43 @@
package inventory
import (
"context"
"errors"
"testing"
"time"
)
// **A heal is begun under the lease and kept with its epoch** (novox/hq to-be 45 §7): refused, and
// nothing written, by a process that may not act; finished whatever the lease says by then; read back
// in the order they were taken.
func TestAHealIsBegunUnderTheLeaseAndFinishedAfter(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
inv.ActsUnder(func(context.Context) (uint64, error) { return 0, errors.New("the lease is not held") })
if _, err := inv.BeginHeal(ctx, Heal{Healer: "H1", ConditionKey: "machine.anchor.sent-not-reported", Act: "asked"}); err == nil {
t.Fatal("a heal was begun by a process that may not act")
}
inv.ActsUnder(func(context.Context) (uint64, error) { return 57, nil })
at := time.Now().Add(-time.Minute).UTC().Truncate(time.Microsecond)
first, err := inv.BeginHeal(ctx, Heal{Healer: "H1", ConditionKey: "machine.anchor.sent-not-reported",
Kind: "sent-not-reported", Act: "asked", At: at})
if err != nil || first.ID == 0 || first.Epoch != 57 || first.BudgetKey != first.ConditionKey ||
first.Outcome != HealActing || !first.At.Equal(at) {
t.Fatalf("begun: %+v %v", first, err)
}
if _, err := inv.BeginHeal(ctx, Heal{Healer: "H3"}); err == nil {
t.Fatal("a heal naming no condition and no act was begun")
}
inv.ActsUnder(func(context.Context) (uint64, error) { return 0, errors.New("lost meanwhile") })
if err := inv.FinishHeal(ctx, first.ID, HealActed, "it reported"); err != nil {
t.Fatalf("what was done could not be recorded once the lease was gone: %v", err)
}
if err := inv.FinishHeal(ctx, first.ID+1000, HealActed, ""); err == nil {
t.Fatal("a heal that was never begun was finished")
}
heals, err := inv.HealsSince(ctx, time.Now().Add(-time.Hour))
if err != nil || len(heals) != 1 || heals[0].Outcome != HealActed || heals[0].Said != "it reported" ||
heals[0].Finished == nil || heals[0].Epoch != 57 {
t.Fatalf("read back: %+v %v", heals, err)
}
}
@@ -0,0 +1,34 @@
-- A known failure heals itself, under a brake, and every repair is said (novox/hq to-be 45 §7,
-- Phase 3, ADR 0227 rule 7).
--
-- Every act a healer takes is a row here, written before the act and finished after it: what the
-- budgets and the mesh-wide brake count, and what `healers` reads back. In the store and not in a
-- bucket on the bus, because the budget must outlive the controller — a controller restarting in a
-- loop must not reset a healer's count each time — and because the bus's user list would have to
-- grant a new bucket before the first heal could be counted.
--
-- Numbered 0070, past 0069, while the consumer-retirement feature (ADR 0230) is built beside this one:
-- a migration it adds takes 0069 if it merges first; one merged after this must be numbered above
-- 0070, because the store refuses to run a lower number than one already run.
create table heal (
id bigserial primary key,
-- The healer's id in the registry: H1, H2, ...
healer text not null,
-- The condition it acted on, and the condition's kind.
condition_key text not null,
kind text not null,
-- What its budget is counted against: the condition, a machine, a holder, a consumer.
budget_key text not null,
-- What it did, in the mesh's words.
act text not null,
-- 'acting' while the act runs; then 'acted', 'failed' or 'escalated' (a budget spent, said).
outcome text not null default 'acting',
said text not null default '',
at timestamptz not null default now(),
finished timestamptz,
-- The lease epoch it acted under (to-be 45 §6); null where none was claimed.
epoch bigint
);
create index heal_at on heal (at);
create index heal_budget on heal (healer, budget_key, at);
+4
View File
@@ -66,6 +66,10 @@ const (
// KeySecretReplaced: a value given to the mesh by hand was replaced with one it made, after its
// module's first good start (novox/hq ADR 0228). Never the value.
KeySecretReplaced = "secret-replaced"
// KeyHealerActed: a healer acted on a condition — what it did and what came of it, or that its budget
// is spent and the condition is the operator's (novox/hq to-be 45 §7). Never a person's act: those
// are the hand-act log's.
KeyHealerActed = "healer-acted"
)
// Applied is what a machine now runs, as the mesh states it.