A VPN client rewrote the laptop's resolver file and nothing said so. The engine now states its machine's networking; the controller keeps it with the machine's health (migration 0077) and raises the rewrite as its own finding naming the writer, the machine's own faults as machine.<m>.network, and what several machines cannot reach once, there. The gate waits on a rewrite it did not make rather than putting back a good build.
989 lines
39 KiB
Go
989 lines
39 KiB
Go
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
|
|
// 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 "<module>@<machine>", the provider its findings are held under (ADR 0240 rule 5).
|
|
heldOn map[string]string
|
|
}
|
|
|
|
// 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 {
|
|
if c.Source == gateProbe || c.Raised.Before(since) {
|
|
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))
|
|
if !aboutIt || c.Source == gateProbe || c.Raised.Before(since) {
|
|
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 == 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
|
|
}
|
|
// **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)
|
|
}
|
|
worst, why := healthGood, ""
|
|
var failing []string
|
|
broken := map[string]bool{}
|
|
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 {
|
|
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 h != healthGood && !slices.Contains(failing, j.module) {
|
|
failing = append(failing, j.module)
|
|
}
|
|
if h == healthBroken {
|
|
broken[j.module] = true
|
|
}
|
|
if h > worst {
|
|
worst, why = h, said
|
|
} else if h == worst && h != healthGood && why == "" {
|
|
why = said
|
|
}
|
|
}
|
|
switch {
|
|
case worst == healthBroken:
|
|
// What broke is put back; what was only not yet healthy beside it is too — they moved together.
|
|
g.Failing = failing
|
|
decide(g, inventory.GateFailed, why, now)
|
|
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 now.Sub(*g.Since) > gateBound {
|
|
decide(g, inventory.GateFailed, fmt.Sprintf("not healthy within %s of its apply: %s", gateBound, why), now)
|
|
}
|
|
default:
|
|
g.Passes++
|
|
g.LastPass, g.Last, g.Failing = &now, "", nil
|
|
if g.Passes >= gatePasses && now.Sub(*g.Since) >= gateSettle {
|
|
decide(g, inventory.GatePassed, fmt.Sprintf("healthy %d times over %s", g.Passes,
|
|
now.Sub(*g.Since).Round(time.Second)), now)
|
|
}
|
|
}
|
|
return g.Verdict, nil
|
|
}
|
|
|
|
// 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: g.Why, 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.<component>.<machine>.rolled-back` for the
|
|
// core (to-be 45 §2), `build.<module>.<machine>.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.<component>.
|
|
// <machine>.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 <id>` 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.Rollback != "" {
|
|
line += "; " + g.Rollback
|
|
}
|
|
return line
|
|
}
|
|
|
|
// witnessObservation is a witness's verdict as a condition (ADR 0236, the host's contract):
|
|
// `core.<component>.<node>.<outcome>`, 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)}
|
|
}
|