diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 7c9a4b5..8a78ae3 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -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. diff --git a/cmd/mesh-controller/conditions.go b/cmd/mesh-controller/conditions.go index 3351436..8795174 100644 --- a/cmd/mesh-controller/conditions.go +++ b/cmd/mesh-controller/conditions.go @@ -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 ` says more; each clears when observation says it is resolved, " + "never by hand — `conditions silence --for --why ` 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 +} diff --git a/cmd/mesh-controller/doctor.go b/cmd/mesh-controller/doctor.go index e705a9b..b7e0d3a 100644 --- a/cmd/mesh-controller/doctor.go +++ b/cmd/mesh-controller/doctor.go @@ -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 } diff --git a/cmd/mesh-controller/healers.go b/cmd/mesh-controller/healers.go new file mode 100644 index 0000000..e77eb2d --- /dev/null +++ b/cmd/mesh-controller/healers.go @@ -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: `.`. +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 +} diff --git a/cmd/mesh-controller/healers_test.go b/cmd/mesh-controller/healers_test.go new file mode 100644 index 0000000..c181095 --- /dev/null +++ b/cmd/mesh-controller/healers_test.go @@ -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) + } +} diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index a7c71ee..cce4f30 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -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 diff --git a/cmd/mesh-controller/probes.go b/cmd/mesh-controller/probes.go index 865e539..6476acf 100644 --- a/cmd/mesh-controller/probes.go +++ b/cmd/mesh-controller/probes.go @@ -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 diff --git a/cmd/mesh-controller/readable.go b/cmd/mesh-controller/readable.go index d544e45..e7a84d4 100644 --- a/cmd/mesh-controller/readable.go +++ b/cmd/mesh-controller/readable.go @@ -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{} diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 54de7a9..5b5e2a4 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -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 != "" { diff --git a/cmd/mesh-controller/seatverbs_schema_test.go b/cmd/mesh-controller/seatverbs_schema_test.go index bf5b907..df4c877 100644 --- a/cmd/mesh-controller/seatverbs_schema_test.go +++ b/cmd/mesh-controller/seatverbs_schema_test.go @@ -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 diff --git a/cmd/mesh-controller/signals.go b/cmd/mesh-controller/signals.go index e29b00d..9b449e5 100644 --- a/cmd/mesh-controller/signals.go +++ b/cmd/mesh-controller/signals.go @@ -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 diff --git a/cmd/mesh-controller/signals_test.go b/cmd/mesh-controller/signals_test.go index 4a60cc0..98c67b7 100644 --- a/cmd/mesh-controller/signals_test.go +++ b/cmd/mesh-controller/signals_test.go @@ -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. diff --git a/cmd/mesh-controller/status.go b/cmd/mesh-controller/status.go index a7c732f..8025552 100644 --- a/cmd/mesh-controller/status.go +++ b/cmd/mesh-controller/status.go @@ -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. diff --git a/cmd/mesh-controller/status_summary.go b/cmd/mesh-controller/status_summary.go index ff91fe6..c9c9085 100644 --- a/cmd/mesh-controller/status_summary.go +++ b/cmd/mesh-controller/status_summary.go @@ -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. diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index 5213b5b..b121bd6 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -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) diff --git a/internal/broker/derived.go b/internal/broker/derived.go index 6244c50..a0bad9c 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -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 diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 2738b1e..b72b8d7 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -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 diff --git a/internal/broker/states_agreement_test.go b/internal/broker/states_agreement_test.go index f4e689a..9d11f88 100644 --- a/internal/broker/states_agreement_test.go +++ b/internal/broker/states_agreement_test.go @@ -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) diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 79d770b..5f92855 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -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", diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 7b9caed..782634a 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -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"] } diff --git a/internal/broker/writers.go b/internal/broker/writers.go index 1638a93..63d0862 100644 --- a/internal/broker/writers.go +++ b/internal/broker/writers.go @@ -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 diff --git a/internal/broker/writers_test.go b/internal/broker/writers_test.go index d1e3972..b4f17d9 100644 --- a/internal/broker/writers_test.go +++ b/internal/broker/writers_test.go @@ -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", diff --git a/internal/catalogue/seats.go b/internal/catalogue/seats.go index 4012cee..5bef725 100644 --- a/internal/catalogue/seats.go +++ b/internal/catalogue/seats.go @@ -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. diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index 4b3b501..025bb7a 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -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 " + diff --git a/internal/conditions/condition.go b/internal/conditions/condition.go index 81864f9..50311e7 100644 --- a/internal/conditions/condition.go +++ b/internal/conditions/condition.go @@ -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 } diff --git a/internal/conditions/store.go b/internal/conditions/store.go index 4c89690..dd5bd6d 100644 --- a/internal/conditions/store.go +++ b/internal/conditions/store.go @@ -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 diff --git a/internal/conditions/tried_test.go b/internal/conditions/tried_test.go new file mode 100644 index 0000000..bba81e9 --- /dev/null +++ b/internal/conditions/tried_test.go @@ -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) + } +} diff --git a/internal/inventory/heals.go b/internal/inventory/heals.go new file mode 100644 index 0000000..7792b44 --- /dev/null +++ b/internal/inventory/heals.go @@ -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 +} diff --git a/internal/inventory/heals_test.go b/internal/inventory/heals_test.go new file mode 100644 index 0000000..8069c58 --- /dev/null +++ b/internal/inventory/heals_test.go @@ -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) + } +} diff --git a/internal/inventory/migrations/0070-a-known-failure-heals-itself.sql b/internal/inventory/migrations/0070-a-known-failure-heals-itself.sql new file mode 100644 index 0000000..f8d10fa --- /dev/null +++ b/internal/inventory/migrations/0070-a-known-failure-heals-itself.sql @@ -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); diff --git a/internal/link/events.go b/internal/link/events.go index 9ef5ee3..cf51d9f 100644 --- a/internal/link/events.go +++ b/internal/link/events.go @@ -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.