diff --git a/cmd/mesh-controller/bus_step.go b/cmd/mesh-controller/bus_step.go index e13854d..94ee7c2 100644 --- a/cmd/mesh-controller/bus_step.go +++ b/cmd/mesh-controller/bus_step.go @@ -218,7 +218,7 @@ func busCommand(ctx context.Context, args []string) error { fmt.Printf("bus upgrade %d: %s %s → %s on %s; streams snapshotted at %s; %s\n", step.ID, b.module, step.From, short(b.to), strings.Join(moving, ", "), where, map[bool]string{true: "reversible: putting the old build back undoes it", false: "NOT reversible: the snapshot is the only way back"}[*reversible]) - sent, err := sendRollout(withBusStep(ctx), open, moving) + sent, err := sendRollout(withBusStep(withScope(ctx, sendScope{person: true})), open, moving) if err != nil { _ = inv.EndBusStep(ctx, step.ID, "failed", "the send was refused: "+err.Error()) return fmt.Errorf("the bus's machine could not be sent its new build: %w — nothing was replaced", err) diff --git a/cmd/mesh-controller/gate.go b/cmd/mesh-controller/gate.go index 6fa4936..4316ac8 100644 --- a/cmd/mesh-controller/gate.go +++ b/cmd/mesh-controller/gate.go @@ -139,7 +139,7 @@ var gatherGateFacts = func(ctx context.Context, open *stores, component string) } else { f.servedErr = errors.New("this process does not serve the mesh, so it cannot ask the bus who serves what") } - if component == lease.ComponentController { + if theLease != nil { h, found, err := theLease.holder(ctx) switch { case err != nil: @@ -147,6 +147,9 @@ var gatherGateFacts = func(ctx context.Context, open *stores, component string) 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 } @@ -263,12 +266,31 @@ func judgeGate(ctx context.Context, open *stores, p *inventory.Plan, module stri 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 @@ -278,8 +300,16 @@ func judgeGate(ctx context.Context, open *stores, p *inventory.Plan, module stri return "", err } worst, why := healthGood, "" - for _, n := range g.Machines { - h, said := judgeHealth(module, g.Component, shelf[module], n, *g.Since, facts) + 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) + 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 == "" { @@ -288,15 +318,17 @@ func judgeGate(ctx context.Context, open *stores, p *inventory.Plan, module stri } 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 == healthNotYet: - g.Passes, g.LastPass, g.Last = 0, nil, why + 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 = &now, "" + 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) @@ -320,6 +352,7 @@ func gatePassed(ctx context.Context, open *stores, p *inventory.Plan, module str 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 { @@ -443,7 +476,7 @@ 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) } - sent, err := sendRollout(ctx, open, g.Machines) + 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])) @@ -553,6 +586,8 @@ func probeGates(ctx context.Context, d *doctor) ([]conditions.Observation, error 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 } diff --git a/cmd/mesh-controller/healers.go b/cmd/mesh-controller/healers.go index e77eb2d..a7471a7 100644 --- a/cmd/mesh-controller/healers.go +++ b/cmd/mesh-controller/healers.go @@ -614,6 +614,9 @@ func lastReportOf(ctx context.Context, inv *inventory.Inventory, node string) (i // planStale is why a plan's wait is superseded or finished, and the state closing it leaves it in; empty // when it is neither, which is not H2's to repair. func planStale(ctx context.Context, inv *inventory.Inventory, p inventory.Plan) (state, why string, err error) { + if p.Release != nil { + return "", "", nil // a release plan walks machines, and its own gate says when it is done (ADR 0236) + } recent, err := inv.RecentPlans(ctx, 50) if err != nil { return "", "", err diff --git a/cmd/mesh-controller/held_back.go b/cmd/mesh-controller/held_back.go index 22da1e3..fc22083 100644 --- a/cmd/mesh-controller/held_back.go +++ b/cmd/mesh-controller/held_back.go @@ -120,6 +120,10 @@ func heldMachines(ctx context.Context, open *stores, names []string) (map[string if err != nil { return nil, err } + f, err := readMoveFacts(ctx, inv) + if err != nil { + return nil, err + } for _, node := range names { plan, _, err := planFor(ctx, open, node) if err != nil { @@ -133,7 +137,25 @@ func heldMachines(ctx context.Context, open *stores, names []string) (map[string if err != nil { return nil, err } - if why := heldBack(node, modules, sent, known, current, plans); len(why) > 0 { + // A rebuild that put the same thing on the machine is no move (ADR 0236): read as the build it + // runs. + same := map[string]string{} + for m, was := range sent { + same[m] = was + if f.identical(m, was, current[m].Commit) { + same[m] = current[m].Commit + } + } + why := heldBack(node, modules, same, known, current, plans) + // And a build no gate has seen is not carried by a send that does not judge it (ADR 0236). + for _, mv := range f.moves(node, modules, same, known, false) { + if planStillToSend(plans, mv.Module, node) == "" { + why = append(why, fmt.Sprintf("%s would move from %s to %s, which has passed no gate yet — a release "+ + "plan sends it, one machine at a time, judged", mv.Module, buildName(mv.From), buildName(mv.To))) + } + } + if len(why) > 0 { + sort.Strings(why) out[node] = why } } diff --git a/cmd/mesh-controller/held_back_test.go b/cmd/mesh-controller/held_back_test.go index bccf4ec..0507b64 100644 --- a/cmd/mesh-controller/held_back_test.go +++ b/cmd/mesh-controller/held_back_test.go @@ -264,7 +264,8 @@ func TestANamedPushLeavesAMachineAPolicyHoldsBack(t *testing.T) { t.Fatalf("the push did not say both:\n%s", said.String()) } - // A policy that rolls out: the laptop is a consequence like any other, and sent. + // A policy that rolls out, and a build no gate has seen: still held — a cascade does not judge it + // (novox/hq ADR 0236). if err := inv.SetUpgradeOf(ctx, "resolver", inventory.Upgrade{RollOut: true}); err != nil { t.Fatal(err) } @@ -273,6 +274,19 @@ func TestANamedPushLeavesAMachineAPolicyHoldsBack(t *testing.T) { compose, d, "", &said); err != nil { t.Fatal(err) } + if len(d.declared) != 0 || !strings.Contains(said.String(), "has passed no gate yet") { + t.Fatalf("a cascade carried a build no gate has seen: %v\n%s", d.declared, said.String()) + } + // Once it passed a gate on some machine, the laptop is a consequence like any other, and sent. + if err := inv.RecordGate(ctx, inventory.GateVerdict{Build: "build-c2", Module: "resolver", Commit: "c2c2c2c2c2", + Machines: []string{"anchor"}, Verdict: inventory.GatePassed}); err != nil { + t.Fatal(err) + } + said.Reset() + if _, err := flushBehind(ctx, open, mustNodes(t, open), map[string]bool{"anchor": true, "spare": true}, + compose, d, "", &said); err != nil { + t.Fatal(err) + } if !reflect.DeepEqual(d.declared, []string{"laptop"}) || digestOfLaptop() == before { t.Fatalf("a rolled-out upgrade's machine was not sent: %v\n%s", d.declared, said.String()) } diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index f8aa3f4..fba0f7b 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -1,6 +1,7 @@ package main import ( + "slices" "context" "crypto/sha256" "encoding/hex" @@ -466,6 +467,28 @@ func pushCommand(ctx context.Context, args []string) error { which, len(asked), strings.Join(asked, ", ")) } + // **A whole-mesh push sends no build a gate has not seen** (novox/hq ADR 0236): a machine where one + // waits is left, named, for the release plan — `push ` still sends one machine by a person's word. + if len(args) == 0 { + f, err := readMoveFacts(ctx, inv) + if err != nil { + return err + } + var kept []string + for _, n := range asked { + moves, err := machineMoves(ctx, open, f, n, false) + if err != nil { + return err + } + if len(moves) > 0 { + fmt.Printf(" %s is left: %d build(s) wait there for a gate, which a release plan sends one machine "+ + "at a time (`upgrade backlog` lists them; `push %s` sends it by name)\n", n, len(moves), n) + continue + } + kept = append(kept, n) + } + asked = kept + } // **The bus is replaced only as a planned step** (novox/hq ADR 0236, to-be 45 §8): a machine whose bus // would move is not sent by a push — named, it is refused; otherwise it is left and said. busKept, err := busHeld(ctx, inv, asked) @@ -1030,7 +1053,16 @@ func sendToEach(ctx context.Context, open *stores, names []string) ([]string, er if err != nil { return nil, err } + added := "" + if behind && !slices.Contains(names, holder) { + added = holder + } names = brokerFirst(names, holder, behind) + // **No build moves on a machine without a gate** (novox/hq ADR 0236): a send outside its scope that + // would carry one is refused, said with what waits; the bus's machine added for its user list is left. + if names, err = ungatedIn(ctx, open, names, added); err != nil { + return nil, err + } // Held from composing to sending (novox/hq ADR 0100); a caller that holds them already — // converge, which flips the node and then sends it — is not made to wait on itself. diff --git a/cmd/mesh-controller/release.go b/cmd/mesh-controller/release.go new file mode 100644 index 0000000..226bad1 --- /dev/null +++ b/cmd/mesh-controller/release.go @@ -0,0 +1,603 @@ +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "slices" + "sort" + "strings" + "time" + + "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/link" +) + +// No build reaches a machine without a gate (novox/hq ADR 0236). +// +// **A send carries the machine's whole declaration.** A plan sending one module to its first machine +// sends every other module whose build moved there too; a cascade, a healer's resend and a whole-mesh +// push do the same. Under the old default (`record`) builds were registered and sent nowhere, so on the +// day the default became `roll` every machine was behind on dozens of builds no gate had seen — and the +// next send of anything would have restarted all of them at once, on every machine. +// +// So **a build moves on a machine only through a send that judges it there**, or a person's: +// +// - a gated send — a plan's first machine, a release plan's machine — carries every move waiting on that +// machine, and its gate judges each of them; a pass is each build's verdict, a failure puts back what +// failed; +// - a module a send exists for may move anywhere it is sent: a policy of *together*, a rollback, a build +// whose gate already passed; +// - a send naming a machine by a person (`push `, the bus step) carries what it carries; +// - every other send — a plan's "rest", a cascade, a healer's, the bus's user list carried — is refused, +// or leaves the machine, while a move there waits for a gate. A rebuild that made the same artifacts +// from the same manifest is no move. +// +// **The release plan** is what walks the waiting moves through: whenever moves wait for a gate that no +// started plan is walking, the mesh opens one — every such machine, one at a time, the control node last, +// each sent and judged before the next. A release plan that fails stops the next one opening on its own: +// a person releases it again (`upgrade release-backlog --why`), after the condition says why. + +// sendScope is what a send may carry that no gate has seen. +type sendScope struct { + // judged are the machines whose every move this send's gate judges. + judged map[string]bool + // modules are those a send exists for, which may move wherever it goes. + modules map[string]bool + // person is a person's act: it carries what it carries. + person bool +} + +type sendScopeKey struct{} + +func withScope(ctx context.Context, s sendScope) context.Context { + return context.WithValue(ctx, sendScopeKey{}, s) +} + +func scopeOf(ctx context.Context) sendScope { + s, _ := ctx.Value(sendScopeKey{}).(sendScope) + return s +} + +// errUngated is a send refused because it would carry a build no gate has seen. +var errUngated = errors.New("a build waits there for its gate") + +// errWalkedElsewhere is a gated send refused because a plan that has started walks a move there. +var errWalkedElsewhere = errors.New("a plan already walking a build there sends it") + +// moveFacts is what tells a move from a rebuild, and a gated build from one no gate has seen. +type moveFacts struct { + current map[string]inventory.CurrentBuild + fps map[string]map[string]string + passed map[string]map[string]bool + plans []inventory.Plan +} + +func readMoveFacts(ctx context.Context, inv *inventory.Inventory) (moveFacts, error) { + var f moveFacts + var err error + if f.current, err = inv.CurrentBuilds(ctx); err != nil { + return f, err + } + if f.fps, err = inv.Fingerprints(ctx); err != nil { + return f, err + } + if f.passed, err = inv.PassedCommits(ctx); err != nil { + return f, err + } + f.plans, err = inv.OpenPlans(ctx) + return f, err +} + +// identical is whether two builds of a module put the same thing on a machine: the same commit, or +// builds that made the same artifacts from the same manifest. +func (f moveFacts) identical(module, a, b string) bool { + if a == b || sameCommit(a, b) { + return true + } + fa := f.fps[module][a] + return fa != "" && fa == f.fps[module][b] +} + +// gated is whether a build of a module from this commit — or one identical to it — passed a gate. +func (f moveFacts) gated(module, commit string) bool { + for c := range f.passed[module] { + if f.identical(module, c, commit) { + return true + } + } + return false +} + +// moves is what a machine's next send would move that no gate has seen: modules whose policy rolls out +// (a recorded one is a person's push), that the machine was sent before, moving to a build not identical +// to the one it runs and that has passed no gate. A machine whose last send's builds are not known names +// none: the cascade holds it whole (ADR 0221). +// +// With all, a move to a build that passed a gate elsewhere counts too: what a gated send carries. +func (f moveFacts) moves(node string, modules []string, sent map[string]string, known, all bool) []inventory.CarriedMove { + if !known { + return nil + } + var out []inventory.CarriedMove + for _, m := range modules { + now := f.current[m] + was, carried := sent[m] + if !now.RollOut || !carried || f.identical(m, was, now.Commit) || (!all && f.gated(m, now.Commit)) { + continue + } + out = append(out, inventory.CarriedMove{Module: m, Node: node, From: was, To: now.Commit}) + } + sort.Slice(out, func(i, j int) bool { return out[i].Module < out[j].Module }) + return out +} + +// walkedBy is the open plan that has started walking a module's build — sent it to a first machine, +// not yet passed — other than to this machine; empty when none does. +func (f moveFacts) walkedBy(module, node string) string { + for _, p := range f.plans { + s, holds := p.Modules[module] + if !p.Open() || !holds || s == nil || s.FirstAt == nil || s.SentAt != nil || slices.Contains(s.First, node) { + continue + } + if s.Gate != nil && s.Gate.Verdict == inventory.GatePassed { + continue + } + return p.ID + } + return "" +} + +// machineMoves is the moves no gate has seen on one machine, read from what it would be sent now. +var machineMoves = func(ctx context.Context, open *stores, f moveFacts, node string, all bool) ([]inventory.CarriedMove, error) { + plan, _, err := planFor(ctx, open, node) + if err != nil { + return nil, nil // it cannot be worked out: the send says why + } + modules := make([]string, 0, len(plan.Modules)) + for _, m := range plan.Modules { + modules = append(modules, m.Module) + } + sent, known, err := open.inventory.SentBuilds(ctx, node) + if err != nil { + return nil, err + } + return f.moves(node, modules, sent, known, all), nil +} + +// ungatedIn refuses a send whose machines would move a build no gate has seen, outside its scope: the +// machines named, and why; the machine holding the bus, added only for its user list, is left instead. +func ungatedIn(ctx context.Context, open *stores, names []string, addedHolder string) ([]string, error) { + scope := scopeOf(ctx) + if scope.person { + return names, nil + } + f, err := readMoveFacts(ctx, open.inventory) + if err != nil { + return nil, err + } + var kept, refused []string + for _, n := range names { + moves, err := machineMoves(ctx, open, f, n, false) + if err != nil { + return nil, err + } + var waiting []string + for _, mv := range moves { + if !scope.judged[n] && !scope.modules[mv.Module] { + waiting = append(waiting, fmt.Sprintf("%s %s → %s", mv.Module, short(mv.From), short(mv.To))) + } + } + switch { + case len(waiting) == 0: + kept = append(kept, n) + case n == addedHolder: + fmt.Printf("%s holds the bus and its user list changed, and it is not sent now: %s wait there for a gate "+ + "— the bus may refuse what was newly granted until it is\n", n, strings.Join(waiting, ", ")) + default: + refused = append(refused, fmt.Sprintf("%s (%s)", n, strings.Join(waiting, ", "))) + } + } + if len(refused) > 0 { + return nil, fmt.Errorf("%w: %s — a release plan sends them, one machine at a time, each judged (novox/hq ADR "+ + "0236); `upgrade backlog` lists them", errUngated, strings.Join(refused, "; ")) + } + return kept, nil +} + +// 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) { + inv := open.inventory + f, err := readMoveFacts(ctx, inv) + if err != nil { + return nil, nil, err + } + moves, err := machineMoves(ctx, open, f, node, true) + if err != nil { + return nil, nil, err + } + for _, mv := range moves { + if own != nil && mv.Module == own.Module { + continue + } + if id := f.walkedBy(mv.Module, node); id != "" { + return nil, nil, fmt.Errorf("%w: %s's build %s waits on %s, which %s is walking", errWalkedElsewhere, + 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) + } + if len(moves) == 0 && own == nil { + return nil, nil, nil + } + for i := range moves { + if moves[i].Build == "" { + if moves[i].Build, err = inv.BuildOf(ctx, moves[i].Module, moves[i].To); err != nil { + return nil, nil, err + } + } + } + sent, err := sendRollout(withScope(ctx, sendScope{judged: map[string]bool{node: true}}), open, []string{node}) + if err != nil { + return nil, nil, err + } + return moves, sent, nil +} + +// passCarried keeps a pass as the verdict of every build the gate judged beside its own module. +func passCarried(ctx context.Context, open *stores, p *inventory.Plan, g *inventory.PlanGate, except string) { + seen := map[string]bool{} + for _, c := range g.Carried { + if c.Module == except || c.Build == "" || seen[c.Build] { + continue + } + seen[c.Build] = true + if err := open.inventory.RecordGate(ctx, inventory.GateVerdict{Build: c.Build, Module: c.Module, Commit: c.To, + Previous: c.From, Plan: p.ID, Machines: []string{c.Node}, Verdict: inventory.GatePassed, Why: g.Why, + Component: coreComponent(c.Module), JudgingFrom: g.Since}); err != nil { + fmt.Printf("%s: %s passed its gate on %s, and the verdict could not be kept: %v\n", p.ID, c.Module, c.Node, err) + } + } +} + +// failCarried puts back every build the gate carried that it found wanting, on the machines that were +// sent it, once each. +func failCarried(ctx context.Context, open *stores, p *inventory.Plan, g *inventory.PlanGate, except string) { + var notes []string + if p.Note != "" { + notes = append(notes, p.Note) + } + done := map[string]bool{except: true} + for _, c := range g.Carried { + if done[c.Module] || (len(g.Failing) > 0 && !slices.Contains(g.Failing, c.Module)) { + continue + } + done[c.Module] = true + machines, err := sentTheBuild(ctx, open, c.Module, c.To) + if err != nil || len(machines) == 0 { + machines = []string{c.Node} + } + state := &inventory.PlanModule{Build: c.Build, Previous: c.From, Commit: c.To} + p.Note = "" + gateFailed(ctx, open, p, c.Module, state, machines, g.Why) + notes = append(notes, p.Note) + } + p.State = inventory.PlanFailed + p.Note = strings.Join(notes, "; ") +} + +// sentTheBuild is every machine running a module that was last sent this build of it. +func sentTheBuild(ctx context.Context, open *stores, module, commit string) ([]string, error) { + running, err := open.inventory.Running(ctx, module) + if err != nil { + return nil, err + } + var out []string + for _, n := range running { + sent, known, err := open.inventory.SentBuilds(ctx, n) + if err != nil { + return nil, err + } + if known && sameCommit(sent[module], commit) { + out = append(out, n) + } + } + return out, nil +} + +// releaseRepository is what a release plan says it is for, where a merge's says its repository. +const releaseRepository = "the builds waiting for a gate" + +// releaseHeard is the machines a release plan may send: heard within their heartbeat's bound. A machine +// away is left, not judged against a bound it cannot meet. A variable so a test says who is heard. +var releaseHeard = func(ctx context.Context, open *stores) (map[string]bool, error) { + if d := doctorFrom; d != nil && d.watchdogs != nil { + return heardMachines(d), nil + } + reports, err := open.inventory.LastReports(ctx) + if err != nil { + return nil, err + } + out := map[string]bool{} + for _, r := range reports { + if r.At != nil && time.Since(*r.At) < 15*time.Minute { + out[r.Node] = true + } + } + return out, nil +} + +// backlogNow is what the newest look found waiting, for `upgrade backlog` and the gate's probe. +var backlogNow struct { + held string + waiting map[string][]inventory.CarriedMove +} + +// waitingMoves is every machine's moves no gate has seen and no started plan walks. +func waitingMoves(ctx context.Context, open *stores, all bool) (map[string][]inventory.CarriedMove, error) { + f, err := readMoveFacts(ctx, open.inventory) + if err != nil { + return nil, err + } + nodes, err := open.inventory.Nodes(ctx) + if err != nil { + return nil, err + } + out := map[string][]inventory.CarriedMove{} + for _, n := range nodes { + moves, err := machineMoves(ctx, open, f, n.Name, all) + if err != nil { + return nil, err + } + for _, mv := range moves { + if f.walkedBy(mv.Module, n.Name) == "" { + out[n.Name] = append(out[n.Name], mv) + } + } + } + return out, nil +} + +// releaseBacklog opens a release plan when builds wait for a gate and none is open; not after a release +// plan failed, until a person releases one (by). Called with the plans held. +func releaseBacklog(ctx context.Context, open *stores, by string) (*inventory.Plan, error) { + inv := open.inventory + plans, err := inv.OpenPlans(ctx) + if err != nil { + return nil, err + } + for _, p := range plans { + if p.Release != nil { + if by != "" { + return nil, fmt.Errorf("%s is already releasing what waits; `plans %s` says where it is", p.ID, p.ID) + } + return nil, nil + } + } + waiting, err := waitingMoves(ctx, open, false) + if err != nil { + return nil, err + } + backlogNow.waiting, backlogNow.held = waiting, "" + if len(waiting) == 0 { + return nil, nil + } + // Opened by what no gate has seen; it walks every machine where anything of the release waits — a + // build that passed on the first machine still goes to the next one by this plan, judged there too. + if waiting, err = waitingMoves(ctx, open, true); err != nil { + return nil, err + } + if by == "" { + recent, err := inv.RecentPlans(ctx, 50) + if err != nil { + return nil, err + } + for _, p := range recent { + if p.Release == nil { + continue + } + if p.State == inventory.PlanFailed { + backlogNow.held = fmt.Sprintf("%s failed (%s); what waits is released again by a person", p.ID, p.Note) + return nil, nil + } + break + } + } + heard, err := releaseHeard(ctx, open) + if err != nil { + return nil, err + } + controllers, err := inv.Running(ctx, catalogue.ControllerSeatName) + if err != nil { + return nil, err + } + var order, last []string + modules := map[string]bool{} + for node, moves := range waiting { + if !heard[node] { + continue + } + for _, mv := range moves { + modules[mv.Module] = true + } + if slices.Contains(controllers, node) { + last = append(last, node) + } else { + order = append(order, node) + } + } + if len(order)+len(last) == 0 { + return nil, nil + } + sort.Strings(order) + sort.Strings(last) + order = append(order, last...) + names := make([]string, 0, len(modules)) + for m := range modules { + names = append(names, m) + } + sort.Strings(names) + now := time.Now().UTC() + p := inventory.Plan{ID: fmt.Sprintf("release-%d", now.UnixNano()), Repository: releaseRepository, + Created: now, State: inventory.PlanRolling, Tiers: [][]string{names}, Modules: map[string]*inventory.PlanModule{}, + Release: &inventory.PlanRelease{Order: order, By: by}, + Note: fmt.Sprintf("%d build(s) wait for a gate on %s; one machine at a time, each judged", len(names), + strings.Join(order, ", "))} + if err := inv.SavePlan(ctx, &p); err != nil { + return nil, err + } + fmt.Printf("%s: %s\n", p.ID, p.Note) + return &p, nil +} + +// advanceRelease takes one step of a release plan: the next machine sent everything waiting there under a +// gate, or the machine being judged judged once more; the plan done when every machine is. +func advanceRelease(ctx context.Context, open *stores, p *inventory.Plan) (bool, error) { + r := p.Release + now := time.Now().UTC() + if r.Gate == nil { + heard, err := releaseHeard(ctx, open) + if err != nil { + return false, err + } + for r.Next < len(r.Order) { + node := r.Order[r.Next] + if !heard[node] { + r.Skipped = append(r.Skipped, node) + r.Next++ + continue + } + moves, sent, err := gatedSend(ctx, open, node, nil) + if errors.Is(err, errWalkedElsewhere) { + note := fmt.Sprintf("waiting before %s: %v", node, err) + changed := p.Note != note + p.Note = note + return changed, nil + } + if err != nil { + return false, fmt.Errorf("sending %s what waits there: %w", node, err) + } + if len(moves) == 0 { + r.Done = append(r.Done, node) + r.Next++ + continue + } + r.Gate = &inventory.PlanGate{Machines: sent, Since: &now, Carried: moves} + p.Note = fmt.Sprintf("sent %s %d build(s) that waited for a gate; judging them there", node, len(moves)) + fmt.Printf("%s: %s\n", p.ID, p.Note) + return true, nil + } + p.State, p.Tier = inventory.PlanDone, len(p.Tiers) + p.Note = fmt.Sprintf("released on %s", orNone(strings.Join(r.Done, ", "))) + if len(r.Skipped) > 0 { + p.Note += "; not heard from, left as they were: " + strings.Join(r.Skipped, ", ") + } + return true, nil + } + g := r.Gate + verdict, err := judgeMoves(ctx, open, g, nil, now) + if err != nil { + return false, err + } + switch verdict { + case "": + note := fmt.Sprintf("judging %s: %s", strings.Join(g.Machines, ", "), gateLine(g)) + changed := p.Note != note + p.Note = note + return changed, nil + case inventory.GatePassed: + passCarried(ctx, open, p, g, "") + r.Done = append(r.Done, firstOf(g.Machines)) + r.Next++ + r.Gate = nil + return true, nil + } + p.Note = "" + failCarried(ctx, open, p, g, "") + 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 +} + +// backlogObservation is what the gate's probe says of a release held after a failure. +func backlogObservation() []conditions.Observation { + if backlogNow.held == "" || len(backlogNow.waiting) == 0 { + return nil + } + n := 0 + var machines []string + for node, moves := range backlogNow.waiting { + n += len(moves) + machines = append(machines, node) + } + sort.Strings(machines) + return []conditions.Observation{{Scope: conditions.ScopeMesh, ID: "release", Token: "held", Kind: "release-held", + Severity: conditions.Warning, Resolver: conditions.ResolverOperator, + Summary: fmt.Sprintf("%d build move(s) on %s wait for a gate and are not released: %s — `upgrade backlog` lists "+ + "them, `upgrade release-backlog --why …` releases them", n, strings.Join(machines, ", "), backlogNow.held)}} +} + +// backlogCommand is `upgrade backlog`, read-only, and `upgrade release-backlog --why`. +func backlogCommand(ctx context.Context, sub string, args []string) error { + set := flag.NewFlagSet("upgrade "+sub, flag.ContinueOnError) + why := addHandActFlags(set) + if rest, err := parseAround(set, args); err != nil { + return err + } else if len(rest) > 0 { + return errors.New("upgrade backlog | upgrade release-backlog --why ") + } + if sub == "release-backlog" { + if err := why.require("upgrade release-backlog"); err != nil { + return err + } + } + open, err := openStores(ctx) + if err != nil { + return err + } + defer open.Close() + if sub == "release-backlog" { + release, err := open.inventory.HoldPlans(ctx, true) + if err != nil { + return err + } + defer release() + why.record(ctx, "upgrade release-backlog", nil) + p, err := releaseBacklog(ctx, open, link.Caller()) + if err != nil { + return err + } + if p == nil { + fmt.Println("nothing waits for a gate on any machine heard from: nothing to release") + return nil + } + fmt.Printf("%s releases it; `plans %s` says where it is\n", p.ID, p.ID) + return nil + } + waiting, err := waitingMoves(ctx, open, false) + if err != nil { + return err + } + if len(waiting) == 0 { + fmt.Println("no build waits for a gate on any machine") + return nil + } + nodes := make([]string, 0, len(waiting)) + for n := range waiting { + nodes = append(nodes, n) + } + sort.Strings(nodes) + for _, n := range nodes { + fmt.Printf("%s: %d build(s) wait for a gate\n", n, len(waiting[n])) + for _, mv := range waiting[n] { + fmt.Printf(" %-28s %s → %s\n", mv.Module, short(mv.From), short(mv.To)) + } + } + return nil +} diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 2a59c22..1adb2fa 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -476,6 +476,10 @@ func advancePlans(ctx context.Context, open *stores) { return } defer release() + // What waits for a gate, released one machine at a time (ADR 0236). + if _, err := releaseBacklog(ctx, open, ""); err != nil { + fmt.Printf("plans: what waits for a gate could not be looked at: %v\n", err) + } advanceHeld(ctx, open) } @@ -531,6 +535,9 @@ func advanceHeld(ctx context.Context, open *stores) { func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, edges []inventory.Edge, rollsOut func(string) bool) (bool, error) { inv := open.inventory + if p.Release != nil { + return advanceRelease(ctx, open, p) + } if p.Tier >= len(p.Tiers) { p.State = inventory.PlanDone fmt.Printf("%s: done — %s at %s, %d tier(s)\n", p.ID, p.Repository, short(p.Commit), len(p.Tiers)) @@ -631,6 +638,9 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, // 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) return true, nil @@ -653,6 +663,9 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, 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) return true, nil @@ -677,7 +690,20 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, } // 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. - sent, err := sendRollout(ctx, open, step.send) + 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) + } 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. @@ -690,6 +716,8 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, } 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", @@ -1059,6 +1087,18 @@ func plansCommand(ctx context.Context, args []string) error { return err } fmt.Printf("%s — %s\n", p.ID, planLineWith(p, now, buildSeatPause(ctx, inv, []inventory.Plan{p}))) + if r := p.Release; r != nil { + // A release plan's walk (ADR 0236): machines done, the one judged, those to come. + fmt.Printf(" machines in order: %s; done: %s; skipped: %s\n", strings.Join(r.Order, ", "), + orNone(strings.Join(r.Done, ", ")), orNone(strings.Join(r.Skipped, ", "))) + if r.Gate != nil { + fmt.Printf(" %s\n", gateLine(r.Gate)) + for _, c := range r.Gate.Carried { + fmt.Printf(" %-22s %s → %s\n", c.Module, short(c.From), short(c.To)) + } + } + return nil + } for i, tier := range p.Tiers { marker := " " if i == p.Tier && p.Open() { diff --git a/cmd/mesh-controller/release_test.go b/cmd/mesh-controller/release_test.go new file mode 100644 index 0000000..11c4c82 --- /dev/null +++ b/cmd/mesh-controller/release_test.go @@ -0,0 +1,257 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "reflect" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/lease" +) + +// The backlog the old default left — builds registered and sent nowhere — is released by a plan, one +// machine at a time, each judged, never by the next send of something else (novox/hq ADR 0236). + +type backlogMesh struct { + open *stores + sent [][]string + broken map[string]bool +} + +// aBacklog is a mesh where `app` moved c1 → c2 on anchor and laptop, `late` moved on laptop only, and +// `same` was rebuilt with the very artifacts it had: two machines heard, every send applied at once. +func aBacklog(t *testing.T) *backlogMesh { + t.Helper() + open := aMesh(t) + ctx := t.Context() + inv := open.inventory + b := &backlogMesh{open: open, broken: map[string]bool{}} + withConditionsInMemory(t) + was := doctorFrom + doctorFrom = nil + t.Cleanup(func() { doctorFrom = was }) + + build := func(module, commit, made string, asked time.Time) { + manifest, _ := json.Marshal(catalogue.Manifest{Module: module, Version: "1"}) + if err := inv.RecordBuild(ctx, inventory.Build{ID: "build-" + module + "-" + commit, Module: module, Commit: commit, + Repository: "novox/mesh-catalog", Path: "modules/" + module, Manifest: manifest, Asked: asked, At: asked, + Made: []inventory.Artifact{{Name: "x", Kind: "bundle", Reference: made}}}); err != nil { + t.Fatal(err) + } + if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: module, Version: "1"}, inventory.Source{ + Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/" + module, BuiltFrom: commit, Head: commit, + Asked: asked}); err != nil { + t.Fatal(err) + } + } + old, now := time.Now().Add(-2*time.Hour), time.Now().Add(-time.Minute) + on := map[string][]string{"app": {"anchor", "laptop"}, "late": {"laptop"}, "same": {"anchor", "laptop"}} + for _, m := range []string{"app", "late", "same"} { + build(m, "c1", "sha256:"+m+"-1", old) + for _, n := range on[m] { + if _, err := inv.Assign(ctx, n, m); err != nil { + t.Fatal(err) + } + } + } + for _, n := range []string{"anchor", "laptop"} { + sent := map[string]string{} + for m, nodes := range on { + for _, x := range nodes { + if x == n { + sent[m] = "c1" + } + } + } + if err := inv.RecordSent(ctx, nodeID(t, open, n), "d-"+n, sent); err != nil { + t.Fatal(err) + } + } + build("app", "c2", "sha256:app-2", now) + build("late", "c2", "sha256:late-2", now) + build("same", "c2", "sha256:same-1", now) // rebuilt, the same bytes + + assigned := func(ctx context.Context, open *stores, f moveFacts, node string, all bool) ([]inventory.CarriedMove, error) { + modules, err := open.inventory.Assigned(ctx, node) + if err != nil { + return nil, err + } + sent, known, err := open.inventory.SentBuilds(ctx, node) + if err != nil { + return nil, err + } + return f.moves(node, modules, sent, known, all), nil + } + wasMoves, wasHeard, wasSend, wasGather := machineMoves, releaseHeard, sendRollout, gatherGateFacts + machineMoves = assigned + releaseHeard = func(context.Context, *stores) (map[string]bool, error) { + return map[string]bool{"anchor": true, "laptop": true}, nil + } + n := 0 + sendRollout = func(ctx context.Context, open *stores, names []string) ([]string, error) { + b.sent = append(b.sent, append([]string(nil), names...)) + current, err := open.inventory.CurrentBuilds(ctx) + if err != nil { + return nil, err + } + for _, node := range names { + modules, _ := open.inventory.Assigned(ctx, node) + carried := map[string]string{} + for _, m := range modules { + carried[m] = current[m].Commit + } + n++ + digest := fmt.Sprintf("d-%s-%d", node, n) + if err := open.inventory.RecordSent(ctx, nodeID(t, open, node), digest, carried); 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 + } + 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{}} + reports, err := open.inventory.LastReports(ctx) + if err != nil { + return f, err + } + for _, r := range reports { + if b.broken[r.Node] { + r.Outcome = inventory.OutcomeFailed + } + f.reports[r.Node] = r + } + return f, nil + } + wasSettle, wasEvery, wasBound := gateSettle, gateEvery, gateBound + gateSettle, gateEvery = 0, time.Hour + t.Cleanup(func() { + machineMoves, releaseHeard, sendRollout, gatherGateFacts = wasMoves, wasHeard, wasSend, wasGather + gateSettle, gateEvery, gateBound = wasSettle, wasEvery, wasBound + }) + return b +} + +func (b *backlogMesh) release(t *testing.T) inventory.Plan { + t.Helper() + plans, err := b.open.inventory.RecentPlans(t.Context(), 10) + if err != nil { + t.Fatal(err) + } + for _, p := range plans { + if p.Release != nil { + return p + } + } + t.Fatal("no release plan") + return inventory.Plan{} +} + +// The backlog goes out one machine at a time: the first judged before the second is sent, each build's +// pass kept; a rebuild with the same bytes is no move at all; and a send outside a gate is refused. +func TestTheBacklogIsReleasedOneMachineAtATimeEachJudged(t *testing.T) { + b := aBacklog(t) + ctx := t.Context() + inv := b.open.inventory + + f, err := readMoveFacts(ctx, inv) + if err != nil { + t.Fatal(err) + } + moves, _ := machineMoves(ctx, b.open, f, "laptop", false) + var names []string + for _, mv := range moves { + names = append(names, mv.Module) + } + if !reflect.DeepEqual(names, []string{"app", "late"}) { + t.Fatalf("laptop waits for %v; a rebuild with the same bytes is no move", names) + } + // Nothing else may carry them: a send that judges nothing is refused. + if _, err := ungatedIn(ctx, b.open, []string{"laptop"}, ""); !errors.Is(err, errUngated) { + t.Fatalf("a send that judges nothing would carry them: %v", err) + } + + advancePlans(ctx, b.open) + if !reflect.DeepEqual(b.sent, [][]string{{"anchor"}}) { + t.Fatalf("sent %v: the first machine alone, before it is judged", b.sent) + } + p := b.release(t) + if !reflect.DeepEqual(p.Release.Order, []string{"anchor", "laptop"}) || p.Release.Gate == nil { + t.Fatalf("the release plan is %+v", p.Release) + } + gateEvery = 0 + for i := 0; i < 4; i++ { + advancePlans(ctx, b.open) + } + if !reflect.DeepEqual(b.sent, [][]string{{"anchor"}, {"laptop"}}) { + t.Fatalf("sent %v", b.sent) + } + p = b.release(t) + if p.State != inventory.PlanDone || !reflect.DeepEqual(p.Release.Done, []string{"anchor", "laptop"}) { + t.Fatalf("the release plan is %s: %s %+v", p.State, p.Note, p.Release) + } + for _, build := range []string{"build-app-c2", "build-late-c2"} { + if v, found, err := inv.GateOf(ctx, build); err != nil || !found || v.Verdict != inventory.GatePassed { + t.Fatalf("%s's pass was not kept: %+v %v %v", build, v, found, err) + } + } + // Nothing waits now, and nothing opens again. + if p, err := releaseBacklog(ctx, b.open, ""); err != nil || p != nil { + t.Fatalf("a release opened with nothing waiting: %+v %v", p, err) + } +} + +// A release that fails on its first machine puts back what it carried there, goes no further, and the +// next one waits for a person, said as a condition; a person releases it with why. +func TestAFailedReleaseIsPutBackAndTheNextWaitsForAPerson(t *testing.T) { + b := aBacklog(t) + ctx := t.Context() + inv := b.open.inventory + gateEvery = 0 + b.broken["anchor"] = true + advancePlans(ctx, b.open) + + p := b.release(t) + if p.State != inventory.PlanFailed { + t.Fatalf("the release is %s: %s", p.State, p.Note) + } + if !reflect.DeepEqual(b.sent, [][]string{{"anchor"}, {"anchor"}}) { + t.Fatalf("sent %v: the first machine, then the put-back to it, and nothing to laptop", b.sent) + } + if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c1" { + t.Fatalf("app is at %s, not put back", current["app"].Commit) + } + if failed, _ := inv.GateFailed(ctx, "build-app-c2"); !failed { + t.Fatal("app's build is not marked failed at its gate") + } + // late still waits on laptop, and is not released on its own after a failure. + if p, err := releaseBacklog(ctx, b.open, ""); err != nil || p != nil { + t.Fatalf("a release opened after a failed one: %+v %v", p, err) + } + if obs := backlogObservation(); len(obs) != 1 || !strings.Contains(obs[0].Summary, "release-backlog") { + t.Fatalf("the held release is not said: %+v", obs) + } + b.broken["anchor"] = false + if err := backlogCommand(ctx, "release-backlog", []string{"--why", "anchor is fixed"}); err != nil { + t.Fatal(err) + } + for i := 0; i < 4; i++ { + advancePlans(ctx, b.open) + } + if p := b.release(t); p.State != inventory.PlanDone || p.Release.By == "" { + t.Fatalf("the person's release is %s (%+v)", p.State, p.Release) + } + if sent, _, _ := inv.SentBuilds(ctx, "laptop"); sent["late"] != "c2" { + t.Fatalf("late was not released to laptop: %v", sent) + } +} diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index d4a04f0..3517e34 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -509,6 +509,16 @@ func (a *verbArguments) commandLine() ([]string, error) { } return argv, nil case "upgrade": + if on("backlog") { + return []string{"upgrade", "backlog"}, nil + } + if on("release-backlog") { + argv := []string{"upgrade", "release-backlog", "--why", str("why")} + if c := str("cause"); c != "" { + argv = append(argv, "--cause", c) + } + return argv, nil + } m, policy := str("module"), str("policy") if m == "" { if policy != "" { diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index 33cce27..515d311 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -74,24 +74,12 @@ func (f following) Upgraded(ctx context.Context, u link.Upgraded) error { u.Module, shortCommit(u.Commit), id, readableList(on)) return nil } - if decision.Together { - fmt.Printf("%s moved to %s; sending %s together\n", - u.Module, shortCommit(u.Commit), readableList(on)) - return askAgainOnGrants(sendTo(ctx, f.open, on)) - } - // One at a time, and stopping at the first that fails. - // - // **Stopping is the point.** The machines are done one after another precisely so that a - // version that breaks the first one does not reach the rest; carrying on past a failure would - // make this the same as sending them together, only slower. - fmt.Printf("%s moved to %s; sending %s one at a time\n", - u.Module, shortCommit(u.Commit), readableList(on)) - for _, node := range on { - if err := sendTo(ctx, f.open, []string{node}); err != nil { - return askAgainOnGrants(fmt.Errorf("%s did not take %s, so the machines after it were left alone: %w", - node, u.Module, err)) - } - } + // **No plan holds it: a release plan sends it** (novox/hq ADR 0236). A build asked outside a plan — a + // `rebuild`, a `replay` registered — was sent here one machine after another without any judging; it + // now waits for a gate like any other build, and the release plan walks it through the machines, one + // at a time, each judged. + fmt.Printf("%s moved to %s outside a plan; a release plan sends it to %s, one machine at a time, each judged "+ + "(`upgrade backlog` lists what waits)\n", u.Module, shortCommit(u.Commit), readableList(on)) return nil } @@ -147,6 +135,10 @@ func isAre(n int) string { // upgradeCommand says what should happen when a module's current version moves (ADR 0162 §3, ADR // 0235): every module's policy and where it comes from; one module's; or a person's choice for one. func upgradeCommand(ctx context.Context, args []string) error { + // What waits for a gate, and releasing it by a person's word (ADR 0236). + if len(args) > 0 && (args[0] == "backlog" || args[0] == "release-backlog") { + return backlogCommand(ctx, args[0], args[1:]) + } set := flag.NewFlagSet("upgrade", flag.ContinueOnError) together := set.Bool("together", false, "send every machine running it at once, instead of one machine first") diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index 0e9b535..d6020d2 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -277,8 +277,12 @@ var ControllerVerbs = []Verb{ "module": "one module", "policy": "with module: roll-out, record or default", "together": "\"true\": with roll-out, every machine at once instead of one machine first", - "why": "with policy: why — required for record, kept and said with the policy", - }, nil, "together")}, + "why": "with policy or release-backlog: why — required for record and for a release, kept", + "backlog": "\"true\": every build waiting for a gate, per machine — what a release plan would send", + "release-backlog": "\"true\": release what waits now, one machine at a time, each judged — after a release " + + "plan failed, the mesh waits for this; with why", + "cause": "with release-backlog: the cause in a word (optional)", + }, nil, "together", "backlog", "release-backlog")}, {Name: "bus", Description: "The bus as a planned step (novox/hq to-be 45 §8, ADR 0236): what a bus upgrade " + "would do — the bus's build on each machine against the one the mesh holds — and how the last step went. " + "With upgrade, start one: a person's act with why, after the streams are snapshotted (snapshot-taken says " + diff --git a/internal/inventory/gate.go b/internal/inventory/gate.go index a07639e..9d32b07 100644 --- a/internal/inventory/gate.go +++ b/internal/inventory/gate.go @@ -46,8 +46,8 @@ type GateVerdict struct { Epoch uint64 } -// RecordGate writes a build's verdict. **A failed build's row is written once**: a second failure for -// the same build is refused with ErrGateKept, which is what keeps a rollback to one per build — the row +// RecordGate writes a build's verdict. A build that passed on one machine may still fail on the next one +// judged; **a failed build's row is written once**: any later verdict for it is refused with ErrGateKept, which is what keeps a rollback to one per build — the row // is written before the rollback's send, and a controller replaced in between finds it. func (i *Inventory) RecordGate(ctx context.Context, v GateVerdict) error { epoch, err := i.actingEpoch(ctx) @@ -64,7 +64,7 @@ func (i *Inventory) RecordGate(ctx context.Context, v GateVerdict) error { on conflict (build) do update set verdict = excluded.verdict, rollback = excluded.rollback, why = excluded.why, previous = excluded.previous, machines = excluded.machines, judged_at = now(), epoch = excluded.epoch - where build_gate.verdict = 'passed' and excluded.verdict = 'passed'`, + where build_gate.verdict = 'passed'`, v.Build, v.Module, v.Commit, v.Previous, v.Plan, v.Machines, v.Verdict, v.Rollback, v.Why, v.Component, v.JudgingFrom, epoch) if err != nil { @@ -245,3 +245,65 @@ func (i *Inventory) BuildFingerprints(ctx context.Context, module string) (map[s } return out, rows.Err() } + +// Fingerprints is BuildFingerprints for every module at once: module → commit → fingerprint. +func (i *Inventory) Fingerprints(ctx context.Context) (map[string]map[string]string, error) { + rows, err := i.store.Pool().Query(ctx, + `select module, commit_hash, made, coalesce(manifest::text, '') from build + where module is not null and failed = '' and commit_hash <> '' order by at desc`) + if err != nil { + return nil, err + } + defer rows.Close() + out := map[string]map[string]string{} + for rows.Next() { + var module, commit, manifest string + var made []byte + if err := rows.Scan(&module, &commit, &made, &manifest); err != nil { + return nil, err + } + if out[module] == nil { + out[module] = map[string]string{} + } + if _, seen := out[module][commit]; seen { + continue + } + sum := sha256.Sum256(append(append(made, 0), []byte(manifest)...)) + out[module][commit] = hex.EncodeToString(sum[:]) + } + return out, rows.Err() +} + +// PassedCommits is, per module, the commits a build of which passed its gate on some machine. +func (i *Inventory) PassedCommits(ctx context.Context) (map[string]map[string]bool, error) { + rows, err := i.store.Pool().Query(ctx, `select module, commit_hash from build_gate where verdict = 'passed'`) + if err != nil { + return nil, err + } + defer rows.Close() + out := map[string]map[string]bool{} + for rows.Next() { + var module, commit string + if err := rows.Scan(&module, &commit); err != nil { + return nil, err + } + if out[module] == nil { + out[module] = map[string]bool{} + } + out[module][commit] = true + } + return out, rows.Err() +} + +// BuildOf is the id of the newest successful build of a module from a commit; empty when none is +// recorded (a manifest handed over by hand). +func (i *Inventory) BuildOf(ctx context.Context, module, commit string) (string, error) { + var id string + err := i.store.Pool().QueryRow(ctx, + `select id from build where module = $1 and commit_hash = $2 and failed = '' order by at desc limit 1`, + module, commit).Scan(&id) + if errors.Is(err, pgx.ErrNoRows) { + return "", nil + } + return id, err +} diff --git a/internal/inventory/migrations/0073-a-build-rolls-out-gated-and-rolls-back.sql b/internal/inventory/migrations/0073-a-build-rolls-out-gated-and-rolls-back.sql index e6b312d..9d57766 100644 --- a/internal/inventory/migrations/0073-a-build-rolls-out-gated-and-rolls-back.sql +++ b/internal/inventory/migrations/0073-a-build-rolls-out-gated-and-rolls-back.sql @@ -65,3 +65,8 @@ create table bus_step ( outcome text not null default '' check (outcome in ('', 'done', 'failed')), found text not null default '' ); + +-- 4. A release plan (ADR 0236): the builds that wait for a gate — the backlog the old default left, and +-- whatever a plan built and did not send — walked through the machines one at a time, each judged +-- before the next. Its walk is kept with the plan. +alter table release_plan add column release jsonb; diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index 7508b60..cbb37f4 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -38,6 +38,9 @@ type Plan struct { Revision int64 `json:"revision"` // Epoch is the controller lease epoch that wrote it last; zero for a write that claimed none. Epoch uint64 `json:"epoch,omitempty"` + // Release is set on a release plan (novox/hq ADR 0236): not a merge's, but the builds waiting for a + // gate, walked through the machines one at a time. + Release *PlanRelease `json:"release,omitempty"` } // ErrPlanMoved is a save against a plan written by somebody else since it was read. @@ -103,6 +106,35 @@ type PlanGate struct { Rollback string `json:"rollback,omitempty"` // Kept says a passing verdict was written to the gate's records. Kept bool `json:"kept,omitempty"` + // Carried is every module whose build moved on the judged machines with the send — the plan's own + // module and whatever else was waiting there for a gate (novox/hq ADR 0236): each is judged here, a + // pass is its verdict too, and one that fails is put back. + Carried []CarriedMove `json:"carried,omitempty"` + // Failing names the modules the last judging found wanting. + Failing []string `json:"failing,omitempty"` +} + +// CarriedMove is one module's build moving on a machine with a gated send. +type CarriedMove struct { + Module string `json:"module"` + Node string `json:"node"` + From string `json:"from,omitempty"` + To string `json:"to"` + Build string `json:"build,omitempty"` +} + +// PlanRelease is a release plan's walk through the machines (novox/hq ADR 0236): every module build +// that waits for a gate, sent one machine at a time, each judged before the next. +type PlanRelease struct { + Order []string `json:"order"` + Next int `json:"next"` + // Gate is the machine being judged; nil between machines. + Gate *PlanGate `json:"gate,omitempty"` + Done []string `json:"done,omitempty"` + // Skipped are the machines not heard from when their turn came, left as they were. + Skipped []string `json:"skipped,omitempty"` + // By is the person who released it, empty when the mesh did. + By string `json:"by,omitempty"` } // The states a plan passes through. @@ -138,6 +170,12 @@ func (i *Inventory) SavePlan(ctx context.Context, p *Plan) error { if err != nil { return err } + var release []byte + if p.Release != nil { + if release, err = json.Marshal(p.Release); err != nil { + return err + } + } // **And how long the tier it left took** (novox/hq to-be 45 Phase 0): measured here, where the // plan moves, in the same transaction as the move, so no save can move a tier unmeasured or // measure one twice. @@ -153,15 +191,16 @@ func (i *Inventory) SavePlan(ctx context.Context, p *Plan) error { var revision int64 err = tx.QueryRow(ctx, `insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, - branch, tier_entered, revision, epoch) - values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10, $11, 1, $13) + branch, tier_entered, revision, epoch, release) + values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10, $11, 1, $13, $14) on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier, tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch, - tier_entered = excluded.tier_entered, revision = release_plan.revision + 1, epoch = excluded.epoch + tier_entered = excluded.tier_entered, revision = release_plan.revision + 1, epoch = excluded.epoch, + release = excluded.release where release_plan.revision = $12 returning revision`, p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch, entered, - p.Revision, epoch).Scan(&revision) + p.Revision, epoch, release).Scan(&revision) if errors.Is(err, pgx.ErrNoRows) { // The row is there and at another revision — moved since this was read, or there already // when this one is new: either way not this writer's to overwrite. (A plan saved before plans @@ -216,7 +255,7 @@ func (i *Inventory) PlanByID(ctx context.Context, id string) (Plan, error) { func (i *Inventory) plans(ctx context.Context, tail string) ([]Plan, error) { rows, err := i.store.Pool().Query(ctx, `select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch, - coalesce(tier_entered, created), revision, coalesce(epoch, 0) + coalesce(tier_entered, created), revision, coalesce(epoch, 0), release from release_plan `+tail) if err != nil { return nil, err @@ -225,12 +264,17 @@ func (i *Inventory) plans(ctx context.Context, tail string) ([]Plan, error) { var out []Plan for rows.Next() { var p Plan - var tiers, modules []byte + var tiers, modules, release []byte var epoch int64 if err := rows.Scan(&p.ID, &p.Repository, &p.Commit, &p.Created, &p.Updated, &p.State, - &p.Tier, &tiers, &modules, &p.Note, &p.Branch, &p.TierEntered, &p.Revision, &epoch); err != nil { + &p.Tier, &tiers, &modules, &p.Note, &p.Branch, &p.TierEntered, &p.Revision, &epoch, &release); err != nil { return nil, err } + if len(release) > 0 { + if err := json.Unmarshal(release, &p.Release); err != nil { + return nil, err + } + } p.Epoch = uint64(epoch) if err := json.Unmarshal(tiers, &p.Tiers); err != nil { return nil, err