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); or a delivery held past its state's bound where mesh-delivery's table lets H2 " + "take a transition (D14, novox/hq ADR 0239)", 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 — or, for a delivery, mesh-delivery's `close`, which " + "reads its walk again and takes only the transition its table names", 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) { if p.Release != nil { return "", "", nil // a release plan walks machines, and its own gate says when it is done (ADR 0236) } 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) { if c.Subject.Scope == conditions.ScopeDelivery { return appliesToAStalledDelivery(ctx, h, c) } 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) { if c.Subject.Scope == conditions.ScopeDelivery { return repairDelivery(ctx, h, c) } 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 }