package main import ( "context" "encoding/json" "errors" "fmt" "slices" "sort" "strings" "time" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/micro" "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/conditions" "github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/lease" "github.com/novox/mesh-controller/internal/link" ) // The gate on a release plan's first machine, and the rollback after it (novox/hq ADR 0236, to-be 45 // §8, ADR 0227 rule 8). // // **"Reported applied" is not enough.** ADR 0218 sent a module to one machine first and the rest once // that machine reported it applied. A build that applies and then does nothing, crashes, serves no // tools, or breaks the machine's word to the mesh passed that test. Now the first machine is judged by // the component's health — the core's health definitions, or a module's own — passing three times over // at least two minutes, within ten minutes of the apply. Only then are the rest sent. // // **A failing gate stops the plan and puts the previous build back there.** The build is marked failed // at its gate — registration refuses it and nothing sends it again on its own — the module's registered // build goes back to the one the first machine ran before, and that machine is sent it: the ordinary // path, again. Once per build: the verdict is written before the send, and a verdict once written is // not written over. Said as a condition, urgent for the core and when it could not be put back, and as // the event `rolled-back`. // The gate's bounds (to-be 45 §8). Variables so a test can judge in a second, not in minutes. var ( // gateSettle is how long the first machine must stay healthy, at the least. gateSettle = 2 * time.Minute // gateBound is how long after the apply the build may take to become healthy. gateBound = 10 * time.Minute // gatePasses is how many consecutive judgings must find it healthy, gateEvery apart at the least. gatePasses = 3 gateEvery = 40 * time.Second ) // gateProbe is the registry's row for the gate's verdicts and the witnesses' rollbacks: its conditions // are raised by a plan as it judges, and kept or cleared by the probe on every run. const gateProbe = "DG" // The kinds a gate raises. const ( kindRolledBack = "rolled-back" kindRollbackFailed = "rollback-failed" ) // coreComponent is the core component a module is, as a witness names it — controller, node-engine, // node-tools — or empty: the core is judged by its health definitions, anything else by its own health. func coreComponent(module string) string { switch module { case catalogue.ControllerSeatName: return lease.ComponentController case hostModule: return lease.ComponentEngine case broker.RuntimeModule: return lease.ComponentNodeTools } return "" } // health is a judging's word on one machine. type health int const ( healthGood health = iota // healthPerson is a wait for a person (novox/hq ADR 0254): everything unhealthy of the module waits for // one person's new login (personWait). The build did what it should; no bound a machine can meet ends // the wait. It counts as a pass, and the gate carries the reading along in its verdict. healthPerson // healthWaiting is not a pass and not a fault: what the module's checks find waits on an unhealthy // provider (ADR 0240 rule 5), so the judging waits — past the bound too — rather than putting back a // build for something it did not do. healthWaiting healthNotYet healthBroken ) // served is what one machine's node tools answered the bus's discovery with. type served struct { runtime bool tools map[string]bool } // gateFacts is what one judging reads, gathered once for every machine it judges. type gateFacts struct { now time.Time reports map[string]inventory.Reported // engines is each machine's node-engine build as it last reported it. engines map[string]string // served is what the bus's discovery answered; servedErr why it could not be asked. served map[string]served servedErr error // open are the open conditions; openErr why they could not be read (nil keeper: not judged). open []conditions.Condition openErr error judged bool // rolledBack is what each machine's witnesses say they decided. rolledBack map[string][]lease.Rollback // holder is who holds the controller lease, for judging the controller. holder *lease.Holder holderErr error // health is each machine's newest health statement (ADR 0240); a machine absent never stated one. // healthErr is why they could not be read. health map[string]inventory.NodeHealth healthErr error // heldOn is, per "@", the provider its findings are held under (ADR 0240 rule 5). heldOn map[string]string // groupsAdded is, per module, whether the move judged puts an account in a group its previous build did // not (issue 318 review): the only move whose wait for a new login is excused. groupsAdded map[string]bool } // gatherGateFacts reads what a judging needs, from the store, the bus and this controller's memory. A // variable so a test can hand a judging its facts. var gatherGateFacts = func(ctx context.Context, open *stores, component string) (gateFacts, error) { inv := open.inventory f := gateFacts{now: time.Now(), reports: map[string]inventory.Reported{}, engines: map[string]string{}, rolledBack: witnessed.all()} reports, err := inv.LastReports(ctx) if err != nil { return f, err } for _, r := range reports { f.reports[r.Node] = r } nodes, err := inv.Nodes(ctx) if err != nil { return f, err } for _, n := range nodes { f.engines[n.Name] = n.HostVersion } // What each machine says of its long-running resources (ADR 0240): unreadable is said, never read as // healthy. f.health, f.healthErr = inv.Healths(ctx) if d := doctorFrom; d != nil { if d.keeper != nil { f.judged = true f.open, f.openErr = d.keeper.Open(ctx) } if d.js != nil { f.served, f.servedErr = servedOnTheBus(ctx, d.js.Conn()) } else { f.servedErr = errors.New("this controller has no bus to ask") } } else { f.servedErr = errors.New("this process does not serve the mesh, so it cannot ask the bus who serves what") } // Whose findings wait on an unhealthy provider (ADR 0240 rule 5): their gates wait, not fail. if f.healthErr == nil && f.openErr == nil { if hold, err := readHolding(ctx, inv, f.open); err == nil { f.heldOn = map[string]string{} for machine := range f.health { for module, p := range hold.heldModules(machine) { f.heldOn[module+"@"+machine] = p.Module + " on " + p.Node } } } } if theLease != nil { h, found, err := theLease.holder(ctx) switch { case err != nil: f.holderErr = err case found: f.holder = &h } if component != lease.ComponentController && f.holderErr != nil { f.holderErr = nil // read only for the controller's own judging } } return f, nil } // judgeHealth is one machine's health for a module's new build, from what one judging read: healthy, // not yet (with what is wanting), or healthBroken — a witness put it back, or the machine refused or failed // what it was sent. Pure. func judgeHealth(module, component string, m catalogue.Manifest, machine string, since time.Time, f gateFacts) (health, string) { // A witness's verdict made since the build was sent: the build failed its health there. One that // could not judge at all (unwitnessed) is not a verdict on the build. for _, r := range f.rolledBack[machine] { if component == "" || r.Component != component || r.Outcome == lease.OutcomeUnwitnessed || r.At.Before(since) { continue } return healthBroken, fmt.Sprintf("the witness on %s judged the %s %s and %s: %s", machine, r.Component, short(r.From), r.Outcome, r.Why) } r, said := f.reports[machine] switch { case !said || r.At == nil || !r.Current: return healthNotYet, fmt.Sprintf("%s has not reported on what it was sent", machine) case r.Outcome == inventory.OutcomeFailed || r.Outcome == inventory.OutcomeRefused: return healthBroken, fmt.Sprintf("%s %s what it was sent", machine, r.Outcome) case r.Outcome != inventory.OutcomeApplied: return healthNotYet, fmt.Sprintf("%s reported %q", machine, r.Outcome) } // **No new condition about it**: about the machine itself, or naming the module on that machine, // raised since the judging began. The gate's own are not evidence about the build. if f.judged { if f.openErr != nil { return healthNotYet, "what is wrong cannot be read, so whether the build made anything wrong is not known: " + firstLine(f.openErr.Error()) } for _, c := range f.open { // A wait for a person's new login, or for a directory used as found to be handed over, is the module's // reading, not a fault raised since the send: the gate reads it from the statement below (ADR 0254, // novox/hq issue 339). if c.Source == gateProbe || c.Raised.Before(since) || c.Kind == kindReloginNeeded || c.Kind == kindUsedAsFound { continue } onIt := c.Subject.Machine == machine || slices.Contains(c.Subject.Also, machine) || (c.Subject.Scope == conditions.ScopeMachine && c.Subject.ID == machine) if !onIt { continue } // A module's own health condition names it whole, its name's dots and all (ADR 0240). ownHealth := c.Subject.Scope == conditions.ScopeModule && c.Subject.ID == module+"."+machine if c.Subject.Scope == conditions.ScopeMachine || ownHealth || slices.Contains(strings.Split(c.Subject.ID, "."), module) { return healthNotYet, fmt.Sprintf("raised since it was sent: %s — %s", c.Key, c.Summary) } } } switch component { case lease.ComponentEngine: // The node-engine has reported its current declaration under its own build. want := deliveredVersions(m) if len(want) > 0 && !slices.Contains(want, f.engines[machine]) { return healthNotYet, fmt.Sprintf("%s's node-engine reports build %s, not the new %s", machine, orNotKnown(f.engines[machine]), strings.Join(want, " or ")) } case lease.ComponentNodeTools: // The node tools are announced and answer. if f.servedErr != nil { return healthNotYet, "whether the node tools answer cannot be asked: " + firstLine(f.servedErr.Error()) } if !f.served[machine].runtime { return healthNotYet, fmt.Sprintf("the node tools on %s do not answer the bus", machine) } case lease.ComponentController: // The new controller holds the lease and says it is ready. switch { case f.holderErr != nil: return healthNotYet, "who holds the controller lease cannot be read: " + firstLine(f.holderErr.Error()) case f.holder == nil: return healthNotYet, "no controller holds the lease" case f.holder.Taken.Before(since.Add(-time.Minute)): return healthNotYet, fmt.Sprintf("the lease is held since %s, by a controller older than the new build", f.holder.Taken.UTC().Format(time.RFC3339)) case f.holder.Health == nil || !f.holder.Health.Ready: why := "the controller holding the lease does not say it is ready" if f.holder.Health != nil && f.holder.Health.Why != "" { why += ": " + f.holder.Health.Why } return healthNotYet, why } default: // A module's tools answer, where it has any and the machine runs the node tools that serve them. if len(m.Tools) > 0 { if f.servedErr != nil { return healthNotYet, "whether its tools are served cannot be asked: " + firstLine(f.servedErr.Error()) } if s := f.served[machine]; s.runtime && !s.tools[module] { return healthNotYet, fmt.Sprintf("the node tools on %s do not serve %s's tools", machine, module) } } // **And what it runs is stated healthy** (ADR 0240 §4): every long-running resource of it on that // machine, in a statement heard since the send. A resource still starting makes the judging wait. if h, why := moduleHealthWord(module, machine, since, f); h != healthGood { return h, why } } return healthGood, "" } // machineWord is what one judging found wrong with a machine itself, apart from its modules: facts is // the judging's facts without those conditions, on what each names among the modules the send moved, // and whole what holds the machine back as a whole. type machineWord struct { facts gateFacts on map[string]string whole string // waiting is what holds the machine on something shown to be another's (ADR 0241): the judging waits. waiting string } // kindCoreBehind is D10's kind: a machine runs core components older than the mesh holds, or has not // yet reported applying the node tools it was last sent. const kindCoreBehind = "core-behind" // aboutTheMachine sorts the conditions raised about a machine itself since a send was made there // (novox/hq issue 281). The gate read every one of them as the module's it was kept on: a machine-level // condition caused by anything else in the send — or by a tier sent one module at a time — failed that // module, at the bound, with a reason that was never about it. // // - one naming a module the send moved is that module's; // - the core being behind (D10) is the core's: of the node-engine or the node tools when the send // moved them — inside their settle window, which is the gate's bound — and otherwise no evidence about // what was sent: a send not yet applied the machine's reports already say, and a newer core build // that no plan sends is not this send's; // - anything else holds the machine back as a whole: everything the send moved there waits on it, and // fails with it, together, at the bound, with that reason. // // Pure. func aboutTheMachine(machine string, moved []string, since time.Time, f gateFacts) machineWord { w := machineWord{facts: f, on: map[string]string{}} if !f.judged || f.openErr != nil { return w } var core []string for _, m := range moved { if c := coreComponent(m); c == lease.ComponentEngine || c == lease.ComponentNodeTools { core = append(core, m) } } kept := make([]conditions.Condition, 0, len(f.open)) for _, c := range f.open { aboutIt := c.Subject.Scope == conditions.ScopeMachine && (c.Subject.ID == machine || c.Subject.Machine == machine || slices.Contains(c.Subject.Also, machine)) // A directory used as found waits for a person, whatever the send did (novox/hq issue 339). if !aboutIt || c.Source == gateProbe || c.Raised.Before(since) || c.Kind == kindUsedAsFound { kept = append(kept, c) continue } said := fmt.Sprintf("raised since it was sent: %s — %s", c.Key, c.Summary) parts := strings.Split(c.Subject.ID, ".") var named []string for _, m := range moved { if slices.Contains(parts, m) { named = append(named, m) } } coreBehind := c.Kind == kindCoreBehind || strings.HasSuffix(c.Key, "."+kindCoreBehind) switch { case len(named) > 0: for _, m := range named { if _, already := w.on[m]; !already { w.on[m] = said } } case coreBehind && len(core) > 0: for _, m := range core { if _, already := w.on[m]; !already { w.on[m] = said + " (its settle window runs to the gate's bound)" } } case coreBehind: // Not what the send moved: said nowhere against it. case c.Kind == kindMachineUnits: // **The machine's own failed units** (issue 315): no module places any of them, so nothing a // send moved is among them — a module's failed unit is that module's own condition. Said // nowhere against the send. case c.Kind == kindNetworkRewritten || c.Kind == kindNetworkUnreachable && c.Subject.ID != machine: // **Shown to be somebody else's** (ADR 0241): another program rewrote the resolver file the send // did not move, or the machine cannot reach another that is down. Nothing the send did; the // judging waits for it rather than putting back a build at the bound. if w.waiting == "" { w.waiting = machine + "'s network: " + said } default: if w.whole == "" { w.whole = machine + " as a whole: " + said } } } w.facts.open = kept return w } // judgeGate takes one judging of a module's first machines and records it in the plan's gate: a pass // counted, a pass missed (and why), or the verdict. Answers the verdict once there is one. func judgeGate(ctx context.Context, open *stores, p *inventory.Plan, module string, state *inventory.PlanModule, running []string, now time.Time) (string, error) { g := state.Gate if g == nil { // Judged from the send: a witness's verdict, a condition, the bound — all counted from when the // first machine was sent the build. start := now if state.FirstAt != nil { start = *state.FirstAt } var machines []string for _, n := range state.First { if slices.Contains(running, n) { machines = append(machines, n) } } g = &inventory.PlanGate{Component: coreComponent(module), Machines: machines, From: state.Previous, To: state.Commit, Since: &start} state.Gate = g } pairs := []judged{} for _, n := range g.Machines { pairs = append(pairs, judged{module: module, node: n}) } return judgeMoves(ctx, open, g, pairs, now) } // judged is one module on one machine, as a gate judges it. type judged struct{ module, node string } // judgeMoves takes one judging of a gate over the modules it judges on their machines — its own, and // everything the send carried (Carried) — and records it: a pass counted, a pass missed (what is // wanting, and which modules), or the verdict. Answers the verdict once there is one. func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs []judged, now time.Time) (string, error) { if g.Verdict != "" { return g.Verdict, nil } if g.LastPass != nil && now.Sub(*g.LastPass) < gateEvery { return "", nil } for _, c := range g.Carried { if !slices.Contains(pairs, judged{module: c.Module, node: c.Node}) { pairs = append(pairs, judged{module: c.Module, node: c.Node}) } } shelf, err := open.inventory.Catalogue(ctx) if err != nil { return "", err } facts, err := gatherGateFacts(ctx, open, g.Component) if err != nil { return "", err } facts.groupsAdded = movesAddingGroups(ctx, open.inventory, g, pairs, shelf) // **What is wrong with a machine itself is the machine's** (novox/hq issue 281): read once for each // machine judged, apart from what is wrong with a module there, and never pinned on the module the // gate happens to be kept on. byMachine := map[string][]string{} for _, j := range pairs { byMachine[j.node] = append(byMachine[j.node], j.module) } words := map[string]machineWord{} for node, moved := range byMachine { words[node] = aboutTheMachine(node, moved, *g.Since, facts) } // **Each module is judged on its own** (novox/hq ADR 0254, issue 318): its reading is the worst of its // machines', its passes are counted apart, and the send's verdict is still one. **Every module of a send // leaves it with a verdict** (issue 318 review): one healthy for the passes the gate asks, after the // settle time, keeps a pass when another beside it fails; one found broken holds the send's verdict until // the modules beside it have theirs, or the bound; and at the verdict, every module that did not pass is // put back with what failed, so none is left on the machine unjudged for other walks to wait on. worst, why := healthGood, "" reading := map[string]health{} waits := map[string]string{} var modules []string for _, j := range pairs { w := words[j.node] h, said := judgeHealth(j.module, coreComponent(j.module), shelf[j.module], j.node, *g.Since, w.facts) if h == healthGood || h == healthPerson { if on, named := w.on[j.module]; named { h, said = healthNotYet, on } else if w.whole != "" { h, said = healthNotYet, w.whole } else if w.waiting != "" { h, said = healthWaiting, w.waiting } } if _, seen := reading[j.module]; !seen { modules = append(modules, j.module) } if h == healthBroken && !slices.Contains(g.Broken, j.module) { g.Broken = append(g.Broken, j.module) if g.BrokenWhy == "" { g.BrokenWhy = said } } if slices.Contains(g.Broken, j.module) { h = healthBroken // found broken once, broken for the rest of the judging } if h > reading[j.module] { reading[j.module] = h } if h == healthPerson { waits[j.module] = joinSaid(waits[j.module], said) } if h > worst { worst, why = h, said } else if h == worst && h > healthPerson && why == "" { why = said } } if g.Healthy == nil { g.Healthy = map[string]int{} } if g.HealthyAt == nil { g.HealthyAt = map[string]time.Time{} } g.Waits = nil var failing []string for _, m := range modules { if reading[m] > healthPerson { failing = append(failing, m) g.Healthy[m] = 0 delete(g.HealthyAt, m) continue } // **A pass is counted only gateEvery after the one before** (issue 318 review): a send judged more // often while another module beside it is not yet healthy counts no faster. if at, counted := g.HealthyAt[m]; !counted || now.Sub(at) >= gateEvery { g.Healthy[m]++ g.HealthyAt[m] = now } if w := waits[m]; w != "" { if g.Waits == nil { g.Waits = map[string]string{} } g.Waits[m] = w } } settled := now.Sub(*g.Since) >= gateSettle passedAlone := func(m string) bool { return settled && reading[m] <= healthPerson && g.Healthy[m] >= gatePasses } // fail decides the send failed: what passed on its own keeps its pass, everything else is put back. fail := func(why string) { g.Passing, g.Failing = nil, nil for _, m := range modules { if passedAlone(m) { g.Passing = append(g.Passing, m) } else { g.Failing = append(g.Failing, m) } } decide(g, inventory.GateFailed, why, now) } pastBound := now.Sub(*g.Since) > gateBound switch { case worst == healthBroken: var judging []string for _, m := range modules { if reading[m] != healthBroken && !passedAlone(m) { judging = append(judging, m) } } if len(judging) == 0 || pastBound { fail(g.BrokenWhy) break } // What broke fails the send, and is put back at once by the caller (putBackBroken); what is beside // it is judged to its own verdict first, within the bound. g.Passes, g.LastPass, g.Failing = 0, nil, failing g.Last = fmt.Sprintf("%s; %s judged to its own verdict before the send's", g.BrokenWhy, strings.Join(judging, ", ")) case worst == healthWaiting: // Waiting on a provider that is unhealthy: not a pass, and not a failure at the bound either — // the provider's own condition says what is wrong (ADR 0240 rule 5). g.Passes, g.LastPass, g.Last, g.Failing = 0, nil, why, failing case worst == healthNotYet: g.Passes, g.LastPass, g.Last, g.Failing = 0, nil, why, failing if pastBound { fail(fmt.Sprintf("not healthy within %s of its apply: %s", gateBound, why)) } default: // Healthy, or waiting for a person (ADR 0254): a pass, the wait carried along in the verdict. g.Passes++ g.LastPass, g.Last, g.Failing = &now, "", nil if g.Passes >= gatePasses && settled { decide(g, inventory.GatePassed, fmt.Sprintf("healthy %d times over %s", g.Passes, now.Sub(*g.Since).Round(time.Second))+waitsSaid(g.Waits), now) } } return g.Verdict, nil } // whyFor is a passing gate's why as one module's verdict says it: the send's, and that module's own wait // for a person, never another's (issue 318 review). func whyFor(g *inventory.PlanGate, module string) string { why, _, _ := strings.Cut(g.Why, waitsPrefix) if w := g.Waits[module]; w != "" { why += waitsPrefix + w } return why } const waitsPrefix = "; and it waits for a person: " // waitsSaid is the waits for a person a passing gate carries, as its verdict says them. func waitsSaid(waits map[string]string) string { if len(waits) == 0 { return "" } modules := make([]string, 0, len(waits)) for m := range waits { modules = append(modules, m) } sort.Strings(modules) var said []string for _, m := range modules { said = append(said, waits[m]) } return waitsPrefix + strings.Join(said, "; ") } func joinSaid(a, b string) string { switch { case a == "": return b case b == "" || strings.Contains(a, b): return a } return a + "; " + b } // movesAddingGroups is, per module a gate judges, whether its move puts an account in a group the build it // moved from did not: read from the builds' manifests. A module whose move is not known adds none. func movesAddingGroups(ctx context.Context, inv *inventory.Inventory, g *inventory.PlanGate, pairs []judged, shelf map[string]catalogue.Manifest) map[string]bool { out := map[string]bool{} for _, j := range pairs { if _, done := out[j.module]; done { continue } from, to := g.From, g.To for _, c := range g.Carried { if c.Module == j.module { from, to = c.From, c.To break } } target, found, err := inv.ManifestAt(ctx, j.module, to) if err != nil || !found { target = shelf[j.module] } // What it moved from is not known — no build named (a plan's own module whose previous build was not // recorded), or no manifest kept for it: no wait is excused, rather than one the move did not bring. if from == "" { out[j.module] = false continue } before, had, err := inv.ManifestAt(ctx, j.module, from) if err != nil || !had { out[j.module] = false continue } out[j.module] = addsAccountGroups(before, had, target) } return out } // decide sets a gate's verdict. func decide(g *inventory.PlanGate, verdict, why string, now time.Time) { g.Verdict, g.Why, g.JudgedAt = verdict, why, &now g.Took = now.Sub(*g.Since).Round(time.Second).String() } // gatePassed keeps a passing build's verdict, so `plans` and the gate's probe can read it. func gatePassed(ctx context.Context, open *stores, p *inventory.Plan, module string, state *inventory.PlanModule) { g := state.Gate err := open.inventory.RecordGate(ctx, inventory.GateVerdict{Build: state.Build, Module: module, Commit: state.Commit, Previous: state.Previous, Plan: p.ID, Machines: g.Machines, Verdict: inventory.GatePassed, Why: whyFor(g, module), Component: g.Component, JudgingFrom: g.Since}) if err != nil && state.Build != "" { fmt.Printf("%s: %s passed its gate, and the verdict could not be kept: %v\n", p.ID, module, err) } passCarried(ctx, open, p, g, module) fmt.Printf("%s: %s passed its gate on %s (%s); the rest are sent\n", p.ID, module, strings.Join(g.Machines, ", "), g.Why) if module == catalogue.ControllerSeatName { carryUserList(ctx, open, p) } } // carryUserList sends the machine holding the bus the user list a new controller composes, once that // controller passed its gate (ADR 0236). The controller's own grants travel in that list, and the old // controller composed the list the plan sent; on 2026-10-06 eight pushes by hand carried a new // controller's grant into it. Not when a build its policy or a plan holds back would go with it (ADR // 0221): then it is said, as a push would say it. func carryUserList(ctx context.Context, open *stores, p *inventory.Plan) { holder, behind, err := brokerBehind(ctx, open, nil) if err != nil || holder == "" || !behind { if err != nil { fmt.Printf("%s: whether the bus's user list is behind the new controller cannot be read: %v\n", p.ID, err) } return } held, err := heldMachines(ctx, open, []string{holder}) if err != nil { fmt.Printf("%s: whether %s may be sent the new user list cannot be read: %v\n", p.ID, holder, err) return } if why, isHeld := held[holder]; isHeld { fmt.Printf("%s: the new controller's user list is not carried to %s, which holds the bus: %s — `push %s` "+ "carries it\n", p.ID, holder, strings.Join(why, "; "), holder) return } if _, err := sendRollout(ctx, open, []string{holder}); err != nil { fmt.Printf("%s: the new controller's user list could not be carried to %s: %v\n", p.ID, holder, err) return } fmt.Printf("%s: carried the new controller's user list to %s, which holds the bus\n", p.ID, holder) } // gateFailed stops the plan at a build that failed its gate and puts the previous build back on the // machines it was judged on — once per build, said as a condition and an event. The plan is saved // before the send: a controller that is itself the build being put back does not outlive it. func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module string, state *inventory.PlanModule, machines []string, why string) { inv := open.inventory now := time.Now().UTC() if state.Gate == nil { state.Gate = &inventory.PlanGate{Component: coreComponent(module), Machines: machines, From: state.Previous, To: state.Commit, Since: state.FirstAt} } g := state.Gate if g.Verdict == "" { if g.Since == nil { g.Since = &now } decide(g, inventory.GateFailed, why, now) } if len(g.Machines) == 0 { g.Machines = machines } state.Why = "failed its gate: " + g.Why p.State = inventory.PlanFailed p.Note = fmt.Sprintf("%s failed its gate on %s in tier %d: %s", module, strings.Join(g.Machines, ", "), p.Tier, g.Why) verdict := inventory.GateVerdict{Build: state.Build, Module: module, Commit: state.Commit, Previous: state.Previous, Plan: p.ID, Machines: g.Machines, Verdict: inventory.GateFailed, Rollback: inventory.RollingBack, Why: g.Why, Component: g.Component, JudgingFrom: g.Since} // **A send that changed nothing of the module there is no verdict on its build** (novox/hq issue // 280). The machine already ran this build, or one that made the same artifacts from the same // manifest: whatever the gate found wanting, this build did not bring it, and there is nothing to // put back. On 2026-10-06 such a module was marked failed, and the rollback looked for an earlier // build of the very commit it had failed — "no build kept" — while the build it had run before was // kept all along. Said, left as it is, never marked. if unchangedBy(ctx, inv, module, state.Previous, state.Commit) { g.Rollback = gateUnchanged state.Why = "stopped with its send; the send changed nothing of it: " + g.Why p.Note = fmt.Sprintf("the send to %s in tier %d failed its gate: %s; %s was left as it was — %s already ran "+ "%s %s, or a build identical to it, before the send, so nothing of it moved and nothing is put back", strings.Join(g.Machines, ", "), p.Tier, g.Why, module, strings.Join(g.Machines, ", "), module, short(state.Commit)) fmt.Printf("%s: %s\n", p.ID, p.Note) return } if state.Build == "" { // A plan from before builds were asked by id: nothing to mark, so nothing is put back by the // mesh — said, for a person. g.Rollback = inventory.NotRolledBack p.Note += "; not put back: the plan does not know which build it sent" sayRollback(ctx, open, module, g, "") return } if err := inv.RecordGate(ctx, verdict); err != nil { if errors.Is(err, inventory.ErrGateKept) { // Already judged and acted on, by this controller before a restart or by another: never twice. if kept, found, _ := inv.GateOf(ctx, state.Build); found { g.Rollback = kept.Rollback } p.Note += "; its rollback was already made once and is not made again" return } g.Rollback = inventory.NotRolledBack p.Note += "; not put back: its verdict could not be kept, and a rollback that cannot be counted is not made — " + err.Error() sayRollback(ctx, open, module, g, "") return } // The previous build: the one the first machine ran, from the build records. notBack := func(why string) { g.Rollback = inventory.NotRolledBack p.Note += "; NOT put back: " + why if err := inv.SetRollback(ctx, state.Build, inventory.NotRolledBack, g.Why+"; not put back: "+why); err != nil { fmt.Printf("%s: how %s's rollback went could not be kept: %v\n", p.ID, module, err) } sayRollback(ctx, open, module, g, why) } if state.Previous == "" { // A first build there: nothing ran before it, so nothing can be put back, and the machine is left // with it — said as that, not as a build the mesh lost. notBack(fmt.Sprintf("this is the first build of %s that %s was sent, or what it was sent before is not known: "+ "there is no earlier build there to put back, so it is left with this one — `unassign` takes it off, "+ "a newer merge replaces it", module, strings.Join(g.Machines, ", "))) return } failed, _, err := inv.BuildByID(ctx, state.Build) if err != nil { notBack("the failed build's record cannot be read: " + err.Error()) return } previous, found, err := inv.PreviousBuild(ctx, module, state.Previous, failed) if err != nil { notBack("the build records cannot be read: " + err.Error()) return } if !found { notBack(fmt.Sprintf("no build of %s from %s, asked before the failed one, is among the %d newest kept to put "+ "back", module, short(state.Previous), inventory.KeptBuilds)) return } if err := inv.RestoreModule(ctx, previous); err != nil { notBack(err.Error()) return } g.Rollback = inventory.RollingBack if err := inv.SavePlan(ctx, p); err != nil { fmt.Printf("%s: the plan could not be kept before %s is put back: %v\n", p.ID, module, err) } // Put back with the rest of its send, in one send per machine (issue 281). if b, batched := ctx.Value(rollbacksKey{}).(*rollbacks); batched { b.pending = append(b.pending, pendingRollback{module: module, state: state, g: g, previous: previous}) b.modules[module] = true for _, n := range g.Machines { if !slices.Contains(b.machines, n) { b.machines = append(b.machines, n) } } return } sent, err := sendRollout(withScope(ctx, sendScope{modules: map[string]bool{module: true}}), open, g.Machines) if err != nil { notBack(fmt.Sprintf("its registered build is back at %s, and sending it to %s was refused: %v — `push %s` "+ "sends it", short(previous.Commit), strings.Join(g.Machines, ", "), err, g.Machines[0])) return } g.Rollback = inventory.RolledBack p.Note += fmt.Sprintf("; put back to %s on %s", short(previous.Commit), strings.Join(sent, ", ")) if err := inv.SetRollback(ctx, state.Build, inventory.RolledBack, g.Why); err != nil { fmt.Printf("%s: how %s's rollback went could not be kept: %v\n", p.ID, module, err) } sayRollback(ctx, open, module, g, "") } // rollbacks is what a failed send puts back, sent together (novox/hq issue 281): a gate that judged one // send judges what it moved as one, and what it found wanting goes back in one send per machine — not // in a send for each module, which is the churn that failed the gate in the first place. type rollbacks struct { modules map[string]bool machines []string pending []pendingRollback } type pendingRollback struct { module string state *inventory.PlanModule g *inventory.PlanGate previous inventory.Build } type rollbacksKey struct{} // batchingRollbacks is a context under which gateFailed registers what it puts back and leaves the send // to sendRollbacks. func batchingRollbacks(ctx context.Context) (context.Context, *rollbacks) { b := &rollbacks{modules: map[string]bool{}} return context.WithValue(ctx, rollbacksKey{}, b), b } // sendRollbacks sends what a failed send put back, once to each machine, and says each module's rollback. func sendRollbacks(ctx context.Context, open *stores, p *inventory.Plan, b *rollbacks) { if len(b.pending) == 0 { return } inv := open.inventory sort.Strings(b.machines) sent, err := sendRollout(withScope(ctx, sendScope{modules: b.modules}), open, b.machines) var back []string for _, r := range b.pending { if err != nil { why := fmt.Sprintf("its registered build is back at %s, and sending it to %s was refused: %v — `push %s` "+ "sends it", short(r.previous.Commit), strings.Join(r.g.Machines, ", "), err, firstOf(r.g.Machines)) r.g.Rollback = inventory.NotRolledBack if err := inv.SetRollback(ctx, r.state.Build, inventory.NotRolledBack, r.g.Why+"; not put back: "+why); err != nil { fmt.Printf("%s: how %s's rollback went could not be kept: %v\n", p.ID, r.module, err) } sayRollback(ctx, open, r.module, r.g, why) continue } r.g.Rollback = inventory.RolledBack back = append(back, r.module+" to "+short(r.previous.Commit)) if err := inv.SetRollback(ctx, r.state.Build, inventory.RolledBack, r.g.Why); err != nil { fmt.Printf("%s: how %s's rollback went could not be kept: %v\n", p.ID, r.module, err) } sayRollback(ctx, open, r.module, r.g, "") } if err != nil { p.Note += fmt.Sprintf("; NOT put back: sending %s was refused: %v", strings.Join(b.machines, ", "), err) return } p.Note += fmt.Sprintf("; put back %s on %s, in one send", strings.Join(back, ", "), strings.Join(sent, ", ")) } // gateUnchanged is a failed gate's word on a module its send changed nothing of (novox/hq issue 281): // left as it was, never marked failed, and so no verdict on its build — `plans retry` asks it again. const gateUnchanged = "unchanged" // unchangedBy says whether a send moving a module from one build to another changed nothing of it: the // same commit, builds made from the same source, or builds that made the same artifacts from the same // manifest. Not known is changed. func unchangedBy(ctx context.Context, inv *inventory.Inventory, module, from, to string) bool { if from == "" || to == "" { return false } if sameCommit(from, to) { return true } // Made from the same source (issue 280), or the same artifacts from the same manifest. f := moveFacts{} var err error if f.srcs, err = inv.SourceFingerprints(ctx); err != nil { return false } if f.fps, err = inv.Fingerprints(ctx); err != nil { return false } return f.identical(module, from, to) } // rolledBackEvent is the body of `rolled-back` (ADR 0236): a contract, like a condition's events. type rolledBackEvent struct { Event string `json:"event"` At time.Time `json:"at"` Module string `json:"module"` Component string `json:"component,omitempty"` Machines []string `json:"machines"` From string `json:"from,omitempty"` To string `json:"to,omitempty"` Why string `json:"why"` // Rollback is rolled-back, or not-rolled-back with NotWhy. Rollback string `json:"rollback"` NotWhy string `json:"not_why,omitempty"` Show string `json:"show"` } // sayRollback raises the gate's condition at once — the probe keeps it from then — and says the event. func sayRollback(ctx context.Context, open *stores, module string, g *inventory.PlanGate, notWhy string) { o := gateObservation(module, g.Component, g.Machines, g.Rollback, g.Why, notWhy, g.To, g.From) fmt.Println(o.Summary) d := doctorFrom if d == nil { return } if d.keeper != nil { o.Source = gateProbe if _, err := d.keeper.Observe(ctx, o); err != nil { fmt.Printf("the gate's condition %s could not be kept: %v\n", o.Key(), err) } } if d.teller == nil { return } body, err := json.Marshal(rolledBackEvent{Event: link.KeyRolledBack, At: time.Now().UTC(), Module: module, Component: g.Component, Machines: g.Machines, From: g.To, To: g.From, Why: g.Why, Rollback: g.Rollback, NotWhy: notWhy, Show: conditions.Condition{Key: o.Key()}.Show()}) if err != nil { return } saying, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() if err := d.teller.PublishSeatEvent(saying, conditions.Seat, link.KeyRolledBack, body); err != nil { fmt.Printf("%s's rollback could NOT be said on the bus: %v\n", module, err) } } // gateObservation is a failed gate as a condition: `core...rolled-back` for the // core (to-be 45 §2), `build...rolled-back` for any other module; `rollback-failed` // when it could not be put back. Urgent for the core and for a build left in place; a warning for a // module put back, which runs what it ran before. func gateObservation(module, component string, machines []string, rollback, why, notWhy, failed, previous string) conditions.Observation { machine := strings.Join(machines, ",") o := conditions.Observation{Scope: conditions.ScopeBuild, ID: module + "." + machine, Kind: kindRolledBack, Token: kindRolledBack, Machine: firstOf(machines), Severity: conditions.Warning, Source: gateProbe, Summary: fmt.Sprintf("%s's build %s failed its gate on %s and was put back to %s: %s", module, short(failed), machine, short(previous), why)} if component != "" { o.Scope, o.ID, o.Severity = conditions.ScopeCore, component+"."+machine, conditions.Urgent } if len(machines) > 1 { o.Also = machines[1:] } if rollback != inventory.RolledBack && rollback != inventory.RollingBack { o.Kind, o.Token, o.Severity = kindRollbackFailed, kindRollbackFailed, conditions.Urgent o.Resolver = conditions.ResolverOperator o.Summary = fmt.Sprintf("%s's build %s failed its gate on %s and was NOT put back: %s — %s", module, short(failed), machine, why, orNone(notWhy)) } return o } // probeGates is DG: no build that failed its gate is left without its condition, and no witness's // rollback goes unsaid. Each module whose newest verdict is a failure keeps its condition; a newer build // that passes clears it. Each core component a machine's witness says it put back is `core.. // .rolled-back`, urgent, until the machine's reports stop saying it. func probeGates(ctx context.Context, d *doctor) ([]conditions.Observation, error) { latest, err := d.open.inventory.LatestGates(ctx) if err != nil { return nil, err } var out []conditions.Observation for _, v := range latest { if v.Verdict != inventory.GateFailed { continue } notWhy := "" if v.Rollback == inventory.NotRolledBack { notWhy = v.Why } out = append(out, gateObservation(v.Module, v.Component, v.Machines, v.Rollback, v.Why, notWhy, v.Commit, v.Previous)) } for node, list := range witnessed.all() { for _, r := range list { out = append(out, witnessObservation(node, r)) } } // A release held after a failed one, while builds still wait for a gate (ADR 0236). out = append(out, backlogObservation()...) return sortedFound(dedupeObservations(out)), nil } // dedupeObservations keeps one observation per key, the first. func dedupeObservations(list []conditions.Observation) []conditions.Observation { seen := map[string]bool{} var out []conditions.Observation for _, o := range list { if seen[o.Key()] { continue } seen[o.Key()] = true out = append(out, o) } return out } // servedOnTheBus asks the bus's discovery what every machine's node tools answer and which modules' // tools are served where. A variable so a test needs no runtime. var servedOnTheBus = func(ctx context.Context, conn *nats.Conn) (map[string]served, error) { inbox := conn.NewRespInbox() sub, err := conn.SubscribeSync(inbox) if err != nil { return nil, err } defer func() { _ = sub.Unsubscribe() }() if err := conn.PublishRequest("$SRV.INFO", inbox, nil); err != nil { return nil, fmt.Errorf("asking the bus who serves what: %w", err) } out := map[string]served{} add := func(node string) served { s, ok := out[node] if !ok { s = served{tools: map[string]bool{}} } return s } deadline := time.Now().Add(discoveryPatience) for time.Now().Before(deadline) { wait, cancel := context.WithTimeout(ctx, discoveryQuiet) msg, err := sub.NextMsgWithContext(wait) cancel() if err != nil { if ctx.Err() != nil { return nil, ctx.Err() } break } var info micro.Info if json.Unmarshal(msg.Data, &info) != nil { continue } if info.Name == broker.RuntimeModule { node := info.Metadata["node"] if node == "" { node = info.ID } s := add(node) s.runtime = true out[node] = s } for _, e := range info.Endpoints { if e.Metadata["kind"] != "tool" || e.Metadata["module"] == "" { continue } node := e.Metadata["node"] if node == "" { node = info.ID } s := add(node) s.tools[e.Metadata["module"]] = true out[node] = s } } return out, nil } // gateLines is a plan's gates as `plans ` says them: the rollout's record. func gateLine(g *inventory.PlanGate) string { if g == nil { return "" } what := "judging" switch g.Verdict { case inventory.GatePassed: what = "passed" case inventory.GateFailed: what = "FAILED" } line := fmt.Sprintf("gate on %s: %s", strings.Join(g.Machines, ", "), what) if g.Component != "" { line += " (core: " + g.Component + ")" } if g.From != "" || g.To != "" { line += fmt.Sprintf(", %s → %s", short(orNone(g.From)), short(g.To)) } if g.Took != "" { line += ", after " + g.Took } if g.Why != "" { line += ": " + g.Why } else if g.Last != "" { line += fmt.Sprintf(" (%d of %d passes; wanting: %s)", g.Passes, gatePasses, g.Last) } else if g.Verdict == "" { line += fmt.Sprintf(" (%d of %d passes)", g.Passes, gatePasses) } if g.Verdict == "" && len(g.Waits) > 0 { line += waitsSaid(g.Waits) } if len(g.Passing) > 0 { line += "; passed on their own and kept: " + strings.Join(g.Passing, ", ") } if g.Rollback != "" { line += "; " + g.Rollback } return line } // witnessObservation is a witness's verdict as a condition (ADR 0236, the host's contract): // `core...`, urgent for rolled-back, not-reversible, restore-failed and // halted, a warning for nothing-to-restore and unwitnessed. func witnessObservation(node string, r lease.Rollback) conditions.Observation { severity := conditions.Warning if lease.Urgent(r.Outcome) { severity = conditions.Urgent } outcome := r.Outcome if outcome == "" { outcome = lease.OutcomeRolledBack } what := map[string]string{ lease.OutcomeRolledBack: "and put back " + short(orNone(r.To)), lease.OutcomeNotReversible: "and left it: it is not reversible", lease.OutcomeNothingToRestore: "and had nothing to put back", lease.OutcomeRestoreFailed: "and could not put the previous build back", lease.OutcomeUnwitnessed: "and could not judge it at all", lease.OutcomeHalted: "and gave up: the build before it fails too", }[outcome] return conditions.Observation{Scope: conditions.ScopeCore, ID: r.Component + "." + node, Kind: kindRolledBack, Token: outcome, Machine: node, Severity: severity, Summary: fmt.Sprintf("the witness on %s judged the %s %s not healthy %s at %s: %s", node, r.Component, short(r.From), what, r.At.UTC().Format(time.RFC3339), r.Why)} }