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.
527 lines
15 KiB
Go
527 lines
15 KiB
Go
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
|
|
}
|