795 lines
28 KiB
Go
795 lines
28 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"io"
|
|
"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 <node>`, 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
|
|
// srcs is, per module, per commit, its build's source fingerprint (novox/hq issue 280).
|
|
srcs 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.srcs, err = inv.SourceFingerprints(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,
|
|
// builds made from the same source (novox/hq issue 280) — the module's tree, the contexts it read, its
|
|
// bases and toolchains, whatever digests an image rebuild made of them — 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
|
|
}
|
|
if sa := f.srcs[module][a]; sa != "" && sa == f.srcs[module][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
|
|
}
|
|
|
|
// unchangedOnEvery is whether a plan's build of a module was made from the source of the build every
|
|
// machine running it was last sent (novox/hq issue 280): no move on any of them. False when no machine
|
|
// runs it, when what one was sent is not known, or when the build stands for itself.
|
|
func unchangedOnEvery(ctx context.Context, inv *inventory.Inventory, module, build string, running []string) (bool, error) {
|
|
if build == "" || len(running) == 0 {
|
|
return false, nil
|
|
}
|
|
same, err := inv.SameSourceCommits(ctx, module, build)
|
|
if err != nil || len(same) == 0 {
|
|
return false, err
|
|
}
|
|
for _, n := range running {
|
|
sent, known, err := inv.SentBuilds(ctx, n)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
was, carried := sent[module]
|
|
if !known || !carried || !inRun(same, was) {
|
|
return false, nil
|
|
}
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
// inRun is whether a commit is one of a run's, however either is abbreviated.
|
|
func inRun(run map[string]bool, commit string) bool {
|
|
for c := range run {
|
|
if sameCommit(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 walk 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.
|
|
//
|
|
// 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 {
|
|
return nil, nil, err
|
|
}
|
|
moves, err := machineMoves(ctx, open, f, node, true)
|
|
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(mv.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)
|
|
}
|
|
}
|
|
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 && len(owns) == 0 {
|
|
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
|
|
}
|
|
}
|
|
}
|
|
sort.Slice(moves, func(i, j int) bool { return moves[i].Module < moves[j].Module })
|
|
sayRecreations(ctx, open, moves)
|
|
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}
|
|
// **What passed on its own keeps its pass** (novox/hq ADR 0254, issue 318): healthy for the passes the
|
|
// gate asks while another module of the send failed, it is neither put back nor left unjudged — a build
|
|
// left without a verdict is one more move every other walk would wait for.
|
|
if kept := keepPassing(ctx, open, p, g, except); len(kept) > 0 {
|
|
notes = append(notes, "passed on their own and kept: "+strings.Join(kept, ", "))
|
|
for _, m := range kept {
|
|
done[m] = 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}
|
|
// 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)
|
|
}
|
|
p.State = inventory.PlanFailed
|
|
p.Note = strings.Join(notes, "; ")
|
|
}
|
|
|
|
// keepPassing keeps a pass as the verdict of every build the gate carried that passed on its own when the
|
|
// send failed (ADR 0254), except the plan's own module; answers the modules kept.
|
|
func keepPassing(ctx context.Context, open *stores, p *inventory.Plan, g *inventory.PlanGate, except string) []string {
|
|
var kept []string
|
|
seen := map[string]bool{}
|
|
for _, c := range g.Carried {
|
|
if c.Module == except || c.Build == "" || seen[c.Build] || !slices.Contains(g.Passing, c.Module) {
|
|
continue
|
|
}
|
|
seen[c.Build] = true
|
|
why := fmt.Sprintf("healthy %d times on its own while the send failed: %s", g.Healthy[c.Module], g.Why)
|
|
if w := g.Waits[c.Module]; w != "" {
|
|
why += "; and it waits for a person: " + w
|
|
}
|
|
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: why,
|
|
Component: coreComponent(c.Module), JudgingFrom: g.Since}); err != nil {
|
|
fmt.Printf("%s: %s passed its gate on %s on its own, and the verdict could not be kept: %v\n", p.ID, c.Module,
|
|
c.Node, err)
|
|
continue
|
|
}
|
|
if !slices.Contains(kept, c.Module) {
|
|
kept = append(kept, c.Module)
|
|
}
|
|
}
|
|
return kept
|
|
}
|
|
|
|
// 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))
|
|
if said := recreationsSaid(moves); said != "" {
|
|
p.Note += "; " + said
|
|
}
|
|
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 = ""
|
|
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
|
|
}
|
|
|
|
// 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 <text>")
|
|
}
|
|
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
|
|
}
|
|
|
|
// sayRecreations says, on each move a send carries, what it does to the module's containers (novox/hq
|
|
// ADR 0245): the build the machine ran against the one it is sent, container by container. Left unsaid
|
|
// for a move whose earlier build is not in the records.
|
|
func sayRecreations(ctx context.Context, open *stores, moves []inventory.CarriedMove) {
|
|
if len(moves) == 0 {
|
|
return
|
|
}
|
|
shelf, err := open.inventory.Catalogue(ctx)
|
|
if err != nil {
|
|
return
|
|
}
|
|
for i, mv := range moves {
|
|
to, held := shelf[mv.Module]
|
|
if !held || mv.From == "" {
|
|
continue
|
|
}
|
|
from, found, err := open.inventory.ManifestAt(ctx, mv.Module, mv.From)
|
|
if err != nil || !found {
|
|
continue
|
|
}
|
|
moves[i].Recreates = catalogue.Recreates(from, to).Say(mv.Module)
|
|
}
|
|
}
|
|
|
|
// sayWhatAPushRecreates says, per machine a push is about to send, every module the send moves and what
|
|
// the move recreates (novox/hq ADR 0245) — the same words a gated send keeps on its plan. A person's push
|
|
// carries every module's new build, whatever its policy and whether or not a gate saw it, so it is the
|
|
// send that most needs to say so: on 2026-10-07 a `push` of one machine moved mail to a new build and
|
|
// recreated one of its containers with a new image, and the answer said nothing about it. A machine whose
|
|
// last send's builds are not known is not said, and neither is a failure to read: the push goes on.
|
|
func sayWhatAPushRecreates(ctx context.Context, open *stores, sending []readyNode, w io.Writer) {
|
|
if len(sending) == 0 {
|
|
return
|
|
}
|
|
f, err := readMoveFacts(ctx, open.inventory)
|
|
if err != nil {
|
|
return
|
|
}
|
|
for _, r := range sending {
|
|
sent, known, err := open.inventory.SentBuilds(ctx, r.node)
|
|
if err != nil || !known {
|
|
continue
|
|
}
|
|
moves := pushMoves(f, r.node, r.declared.Builds, sent)
|
|
if len(moves) == 0 {
|
|
continue
|
|
}
|
|
sayRecreations(ctx, open, moves)
|
|
fmt.Fprintf(w, "%s moves %d module(s):\n", r.node, len(moves))
|
|
for _, mv := range moves {
|
|
said := mv.Recreates
|
|
if said == "" {
|
|
said = "recreates none of its containers, or what it ran is not in the records"
|
|
}
|
|
fmt.Fprintf(w, " %-28s %s → %s: %s\n", mv.Module, short(mv.From), short(mv.To), said)
|
|
}
|
|
}
|
|
}
|
|
|
|
// pushMoves is what a send carrying these builds moves on a machine last sent `sent`: every module it ran
|
|
// before at a build not identical to the one it is now sent. A module new to the machine is not a move.
|
|
func pushMoves(f moveFacts, node string, carries, sent map[string]string) []inventory.CarriedMove {
|
|
var out []inventory.CarriedMove
|
|
for m, to := range carries {
|
|
was, ran := sent[m]
|
|
if !ran || was == "" || to == "" || f.identical(m, was, to) {
|
|
continue
|
|
}
|
|
out = append(out, inventory.CarriedMove{Module: m, Node: node, From: was, To: to})
|
|
}
|
|
sort.Slice(out, func(i, j int) bool { return out[i].Module < out[j].Module })
|
|
return out
|
|
}
|
|
|
|
// recreationsSaid is every recreation a send's moves say, joined: what a plan's note carries.
|
|
func recreationsSaid(moves []inventory.CarriedMove) string {
|
|
var said []string
|
|
for _, mv := range moves {
|
|
if mv.Recreates != "" {
|
|
said = append(said, fmt.Sprintf("on %s %s", mv.Node, mv.Recreates))
|
|
}
|
|
}
|
|
return strings.Join(said, "; ")
|
|
}
|