package main import ( "context" "encoding/json" "errors" "flag" "fmt" "slices" "sort" "strings" "time" "github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/link" ) // A merge produces a tiered plan the mesh keeps (novox/hq ADR 0162). // // The handler that hears the merge computes the plan from the catalogue's one dependency relation, // writes it to the store, asks the first tier and returns — the receive loop is never held by a // build. Every outcome taken in advances the plan it belongs to; a ticker advances what outcomes // alone cannot (a tier waiting for machines to report); a controller replaced mid-plan finds the // plan where it left it. // planWaitBound is how long a plan may wait on one thing before `status` names it red. const planWaitBound = 30 * time.Minute // tiersOf sorts a set of modules into tiers along the ordering edges among them: tier 0 depends // on nothing else in the set, tier 1 only on tier 0, and so on. An edge to a module outside the set says // nothing about the order inside it. A cycle — which the catalogue should never produce — puts // what remains in one last tier rather than losing it, and is said by the caller. func tiersOf(set []string, edges []inventory.Edge) [][]string { in := map[string]bool{} for _, m := range set { in[m] = true } deps := map[string]map[string]bool{} for _, m := range set { deps[m] = map[string]bool{} } // The build seat's holders follow the controller that defines their worker (EdgeWorkerOf, // novox/hq issue 206), so the built-by edge from that controller to such a holder yields: the // controller is built by whichever build machine is running, as the runtime image always was. worker := map[string]map[string]bool{} for _, e := range edges { if e.Kind == inventory.EdgeWorkerOf && in[e.From] && in[e.To] { if worker[e.To] == nil { worker[e.To] = map[string]bool{} } worker[e.To][e.From] = true } } for _, e := range edges { // A code dependency — B packages A's source — rebuilds B with A, in the same tier: B's // build needs nothing of A's first. The other kinds order: stands-on and declared after // the base is built, built-by after the build machine is built and running — except for // what the build machine itself stands on, and for the controller whose worker the build // machine binds. The runtime image is built by the builder and the builder is built on the // runtime image; the image comes first, built by the builder that is running. if !in[e.From] || !in[e.To] || e.From == e.To || e.Kind == inventory.EdgePackages { continue } if e.Kind == inventory.EdgeBuiltBy && (isBaseOf(e.From, e.To, edges, in) || worker[e.From][e.To]) { continue } deps[e.From][e.To] = true } placed := map[string]bool{} var tiers [][]string for len(placed) < len(set) { var tier []string for _, m := range set { if placed[m] { continue } free := true for d := range deps[m] { if !placed[d] { free = false break } } if free { tier = append(tier, m) } } if len(tier) == 0 { // A cycle: everything left, together, and the caller says so. for _, m := range set { if !placed[m] { tier = append(tier, m) } } } sort.Strings(tier) for _, m := range tier { placed[m] = true } tiers = append(tiers, tier) } return tiers } // isBaseOf says whether `to` stands on `base`, directly or through other bases in the set, along // the build edges alone. func isBaseOf(base, to string, edges []inventory.Edge, in map[string]bool) bool { seen := map[string]bool{} var walk func(string) bool walk = func(m string) bool { if m == base { return true } if seen[m] { return false } seen[m] = true for _, e := range edges { if e.From == m && in[e.To] && (e.Kind == inventory.EdgeStandsOn || e.Kind == inventory.EdgeDeclared) && walk(e.To) { return true } } return false } return walk(to) } // reachableFrom is the moved modules plus everything that depends on them, through every layer: // what a merge rebuilds. Along the code and build edges only: a module *built by* the build machine // is not changed by a new build machine, so a built-by edge orders and gates a plan and never // widens it — the first plan of 2026-10-01 took the whole catalogue along for a controller change. func reachableFrom(moved []string, edges []inventory.Edge) []string { in := map[string]bool{} for _, m := range moved { in[m] = true } for grew := true; grew; { grew = false for _, e := range edges { // Built-by and worker-of order a plan; neither widens it. A new build machine changes // nothing it builds, and a new controller changes nothing about the holder it orders — // what packages the controller's source is already a code edge. if e.Kind == inventory.EdgeBuiltBy || e.Kind == inventory.EdgeWorkerOf { continue } if in[e.To] && !in[e.From] { in[e.From] = true grew = true } } } out := make([]string, 0, len(in)) for m := range in { out = append(out, m) } sort.Strings(out) return out } // hasCycle says whether the tiers' last tier holds modules that still depend on each other. func hasCycle(tiers [][]string, edges []inventory.Edge) bool { if len(tiers) == 0 { return false } last := map[string]bool{} for _, m := range tiers[len(tiers)-1] { last[m] = true } for _, e := range edges { if last[e.From] && last[e.To] { return true } } return false } // planFor is the plan a merge produces: the moved modules and everything reachable from them, // tiered, with the merge it answers. func planOfMerge(m link.SourceMoved, moved []string, edges []inventory.Edge) inventory.Plan { set := reachableFrom(moved, edges) tiers := tiersOf(set, edges) modules := map[string]*inventory.PlanModule{} for _, name := range set { modules[name] = &inventory.PlanModule{} } return inventory.Plan{ ID: fmt.Sprintf("plan-%d", time.Now().UnixNano()), Repository: m.Owner + "/" + m.Repo, Branch: m.Base, Commit: m.Commit, Created: time.Now().UTC(), State: inventory.PlanBuilding, Tiers: tiers, Modules: modules, } } // supersededBy is what a newer plan takes over from the open plans it supersedes (novox/hq issue // 254, ADR 0218): the modules they had not finished, and those plans closed as superseded. // // **A merge looked at no plan but its own.** Two merges of one repository a few minutes apart were // two open plans asking for the same modules, each sending machines what it built; and a plan that // would never move again — waiting on a report that could not come, at 97b1b2b — stayed open for // ever beside the newer ones, read as work in progress by everyone who looked. The newer merge is the // newer intent for that repository and branch, so its plan takes over: every open plan of the same // repository and branch **created before it** — by the time the plans were made, never by comparing // commits, which have no order of their own — gives up the modules it had not built, and those are // planned again in the newer plan beside what the newer merge moved. // // "Not built" is a module not yet asked, or asked and not answered; **and a module built and not // yet sent to its machines**, where its policy rolls it out: closed, the older plan would never send // it, and the catalogue announces no move for a rebuild (issue 189), so the newer plan builds and // sends it. A build the older plan asked still finishes and registers as any build does — ordered by // when it was asked (issue 219), so the newer plan's ask, made later, is the one that stands. // // A plan with no branch recorded is from before branches were kept, and is superseded by the next // plan of its repository: what it had not built is folded in, so nothing is lost by it. func supersededBy(newer inventory.Plan, open []inventory.Plan, rollsOut func(string) bool) ([]string, []inventory.Plan) { folded := map[string]bool{} var closed []inventory.Plan for _, old := range open { if old.ID == newer.ID || !old.Open() || !strings.EqualFold(old.Repository, newer.Repository) || (old.Branch != "" && old.Branch != newer.Branch) || !old.Created.Before(newer.Created) { continue } var took []string for name, s := range old.Modules { if s != nil && s.State == planDeleted { continue // deleted at its source: nothing to plan again } if s == nil || s.State != "built" || (s.SentAt == nil && rollsOut(name)) { folded[name] = true took = append(took, name) } } sort.Strings(took) old.State = inventory.PlanSuperseded old.Note = fmt.Sprintf("superseded at tier %d by %s (%s at %s)", old.Tier, newer.ID, newer.Repository, short(newer.Commit)) if len(took) > 0 { old.Note += "; " + strings.Join(took, ", ") + " planned there again" } closed = append(closed, old) } out := make([]string, 0, len(folded)) for name := range folded { out = append(out, name) } sort.Strings(out) return out, closed } // gates is what the next tier needs running from this one: a module of the tier that a later // tier is built by — the runtime dependency — and whose policy rolls it out, must be applied by // the machines running it before the next tier is asked. A base an image stands on need only be // built; a source another module packages need not even be that. func gates(p inventory.Plan, edges []inventory.Edge, rollsOut func(string) bool) []string { if p.Tier >= len(p.Tiers) { return nil } inTier := map[string]bool{} all := map[string]bool{} for _, tier := range p.Tiers { for _, m := range tier { all[m] = true } } for _, m := range p.Tiers[p.Tier] { inTier[m] = true } later := map[string]bool{} for _, tier := range p.Tiers[p.Tier+1:] { for _, m := range tier { later[m] = true } } seen := map[string]bool{} var out []string for _, e := range edges { if later[e.From] && inTier[e.To] && e.Kind == inventory.EdgeBuiltBy && !seen[e.To] && rollsOut(e.To) && !isBaseOf(e.From, e.To, edges, all) { seen[e.To] = true out = append(out, e.To) } } sort.Strings(out) return out } // applied says whether every machine running the module has reported since the module was built. // appliedEach is applied with a moment of its own for each machine: the reports that count are the ones // after that machine was sent the build (novox/hq issue 256). func appliedEach(module string, since func(node string) time.Time, running []string, reports []inventory.Reported) (bool, []string) { at := map[string]*time.Time{} for _, r := range reports { at[r.Node] = r.At } var waiting []string for _, n := range running { if t := at[n]; t == nil || t.Before(since(n)) { waiting = append(waiting, n) } } return len(waiting) == 0, waiting } func applied(module string, builtAt time.Time, running []string, reports []inventory.Reported) (bool, []string) { at := map[string]*time.Time{} for _, r := range reports { at[r.Node] = r.At } var waiting []string for _, n := range running { if t := at[n]; t == nil || t.Before(builtAt) { waiting = append(waiting, n) } } return len(waiting) == 0, waiting } // askTier asks the build machine for every module of the tier, and marks each asked. A module // the catalogue no longer holds, or whose ask could not be made, is a failure of the plan: a tier // half asked is a tier that will never complete. func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) error { entries, err := inv.Catalogued(ctx) if err != nil { return err } byName := map[string]inventory.Entry{} for _, e := range entries { byName[e.Manifest.Module] = e } for _, name := range p.Tiers[p.Tier] { askModule(ctx, p, name, byName) } return nil } // askABuild is how a plan asks for one build, not waited for, and learns the id it asked under. A // variable so a test of what a plan does around an ask needs no build machine. var askABuild = func(ctx context.Context, source buildSource, path, ref string) (string, error) { return buildOneAsked(ctx, source, path, ref, 0, false) } // askModule asks the build machine for one module of a plan and marks it asked, with the id it was // asked under (novox/hq ADR 0219) — or failed, with the plan, when it could not be asked. func askModule(ctx context.Context, p *inventory.Plan, name string, byName map[string]inventory.Entry) { now := time.Now().UTC() state := p.Modules[name] if state == nil { state = &inventory.PlanModule{} p.Modules[name] = state } e, known := byName[name] if !known { state.State = "failed" state.Why = "no longer in the catalogue" p.State = inventory.PlanFailed p.Note = name + " is no longer in the catalogue" return } source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat} fmt.Printf(" tier %d: ", p.Tier) // The branch it follows, never a commit a build once named (novox/hq 04-ISSUES/215). id, err := askABuild(ctx, source, e.Source.Path, followedBranch(e.Source.Ref)) if err != nil { state.State = "failed" state.Why = err.Error() p.State = inventory.PlanFailed p.Note = fmt.Sprintf("%s could not be asked for: %v", name, err) return } state.State = "asked" state.AskedAt = &now state.Build = id state.Why, state.Commit, state.BuiltAt = "", "", nil } // planBuilt marks a module built (or failed) in every open plan whose current tier holds it, and // advances what that completes. Called from the daemon's take-in of every outcome. // // **Only a build asked at or after the plan's ask is its outcome** (novox/hq 04-ISSUES/219). Two // plans a few minutes apart both ask for a module; the earlier plan's build, finishing late, is not // the later plan's answer — it stood on the bases from before the later plan's merge, and taking it // would send machines, and the next tier, what the later merge replaced. asked is zero when the // build's request time is not known, and such an outcome is taken as before. func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time, id string) { inv := open.inventory // One controller works the plans at a time (novox/hq issue 213); an outcome waits its turn rather // than write over what the holder is about to save. Not taken, it is still in the build records, // which the holder settles the plan from (issue 214). release, err := inv.HoldPlans(ctx, true) if err != nil { fmt.Printf("plans: %s's outcome is left to the build records: %v\n", module, err) return } defer release() plans, err := inv.OpenPlans(ctx) if err != nil { fmt.Printf("plans: cannot read them: %v\n", err) return } now := time.Now().UTC() for i := range plans { p := &plans[i] if p.Tier >= len(p.Tiers) { continue } inTier := false for _, m := range p.Tiers[p.Tier] { if m == module { inTier = true } } if !inTier { continue } state := p.Modules[module] if state == nil { state = &inventory.PlanModule{} p.Modules[module] = state } // **The plan's own ask is its outcome, by id** (novox/hq ADR 0219); another build of the module // is, as before, when it was asked at or after the plan's ask (issue 219). // **A module asked under an id is answered by that id's outcome and no other** (novox/hq ADR // 0219): a replay, a rebuild beside the plan, or an older ask finishing late is somebody else's // build, made from other source, and settling the plan with it would send that. A plan from // before ids were kept is matched as it was: by when the build was asked (issue 219). if state.Build != "" { if state.Build != id { continue } } else if askedBefore(asked, state.AskedAt) { continue } if failed != "" && deletedAtSource(failed) { // Deleted at its source by the merge, not broken (novox/hq ADR 0236): the plan goes on. state.State, state.Why = planDeleted, "deleted at its source: "+firstLine(failed) forgetDeleted(ctx, inv, module, p) } else if failed != "" { state.State = "failed" state.Why = failed p.State = inventory.PlanFailed p.Note = fmt.Sprintf("%s failed to build in tier %d", module, p.Tier) sayUnsent(p, func(m string) bool { u, err := inv.UpgradeOf(ctx, m) return err == nil && u.RollOut }) } else { state.State = "built" state.BuiltAt = &now state.Commit = commit } if err := inv.SavePlan(ctx, p); err != nil { fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err) continue } if p.State == inventory.PlanFailed { fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note) } } advanceHeld(ctx, open) } // advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built // and whose gates are applied gives way to the next; the last tier done is the plan done. Called // after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes. // // **One controller at a time** (novox/hq issue 213). A plan is read, changed and saved whole; two // controllers — the old and the new while a machine hands its controller over — would each ask a // tier the other had just asked. Taken without waiting: whoever holds the plans is moving them. func advancePlans(ctx context.Context, open *stores) { release, err := open.inventory.HoldPlans(ctx, false) if err != nil { if !errors.Is(err, inventory.ErrPlansBusy) { fmt.Printf("plans: cannot hold them: %v\n", err) } 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) } // advanceHeld is advancePlans for a caller already holding the plans. func advanceHeld(ctx context.Context, open *stores) { inv := open.inventory plans, err := inv.OpenPlans(ctx) if err != nil { fmt.Printf("plans: cannot read them: %v\n", err) return } if len(plans) == 0 { return } edges, err := inv.Dependencies(ctx) if err != nil { fmt.Printf("plans: cannot read the dependencies: %v\n", err) return } rollsOut := func(module string) bool { u, err := inv.UpgradeOf(ctx, module) return err == nil && u.RollOut } for i := range plans { p := &plans[i] for p.Open() { // **Saved only when the step changed it** (novox/hq issue 296). Every step was saved, and // the steps run on every tick and after every build outcome on the mesh: a plan whose tier // stood still for half an hour was written every few seconds, its revision moved, `plan-moved` // said on the bus, and its last save read as when it began building. before := planSnapshot(*p) moved, err := advanceOnce(ctx, open, p, edges, rollsOut) if err != nil { // Kept in the plan, so `plans` says why it has not moved rather than the log alone; // the state is left as it was and the step is tried again on the next tick. Said and // kept when it is new: the same refusal on every tick is one fact, not one per tick. p.Note = "tier " + fmt.Sprint(p.Tier) + ": " + err.Error() + " — tried again" if planSnapshot(*p) != before { fmt.Printf("%s: %v\n", p.ID, err) if err := inv.SavePlan(ctx, p); err != nil { fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err) } } break } if p.State == inventory.PlanFailed { sayUnsent(p, rollsOut) } if planSnapshot(*p) != before { if err := inv.SavePlan(ctx, p); err != nil { fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err) break } } if !moved { break } } } } // planSnapshot is a plan as a save would write it, to compare before and after a step: what the store // stamps itself (when it was saved, its revision and epoch, when its tier was entered) is left out. func planSnapshot(p inventory.Plan) string { p.Updated, p.Revision, p.Epoch, p.TierEntered = time.Time{}, 0, 0, time.Time{} b, err := json.Marshal(p) if err != nil { // Unreadable is never the same: the plan is saved, as every step was before. return fmt.Sprintf("unmarshallable %p %v", &p, err) } return string(b) } // advanceOnce takes one step of one plan and says whether anything changed. 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) } // **A walk that waits for its delivery's word asks nothing** (novox/hq ADR 0239): nothing of it is // registered before its turn, so no other send carries it to a machine. if p.Waiting() { note := waitingNote(*p) changed := p.Note != note p.Note = note return changed, nil } 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)) return true, nil } tier := p.Tiers[p.Tier] // Not yet asked: ask. unasked := 0 for _, m := range tier { if s := p.Modules[m]; s == nil || s.State == "" { unasked++ } } if unasked == len(tier) { if err := askTier(ctx, inv, p); err != nil { return false, err } return true, nil } // **Asked: settle from the build records first** (novox/hq 04-ISSUES/214). An outcome is taken // in by whichever controller hears it, and a merge to the controller's own repository replaces // the controller in its first tier: the build that produced the new one is recorded, and the // plan never hears it. The record is the fact; a build recorded after the ask is that tier's // outcome, whoever was listening. recorded := map[string][]inventory.Build{} byID := map[string]inventory.Build{} for _, m := range tier { if s := p.Modules[m]; s != nil && s.State == "asked" { builds, err := inv.Builds(ctx, m, 5) if err != nil { return false, err } recorded[m] = builds // Its own ask's record, by id — found even when the outcome named no module (ADR 0219). if s.Build != "" { b, found, err := inv.BuildByID(ctx, s.Build) if err != nil { return false, err } if found { byID[s.Build] = b } } } } if settleFromRecords(p, tier, recorded, byID) { return true, nil } // Asked: wait for every build. var latest time.Time for _, m := range tier { s := p.Modules[m] if s == nil || (s.State != "built" && s.State != planDeleted) { return false, nil } if s.BuiltAt != nil && s.BuiltAt.After(latest) { latest = *s.BuiltAt } } // Built: send every module of the tier whose policy rolls out, once, to the machines running // it — whether or not its source commit moved. A dependent rebuilt because its base moved, or // a module that packages another repository's source, keeps its commit; the catalogue announces // no move for it and its machines would keep the old image until somebody pushed (novox/hq // issue 189). A module whose policy records is built and left, as its policy says. // // **One machine first, unless the module's policy says together** (novox/hq issue 249, ADR // 0218). The plan sent every machine running the module at once, and the operator's policy — // one at a time, stopping at the first that fails, which an announced upgrade honours — was not // read here at all: a module whose new declarations broke it broke everywhere in the same // 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) { continue } running, err := inv.Running(ctx, m) 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. if state.FirstAt == nil { if same, err := unchangedOnEvery(ctx, inv, m, state.Build, running); err != nil { return false, err } else if same { now := time.Now().UTC() state.SentAt = &now state.Why = "no move: made from the source every machine running it runs" fmt.Printf("%s: %s built from the source every machine running it runs: no move, nothing sent\n", p.ID, m) continue } } policy, err := inv.UpgradeOf(ctx, m) if err != nil { return false, err } var reports []inventory.Reported if !policy.Together { // Read for the choice of the first machine as well as for its report. if reports, err = inv.LastReports(ctx); err != nil { return false, err } } 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(m, 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 // 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 != "": 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 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 } switch verdict { case "": if putBackBroken(ctx, open, p, state.Gate, m, state) { return true, nil } pending = append(pending, fmt.Sprintf("%s judged on %s: %s", m, strings.Join(state.Gate.Machines, ", "), gateLine(state.Gate))) continue case inventory.GateFailed: failFirstSend(ctx, open, p, m, state, state.Gate.Machines, state.Gate.Why, step.rest) return true, nil case inventory.GatePassed: if !state.Gate.Kept { gatePassed(ctx, open, p, m, state) state.Gate.Kept = true } } } switch { case len(step.send) == 0: // No machine runs it: nothing to send, and nothing to wait for. state.SentAt = &now continue case step.first: firstOn[step.send[0]] = append(firstOn[step.send[0]], m) continue } 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) } 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) } return true, nil } } // 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 { note := "tier " + fmt.Sprint(p.Tier) + " built; waiting for " + strings.Join(pending, "; ") + " to report it applied before the rest are sent" changed := p.State != inventory.PlanRolling || p.Note != note p.State = inventory.PlanRolling p.Note = note return changed, nil } // And wait for what the next tier needs running. needed := gates(*p, edges, rollsOut) if len(needed) > 0 { reports, err := inv.LastReports(ctx) if err != nil { return false, err } var waiting []string for _, m := range needed { running, err := inv.Running(ctx, m) if err != nil { return false, err } state := p.Modules[m] if state == nil { state = &inventory.PlanModule{} p.Modules[m] = state } if state.State == planDeleted { continue } // The plan sends what it waits for. A rebuild from the same source commit is not a // move the catalogue announces — the build machine rebuilt for a controller change // is one — so the roll-out that opens this gate is the plan's to make, once, and // the reports that open it are the ones after the send. since := latest if state.BuiltAt != nil { since = *state.BuiltAt } // **Each machine from its own send** (novox/hq issue 256). With one machine first (ADR 0218) // a module is sent twice — the first machine, then the rest — and SentAt is the second. // Asked of every machine, the first machine's report, made between the two sends, read as // older than the build, and the gate waited for a report it already had, for ever. sinceFor := func(node string) time.Time { at := since sent := state.SentAt for _, n := range state.First { if n == node { sent = state.FirstAt } } if sent != nil && sent.After(at) { at = *sent } return at } if ok, on := appliedEach(m, sinceFor, running, reports); !ok { waiting = append(waiting, fmt.Sprintf("%s on %s", m, strings.Join(on, ", "))) } } if len(waiting) > 0 { note := "tier " + fmt.Sprint(p.Tier) + " built; waiting for " + strings.Join(waiting, "; ") + " to be applied" changed := p.State != inventory.PlanRolling || p.Note != note p.State = inventory.PlanRolling p.Note = note return changed, nil } } p.Tier++ p.State = inventory.PlanBuilding p.Note = "" if p.Tier < len(p.Tiers) { fmt.Printf("%s: tier %d done; asking tier %d: %s\n", p.ID, p.Tier-1, p.Tier, strings.Join(p.Tiers[p.Tier], ", ")) } 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, ", ")) // What the send recreates, said with it (novox/hq ADR 0245). if said := recreationsSaid(carried); said != "" { p.Note += "; " + said fmt.Printf("%s: %s\n", p.ID, said) } 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(module string, 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+": "+whyFor(g, module), 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 && slices.Contains(g.Returned, module) { // Put back at once when it broke: its rollback was made then, and is not made again. state.Why = "put back when it broke; its send failed: " + g.Why p.State = inventory.PlanFailed p.Note = fmt.Sprintf("the send to %s in tier %d failed its gate: %s; %s was put back when it broke", strings.Join(g.Machines, ", "), p.Tier, g.Why, module) } else if g == nil || g.Verdict == "" || len(g.Failing) == 0 || slices.Contains(g.Failing, module) { gateFailed(batched, open, p, module, state, machines, why) } else if slices.Contains(g.Passing, module) && state.Build != "" { // **The module the gate is kept on keeps its own pass too** (issue 318 review): healthy for the passes // asked while another of its send failed, its build has a verdict, and other walks do not wait on it. if err := open.inventory.RecordGate(ctx, inventory.GateVerdict{Build: state.Build, Module: module, Commit: state.Commit, Previous: state.Previous, Plan: p.ID, Machines: g.Machines, Verdict: inventory.GatePassed, Why: passedAloneWhy(g, module), Component: g.Component, JudgingFrom: g.Since}); err != nil { fmt.Printf("%s: %s passed its gate on its own, and the verdict could not be kept: %v\n", p.ID, module, err) } state.Why = "passed on its own; 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 passed on its own and is kept", strings.Join(g.Machines, ", "), p.Tier, g.Why, module) } 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. send []string first bool // waiting names the first machines whose report the rest wait for. waiting string // failed says how a first machine did not take it; rest is what is then left alone. failed string rest []string } // nextRollout is the next step of one module's rollout in a plan (novox/hq issue 249, ADR 0218). // // Together, every machine running it at once, as the policy says. Otherwise one machine first — the // first by name among those that have reported within the bound, so a laptop that is away is not // the one the rest wait on; the first by name when none has; the same choice on every controller and // every resume — and the rest once each machine the first send reached reports, about the // declaration it was last sent, that it applied it. Compared by the store's own record of what was // sent (Reported.Current), never by this controller's clock against the machine's. // // **A first machine that fails, refuses, or does not report within the bound stops the rollout // there** (ADR 0218 §2), naming the machine; the rest are not sent. A wait with no end is not a // rollout: it held the plan open for ever, read as work in progress (issue 254). func nextRollout(s inventory.PlanModule, running []string, together bool, reports []inventory.Reported, now time.Time, bound time.Duration) rolloutStep { if len(running) == 0 { return rolloutStep{} } if together { return rolloutStep{send: running} } byNode := map[string]inventory.Reported{} for _, r := range reports { byNode[r.Node] = r } if s.FirstAt == nil { sorted := append([]string{}, running...) sort.Strings(sorted) for _, n := range sorted { if r, said := byNode[n]; said && r.At != nil && now.Sub(*r.At) <= bound { return rolloutStep{send: []string{n}, first: true} } } return rolloutStep{send: sorted[:1], first: true} } sentFirst := map[string]bool{} for _, n := range s.First { sentFirst[n] = true } var rest []string for _, n := range running { if !sentFirst[n] { rest = append(rest, n) } } var waiting, failed []string for _, n := range s.First { r, said := byNode[n] // Only a report about what it was last sent says anything about this build. if !said || r.At == nil || !r.Current { waiting = append(waiting, n) continue } switch r.Outcome { case inventory.OutcomeApplied: case inventory.OutcomeFailed, inventory.OutcomeRefused: failed = append(failed, n+" "+r.Outcome+" what it was sent") default: waiting = append(waiting, n) } } if len(failed) > 0 { return rolloutStep{failed: strings.Join(failed, "; "), rest: rest} } if len(waiting) > 0 { if now.Sub(*s.FirstAt) > bound { return rolloutStep{failed: fmt.Sprintf("%s did not report it applied within %s", strings.Join(waiting, ", "), bound), rest: rest} } return rolloutStep{waiting: strings.Join(waiting, ", ")} } return rolloutStep{send: rest} } // sayUnsent adds to an ended plan's note the modules it built and never sent (novox/hq issue 249). // An announced move of a module a plan held was left to that plan; a plan that ends without // sending it — failed elsewhere, or closed by hand — would leave its machines behind with nothing // saying so. A module whose rollout stopped at its first machine is not among them: that stop was // the point. Said once. func sayUnsent(p *inventory.Plan, rollsOut func(string) bool) { var unsent []string for name, s := range p.Modules { if s != nil && s.State == "built" && s.SentAt == nil && s.FirstAt == nil && rollsOut(name) { unsent = append(unsent, name) } } if len(unsent) == 0 || strings.Contains(p.Note, "built and never sent") { return } sort.Strings(unsent) p.Note += "; built and never sent: " + strings.Join(unsent, ", ") + " — `push --behind` sends them" } // planTicker advances open plans on a timer, for the steps outcomes alone cannot take. func planTicker(ctx context.Context, open *stores) { advancePlans(ctx, open) tick := time.NewTicker(30 * time.Second) defer tick.Stop() for { select { case <-ctx.Done(): return case <-tick.C: advancePlans(ctx, open) } } } // planLine is one plan as `status` says it, given the least bound of a tier. func planLine(p inventory.Plan, now time.Time) string { return planLineWith(p, now, pauseView{}, tierAtLeast) } // inTierSince is when a plan began waiting on what it waits on now, which `plans`, `status` and the // watchdog of a plan's progress (S3) all count from (novox/hq issue 296): a walk waiting for its // delivery's word from when it was opened, as S16 counts; any other plan from when it entered its tier, // or, when the delivery let it go after that, from the go. Never from the plan's last save: a plan is // saved for many reasons while its tier stands still, and a clock reset by each save said a build queued // for minutes had been building for a few seconds. A plan read without the stamp (one from before it // was kept) counts from its last save, the only time it has. func inTierSince(p inventory.Plan) time.Time { if p.Waiting() { return p.Created } since := p.TierEntered if since.IsZero() { since = p.Updated } if d := p.Delivery; d != nil && d.Go != nil && d.Go.After(since) { since = *d.Go } return since } // planLineWith is planLine knowing whether the build seat is paused (novox/hq ADR 0219): a plan // waiting on builds nobody will take until a person resumes the seat says so, and is not late. bound is // the plan's tier bound, the one its stalled condition is raised at (tierBounds): LATE is that condition // said on the line (novox/hq issue 296). func planLineWith(p inventory.Plan, now time.Time, pause pauseView, bound time.Duration) string { where := fmt.Sprintf("tier %d of %d", min(p.Tier+1, len(p.Tiers)), len(p.Tiers)) switch p.State { case inventory.PlanDone: return fmt.Sprintf("%s %s done, %d tier(s)", p.Repository, short(p.Commit), len(p.Tiers)) case inventory.PlanFailed: return fmt.Sprintf("%s %s FAILED at %s: %s", p.Repository, short(p.Commit), where, p.Note) case inventory.PlanSuperseded: return fmt.Sprintf("%s %s %s", p.Repository, short(p.Commit), p.Note) } since := now.Sub(inTierSince(p)).Round(time.Second) if p.Waiting() { // Waiting for its delivery's word is no lateness of the walk's (novox/hq ADR 0239). return fmt.Sprintf("%s %s %s, %s, for %s", p.Repository, short(p.Commit), where, waitingNote(p), since) } if waiting, paused := pausedWaiting(p, pause, now); paused { return fmt.Sprintf("%s %s %s, %s", p.Repository, short(p.Commit), where, waiting) } late := "" if since > bound { late = " — LATE" } what := "building" if p.State == inventory.PlanRolling { what = p.Note } return fmt.Sprintf("%s %s %s, %s for %s%s", p.Repository, short(p.Commit), where, what, since, late) } // planFailedBuild marks the module a failed build was for when the result names no module: by the // repository and path the plan's modules were asked at. // // **By the id first** (novox/hq ADR 0219): a plan keeps the id it asked each module under, so an // outcome that never learnt its module's name — cancelled, killed, failed at the clone — is matched // to the module it was asked for exactly. Repository and path remain for a plan from before ids // were kept. func planFailedBuild(ctx context.Context, open *stores, result link.BuildResult) { asked, _ := link.BuildAskedAt(result.ID) if plans, err := open.inventory.OpenPlans(ctx); err == nil { if module := moduleAskedAs(plans, result.ID); module != "" { planBuilt(ctx, open, module, result.Commit, result.Failed, asked, result.ID) return } } entries, err := open.inventory.Catalogued(ctx) if err != nil { return } for _, e := range entries { if repositoryMatches(e.Source.Repository, result.Repository) && e.Source.Path == result.Path { planBuilt(ctx, open, e.Manifest.Module, result.Commit, result.Failed, asked, result.ID) return } } } // moduleAskedAs is the module an open plan's current tier asked for under this id, or nothing. func moduleAskedAs(plans []inventory.Plan, id string) string { if id == "" { return "" } for _, p := range plans { if p.Tier >= len(p.Tiers) { continue } for _, m := range p.Tiers[p.Tier] { if s := p.Modules[m]; s != nil && s.Build == id { return m } } } return "" } func repositoryMatches(a, b string) bool { trim := func(s string) string { return strings.ToLower(strings.TrimSuffix(s, ".git")) } return trim(a) == trim(b) || strings.HasSuffix(trim(a), "/"+trim(b)) || strings.HasSuffix(trim(b), "/"+trim(a)) } // planStatus is one plan as `status --json` says it. type planStatus struct { ID string `json:"id"` Repository string `json:"repository"` Commit string `json:"commit"` State string `json:"state"` Tier int `json:"tier"` Tiers int `json:"tiers"` Waiting string `json:"waiting,omitempty"` Since time.Time `json:"since"` Late bool `json:"late"` } func planStatuses(plans []inventory.Plan, now time.Time, pause pauseView, bounds tierBounds) []planStatus { out := make([]planStatus, 0, len(plans)) for _, p := range plans { ps := planStatus{ID: p.ID, Repository: p.Repository, Commit: p.Commit, State: p.State, Tier: p.Tier, Tiers: len(p.Tiers), Since: p.Updated} if p.Open() { ps.Since = inTierSince(p) ps.Waiting = p.Note if ps.Waiting == "" { ps.Waiting = "builds of tier " + fmt.Sprint(p.Tier) } // A walk waiting for its delivery's word is no tier late (ADR 0239): S16 is its watchdog. ps.Late = !p.Waiting() && now.Sub(ps.Since) > bounds.of(p.Repository) // Paused is a person's decision, not lateness (novox/hq ADR 0219). if waiting, paused := pausedWaiting(p, pause, now); paused { ps.Waiting, ps.Late = waiting, false } } out = append(out, ps) } return out } // openPlans is the open plans among the recent ones, and how many have waited past the bound. func openPlans(plans []inventory.Plan, now time.Time, pause pauseView, bounds tierBounds) ([]inventory.Plan, int) { var open []inventory.Plan late := 0 for _, p := range plans { if p.Open() { open = append(open, p) if _, paused := pausedWaiting(p, pause, now); paused || p.Waiting() { continue } if now.Sub(inTierSince(p)) > bounds.of(p.Repository) { late++ } } } return open, late } // plansCommand says what the last merges produced and where each stands; given an id, one plan // tier by tier with every module's state. func plansCommand(ctx context.Context, args []string) error { set := flag.NewFlagSet("plans", flag.ContinueOnError) limit := set.Int("n", 10, "how many to show") whatIf := set.String("what-if", "", "owner/repository: the plan a merge there would produce, saving nothing — with --paths or --modules") paths := set.String("paths", "", "the files the merge would change, comma-separated, from the repository's root") modules := set.String("modules", "", "or the modules it would change, comma-separated") moduleDirs := set.String("module-dirs", "", "with --paths: the directories holding a module.json at the commit, "+ "comma-separated, as the forge's announcer says them (novox/hq issue 278)") // Ending a plan by hand is a repair, and says why (novox/hq to-be 45 §7). why := addHandActFlags(set) positionals, err := parseAround(set, args) if err != nil { return err } if len(positionals) == 2 && (positionals[0] == "stop" || positionals[0] == "close" || positionals[0] == "go") { // Refused before anything is opened: a repair by hand says why. if err := why.require("plans " + positionals[0]); err != nil { return err } } open, err := openStores(ctx) if err != nil { return err } defer open.Close() inv := open.inventory now := time.Now() if len(positionals) == 1 { p, err := inv.PlanByID(ctx, positionals[0]) if err != nil { return err } bounds := readTierBounds(ctx, inv, now) fmt.Printf("%s — %s\n", p.ID, planLineWith(p, now, buildSeatPause(ctx, inv, []inventory.Plan{p}), bounds.of(p.Repository))) 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)) if c.Recreates != "" { fmt.Printf(" %-22s %s\n", "", c.Recreates) } } } return nil } for i, tier := range p.Tiers { marker := " " if i == p.Tier && p.Open() { marker = ">" } fmt.Printf("%s tier %d\n", marker, i) for _, m := range tier { s := p.Modules[m] state := "not yet asked" if s != nil && s.State != "" { state = s.State if s.Commit != "" { state += " from " + short(s.Commit) } if s.Build != "" && s.State != "built" { state += " (" + s.Build + ")" } if s.Why != "" { state += ": " + s.Why } } 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. 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)) } // What the send to its first machine did to its containers (ADR 0245). if s != nil && s.Gate != nil { for _, c := range s.Gate.Carried { if c.Recreates != "" { fmt.Printf(" %-22s on %s %s\n", "", c.Node, c.Recreates) } } } } } return nil } if *whatIf != "" { return planWhatIf(ctx, inv, *whatIf, splitList(*paths), splitList(*modules), moduleDirsOf(*moduleDirs)) } // `go` (novox/hq ADR 0239): a person's word in place of the delivery's owner, for a walk that waits for it // — mesh-delivery down, or a person who will not wait. Recorded in the hand-act log with why. if len(positionals) == 2 && positionals[0] == "go" { why.record(ctx, "plans go", positionals[1:]) var p inventory.Plan if err := sayingOnTheBus(ctx, func() (err error) { p, err = letGo(ctx, inv, positionals[1], "a person ("+whoAsked()+")", strings.TrimSpace(*why.why)) return err }); err != nil { return err } fmt.Printf("%s (%s at %s) is let go by hand; its first tier is asked at the next pass\n", p.ID, p.Repository, short(p.Commit)) return nil } // `retry` (novox/hq ADR 0219): a failed plan's failed builds asked again, and the plan goes on. if len(positionals) == 2 && positionals[0] == "retry" { said, err := retryPlan(ctx, open, positionals[1]) if err != nil { return err } fmt.Println(said) return nil } // `stop`, or `close` (novox/hq issue 254): a person ending a plan that will not move again — one // waiting on a report that cannot come — so it stops reading as work in progress. Marked failed // with who ended it; what it asked still builds and registers. if len(positionals) == 2 && (positionals[0] == "stop" || positionals[0] == "close") { how := "stopped" if positionals[0] == "close" { how = "closed" } release, err := inv.HoldPlans(ctx, true) if err != nil { return err } defer release() p, err := inv.PlanByID(ctx, positionals[1]) if err != nil { return err } if !p.Open() { return fmt.Errorf("%s is already %s", p.ID, p.State) } why.record(ctx, "plans "+positionals[0], positionals[1:]) p.State = inventory.PlanFailed p.Note = how + " by hand at tier " + fmt.Sprint(p.Tier) + ": " + strings.TrimSpace(*why.why) sayUnsent(&p, func(m string) bool { u, err := inv.UpgradeOf(ctx, m) return err == nil && u.RollOut }) if err := inv.SavePlan(ctx, &p); err != nil { return err } fmt.Printf("%s %s at tier %d of %d; what was asked still builds and registers, nothing further is asked\n", p.ID, how, p.Tier, len(p.Tiers)) return nil } plans, err := inv.RecentPlans(ctx, *limit) if err != nil { return err } if len(plans) == 0 { fmt.Println("no merge has produced a plan yet") return nil } pause := buildSeatPause(ctx, inv, plans) bounds := readTierBounds(ctx, inv, now) for _, p := range plans { fmt.Printf("%-28s %s\n", p.ID, planLineWith(p, now, pause, bounds.of(p.Repository))) } return nil } // planWhatIf is the plan a merge would produce, computed the way the merge handler computes one // and saved nowhere: the modules the repository's changed files touch (or the modules named), what // packages their source, everything reachable from them, in tiers. For reading before merging. func planWhatIf(ctx context.Context, inv *inventory.Inventory, repository string, paths, modules []string, moduleDirs *[]string) error { owner, repo, found := strings.Cut(repository, "/") if !found { return fmt.Errorf("--what-if takes owner/repository, not %q", repository) } m := link.SourceMoved{Owner: owner, Repo: repo, Base: "main", Commit: "what-if", Paths: paths} if moduleDirs != nil { m.ModuleDirs, m.ModuleDirsSaid = *moduleDirs, true } entries, err := inv.Catalogued(ctx) if err != nil { return err } read, err := inv.ReadRepositories(ctx) if err != nil { return err } var from, packaging []inventory.Entry named := map[string]bool{} for _, name := range modules { named[name] = true } if len(named) == 0 { // The changed files: the planner's own answer, as the merge handler, the merge gate and a pull // request's check ask it (novox/hq ADR 0238). r := reachOfMerge(m, entries, read, nil) from, packaging = append(append([]inventory.Entry{}, r.Touched...), r.Deleted...), r.Packaging } else { for _, e := range entries { switch { case named[e.Manifest.Module]: from = append(from, e) case readsFrom(read[e.Manifest.Module], m): packaging = append(packaging, e) } } } moved := append(append([]inventory.Entry{}, from...), packaging...) if len(moved) == 0 { fmt.Printf("a merge of %s changing %s would build nothing the mesh holds\n", repository, orNone(strings.Join(append(paths, modules...), ", "))) return nil } edges, err := inv.Dependencies(ctx) if err != nil { return err } var names []string for _, e := range moved { names = append(names, e.Manifest.Module) } p := planOfMerge(m, names, edges) fmt.Printf("a merge of %s would build %d module(s) in %d tier(s):\n", repository, len(p.Modules), len(p.Tiers)) rolls := map[string]string{} for i, tier := range p.Tiers { fmt.Printf(" tier %d\n", i) for _, name := range tier { how := "built; its policy records, so nothing is sent" if u, err := inv.UpgradeOf(ctx, name); err == nil && u.RollOut { running, _ := inv.Running(ctx, name) reports, _ := inv.LastReports(ctx) how = "built, then sent to " + orNone(strings.Join(running, ", ")) // One machine first unless the policy says together (novox/hq issue 249). if first := nextRollout(inventory.PlanModule{}, running, u.Together, reports, time.Now(), planWaitBound); first.first && len(running) > 1 { how = fmt.Sprintf("built, then sent to %s first and to the rest once it has applied it", first.send[0]) } rolls[name] = how } fmt.Printf(" %-22s %s\n", name, how) } } if hasCycle(p.Tiers, edges) { fmt.Println(" the last tier depends on itself and would be built together, in no order") } if len(packaging) > 0 { var also []string for _, e := range packaging { also = append(also, e.Manifest.Module) } fmt.Printf(" %s package source from %s, so they are rebuilt without their own source moving\n", strings.Join(also, ", "), repository) } return nil } // moduleDirsOf is --module-dirs as an announcer would say it: nothing when not given, so the what-if // reads as an announcement from an announcer that does not say. func moduleDirsOf(s string) *[]string { if strings.TrimSpace(s) == "" { return nil } dirs := splitList(s) return &dirs } func splitList(s string) []string { var out []string for _, part := range strings.Split(s, ",") { if part = strings.TrimSpace(part); part != "" { out = append(out, part) } } return out } // settleFromRecords marks every module of the tier still `asked` built — or failed — from a build // recorded after it was asked, and says whether it changed anything (novox/hq 04-ISSUES/214). // Newest first, as Builds answers: the first record after the ask is the outcome of that ask. // // The record of the plan's own ask, by its id, is that outcome before anything else (novox/hq ADR // 0219): the plan asked under it, and a failure recorded without a module — cancelled, killed — is // found by nothing else. func settleFromRecords(p *inventory.Plan, tier []string, recorded map[string][]inventory.Build, byID map[string]inventory.Build) bool { changed := false for _, m := range tier { s := p.Modules[m] if s == nil || s.State != "asked" || s.AskedAt == nil { continue } var outcome *inventory.Build for i := range recorded[m] { // Asked under an id, only that id's record answers (ADR 0219), and it is looked up below. if s.Build != "" { break } b := recorded[m][i] if b.At.Before(*s.AskedAt) { break } // Recorded after the ask and asked before it: an earlier ask's late outcome, not this // one's (novox/hq 04-ISSUES/219). if askedBefore(b.Asked, s.AskedAt) { continue } outcome = &b } if own, found := byID[s.Build]; s.Build != "" && found { outcome = &own } if outcome == nil { continue } at := outcome.At if !outcome.Worked() && deletedAtSource(outcome.Failed) { s.State, s.Why = planDeleted, "deleted at its source: "+firstLine(outcome.Failed) fmt.Printf("%s: %s was deleted at its source (%s): not built, and the plan goes on\n", p.ID, m, outcome.ID) changed = true continue } if outcome.Worked() { s.State = "built" s.BuiltAt = &at s.Commit = outcome.Commit } else { s.State = "failed" s.Why = outcome.Failed p.State = inventory.PlanFailed p.Note = fmt.Sprintf("%s failed to build in tier %d", m, p.Tier) } fmt.Printf("%s: %s settled from the build records as %s (%s)\n", p.ID, m, s.State, outcome.ID) changed = true } return changed } // askedBefore is whether a build asked at asked was asked before a plan asked for its module — and // so is not that plan's outcome (novox/hq 04-ISSUES/219). False when either time is not known. func askedBefore(asked time.Time, planAsked *time.Time) bool { return !asked.IsZero() && planAsked != nil && asked.Before(*planAsked) } // firstRunning is the machines sent first that run the module — not the machine holding the bus, sent // with them only for its user list. func firstRunning(first, running []string) []string { var out []string for _, n := range first { if slices.Contains(running, n) { out = append(out, n) } } if len(out) == 0 { return first } return out } // planDeleted is a plan's module that the merge deleted at its source: not built, not sent, and no // failure of the plan (novox/hq ADR 0236). const planDeleted = "deleted" // forgetDeleted forgets a module a plan found deleted at its source, where nothing holds it, and says // it otherwise. func forgetDeleted(ctx context.Context, inv *inventory.Inventory, module string, p *inventory.Plan) { switch err := inv.ForgetModule(ctx, module); { case err == nil: fmt.Printf("%s: %s was deleted at its source: forgotten, and the plan goes on\n", p.ID, module) default: fmt.Printf("%s: %s was deleted at its source and is not built; it is kept until a person unassigns it and "+ "`module forget %s`: %v\n", p.ID, module, module, firstLine(err.Error())) } }