From d6e0a8250afb4d6a32c1dc2c698d741eada3ecea Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 6 Oct 2026 22:20:22 +0200 Subject: [PATCH] Send a plan's tier to each machine once, and blame no module for its machine (hq issue 281) --- cmd/mesh-controller/gate.go | 223 ++++++++++++++- cmd/mesh-controller/plan_retry.go | 118 +++++++- cmd/mesh-controller/release.go | 32 ++- cmd/mesh-controller/release_plan.go | 235 ++++++++++++---- cmd/mesh-controller/tier_send_test.go | 375 ++++++++++++++++++++++++++ internal/inventory/plans.go | 4 + 6 files changed, 911 insertions(+), 76 deletions(-) create mode 100644 cmd/mesh-controller/tier_send_test.go diff --git a/cmd/mesh-controller/gate.go b/cmd/mesh-controller/gate.go index 4316ac8..30096a6 100644 --- a/cmd/mesh-controller/gate.go +++ b/cmd/mesh-controller/gate.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "slices" + "sort" "strings" "time" @@ -244,6 +245,86 @@ func judgeHealth(module, component string, m catalogue.Manifest, machine string, 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 +} + +// 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. + 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, @@ -299,11 +380,30 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs 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 { - h, said := judgeHealth(j.module, coreComponent(j.module), shelf[j.module], j.node, *g.Since, facts) + 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 + } + } if h != healthGood && !slices.Contains(failing, j.module) { failing = append(failing, j.module) } @@ -418,6 +518,21 @@ func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module str 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. @@ -450,8 +565,11 @@ func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module str sayRollback(ctx, open, module, g, why) } if state.Previous == "" { - notBack(fmt.Sprintf("%s had not been sent %s before, or what it was sent is not known — `unassign` takes it "+ - "off, or a newer merge replaces it", strings.Join(g.Machines, ", "), module)) + // 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) @@ -465,7 +583,8 @@ func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module str return } if !found { - notBack(fmt.Sprintf("no build of %s from %s is still kept to put back", module, short(state.Previous))) + 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 { @@ -476,6 +595,17 @@ func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module str 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` "+ @@ -490,6 +620,91 @@ func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module str 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"` diff --git a/cmd/mesh-controller/plan_retry.go b/cmd/mesh-controller/plan_retry.go index b5c0002..dc8a2b4 100644 --- a/cmd/mesh-controller/plan_retry.go +++ b/cmd/mesh-controller/plan_retry.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "os" + "slices" "sort" "strings" "time" @@ -181,7 +182,7 @@ func retryRefusal(p inventory.Plan, plans []inventory.Plan) error { // **A build that failed its gate is not sent again** (novox/hq ADR 0236): it was put back on its first // machine, and retrying would judge the build the mesh put back, or send the failed one by hand. for _, m := range stopped { - if g := p.Modules[m].Gate; g != nil && g.Verdict == inventory.GateFailed { + if g := p.Modules[m].Gate; g != nil && g.Verdict == inventory.GateFailed && !noVerdictOnItsBuild(g) { return fmt.Errorf("%s failed its gate on %s (%s) and was put back: a build that failed its gate is not "+ "sent again — a newer merge, or `rebuild %s`, makes a new build, judged at the gate again", m, strings.Join(g.Machines, ", "), g.Why, m) @@ -269,6 +270,9 @@ func retryPlan(ctx context.Context, open *stores, id string) (string, error) { if err := retryRefusal(p, plans); err != nil { return "", err } + if again := unjudgedAtGate(p); len(failedIn(p)) == 0 && len(again) > 0 { + return retryTierWhole(ctx, open, &p, again) + } if len(failedIn(p)) == 0 { return retryRollouts(ctx, open, &p) } @@ -388,17 +392,34 @@ func planHolding(module string, openPlans, recent []inventory.Plan) (inventory.P // records that send as the first anew, and sets the plan rolling: from there it goes on as the plan // would have — the rest sent once those report they applied it, the next tier after (ADR 0218). func retryRollouts(ctx context.Context, open *stores, p *inventory.Plan) (string, error) { - var said []string - for _, m := range stoppedRollouts(*p) { - s := p.Modules[m] - sent, err := sendRollout(ctx, open, s.First) - if err != nil { - return "", fmt.Errorf("%s could not be sent to %s again, so %s stays failed: %w", - m, strings.Join(s.First, ", "), p.ID, err) + // Every machine once, for all the modules stopped there (novox/hq issue 281). + stopped := stoppedRollouts(*p) + var machines []string + for _, m := range stopped { + for _, n := range p.Modules[m].First { + if !slices.Contains(machines, n) { + machines = append(machines, n) + } } - now := time.Now().UTC() - s.First, s.FirstAt, s.Why = sent, &now, "" - said = append(said, m+" to "+strings.Join(sent, ", ")) + } + sort.Strings(machines) + sent, err := sendRollout(ctx, open, machines) + if err != nil { + return "", fmt.Errorf("%s could not be sent to %s again, so %s stays failed: %w", + strings.Join(stopped, ", "), strings.Join(machines, ", "), p.ID, err) + } + now := time.Now().UTC() + var said []string + for _, m := range stopped { + s := p.Modules[m] + var again []string + for _, n := range sent { + if slices.Contains(s.First, n) { + again = append(again, n) + } + } + s.First, s.FirstAt, s.Why = again, &now, "" + said = append(said, m+" to "+strings.Join(again, ", ")) } p.State = inventory.PlanRolling p.Note = fmt.Sprintf("tier %d retried by hand; sent %s first again", p.Tier, strings.Join(said, "; ")) @@ -408,3 +429,78 @@ func retryRollouts(ctx context.Context, open *stores, p *inventory.Plan) (string return fmt.Sprintf("%s retried at tier %d of %d: sent %s first again; the rest follow once it reports it "+ "applied, as the plan would have", p.ID, p.Tier, len(p.Tiers), strings.Join(said, "; ")), nil } + +// noVerdictOnItsBuild says a failed gate said nothing about the module's build (novox/hq issue 281): its +// send changed nothing of the module there — the machine already ran that build, carried there by an +// earlier send of the same tier, or one identical to it. Such a module was blamed for its machine. +func noVerdictOnItsBuild(g *inventory.PlanGate) bool { + return g.Rollback == gateUnchanged || (g.From != "" && sameCommit(g.From, g.To)) +} + +// unjudgedAtGate is the modules of a plan's current tier stopped at a gate that was no verdict on their +// build, sorted. +func unjudgedAtGate(p inventory.Plan) []string { + if p.Tier >= len(p.Tiers) { + return nil + } + var out []string + for _, m := range p.Tiers[p.Tier] { + if s := p.Modules[m]; s != nil && s.Gate != nil && s.Gate.Verdict == inventory.GateFailed && noVerdictOnItsBuild(s.Gate) { + out = append(out, m) + } + } + sort.Strings(out) + return out +} + +// retryTierWhole retries a plan stopped at a gate that judged no build of the module it stopped on +// (issue 281): that module is asked again under a new id — its old build may be marked failed, and a new +// verdict is what takes the gate's condition away — and every module of the tier sent first and never +// passed is sent again, the tier whole, one send per machine, judged again. What passed stays passed. +func retryTierWhole(ctx context.Context, open *stores, p *inventory.Plan, again []string) (string, error) { + inv := open.inventory + entries, err := inv.Catalogued(ctx) + if err != nil { + return "", err + } + byName := map[string]inventory.Entry{} + for _, e := range entries { + byName[e.Manifest.Module] = e + } + var resent []string + for _, m := range p.Tiers[p.Tier] { + s := p.Modules[m] + if s == nil || s.FirstAt == nil || s.SentAt != nil || slices.Contains(again, m) || + (s.Gate != nil && s.Gate.Verdict == inventory.GatePassed) { + continue + } + sendAgain(s) + resent = append(resent, m) + } + var asked []string + for _, m := range again { + sendAgain(p.Modules[m]) + askModule(ctx, p, m, byName) + if s := p.Modules[m]; s.State == "asked" { + asked = append(asked, m+" as "+s.Build) + } + } + if p.State == inventory.PlanFailed && len(failedIn(*p)) == 0 { + p.State = inventory.PlanBuilding + p.Note = fmt.Sprintf("tier %d retried by hand: %s asked again; %s sent again with the tier", p.Tier, + strings.Join(again, ", "), orNone(strings.Join(resent, ", "))) + } + if err := inv.SavePlan(ctx, p); err != nil { + return "", err + } + if p.State != inventory.PlanBuilding { + return "", fmt.Errorf("%s could not be resumed: %s", p.ID, p.Note) + } + return fmt.Sprintf("%s retried at tier %d of %d: asked %s; %s sent again with the tier, one send per machine, "+ + "judged again", p.ID, p.Tier, len(p.Tiers), strings.Join(asked, ", "), orNone(strings.Join(resent, ", "))), nil +} + +// sendAgain forgets a module's first send, so its plan sends it again. +func sendAgain(s *inventory.PlanModule) { + s.First, s.FirstAt, s.Gate, s.GatedBy, s.Previous, s.Why = nil, nil, nil, "", "", "" +} diff --git a/cmd/mesh-controller/release.go b/cmd/mesh-controller/release.go index 579aa0e..2d16ce1 100644 --- a/cmd/mesh-controller/release.go +++ b/cmd/mesh-controller/release.go @@ -255,7 +255,11 @@ func ungatedIn(ctx context.Context, open *stores, names []string, addedHolder st // gatedSend sends one machine everything waiting there, under a gate that judges it all: what moved is // answered, with the build each moved to, for the gate to judge and to put back. Refused while a plan that // has started walks one of those builds elsewhere: that plan sends it here once its gate passed. -func gatedSend(ctx context.Context, open *stores, node string, own *inventory.CarriedMove) ([]inventory.CarriedMove, []string, error) { +// +// owns are the moves the send exists for — every module of a plan's tier whose first machine this is, in +// one send (novox/hq issue 281): a send carries the machine's whole declaration (ADR 0221), so a send per +// module was the same declaration sent again and again, each one setting aside the one before. +func gatedSend(ctx context.Context, open *stores, node string, owns []inventory.CarriedMove) ([]inventory.CarriedMove, []string, error) { inv := open.inventory f, err := readMoveFacts(ctx, inv) if err != nil { @@ -265,8 +269,11 @@ func gatedSend(ctx context.Context, open *stores, node string, own *inventory.Ca if err != nil { return nil, nil, err } + own := func(module string) bool { + return slices.ContainsFunc(owns, func(o inventory.CarriedMove) bool { return o.Module == module }) + } for _, mv := range moves { - if own != nil && mv.Module == own.Module { + if own(mv.Module) { continue } if id := f.walkedBy(mv.Module, node); id != "" { @@ -274,10 +281,16 @@ func gatedSend(ctx context.Context, open *stores, node string, own *inventory.Ca mv.Module, short(mv.To), node, id) } } - if own != nil && !slices.ContainsFunc(moves, func(mv inventory.CarriedMove) bool { return mv.Module == own.Module }) { - moves = append(moves, *own) + for _, o := range owns { + i := slices.IndexFunc(moves, func(mv inventory.CarriedMove) bool { return mv.Module == o.Module }) + switch { + case i < 0: + moves = append(moves, o) + case moves[i].Build == "": + moves[i].Build = o.Build + } } - if len(moves) == 0 && own == nil { + if len(moves) == 0 && len(owns) == 0 { return nil, nil, nil } for i := range moves { @@ -287,6 +300,7 @@ func gatedSend(ctx context.Context, open *stores, node string, own *inventory.Ca } } } + sort.Slice(moves, func(i, j int) bool { return moves[i].Module < moves[j].Module }) sent, err := sendRollout(withScope(ctx, sendScope{judged: map[string]bool{node: true}}), open, []string{node}) if err != nil { return nil, nil, err @@ -328,6 +342,10 @@ func failCarried(ctx context.Context, open *stores, p *inventory.Plan, g *invent machines = []string{c.Node} } state := &inventory.PlanModule{Build: c.Build, Previous: c.From, Commit: c.To} + // A module of the plan's tier sent in the same send (issue 281) keeps its own record of it. + if s := p.Modules[c.Module]; s != nil && except != "" && s.GatedBy == except && s.Build == c.Build { + state = s + } p.Note = "" gateFailed(ctx, open, p, c.Module, state, machines, g.Why) notes = append(notes, p.Note) @@ -563,7 +581,9 @@ func advanceRelease(ctx context.Context, open *stores, p *inventory.Plan) (bool, return true, nil } p.Note = "" - failCarried(ctx, open, p, g, "") + batched, back := batchingRollbacks(ctx) + failCarried(batched, open, p, g, "") + sendRollbacks(ctx, open, p, back) p.Note = fmt.Sprintf("failed its gate on %s: %s — %s", strings.Join(g.Machines, ", "), g.Why, p.Note) fmt.Printf("%s: %s\n", p.ID, p.Note) return true, nil diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 4a03df5..2e3e611 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -610,7 +610,20 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, // minute. Now the first machine is sent, the plan records it and waits for that machine's report // after the send to say it applied what it was sent; only then are the rest sent. A first machine // that fails or refuses stops the module's rollout and the plan with it, the rest untouched. + // + // **One send per machine for the whole tier** (novox/hq issue 281). A send carries the machine's + // whole declaration (ADR 0221), and the plan sent the first machine once per module: on 2026-10-06 a + // tier of seventy modules sent one laptop a new declaration every few seconds, the node-engine set + // aside one for the next, the node tools never settled, and the gate blamed a module for the machine + // being behind. Now every module of the tier whose first machine is the same is sent there in one + // send, judged there by one gate — kept on the first of them, the others gated by it — and the rest + // of the tier's machines are sent once, together, when the gates have passed. var pending []string + firstOn := map[string][]string{} + runningOf := map[string][]string{} + var rest []string + restTo := map[string][]string{} + together := map[string]bool{} for _, m := range tier { state := p.Modules[m] if state == nil || state.SentAt != nil || state.State == planDeleted || !rollsOut(m) { @@ -620,6 +633,7 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, if err != nil { return false, err } + runningOf[m] = running // **A build whose source is unchanged is never a move** (novox/hq issue 280): made from the source // of the build every machine running it was sent, it is registered with that build's artifacts — // nothing to send, and nothing for a gate to judge. @@ -647,25 +661,37 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, } now := time.Now().UTC() step := nextRollout(*state, running, policy.Together, reports, now, planWaitBound) + // **Sent with others, judged with them** (issue 281): the gate of the send that carried it is its + // verdict on its first machine. A failure there stopped the plan already. + if state.GatedBy != "" { + lead := p.Modules[state.GatedBy] + if lead == nil || lead.Gate == nil || lead.Gate.Verdict == "" { + pending = append(pending, fmt.Sprintf("%s judged with %s on %s", m, state.GatedBy, strings.Join(state.First, ", "))) + continue + } + if lead.Gate.Verdict != inventory.GatePassed { + continue + } + passedWith(state, state.GatedBy, lead.Gate) + } switch { case step.failed != "": // The first machine refused or failed what it was sent, or never said: the gate failed, and - // the build is put back there (novox/hq ADR 0236); the rest are left as they were. - gateFailed(ctx, open, p, m, state, firstRunning(state.First, running), step.failed) - if state.Gate != nil && len(state.Gate.Carried) > 0 { - failCarried(ctx, open, p, state.Gate, m) - } - p.Note += fmt.Sprintf("; %s left as it was", orNone(strings.Join(step.rest, ", "))) - fmt.Printf("%s: %s\n", p.ID, p.Note) + // what it carried is put back there (novox/hq ADR 0236); the rest are left as they were. + failFirstSend(ctx, open, p, m, state, firstRunning(state.First, running), step.failed, step.rest) return true, nil case step.waiting != "": - pending = append(pending, fmt.Sprintf("%s on %s, sent first at %s", m, step.waiting, state.FirstAt.Local().Format("15:04"))) + at := "" + if state.FirstAt != nil { + at = state.FirstAt.Local().Format("15:04") + } + pending = append(pending, fmt.Sprintf("%s on %s, sent first at %s", m, step.waiting, at)) continue } // **The gate** (novox/hq ADR 0236, to-be 45 §8): the first machine reported the build applied; // it is judged by its health before anything else is sent — the rest, or, where it is the only // machine, the plan's next step. - if !policy.Together && state.FirstAt != nil && len(state.First) > 0 { + if state.GatedBy == "" && !policy.Together && state.FirstAt != nil && len(state.First) > 0 { verdict, err := judgeGate(ctx, open, p, m, state, running, now) if err != nil { return false, err @@ -676,12 +702,7 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, strings.Join(state.Gate.Machines, ", "), gateLine(state.Gate))) continue case inventory.GateFailed: - gateFailed(ctx, open, p, m, state, state.Gate.Machines, state.Gate.Why) - if len(state.Gate.Carried) > 0 { - failCarried(ctx, open, p, state.Gate, m) - } - p.Note += fmt.Sprintf("; %s left as it was", orNone(strings.Join(step.rest, ", "))) - fmt.Printf("%s: %s\n", p.ID, p.Note) + failFirstSend(ctx, open, p, m, state, state.Gate.Machines, state.Gate.Why, step.rest) return true, nil case inventory.GatePassed: if !state.Gate.Kept { @@ -695,51 +716,66 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, // No machine runs it: nothing to send, and nothing to wait for. state.SentAt = &now continue - } - // Read before the send, which records what it carries. - var before map[string]string - var beforeKnown bool - if step.first { - before, beforeKnown, _ = inv.SentBuilds(ctx, step.send[0]) - } - // What sendToEach answers, not what was asked: the machine holding the bus is sent before - // the first when its user list must change (issue 249), and the plan waits for it too. - var sent []string - var carried []inventory.CarriedMove - switch { case step.first: - // **A gated send** (ADR 0236): everything waiting on the first machine goes with the build, - // and the gate judges all of it there. - own := inventory.CarriedMove{Module: m, Node: step.send[0], From: before[m], To: state.Commit, Build: state.Build} - carried, sent, err = gatedSend(ctx, open, step.send[0], &own) - case policy.Together: - sent, err = sendRollout(withScope(ctx, sendScope{modules: map[string]bool{m: true}}), open, step.send) - default: - // The rest, after the gate passed: nothing else may move with it that no gate has seen. - sent, err = sendRollout(ctx, open, step.send) + firstOn[step.send[0]] = append(firstOn[step.send[0]], m) + continue } - if err != nil { - // Not marked sent, so the next step tries again (issue 249): a grant that could not be - // issued is a send that did not happen. - return false, fmt.Errorf("sending %s to %s after tier %d: %w", m, strings.Join(step.send, ", "), p.Tier, err) + rest = append(rest, m) + restTo[m] = step.send + together[m] = policy.Together + } + // The first machines: one send each for everything of the tier that goes there first, one machine + // per step — the plan is kept between them, so the next machine's send sees what this one walks. + if len(firstOn) > 0 { + machines := make([]string, 0, len(firstOn)) + for n := range firstOn { + machines = append(machines, n) } - if step.first { - // What the first machine ran before: what a failed gate puts back (ADR 0236). - if beforeKnown { - state.Previous = before[m] + sort.Strings(machines) + for _, node := range machines { + err := firstSend(ctx, open, p, node, firstOn[node], runningOf) + if errors.Is(err, errWalkedElsewhere) { + // A build this machine would carry is being judged elsewhere: it goes once that gate passed. + pending = append(pending, fmt.Sprintf("%s for %s: %v", node, strings.Join(firstOn[node], ", "), err)) + continue + } + if err != nil { + // Not marked sent, so the next step tries again (issue 249): a grant that could not be + // issued is a send that did not happen. + return false, fmt.Errorf("sending %s to %s first after tier %d: %w", strings.Join(firstOn[node], ", "), + node, p.Tier, err) } - state.First = sent - state.FirstAt = &now - state.Gate = &inventory.PlanGate{Component: coreComponent(m), Machines: firstRunning(sent, running), - From: state.Previous, To: state.Commit, Since: &now, Carried: carried} - p.State = inventory.PlanRolling - p.Note = fmt.Sprintf("tier %d built; sent %s to %s first", p.Tier, m, strings.Join(sent, ", ")) - fmt.Printf("%s: tier %d built; sent %s to %s first, the rest once it reports it applied\n", - p.ID, p.Tier, m, strings.Join(sent, ", ")) return true, nil } - state.SentAt = &now - fmt.Printf("%s: tier %d built; sent %s to %s\n", p.ID, p.Tier, m, strings.Join(sent, ", ")) + } + // The rest, after their gates passed, and what rolls out together: every machine once. + if len(rest) > 0 { + var machines []string + scope := sendScope{modules: map[string]bool{}} + for _, m := range rest { + for _, n := range restTo[m] { + if !slices.Contains(machines, n) { + machines = append(machines, n) + } + } + if together[m] { + // What its policy rolls out together may move wherever it goes; anything else no gate has + // seen still stops the send. + scope.modules[m] = true + } + } + sort.Strings(machines) + sent, err := sendRollout(withScope(ctx, scope), open, machines) + if err != nil { + return false, fmt.Errorf("sending %s to %s after tier %d: %w", strings.Join(rest, ", "), + strings.Join(machines, ", "), p.Tier, err) + } + now := time.Now().UTC() + for _, m := range rest { + p.Modules[m].SentAt = &now + } + fmt.Printf("%s: tier %d built; sent %s to %s, one send each\n", p.ID, p.Tier, strings.Join(rest, ", "), + strings.Join(sent, ", ")) return true, nil } if len(pending) > 0 { @@ -817,6 +853,91 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, return true, nil } +// firstSend sends one machine, in one send, every module of the plan's tier whose first machine it is +// (novox/hq issue 281), and records the send on each: what that machine ran of it before — read once, +// before the send, so a module the same send carries is never read as already moved — the machines it +// went to, and the gate. The gate is kept on the first of them and judges everything the send moved; +// the others are gated by it. +func firstSend(ctx context.Context, open *stores, p *inventory.Plan, node string, modules []string, + runningOf map[string][]string) error { + inv := open.inventory + before, beforeKnown, _ := inv.SentBuilds(ctx, node) + owns := make([]inventory.CarriedMove, 0, len(modules)) + for _, m := range modules { + s := p.Modules[m] + owns = append(owns, inventory.CarriedMove{Module: m, Node: node, From: before[m], To: s.Commit, Build: s.Build}) + } + // A gated send (ADR 0236): everything waiting on the machine goes with the tier, and the gate judges + // all of it there. + carried, sent, err := gatedSend(ctx, open, node, owns) + if err != nil { + return err + } + now := time.Now().UTC() + lead := modules[0] + for _, m := range modules { + s := p.Modules[m] + // What the first machine ran before: what a failed gate puts back (ADR 0236). + if beforeKnown { + s.Previous = before[m] + } + s.First, s.FirstAt = sent, &now + s.Gate = &inventory.PlanGate{Component: coreComponent(m), Machines: firstRunning(sent, runningOf[m]), + From: s.Previous, To: s.Commit, Since: &now} + s.GatedBy = "" + if m == lead { + s.Gate.Carried = carried + } else { + s.GatedBy = lead + } + } + p.State = inventory.PlanRolling + p.Note = fmt.Sprintf("tier %d built; sent %s to %s first, in one send", p.Tier, strings.Join(modules, ", "), + strings.Join(sent, ", ")) + fmt.Printf("%s: tier %d built; sent %d module(s) to %s first in one send (%s), the rest once its gate passes\n", + p.ID, p.Tier, len(modules), strings.Join(sent, ", "), strings.Join(modules, ", ")) + return nil +} + +// passedWith keeps, on a module sent with others, the verdict of the gate that judged the send: the gate +// on the first of them passed, and its pass was kept for every build it carried (passCarried). +func passedWith(s *inventory.PlanModule, lead string, g *inventory.PlanGate) { + if s.Gate == nil { + s.Gate = &inventory.PlanGate{Machines: g.Machines, From: s.Previous, To: s.Commit, Since: g.Since} + } + if s.Gate.Verdict != "" { + return + } + s.Gate.Verdict, s.Gate.Why, s.Gate.JudgedAt, s.Gate.Took = g.Verdict, "judged with "+lead+": "+g.Why, g.JudgedAt, g.Took + s.Gate.Passes, s.Gate.Kept = g.Passes, true +} + +// failFirstSend stops the plan at a first send whose gate failed, or whose machine refused or failed +// what it was sent: what the gate found wanting is put back there, each once — the module the gate is +// kept on only when it is among them (issue 281: the module whose gate it is was blamed for whatever +// the send carried) — and the machines after it are left as they were. +func failFirstSend(ctx context.Context, open *stores, p *inventory.Plan, module string, state *inventory.PlanModule, + machines []string, why string, rest []string) { + batched, back := batchingRollbacks(ctx) + defer func() { + sendRollbacks(ctx, open, p, back) + p.Note += fmt.Sprintf("; %s left as it was", orNone(strings.Join(rest, ", "))) + fmt.Printf("%s: %s\n", p.ID, p.Note) + }() + g := state.Gate + if g == nil || g.Verdict == "" || len(g.Failing) == 0 || slices.Contains(g.Failing, module) { + gateFailed(batched, open, p, module, state, machines, why) + } else { + state.Why = "not found wanting; stopped with the send that carried it: " + g.Why + p.State = inventory.PlanFailed + p.Note = fmt.Sprintf("the send to %s in tier %d failed its gate: %s; %s was not found wanting and is left as it is", + strings.Join(g.Machines, ", "), p.Tier, g.Why, module) + } + if state.Gate != nil && len(state.Gate.Carried) > 0 { + failCarried(batched, open, p, state.Gate, module) + } +} + // rolloutStep is what a plan does next with one built module's machines (novox/hq issue 249). type rolloutStep struct { // send is the machines to send now; first, whether they are the first machine's send. @@ -1139,7 +1260,11 @@ func plansCommand(ctx context.Context, args []string) error { fmt.Printf(" %-22s %s\n", m, state) // The rollout's record at its gate (novox/hq to-be 45 §8): first machine, from and to, // verdict, time to it, rolled back or not. - if s != nil && s.Gate != nil { + switch { + case s != nil && s.GatedBy != "" && (s.Gate == nil || s.Gate.Verdict == ""): + // Sent with others, judged by the gate of that send (issue 281). + fmt.Printf(" %-22s judged with %s, in the send to %s\n", "", s.GatedBy, strings.Join(s.First, ", ")) + case s != nil && s.Gate != nil: fmt.Printf(" %-22s %s\n", "", gateLine(s.Gate)) } } diff --git a/cmd/mesh-controller/tier_send_test.go b/cmd/mesh-controller/tier_send_test.go new file mode 100644 index 0000000..dc93fbb --- /dev/null +++ b/cmd/mesh-controller/tier_send_test.go @@ -0,0 +1,375 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "reflect" + "strings" + "testing" + "time" + + "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" +) + +// A plan's tier is sent to each machine once, and judged there by one gate (novox/hq issue 281). +// +// On 2026-10-06 a merge rebuilt about seventy modules in one tier, and the plan sent the first machine +// one module at a time — one whole declaration per module, seconds apart. The node-engine set aside one +// declaration for the next, the node tools never settled, the mesh said the machine's core was behind, +// and the gate failed the module it happened to be judging for that machine-level condition; it then +// could not put the module back, having read as "the build it ran before" the very build an earlier send +// of the same tier had already carried there. + +// tierMesh is a mesh with every module of one tier running on anchor and laptop at build c1, a newer +// build c2 of each registered, a plan whose one tier built them all, and every send recorded and +// answered: the machine applies what it is sent and reports it. A judging reads the store's reports, the +// node tools as answering everywhere, and the conditions the test opens. +type tierMesh struct { + open *stores + modules []string + sent [][]string + conds []conditions.Condition + keeper *conditions.Keeper +} + +func aTierMesh(t *testing.T, modules ...string) *tierMesh { + t.Helper() + open := aMesh(t) + ctx := t.Context() + inv := open.inventory + tm := &tierMesh{open: open, modules: modules} + tm.keeper, _ = withConditionsInMemory(t) + was := doctorFrom + doctorFrom = &doctor{open: open, keeper: tm.keeper, teller: &conditions.Told{}} + t.Cleanup(func() { doctorFrom = was }) + + register := func(m, commit string, asked time.Time) { + if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: m, Version: commit}, + inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/" + m, BuiltFrom: commit, + Head: commit, Asked: asked}); err != nil { + t.Fatal(err) + } + } + first := map[string]string{} + for _, m := range modules { + for i, commit := range []string{"c1", "c2"} { + at := time.Now().Add(-2 * time.Hour) + if commit == "c2" { + at = time.Now().Add(-time.Minute) + } + b := inventory.Build{ID: fmt.Sprintf("build-%s-%d", m, i+1), Module: m, Commit: commit, + Repository: "novox/mesh-catalog", Path: "modules/" + m, Asked: at, At: at} + b.Manifest, _ = json.Marshal(catalogue.Manifest{Module: m, Version: commit}) + if err := inv.RecordBuild(ctx, b); err != nil { + t.Fatal(err) + } + } + register(m, "c1", time.Now().Add(-2*time.Hour)) + first[m] = "c1" + for _, n := range []string{"anchor", "laptop"} { + if _, err := inv.Assign(ctx, n, m); err != nil { + t.Fatal(err) + } + } + } + for _, n := range []string{"anchor", "laptop"} { + if err := inv.RecordSent(ctx, nodeID(t, open, n), "d-"+n+"-c1", first); err != nil { + t.Fatal(err) + } + } + for _, m := range modules { + register(m, "c2", time.Now().Add(-time.Minute)) + } + + n := 0 + wasSend := sendRollout + sendRollout = func(ctx context.Context, open *stores, names []string) ([]string, error) { + tm.sent = append(tm.sent, append([]string(nil), names...)) + current, err := open.inventory.CurrentBuilds(ctx) + if err != nil { + return nil, err + } + builds := map[string]string{} + for _, m := range tm.modules { + builds[m] = current[m].Commit + } + for _, node := range names { + n++ + digest := fmt.Sprintf("d-%s-%d", node, n) + if err := open.inventory.RecordSent(ctx, nodeID(t, open, node), digest, builds); err != nil { + return nil, err + } + if _, err := open.inventory.RecordDoing(ctx, nodeID(t, open, node), inventory.Doing{Node: node, + Outcome: inventory.OutcomeApplied, Declared: digest, Applied: 1, At: time.Now()}); err != nil { + return nil, err + } + } + return names, nil + } + t.Cleanup(func() { sendRollout = wasSend }) + + wasGather := gatherGateFacts + gatherGateFacts = func(ctx context.Context, open *stores, component string) (gateFacts, error) { + f := gateFacts{now: time.Now(), reports: map[string]inventory.Reported{}, engines: map[string]string{}, + rolledBack: map[string][]lease.Rollback{}, served: map[string]served{}, judged: true, + open: append([]conditions.Condition(nil), tm.conds...)} + reports, err := open.inventory.LastReports(ctx) + if err != nil { + return f, err + } + for _, r := range reports { + f.reports[r.Node] = r + f.served[r.Node] = served{runtime: true, tools: map[string]bool{}} + } + return f, nil + } + t.Cleanup(func() { gatherGateFacts = wasGather }) + + wasSettle, wasEvery, wasBound := gateSettle, gateEvery, gateBound + // One judging per advance until a test says otherwise: a pass waits gateEvery for the next. + gateSettle, gateEvery = 0, time.Hour + t.Cleanup(func() { gateSettle, gateEvery, gateBound = wasSettle, wasEvery, wasBound }) + + built := time.Now().UTC() + states := map[string]*inventory.PlanModule{} + for _, m := range modules { + states[m] = &inventory.PlanModule{State: "built", BuiltAt: &built, Commit: "c2", Build: "build-" + m + "-2"} + } + plan := inventory.Plan{ID: "plan-tier", Repository: "novox/mesh-catalog", Branch: "main", Commit: "c2", + Created: built, State: inventory.PlanBuilding, Tiers: [][]string{modules}, Modules: states} + if err := inv.SavePlan(ctx, &plan); err != nil { + t.Fatal(err) + } + return tm +} + +func (tm *tierMesh) plan(t *testing.T) inventory.Plan { + t.Helper() + p, err := tm.open.inventory.PlanByID(t.Context(), "plan-tier") + if err != nil { + t.Fatal(err) + } + return p +} + +// raise opens a condition about a machine itself, raised now: after the send. +func (tm *tierMesh) raise(machine, kind, summary string) { + tm.conds = append(tm.conds, conditions.Condition{Key: "machine." + machine + "." + kind, Kind: kind, + Subject: conditions.Subject{Scope: conditions.ScopeMachine, ID: machine, Machine: machine}, + Raised: time.Now().Add(time.Second), Summary: summary, Source: "D10"}) +} + +// The incident's tier, small: a dozen modules whose first machine is the same are sent there in one +// send, judged by one gate, and the rest of their machines are sent once — two sends, not two dozen. +func TestATierOfManyModulesIsSentToEachMachineOnce(t *testing.T) { + var modules []string + for i := 1; i <= 12; i++ { + modules = append(modules, fmt.Sprintf("m%02d", i)) + } + tm := aTierMesh(t, modules...) + ctx := t.Context() + + advancePlans(ctx, tm.open) + if !reflect.DeepEqual(tm.sent, [][]string{{"anchor"}}) { + t.Fatalf("sent %v: the first machine once, for the whole tier", tm.sent) + } + p := tm.plan(t) + lead := p.Modules["m01"] + if lead.Gate == nil || len(lead.Gate.Carried) != len(modules) || lead.GatedBy != "" { + t.Fatalf("the gate on the first module does not judge the whole send: %+v", lead.Gate) + } + for _, m := range modules { + s := p.Modules[m] + if s.FirstAt == nil || !reflect.DeepEqual(s.First, []string{"anchor"}) || s.Previous != "c1" { + t.Fatalf("%s's first send is not recorded, or what anchor ran before is not c1: %+v", m, s) + } + if m != "m01" && s.GatedBy != "m01" { + t.Fatalf("%s is not gated by the send's gate: %+v", m, s) + } + } + + gateEvery = 0 + for i := 0; i < 6; i++ { + advancePlans(ctx, tm.open) + } + p = tm.plan(t) + if !reflect.DeepEqual(tm.sent, [][]string{{"anchor"}, {"laptop"}}) { + t.Fatalf("sent %v: one send per machine for the tier", tm.sent) + } + if p.State != inventory.PlanDone { + t.Fatalf("the plan is %s: %s", p.State, p.Note) + } + for _, m := range modules { + if g := p.Modules[m].Gate; g == nil || g.Verdict != inventory.GatePassed { + t.Fatalf("%s's gate: %+v", m, g) + } + if v, found, err := tm.open.inventory.GateOf(ctx, "build-"+m+"-2"); err != nil || !found || v.Verdict != inventory.GatePassed { + t.Fatalf("%s's pass was not kept: %+v %v %v", m, v, found, err) + } + } +} + +// The core behind about the node tools, while the node tools are themselves among what was sent and +// settling, is the node tools' — inside their settle window, the gate's bound — and no failure of the +// modules sent with them: at the bound only the node tools are put back. +func TestACoreBehindWhileTheNodeToolsSettleFailsNoOtherModule(t *testing.T) { + tm := aTierMesh(t, "app1", "app2", broker.RuntimeModule) + ctx := t.Context() + inv := tm.open.inventory + + advancePlans(ctx, tm.open) + if len(tm.sent) != 1 { + t.Fatalf("sent %v", tm.sent) + } + tm.raise("anchor", kindCoreBehind, "anchor runs the node tools it was last sent, not yet reported applied, "+ + "and no plan is rolling them out") + gateEvery = 0 + advancePlans(ctx, tm.open) + g := tm.plan(t).Modules["app1"].Gate + if g.Verdict != "" || !reflect.DeepEqual(g.Failing, []string{broker.RuntimeModule}) { + t.Fatalf("inside the settle window the gate is %+v: only the node tools wait on it", g) + } + + gateBound = -time.Second + advancePlans(ctx, tm.open) + p := tm.plan(t) + g = p.Modules["app1"].Gate + if p.State != inventory.PlanFailed || g.Verdict != inventory.GateFailed || + !reflect.DeepEqual(g.Failing, []string{broker.RuntimeModule}) { + t.Fatalf("the plan is %s (%s), the gate %+v", p.State, p.Note, g) + } + current, _ := inv.CurrentBuilds(ctx) + if current[broker.RuntimeModule].Commit != "c1" { + t.Fatalf("the node tools were not put back: %s", current[broker.RuntimeModule].Commit) + } + for _, m := range []string{"app1", "app2"} { + if current[m].Commit != "c2" { + t.Fatalf("%s was put back for the node tools' condition: %s", m, current[m].Commit) + } + if failed, _ := inv.GateFailed(ctx, "build-"+m+"-2"); failed { + t.Fatalf("%s's build was marked failed for the node tools' condition", m) + } + } + if !strings.Contains(p.Note, "app1 was not found wanting") { + t.Fatalf("the plan blames the module its gate was kept on: %s", p.Note) + } +} + +// A machine-level condition about anything else holds the machine back as a whole: everything the send +// moved there fails together, once, with that reason — not one module for it. +func TestAMachineUnhealthyAsAWholeFailsTheWholeSendTogether(t *testing.T) { + tm := aTierMesh(t, "app1", "app2", "app3") + ctx := t.Context() + inv := tm.open.inventory + advancePlans(ctx, tm.open) + tm.raise("anchor", "disk-full", "anchor's root file system is full") + gateEvery, gateBound = 0, -time.Second + advancePlans(ctx, tm.open) + p := tm.plan(t) + g := p.Modules["app1"].Gate + if p.State != inventory.PlanFailed || !strings.Contains(g.Why, "anchor as a whole") || + !reflect.DeepEqual(g.Failing, []string{"app1", "app2", "app3"}) { + t.Fatalf("the plan is %s (%s), the gate %+v", p.State, p.Note, g) + } + current, _ := inv.CurrentBuilds(ctx) + for _, m := range []string{"app1", "app2", "app3"} { + if current[m].Commit != "c1" { + t.Fatalf("%s was not put back with the send: %s", m, current[m].Commit) + } + } + if !reflect.DeepEqual(tm.sent, [][]string{{"anchor"}, {"anchor"}}) { + t.Fatalf("sent %v: the first send and one send putting all of it back, and nothing to laptop", tm.sent) + } + if !strings.Contains(p.Note, "put back app1 to c1, app2 to c1, app3 to c1 on anchor, in one send") { + t.Fatalf("the plan does not say what was put back: %s", p.Note) + } +} + +// The machine-level conditions, sorted: one naming a module is that module's; the core behind is the +// core's when the send moved it and nobody's when it did not; anything else is the machine's as a whole. +func TestTheConditionsAboutAMachineItself(t *testing.T) { + since := time.Now().Add(-time.Minute) + about := func(kind, id string) conditions.Condition { + return conditions.Condition{Key: "machine." + id + "." + kind, Kind: kind, Raised: time.Now(), + Subject: conditions.Subject{Scope: conditions.ScopeMachine, ID: id, Machine: "anchor"}} + } + f := gateFacts{judged: true, open: []conditions.Condition{about(kindCoreBehind, "anchor")}} + if w := aboutTheMachine("anchor", []string{"app", "blueman"}, since, f); w.whole != "" || len(w.on) != 0 || len(w.facts.open) != 0 { + t.Errorf("the core behind, with no core component sent, counted against what was sent: %+v", w) + } + if w := aboutTheMachine("anchor", []string{"app", broker.RuntimeModule}, since, f); w.whole != "" || + w.on[broker.RuntimeModule] == "" || w.on["app"] != "" { + t.Errorf("the core behind, with the node tools sent, is not theirs alone: %+v", w) + } + f.open = []conditions.Condition{about("x", "anchor.app")} + if w := aboutTheMachine("anchor", []string{"app", "other"}, since, f); w.on["app"] == "" || w.on["other"] != "" || w.whole != "" { + t.Errorf("a condition naming a module is not that module's: %+v", w) + } + f.open = []conditions.Condition{about("unreachable", "anchor")} + if w := aboutTheMachine("anchor", []string{"app"}, since, f); !strings.Contains(w.whole, "anchor as a whole") { + t.Errorf("a condition about the machine is not the machine's as a whole: %+v", w) + } + f.open[0].Raised = since.Add(-time.Hour) + if w := aboutTheMachine("anchor", []string{"app"}, since, f); w.whole != "" || len(w.facts.open) != 1 { + t.Errorf("a condition older than the send counted: %+v", w) + } +} + +// A module its first send changed nothing of — an earlier send already carried this very build there — +// fails a gate as no verdict on its build: left as it was, never marked failed, nothing said as a +// rollback that failed; and `plans retry` asks it again rather than refusing a build that was never +// judged (the incident's "NOT put back: no build kept"). +func TestAModuleItsSendChangedNothingOfIsLeftAsItWasAndRetried(t *testing.T) { + g := aGateMesh(t) + ctx := t.Context() + inv := g.open.inventory + // An earlier send carried c2 to anchor before the plan's own. + if err := inv.RecordSent(ctx, nodeID(t, g.open, "anchor"), "d-anchor-early", map[string]string{"app": "c2"}); err != nil { + t.Fatal(err) + } + advancePlans(ctx, g.open) + if p := g.plan(t); p.Modules["app"].Previous != "c2" { + t.Fatalf("what anchor ran before the send: %+v", p.Modules["app"]) + } + g.health["anchor"] = healthNotYet + gateEvery, gateBound = 0, -time.Second + advancePlans(ctx, g.open) + + p := g.plan(t) + gate := p.Modules["app"].Gate + if p.State != inventory.PlanFailed || gate.Verdict != inventory.GateFailed || gate.Rollback != gateUnchanged || + !strings.Contains(p.Note, "app was left as it was") { + t.Fatalf("the plan is %s (%s), the gate %+v", p.State, p.Note, gate) + } + if failed, _ := inv.GateFailed(ctx, "build-2"); failed { + t.Fatal("a build its send changed nothing of was marked failed") + } + if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c2" { + t.Fatalf("the module was moved: %s", current["app"].Commit) + } + open, _ := g.keeper.Open(ctx) + for _, c := range open { + if c.Kind == kindRollbackFailed || c.Kind == kindRolledBack { + t.Fatalf("said as a rollback: %+v", c) + } + } + + was := askABuild + askABuild = func(context.Context, buildSource, string, string) (string, error) { return "build-3", nil } + t.Cleanup(func() { askABuild = was }) + g.health["anchor"] = healthGood + said, err := retryPlan(ctx, g.open, "plan-gate") + if err != nil { + t.Fatalf("a plan stopped at a gate that judged no build was not retried: %v", err) + } + p = g.plan(t) + s := p.Modules["app"] + if p.State != inventory.PlanBuilding || s.State != "asked" || s.Build != "build-3" || s.FirstAt != nil || s.Gate != nil { + t.Fatalf("retried as %q: the plan is %s, app %+v", said, p.State, s) + } +} diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index cbb37f4..0e0e758 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -78,6 +78,10 @@ type PlanModule struct { // Gate is the new build's judging on its first machine (novox/hq ADR 0236, to-be 45 §8), kept so a // controller replaced mid-judging resumes it, and read back through `plans` as the rollout's record. Gate *PlanGate `json:"gate,omitempty"` + // GatedBy names the module of the same tier whose gate judges this one on its first machine: they + // went there in one send (novox/hq issue 281), and one gate judges what one send moved. Empty for + // the module the gate is kept on, and for a plan from before tiers were sent whole. + GatedBy string `json:"gated_by,omitempty"` } // PlanGate is one module's rollout record at its gate (to-be 45 §8): the component, the first machine,