From a69a8322487bbacebad015de75f3c6dc3102d1f3 Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 1 Oct 2026 17:53:55 +0200 Subject: [PATCH 1/2] A merge produces a tiered plan the mesh keeps (hq ADR 0162) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A module's dependencies are one relation in the catalogue — stands-on, packages, built-by, declared — answered by one call. A merge takes what moved and everything reachable from it, sorts the set into tiers (a code dependency in the same tier, a build dependency after its base is built, a runtime dependency after the build machine is built and running; the build machine's own base comes first, built by the one that runs), writes the plan to the store, asks the first tier and returns. Every outcome advances the plan; a ticker advances what outcomes cannot; a controller replaced mid-plan resumes it. status lists open plans and names one that has waited too long. --- cmd/mesh-controller/build.go | 2 + cmd/mesh-controller/main.go | 9 + cmd/mesh-controller/push.go | 3 + cmd/mesh-controller/readable.go | 5 +- cmd/mesh-controller/release_plan.go | 526 ++++++++++++++++++ cmd/mesh-controller/release_plan_test.go | 109 ++++ cmd/mesh-controller/status.go | 16 + cmd/mesh-controller/upgrades.go | 59 +- internal/inventory/dependencies.go | 124 +++++ internal/inventory/dependencies_test.go | 64 +++ .../0052-a-merge-is-a-plan-the-mesh-keeps.sql | 19 + internal/inventory/plans.go | 123 ++++ 12 files changed, 1023 insertions(+), 36 deletions(-) create mode 100644 cmd/mesh-controller/release_plan.go create mode 100644 cmd/mesh-controller/release_plan_test.go create mode 100644 internal/inventory/dependencies.go create mode 100644 internal/inventory/dependencies_test.go create mode 100644 internal/inventory/migrations/0052-a-merge-is-a-plan-the-mesh-keeps.sql create mode 100644 internal/inventory/plans.go diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 0b1d0ab..fef3b53 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -581,6 +581,8 @@ type answers struct { // pair that answers "has it caught up", which waiting alone cannot (the sent digest is // recorded at send, not at apply). reported []inventory.Reported + // plans is what the last merges produced and where each stands (novox/hq ADR 0162). + plans []inventory.Plan // refused is why a machine cannot be worked out at all, by name. A different thing from every // other answer here: those are about a machine that was told something, and this is about one // that cannot be told anything — it never reaches waiting, because nothing was computed for it diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index 4f6ce60..70346c6 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -253,12 +253,21 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error { switch { case err != nil && result.Failed != "": fmt.Printf("%s: %v\n", result.ID, err) + if result.Module != "" { + planBuilt(ctx, b.inv, result.Module, result.Commit, result.Failed) + } else { + planFailedBuild(ctx, b.inv, result) + } return nil case err != nil: fmt.Printf("%s: heard and recorded, and not registered: %v\n", result.ID, err) + if manifest.Module != "" { + planBuilt(ctx, b.inv, manifest.Module, result.Commit, err.Error()) + } return nil } fmt.Printf("%s: %s %s registered, built on %s from %s\n", result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit)) + planBuilt(ctx, b.inv, manifest.Module, result.Commit, "") return nil } diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 76f56bc..1909add 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -114,6 +114,9 @@ func serve(ctx context.Context) error { // And build results nobody was waiting for. A build triggered any other way than `build` // would otherwise be reported into the void, which is the same as not reporting it. server.Records(builds{inv}) + // Open plans move on a timer as well as on outcomes (novox/hq ADR 0162): a tier waiting for + // machines to report moves when they have, and a plan left by a replaced controller resumes. + go planTicker(ctx, inv) // And what the catalogue decided a build meant. The builder's own result is already handled // above; this is the other half — the control plane is the only one of the three that knows // which machines run the thing, so it is the one that acts (novox/hq ADR 0072). diff --git a/cmd/mesh-controller/readable.go b/cmd/mesh-controller/readable.go index 0c4d0fb..e08ed50 100644 --- a/cmd/mesh-controller/readable.go +++ b/cmd/mesh-controller/readable.go @@ -39,6 +39,9 @@ type meshStatus struct { // whose is older is still working — and Waiting cannot tell those apart, because the sent // digest is recorded at send, not at apply. Reported []machineReported `json:"reported"` + // Plans is what the last merges produced and where each stands (novox/hq ADR 0162): the + // open ones first, each saying its tier, what it waits for, and whether it has waited too long. + Plans []planStatus `json:"plans"` // Unresolved is every machine that cannot be worked out at all, with what the mesh said when // it tried. **A machine here is in none of the lists above**: nothing was computed for it, so // there is nothing to compare it against and nothing it can be behind — which is why a @@ -152,7 +155,7 @@ func statusAsJSON(asked answers) ([]byte, error) { out := meshStatus{Machines: len(nodes), Wrong: []machineDoing{}, Quiet: []machineQuiet{}, Behind: []moduleBehind{}, Waiting: []machineWaiting{}, Reported: []machineReported{}, Unresolved: []machineUnresolved{}, - Network: asked.network, Adopted: adoptedNodes(nodes)} + Network: asked.network, Adopted: adoptedNodes(nodes), Plans: planStatuses(asked.plans, time.Now())} // In a stated order, so two readings of an unchanged mesh are the same document. untakenNodes := make([]string, 0, len(asked.untaken)) for name := range asked.untaken { diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go new file mode 100644 index 0000000..550a3e9 --- /dev/null +++ b/cmd/mesh-controller/release_plan.go @@ -0,0 +1,526 @@ +package main + +import ( + "context" + "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{} + } + 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. 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, which is the only one there could be. + 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) { + 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. +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 { + 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) + if err := buildOne(ctx, source, e.Source.Path, 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. +func planBuilt(ctx context.Context, inv *inventory.Inventory, module, commit, failed string) { + 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 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) + } + } + advancePlans(ctx, inv) +} + +// 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. +func advancePlans(ctx context.Context, inv *inventory.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, inv, 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, inv *inventory.Inventory, p *inventory.Plan, + edges []inventory.Edge, rollsOut func(string) bool) (bool, error) { + 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: 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: 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 + } + builtAt := latest + if s := p.Modules[m]; s != nil && s.BuiltAt != nil { + builtAt = *s.BuiltAt + } + if ok, on := applied(m, builtAt, 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, inv *inventory.Inventory) { + advancePlans(ctx, inv) + tick := time.NewTicker(30 * time.Second) + defer tick.Stop() + for { + select { + case <-ctx.Done(): + return + case <-tick.C: + advancePlans(ctx, inv) + } + } +} + +// 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, inv *inventory.Inventory, result link.BuildResult) { + entries, err := inv.Catalogued(ctx) + if err != nil { + return + } + for _, e := range entries { + if repositoryMatches(e.Source.Repository, result.Repository) && e.Source.Path == result.Path { + planBuilt(ctx, inv, e.Manifest.Module, result.Commit, result.Failed) + 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 +} diff --git a/cmd/mesh-controller/release_plan_test.go b/cmd/mesh-controller/release_plan_test.go new file mode 100644 index 0000000..554f378 --- /dev/null +++ b/cmd/mesh-controller/release_plan_test.go @@ -0,0 +1,109 @@ +package main + +import ( + "testing" + "time" + + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// A merge produces a tiered plan (novox/hq ADR 0162): what moved and everything reachable from it, +// sorted so a tier depends only on earlier ones — with the three kinds of dependency told apart. +func TestAMergeIsPlannedInTiersAlongTheThreeKindsOfDependency(t *testing.T) { + edges := []inventory.Edge{ + // build dependencies: images on the runtime, a plugin on one of them + {From: "shop", To: "mesh-tools", Kind: inventory.EdgeStandsOn}, + {From: "postgres", To: "mesh-tools", Kind: inventory.EdgeStandsOn}, + {From: "shop-plugin", To: "shop", Kind: inventory.EdgeDeclared}, + // a code dependency: the proxy packages the controller's source — same tier + {From: "route-proxy", To: "mesh-controller", Kind: inventory.EdgePackages}, + // runtime dependencies: everything source-built is built by the builder + {From: "shop", To: "builder", Kind: inventory.EdgeBuiltBy}, + {From: "postgres", To: "builder", Kind: inventory.EdgeBuiltBy}, + {From: "shop-plugin", To: "builder", Kind: inventory.EdgeBuiltBy}, + {From: "route-proxy", To: "builder", Kind: inventory.EdgeBuiltBy}, + {From: "mesh-controller", To: "builder", Kind: inventory.EdgeBuiltBy}, + {From: "mesh-tools", To: "builder", Kind: inventory.EdgeBuiltBy}, + {From: "builder", To: "mesh-tools", Kind: inventory.EdgeStandsOn}, + {From: "unrelated", To: "alpine", Kind: inventory.EdgeStandsOn}, + } + // The runtime image moved: everything on it, and what is built by what is on it. + set := reachableFrom([]string{"mesh-tools"}, edges) + want := []string{"builder", "mesh-controller", "mesh-tools", "postgres", "route-proxy", "shop", "shop-plugin"} + if len(set) != len(want) { + t.Fatalf("reachable from the runtime: %v, want %v", set, want) + } + tiers := tiersOf(set, edges) + pos := map[string]int{} + for i, tier := range tiers { + for _, m := range tier { + pos[m] = i + } + } + if pos["mesh-tools"] != 0 || pos["builder"] != 1 { + t.Fatalf("the runtime then the builder: %v", tiers) + } + if !(pos["shop"] > pos["builder"] && pos["postgres"] > pos["builder"] && pos["mesh-controller"] > pos["builder"]) { + t.Fatalf("what the builder builds comes after the builder: %v", tiers) + } + if pos["shop-plugin"] <= pos["shop"] { + t.Fatalf("a plugin after what it is declared on: %v", tiers) + } + if pos["route-proxy"] != pos["mesh-controller"] { + t.Fatalf("a code dependency is rebuilt in the same tier as its source: %v", tiers) + } + if hasCycle(tiers, edges) { + t.Fatalf("no cycle here: %v", tiers) + } + + // The controller alone moved: the proxy with it, nothing else. + small := reachableFrom([]string{"mesh-controller"}, edges) + if len(small) != 2 { + t.Fatalf("a controller merge rebuilds the controller and what packages it: %v", small) + } + if tiers := tiersOf(small, edges); len(tiers) != 1 { + t.Fatalf("both in one tier: %v", tiers) + } + + // Only a runtime dependency gates on deployment, and only when the module rolls out. + p := planOfMerge(link.SourceMoved{Owner: "novox", Repo: "mesh-tools", Commit: "abc"}, []string{"mesh-tools"}, edges) + p.Tier = pos["builder"] + rollsOut := func(m string) bool { return m == "builder" } + if g := gates(p, edges, rollsOut); len(g) != 1 || g[0] != "builder" { + t.Fatalf("the builder gates the tier after it: %v", g) + } + p.Tier = 0 + if g := gates(p, edges, rollsOut); len(g) != 0 { + t.Fatalf("the runtime image is a build dependency and gates nothing: %v", g) + } + if g := gates(p, edges, func(string) bool { return false }); len(g) != 0 { + t.Fatalf("a module that only records its upgrade gates nothing: %v", g) + } +} + +// A gate is open once every machine running the module has reported after it was built. +func TestAGateOpensWhenTheMachinesHaveReportedSinceTheBuild(t *testing.T) { + built := time.Date(2026, 10, 1, 15, 0, 0, 0, time.UTC) + before, after := built.Add(-time.Minute), built.Add(time.Minute) + reports := []inventory.Reported{{Node: "anchor", At: &after}, {Node: "home-server", At: &before}} + ok, waiting := applied("builder", built, []string{"anchor", "home-server"}, reports) + if ok || len(waiting) != 1 || waiting[0] != "home-server" { + t.Fatalf("one machine has not reported since the build: ok=%v waiting=%v", ok, waiting) + } + if ok, _ := applied("builder", built, []string{"anchor"}, reports); !ok { + t.Fatal("the machine that reported after the build holds the gate open") + } + if ok, _ := applied("builder", built, nil, reports); !ok { + t.Fatal("a module running nowhere gates nothing") + } +} + +// A cycle is not lost: what remains is one last tier, and the caller says so. +func TestACycleIsOneLastTierAndSaidSo(t *testing.T) { + edges := []inventory.Edge{{From: "a", To: "b", Kind: inventory.EdgeStandsOn}, {From: "b", To: "a", Kind: inventory.EdgeStandsOn}} + tiers := tiersOf([]string{"a", "b"}, edges) + if len(tiers) != 1 || len(tiers[0]) != 2 || !hasCycle(tiers, edges) { + t.Fatalf("a cycle should be one tier of two, said: %v", tiers) + } +} diff --git a/cmd/mesh-controller/status.go b/cmd/mesh-controller/status.go index 34794b2..7fe09df 100644 --- a/cmd/mesh-controller/status.go +++ b/cmd/mesh-controller/status.go @@ -138,6 +138,18 @@ func printStatus(asked answers) error { len(quiet), strings.Join(said, "\n ")) } + if open, late := openPlans(asked.plans); len(open) > 0 { + fmt.Printf("%d plan(s) open", len(open)) + if late > 0 { + fmt.Printf(", %d waiting past %s", late, planWaitBound) + } + fmt.Println(":") + for _, p := range open { + fmt.Printf(" %s\n", planLine(p, time.Now())) + } + fmt.Println() + } + if len(behind) > 0 { var names []string for m := range behind { @@ -342,6 +354,10 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) { if err != nil { return answers{}, err } + out.plans, err = inv.RecentPlans(ctx, 5) + if err != nil { + return answers{}, err + } // And which machines are not running what the mesh would send them. The same question as a // module being behind its source, one level down: that one says the catalogue is out of date, diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index 091ce02..0aa9ffd 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -295,31 +295,31 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error { "built from\n", m.Owner, m.Repo, m.Base, m.Commit) return nil } - against, err := inv.BuiltAgainst(ctx) + // A merge produces a plan the mesh keeps (novox/hq ADR 0162): what moved and everything that + // depends on it, along the catalogue's one dependency relation, sorted into tiers. The plan is + // written before any build is asked; the first tier is asked; this returns. Outcomes advance it. + edges, err := inv.Dependencies(ctx) if err != nil { return notNow(err) } - // And whatever stands on what moved. A base rebuilt without its dependents is a mesh half on - // the old image until somebody remembers to ask — on 2026-10-01 forty-two modules, twice, by - // hand (novox/hq issue 186). The relation is the one `build --on` reads; this is the same - // rebuild, asked by the merge that made it necessary, in base order. - standing := dependentsOf(moved, entries, against) - if len(standing) > 0 { - var on []string - for _, e := range standing { - on = append(on, e.Manifest.Module) - } - fmt.Printf(" %d module(s) stand on what moved and are rebuilt with it: %s\n", - len(standing), strings.Join(on, ", ")) - moved = append(moved, standing...) + var movedNames []string + for _, e := range moved { + movedNames = append(movedNames, e.Manifest.Module) } - ordered := orderByBases(moved, against) - names := make([]string, 0, len(ordered)) - for _, e := range ordered { - names = append(names, e.Manifest.Module) + plan := planOfMerge(m, movedNames, edges) + if hasCycle(plan.Tiers, edges) { + fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n", + strings.Join(plan.Tiers[len(plan.Tiers)-1], ", ")) } - fmt.Printf("%s/%s merged into %s (%.8s); building %s\n", - m.Owner, m.Repo, m.Base, m.Commit, strings.Join(names, ", ")) + if err := inv.SavePlan(ctx, plan); err != nil { + return notNow(err) + } + var tiers []string + for i, t := range plan.Tiers { + tiers = append(tiers, fmt.Sprintf("%d: %s", i, strings.Join(t, ", "))) + } + fmt.Printf("%s/%s merged into %s (%.8s); plan %s, %d module(s) in %d tier(s)\n %s\n", + m.Owner, m.Repo, m.Base, m.Commit, plan.ID, len(plan.Modules), len(plan.Tiers), strings.Join(tiers, "\n ")) if len(packaging) > 0 { var also []string for _, e := range packaging { @@ -328,22 +328,11 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error { fmt.Printf(" %s package source from it, so they are rebuilt and their own source record "+ "is left where it is\n", strings.Join(also, ", ")) } - var failed []string - for _, e := range ordered { - source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat} - if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, 20*time.Minute); err != nil { - fmt.Printf(" %s: %v\n", e.Manifest.Module, err) - failed = append(failed, e.Manifest.Module) - // A base that failed is a reason to stop: what stands on it would be built against - // the old one, and report success (novox/hq 04-ISSUES/131). - if standsOn(ordered, e.Manifest.Module, against) { - fmt.Printf(" stopping: %s is a base of what was still to build\n", e.Manifest.Module) - break - } - } + if err := askTier(ctx, inv, &plan); err != nil { + return notNow(err) } - if len(failed) > 0 { - fmt.Printf("%d of %d not built: %s\n", len(failed), len(ordered), strings.Join(failed, ", ")) + if err := inv.SavePlan(ctx, plan); err != nil { + return notNow(err) } return nil } diff --git a/internal/inventory/dependencies.go b/internal/inventory/dependencies.go new file mode 100644 index 0000000..0633526 --- /dev/null +++ b/internal/inventory/dependencies.go @@ -0,0 +1,124 @@ +package inventory + +import ( + "context" + "sort" + "strings" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// The kinds of edge in the catalogue's one dependency relation (novox/hq ADR 0162). +const ( + // EdgeStandsOn: the module's artifact is built on the other's. + EdgeStandsOn = "stands-on" + // EdgePackages: the module's build reads the other's repository. + EdgePackages = "packages" + // EdgeBuiltBy: the module is built by the holder of the build-machine seat. + EdgeBuiltBy = "built-by" + // EdgeDeclared: the manifest's own `build.on`. + EdgeDeclared = "declared" +) + +// Edge is one dependency: From depends on To, in the way Kind says. +type Edge struct { + From string `json:"from"` + To string `json:"to"` + Kind string `json:"kind"` +} + +// Dependencies is the catalogue's dependency relation, whole: every module the mesh holds, with +// an edge to each module it depends on and the kind of dependency on the edge. One answer, so +// nothing else computes an edge (novox/hq ADR 0162) — the merge handler, `build --on` and the +// overview all read this. +// +// Four sources, one relation: a manifest's `build.on`; the artifacts the latest build was made +// against (an `artifact-store:///…` reference is an edge to that module); the repositories +// the latest build read (an edge to the module whose source that is); and the build machine, which +// every source-built module is built by. +func (i *Inventory) Dependencies(ctx context.Context) ([]Edge, error) { + entries, err := i.Catalogued(ctx) + if err != nil { + return nil, err + } + against, err := i.BuiltAgainst(ctx) + if err != nil { + return nil, err + } + read, err := i.ReadRepositories(ctx) + if err != nil { + return nil, err + } + return dependenciesOf(entries, against, read), nil +} + +// dependenciesOf is Dependencies over what was read, so a test can hand it a catalogue. +func dependenciesOf(entries []Entry, against map[string][]string, read map[string][]ReadRepository) []Edge { + known := map[string]bool{} + byRepository := map[string][]string{} + var builders []string + for _, e := range entries { + name := e.Manifest.Module + known[name] = true + if r := repositoryKey(e.Source.Repository); r != "" { + byRepository[r] = append(byRepository[r], name) + } + if e.Manifest.ClaimsSeat("mesh-build-machine") { + builders = append(builders, name) + } + } + seen := map[Edge]bool{} + var out []Edge + add := func(from, to, kind string) { + if from == to || !known[to] { + return + } + e := Edge{From: from, To: to, Kind: kind} + if !seen[e] { + seen[e] = true + out = append(out, e) + } + } + for _, e := range entries { + name := e.Manifest.Module + if e.Manifest.Build != nil { + for _, on := range e.Manifest.Build.On { + if on.Module != "" { + add(name, on.Module, EdgeDeclared) + } + } + } + for _, ref := range against[name] { + if rest, ok := strings.CutPrefix(ref, catalogue.ArtifactStoreScheme); ok { + if base, _, found := strings.Cut(rest, "/"); found { + add(name, base, EdgeStandsOn) + } + } + } + for _, r := range read[name] { + for _, other := range byRepository[repositoryKey(r.Repository)] { + add(name, other, EdgePackages) + } + } + if e.Source.Repository != "" { + for _, b := range builders { + add(name, b, EdgeBuiltBy) + } + } + } + sort.Slice(out, func(a, b int) bool { + if out[a].From != out[b].From { + return out[a].From < out[b].From + } + if out[a].To != out[b].To { + return out[a].To < out[b].To + } + return out[a].Kind < out[b].Kind + }) + return out +} + +// repositoryKey is a repository as compared: lower-cased, without a trailing `.git`. +func repositoryKey(repository string) string { + return strings.ToLower(strings.TrimSuffix(strings.TrimSpace(repository), ".git")) +} diff --git a/internal/inventory/dependencies_test.go b/internal/inventory/dependencies_test.go new file mode 100644 index 0000000..198e193 --- /dev/null +++ b/internal/inventory/dependencies_test.go @@ -0,0 +1,64 @@ +package inventory + +import ( + "testing" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// The catalogue's one dependency relation (novox/hq ADR 0162): four kinds of edge from one call. +func TestDependenciesAreOneRelationWithTheirKinds(t *testing.T) { + entry := func(name, repository string) Entry { + return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}} + } + builder := entry("builder", "http://forge/novox/mesh-catalog.git") + builder.Manifest.Claims = []catalogue.Claim{{Name: "mesh-build-machine", Scope: catalogue.ScopeMesh}} + plugin := entry("shop-plugin", "http://forge/novox/mesh-catalog.git") + plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}} + entries := []Entry{ + entry("mesh-tools", "http://forge/novox/mesh-tools.git"), + entry("shop", "http://forge/novox/mesh-catalog.git"), + plugin, + builder, + entry("mesh-controller", "http://forge/novox/mesh-controller.git"), + entry("route-proxy", "http://forge/novox/mesh-catalog.git"), + {Manifest: catalogue.Manifest{Module: "hand-made"}}, + } + against := map[string][]string{ + "shop": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"}, + "builder": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"}, + } + read := map[string][]ReadRepository{ + "route-proxy": {{Repository: "http://forge/novox/mesh-controller.git", Ref: "main"}}, + } + got := dependenciesOf(entries, against, read) + has := func(from, to, kind string) bool { + for _, e := range got { + if e == (Edge{From: from, To: to, Kind: kind}) { + return true + } + } + return false + } + for _, want := range []Edge{ + {"shop", "mesh-tools", EdgeStandsOn}, + {"builder", "mesh-tools", EdgeStandsOn}, + {"shop-plugin", "shop", EdgeDeclared}, + {"route-proxy", "mesh-controller", EdgePackages}, + {"shop", "builder", EdgeBuiltBy}, + {"mesh-controller", "builder", EdgeBuiltBy}, + {"mesh-tools", "builder", EdgeBuiltBy}, + } { + if !has(want.From, want.To, want.Kind) { + t.Errorf("missing %+v in %+v", want, got) + } + } + if has("builder", "builder", EdgeBuiltBy) { + t.Error("the builder is not built by itself") + } + for _, e := range got { + if e.From == "hand-made" { + t.Errorf("a module with no source depends on nothing: %+v", e) + } + } +} diff --git a/internal/inventory/migrations/0052-a-merge-is-a-plan-the-mesh-keeps.sql b/internal/inventory/migrations/0052-a-merge-is-a-plan-the-mesh-keeps.sql new file mode 100644 index 0000000..3d1b4a0 --- /dev/null +++ b/internal/inventory/migrations/0052-a-merge-is-a-plan-the-mesh-keeps.sql @@ -0,0 +1,19 @@ +-- A merge produces a tiered plan the mesh keeps (novox/hq ADR 0162): what the merge changed and +-- everything standing on it, sorted into tiers, each module's state, and the tier the plan is at. +-- Kept so a controller replaced mid-plan resumes it, and so `status` can say what a merge still +-- waits for. +create table release_plan ( + id text primary key, + repository text not null, + commit_hash text not null, + created timestamptz not null default now(), + updated timestamptz not null default now(), + -- building: a tier's builds are asked; rolling: the tier is built and the machines are applying + -- what a later tier needs running; done; failed. + state text not null, + tier int not null default 0, + tiers jsonb not null, + modules jsonb not null, + note text not null default '' +); +create index release_plan_open on release_plan (created) where state in ('building', 'rolling'); diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go new file mode 100644 index 0000000..1f1a367 --- /dev/null +++ b/internal/inventory/plans.go @@ -0,0 +1,123 @@ +package inventory + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" +) + +// A Plan is what a merge produces (novox/hq ADR 0162): the modules it changed and everything +// standing on them, sorted into tiers, each module's state, and the tier the plan is at. Kept in +// the store so a controller replaced mid-plan resumes it, and so `status` can say what a merge +// still waits for. +type Plan struct { + ID string `json:"id"` + Repository string `json:"repository"` + Commit string `json:"commit"` + Created time.Time `json:"created"` + Updated time.Time `json:"updated"` + State string `json:"state"` + Tier int `json:"tier"` + Tiers [][]string `json:"tiers"` + Modules map[string]*PlanModule `json:"modules"` + Note string `json:"note,omitempty"` +} + +// PlanModule is one module's state within a plan. +type PlanModule struct { + // State: asked, built, failed; empty for a module whose tier has not been asked yet. + State string `json:"state,omitempty"` + AskedAt *time.Time `json:"asked_at,omitempty"` + BuiltAt *time.Time `json:"built_at,omitempty"` + Commit string `json:"commit,omitempty"` + Why string `json:"why,omitempty"` +} + +// The states a plan passes through. +const ( + PlanBuilding = "building" + PlanRolling = "rolling" + PlanDone = "done" + PlanFailed = "failed" +) + +// Open says whether the plan is still being worked. +func (p Plan) Open() bool { return p.State == PlanBuilding || p.State == PlanRolling } + +// SavePlan writes a plan, new or changed, whole: the plan is small and read as one thing. +func (i *Inventory) SavePlan(ctx context.Context, p Plan) error { + tiers, err := json.Marshal(p.Tiers) + if err != nil { + return err + } + modules, err := json.Marshal(p.Modules) + if err != nil { + return err + } + _, err = i.store.Pool().Exec(ctx, + `insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note) + values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9) + on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier, + tiers = excluded.tiers, modules = excluded.modules, note = excluded.note`, + p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note) + return err +} + +// OpenPlans is every plan still being worked, oldest first. +func (i *Inventory) OpenPlans(ctx context.Context) ([]Plan, error) { + return i.plans(ctx, `where state in ('building', 'rolling') order by created`) +} + +// RecentPlans is the last few plans, newest first, open or not — what the overview shows. +func (i *Inventory) RecentPlans(ctx context.Context, limit int) ([]Plan, error) { + return i.plans(ctx, fmt.Sprintf(`order by created desc limit %d`, limit)) +} + +// PlanByID is one plan. +func (i *Inventory) PlanByID(ctx context.Context, id string) (Plan, error) { + plans, err := i.plans(ctx, `where id = '`+id+`'`) + if err != nil { + return Plan{}, err + } + if len(plans) == 0 { + return Plan{}, fmt.Errorf("no plan %s", id) + } + return plans[0], nil +} + +func (i *Inventory) plans(ctx context.Context, tail string) ([]Plan, error) { + rows, err := i.store.Pool().Query(ctx, + `select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note + from release_plan `+tail) + if err != nil { + return nil, err + } + defer rows.Close() + var out []Plan + for rows.Next() { + var p Plan + var tiers, modules []byte + if err := rows.Scan(&p.ID, &p.Repository, &p.Commit, &p.Created, &p.Updated, &p.State, + &p.Tier, &tiers, &modules, &p.Note); err != nil { + return nil, err + } + if err := json.Unmarshal(tiers, &p.Tiers); err != nil { + return nil, err + } + if err := json.Unmarshal(modules, &p.Modules); err != nil { + return nil, err + } + if p.Modules == nil { + p.Modules = map[string]*PlanModule{} + } + out = append(out, p) + } + if errors.Is(rows.Err(), pgx.ErrNoRows) { + return nil, nil + } + return out, rows.Err() +} From fdf338dbd196a4aaa43658b79e07e4b3f3cf8e3b Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 1 Oct 2026 17:56:06 +0200 Subject: [PATCH 2/2] A plan is kept, advanced and resumed from the store: the test --- internal/inventory/plans_test.go | 48 ++++++++++++++++++++++++++++++++ 1 file changed, 48 insertions(+) create mode 100644 internal/inventory/plans_test.go diff --git a/internal/inventory/plans_test.go b/internal/inventory/plans_test.go new file mode 100644 index 0000000..a9153a6 --- /dev/null +++ b/internal/inventory/plans_test.go @@ -0,0 +1,48 @@ +package inventory + +import ( + "testing" + "time" +) + +// A plan is a record the mesh keeps and resumes (novox/hq ADR 0162): written whole, read back open, +// advanced, and gone from the open ones when done. +func TestAPlanIsKeptAdvancedAndResumedFromTheStore(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + p := Plan{ID: "plan-1", Repository: "novox/mesh-tools", Commit: "abc", Created: time.Now().UTC(), + State: PlanBuilding, Tiers: [][]string{{"mesh-tools"}, {"builder"}, {"shop"}}, + Modules: map[string]*PlanModule{"mesh-tools": {}, "builder": {}, "shop": {}}} + if err := inv.SavePlan(ctx, p); err != nil { + t.Fatal(err) + } + open, err := inv.OpenPlans(ctx) + if err != nil || len(open) != 1 || open[0].ID != "plan-1" || len(open[0].Tiers) != 3 { + t.Fatalf("the plan was not kept whole: %v %+v", err, open) + } + // Another controller picks it up where it was left: a tier advanced and a module built. + now := time.Now().UTC() + resumed := open[0] + resumed.Tier = 1 + resumed.Modules["mesh-tools"].State = "built" + resumed.Modules["mesh-tools"].BuiltAt = &now + resumed.State = PlanRolling + resumed.Note = "tier 0 built; waiting for builder on anchor to be applied" + if err := inv.SavePlan(ctx, resumed); err != nil { + t.Fatal(err) + } + again, err := inv.PlanByID(ctx, "plan-1") + if err != nil || again.Tier != 1 || again.Modules["mesh-tools"].State != "built" || again.State != PlanRolling { + t.Fatalf("the advanced plan did not come back as left: %v %+v", err, again) + } + again.State = PlanDone + if err := inv.SavePlan(ctx, again); err != nil { + t.Fatal(err) + } + if open, _ = inv.OpenPlans(ctx); len(open) != 0 { + t.Fatalf("a done plan is not open: %+v", open) + } + if recent, _ := inv.RecentPlans(ctx, 5); len(recent) != 1 || recent[0].State != PlanDone { + t.Fatalf("a done plan is still among the recent ones: %+v", recent) + } +}