The controller's machine moves it from the container to a process by starting the process first and removing the container once the process is up (mesh-host's `replaces`). For that moment two controllers share the store and the bus. Checked what each does: - the seat's verbs: a queue group per seat, each call answered once. Safe. - the controller's consumers on CONTROL and EVENTS: push consumers with no delivery group, so the second bind is refused with "consumer is already bound" and serve exited. The process would restart for ever, the host would never see it up, and the container would never go. The second controller now stands by and binds when the first lets go (tested on a real bus; fails without the change). - plans: read, changed and saved whole by the 30s timer, by build outcomes, by a merge and by `plans stop`. Two timers would each ask a tier the other had just asked. Working the plans now takes a session-level advisory lock on the inventory: the timer skips while another holds it, the other paths wait for it. Build asks happen only inside plan work and are covered by the same lock.
876 lines
27 KiB
Go
876 lines
27 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"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,
|
|
Commit: m.Commit,
|
|
Created: time.Now().UTC(),
|
|
State: inventory.PlanBuilding,
|
|
Tiers: tiers,
|
|
Modules: modules,
|
|
}
|
|
}
|
|
|
|
// 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.
|
|
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
|
|
}
|
|
now := time.Now().UTC()
|
|
for _, name := range p.Tiers[p.Tier] {
|
|
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"
|
|
continue
|
|
}
|
|
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).
|
|
if err := buildOne(ctx, source, e.Source.Path, followedBranch(e.Source.Ref), 0); 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)
|
|
continue
|
|
}
|
|
state.State = "asked"
|
|
state.AskedAt = &now
|
|
}
|
|
return 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) {
|
|
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
|
|
}
|
|
if askedBefore(asked, state.AskedAt) {
|
|
continue
|
|
}
|
|
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)
|
|
} 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()
|
|
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() {
|
|
moved, err := advanceOnce(ctx, open, p, edges, rollsOut)
|
|
if err != nil {
|
|
fmt.Printf("%s: %v\n", p.ID, err)
|
|
break
|
|
}
|
|
if err := inv.SavePlan(ctx, *p); err != nil {
|
|
fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err)
|
|
break
|
|
}
|
|
if !moved {
|
|
break
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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.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{}
|
|
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
|
|
}
|
|
}
|
|
if settleFromRecords(p, tier, recorded) {
|
|
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" {
|
|
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.
|
|
for _, m := range tier {
|
|
state := p.Modules[m]
|
|
if state == nil || state.SentAt != nil || !rollsOut(m) {
|
|
continue
|
|
}
|
|
running, err := inv.Running(ctx, m)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
now := time.Now().UTC()
|
|
state.SentAt = &now
|
|
if len(running) == 0 {
|
|
continue
|
|
}
|
|
if err := sendTo(ctx, open, running); err != nil {
|
|
return false, fmt.Errorf("sending %s to %s after tier %d: %w", m, strings.Join(running, ", "), p.Tier, err)
|
|
}
|
|
fmt.Printf("%s: tier %d built; sent %s to %s\n", p.ID, p.Tier, m, strings.Join(running, ", "))
|
|
return true, 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
|
|
}
|
|
// 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
|
|
}
|
|
if state.SentAt != nil && state.SentAt.After(since) {
|
|
since = *state.SentAt
|
|
}
|
|
if ok, on := applied(m, since, 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
|
|
}
|
|
|
|
// 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.
|
|
func planLine(p inventory.Plan, now time.Time) 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)
|
|
}
|
|
since := now.Sub(p.Updated).Round(time.Minute)
|
|
late := ""
|
|
if since > planWaitBound {
|
|
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.
|
|
func planFailedBuild(ctx context.Context, open *stores, result link.BuildResult) {
|
|
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 {
|
|
asked, _ := link.BuildAskedAt(result.ID)
|
|
planBuilt(ctx, open, e.Manifest.Module, result.Commit, result.Failed, asked)
|
|
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) []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.Waiting = p.Note
|
|
if ps.Waiting == "" {
|
|
ps.Waiting = "builds of tier " + fmt.Sprint(p.Tier)
|
|
}
|
|
ps.Late = now.Sub(p.Updated) > planWaitBound
|
|
}
|
|
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) ([]inventory.Plan, int) {
|
|
var open []inventory.Plan
|
|
late := 0
|
|
for _, p := range plans {
|
|
if p.Open() {
|
|
open = append(open, p)
|
|
if time.Since(p.Updated) > planWaitBound {
|
|
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")
|
|
positionals, err := parseAround(set, args)
|
|
if 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
|
|
}
|
|
fmt.Printf("%s — %s\n", p.ID, planLine(p, now))
|
|
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.Why != "" {
|
|
state += ": " + s.Why
|
|
}
|
|
}
|
|
fmt.Printf(" %-22s %s\n", m, state)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
if *whatIf != "" {
|
|
return planWhatIf(ctx, inv, *whatIf, splitList(*paths), splitList(*modules))
|
|
}
|
|
if len(positionals) == 2 && positionals[0] == "stop" {
|
|
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)
|
|
}
|
|
p.State = inventory.PlanFailed
|
|
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
|
|
release, err := inv.HoldPlans(ctx, true)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer release()
|
|
if p, err = inv.PlanByID(ctx, positionals[1]); err != nil {
|
|
return err
|
|
}
|
|
if !p.Open() {
|
|
return fmt.Errorf("%s is already %s", p.ID, p.State)
|
|
}
|
|
p.State = inventory.PlanFailed
|
|
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
|
|
if err := inv.SavePlan(ctx, p); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("%s stopped at tier %d of %d; what was asked still builds and registers, nothing further is asked\n",
|
|
p.ID, 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
|
|
}
|
|
for _, p := range plans {
|
|
fmt.Printf("%-28s %s\n", p.ID, planLine(p, now))
|
|
}
|
|
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) 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}
|
|
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
|
|
}
|
|
for _, e := range entries {
|
|
switch {
|
|
case named[e.Manifest.Module]:
|
|
from = append(from, e)
|
|
case len(named) == 0 && sourceIs(e.Source, m):
|
|
from = append(from, e)
|
|
case readsFrom(read[e.Manifest.Module], m):
|
|
packaging = append(packaging, e)
|
|
}
|
|
}
|
|
if len(named) == 0 {
|
|
from = whatTheMergeTouched(from, entries, m)
|
|
}
|
|
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)
|
|
how = "built, then sent to " + orNone(strings.Join(running, ", "))
|
|
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
|
|
}
|
|
|
|
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.
|
|
func settleFromRecords(p *inventory.Plan, tier []string, recorded 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] {
|
|
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 outcome == nil {
|
|
continue
|
|
}
|
|
at := outcome.At
|
|
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)
|
|
}
|