Three faults the mesh's own logs showed this morning. A handler that outlives the acknowledgement window was handed its message again while it was still working: acting on a merge builds modules, minutes against a thirty-second window, so one merge ran the whole catalogue five times over. The transport now says the work is in progress while it runs, which is where the window belongs. Everything the mesh hands out — a token, a membership, a person's credential — took its address from the enrolment setting, which on a mesh that has moved still names the broker it moved from: the first person issued after the move was handed the retired broker's port. There is one bus, and its address is the one the control plane is connected to. And `operator issue` documented an argument order its parser refused.
525 lines
19 KiB
Go
525 lines
19 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/novox/mesh-controller/internal/catalogue"
|
|
"github.com/novox/mesh-controller/internal/inventory"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// following acts on what the catalogue announces.
|
|
//
|
|
// **It holds the stores, not a copy of the decision.** What to do about an upgrade is read when
|
|
// one arrives, so changing it takes effect on the next upgrade rather than on the next restart of
|
|
// the control plane.
|
|
type following struct{ open *stores }
|
|
|
|
// Upgraded sends the machines running a module the version the catalogue now considers current —
|
|
// or records that they are behind, which is the default and needs no record.
|
|
//
|
|
// **Recording is not a second code path.** A machine that is not running what the mesh would send
|
|
// it is already something the mesh notices and reports; that is what `status` and `push --behind`
|
|
// are built on. So "record it" is the absence of an action, and the only thing this has to decide
|
|
// is whether to act.
|
|
func (f following) Upgraded(ctx context.Context, u link.Upgraded) error {
|
|
inv := f.open.inventory
|
|
|
|
// The store read first, and an outage there said as one, so the announcement is held and asked
|
|
// again (novox/hq issue 083). Only here: a push that fails further down is not asked again.
|
|
decision, err := inv.UpgradeOf(ctx, u.Module)
|
|
if err != nil {
|
|
return notNow(err)
|
|
}
|
|
on, err := inv.Running(ctx, u.Module)
|
|
if err != nil {
|
|
return notNow(err)
|
|
}
|
|
if len(on) == 0 {
|
|
fmt.Printf("%s moved to %s; no machine runs it\n", u.Module, shortCommit(u.Commit))
|
|
return nil
|
|
}
|
|
if !decision.RollOut {
|
|
// Named rather than counted, and said even though nothing happens: an upgrade that was
|
|
// deliberately not rolled out and an upgrade that was never noticed look identical in a
|
|
// log that only speaks when it acts.
|
|
fmt.Printf("%s moved to %s; %s %s behind it, and this mesh records upgrades rather than "+
|
|
"rolling them out — `push --behind` when you want them\n",
|
|
u.Module, shortCommit(u.Commit), readableList(on), isAre(len(on)))
|
|
return nil
|
|
}
|
|
|
|
if decision.Together {
|
|
fmt.Printf("%s moved to %s; sending %s together\n",
|
|
u.Module, shortCommit(u.Commit), readableList(on))
|
|
return sendTo(ctx, f.open, on)
|
|
}
|
|
// One at a time, and stopping at the first that fails.
|
|
//
|
|
// **Stopping is the point.** The machines are done one after another precisely so that a
|
|
// version that breaks the first one does not reach the rest; carrying on past a failure would
|
|
// make this the same as sending them together, only slower.
|
|
fmt.Printf("%s moved to %s; sending %s one at a time\n",
|
|
u.Module, shortCommit(u.Commit), readableList(on))
|
|
for _, node := range on {
|
|
if err := sendTo(ctx, f.open, []string{node}); err != nil {
|
|
return fmt.Errorf("%s did not take %s, so the machines after it were left alone: %w",
|
|
node, u.Module, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// readableList names machines the way a sentence does, because this is read by a person deciding
|
|
// whether an upgrade went where they expected.
|
|
func readableList(names []string) string {
|
|
switch len(names) {
|
|
case 0:
|
|
return "nothing"
|
|
case 1:
|
|
return names[0]
|
|
case 2:
|
|
return names[0] + " and " + names[1]
|
|
}
|
|
return strings.Join(names[:len(names)-1], ", ") + " and " + names[len(names)-1]
|
|
}
|
|
|
|
func shortCommit(commit string) string {
|
|
if len(commit) > 8 {
|
|
return commit[:8]
|
|
}
|
|
return commit
|
|
}
|
|
|
|
func isAre(n int) string {
|
|
if n == 1 {
|
|
return "is"
|
|
}
|
|
return "are"
|
|
}
|
|
|
|
// upgradeCommand says what should happen when a module's current version moves.
|
|
func upgradeCommand(ctx context.Context, args []string) error {
|
|
set := flag.NewFlagSet("upgrade", flag.ContinueOnError)
|
|
together := set.Bool("together", false,
|
|
"send every machine running it at once, instead of one after another")
|
|
positionals, err := parseAround(set, args)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(positionals) == 0 {
|
|
return errors.New("upgrade <module> [roll-out|record] [--together]")
|
|
}
|
|
module := positionals[0]
|
|
|
|
open, err := openStores(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer open.Close()
|
|
inv := open.inventory
|
|
|
|
if len(positionals) == 1 {
|
|
decision, err := inv.UpgradeOf(ctx, module)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Println(sayUpgrade(module, decision))
|
|
return nil
|
|
}
|
|
|
|
var decision inventory.Upgrade
|
|
switch positionals[1] {
|
|
case "roll-out":
|
|
decision = inventory.Upgrade{RollOut: true, Together: *together}
|
|
case "record":
|
|
if *together {
|
|
// Refused rather than ignored: --together only means anything for a roll-out, and
|
|
// accepting it here would store a preference that never applies and looks like it does.
|
|
return errors.New("`--together` says how to roll out, so it cannot be given with " +
|
|
"`record`, which is the choice not to")
|
|
}
|
|
decision = inventory.Upgrade{}
|
|
default:
|
|
return fmt.Errorf("upgrade <module> roll-out|record — not %q", positionals[1])
|
|
}
|
|
if err := inv.SetUpgradeOf(ctx, module, decision); err != nil {
|
|
return err
|
|
}
|
|
fmt.Println(sayUpgrade(module, decision))
|
|
return nil
|
|
}
|
|
|
|
func sayUpgrade(module string, u inventory.Upgrade) string {
|
|
if !u.RollOut {
|
|
return fmt.Sprintf("when %s moves, the mesh records it and the machines running it are "+
|
|
"reported as behind", module)
|
|
}
|
|
if u.Together {
|
|
return fmt.Sprintf("when %s moves, every machine running it is sent the new version "+
|
|
"together", module)
|
|
}
|
|
return fmt.Sprintf("when %s moves, the machines running it are sent the new version one at "+
|
|
"a time, stopping at the first that fails", module)
|
|
}
|
|
|
|
// Announceable is every build this mesh recorded, in the shape the builder announces one.
|
|
//
|
|
// **The catalogue asks for this when it starts, and the answer is the graph's foundation**
|
|
// (novox/hq 04-ISSUES/050). A durable queue keeps what arrived after it existed, so a running
|
|
// catalogue misses nothing — but the modules built before it first ran were announced to a queue
|
|
// that did not exist, and on a fresh mesh those are always the same three: the shared base, the
|
|
// store the catalogue runs on, and the catalogue itself.
|
|
//
|
|
// **Announced as fetchable, recorded as what it is** (novox/hq 04-ISSUES/102). A build is
|
|
// recorded by digest and path; the catalogue hears the builder's own announcements, which name
|
|
// the store's address, so a replay composes the address back in — the store's address as the
|
|
// network reaches it NOW, which is the whole point of not having recorded the old one. With no
|
|
// store on the network yet, the recorded form goes as it is.
|
|
func (f following) Announceable(ctx context.Context) ([]link.Announcement, error) {
|
|
builds, err := f.open.inventory.Announceable(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
address, err := whereTheStoreIs(ctx, f.open.inventory, "")
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out := make([]link.Announcement, 0, len(builds))
|
|
for _, b := range builds {
|
|
a := link.Announcement{
|
|
Module: b.Module, Commit: b.Commit, Repository: b.Repository,
|
|
Path: b.Path, Ref: b.Ref, Against: b.Against,
|
|
}
|
|
if len(b.Manifest) > 0 {
|
|
a.Manifest = b.Manifest
|
|
if address != "" {
|
|
a.Manifest = routedManifest(b.Manifest, b.Made, address)
|
|
}
|
|
}
|
|
for _, made := range routedArtifacts(b.Made, address) {
|
|
a.Made = append(a.Made, link.MadeArtifact{
|
|
Name: made.Name, Kind: made.Kind, Reference: made.Reference,
|
|
})
|
|
}
|
|
out = append(out, a)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// notNow marks a store that could not be read right now, so the announcement is held rather than
|
|
// lost; anything else is returned as it was.
|
|
func notNow(err error) error {
|
|
if inventory.Unreachable(err) {
|
|
return fmt.Errorf("%w: %w", link.ErrTryAgain, err)
|
|
}
|
|
return err
|
|
}
|
|
|
|
// SourceMoved is the forge announcing a merge: every module recorded as built from that
|
|
// repository and branch is marked as moved to the merge commit, and built — bases first, so a
|
|
// module that stands on another's artifact is built after it and not against the old one
|
|
// (novox/hq 04-ISSUES/131). Nothing is pushed here: what a finished build does to the machines
|
|
// running the module is the upgrade's decision, taken when the catalogue announces it.
|
|
func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
|
|
inv := f.open.inventory
|
|
entries, err := inv.Catalogued(ctx)
|
|
if err != nil {
|
|
return notNow(err)
|
|
}
|
|
read, err := inv.ReadRepositories(ctx)
|
|
if err != nil {
|
|
return notNow(err)
|
|
}
|
|
|
|
// Two kinds of module are affected by one merge, and they are affected differently.
|
|
//
|
|
// A module **built from** this repository and branch has moved: the mesh records the new commit
|
|
// as what its source now has, and only what the merge actually changed is rebuilt. A module that
|
|
// only **packages source from** it has not moved — its own source is somewhere else, at the
|
|
// commit it already records — so it is rebuilt and its record left alone. Writing this commit as
|
|
// its source would make it permanently behind a repository its manifest does not come from.
|
|
var from, packaging []inventory.Entry
|
|
already := 0
|
|
for _, e := range entries {
|
|
switch {
|
|
case sourceIs(e.Source, m):
|
|
if e.Source.BuiltFrom == m.Commit {
|
|
already++
|
|
continue
|
|
}
|
|
// **A merge older than the last look at the source is history, not a move.** The forge
|
|
// announces what it finds merged, and an old merge surfacing late would otherwise move
|
|
// the recorded head backwards and rebuild everything built from that repository, once
|
|
// per old merge (2026-09-28).
|
|
if isHistory(m.MergedAt, e.Source.Seen) {
|
|
continue
|
|
}
|
|
from = append(from, e)
|
|
case readsFrom(read[e.Manifest.Module], m):
|
|
packaging = append(packaging, e)
|
|
}
|
|
}
|
|
if len(from) == 0 && len(packaging) == 0 {
|
|
// "Already built from it" and "nothing reads it" are different facts, and reading the first
|
|
// as the second sends somebody looking for a broken trigger when the mesh is up to date.
|
|
if already > 0 {
|
|
fmt.Printf("%s/%s merged into %s (%.8s); %d module(s) the mesh holds are already built "+
|
|
"from it\n", m.Owner, m.Repo, m.Base, m.Commit, already)
|
|
return nil
|
|
}
|
|
fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds reads it\n",
|
|
m.Owner, m.Repo, m.Base, m.Commit)
|
|
return nil
|
|
}
|
|
// The same judgement for the packaging kind, against the newest look at that repository by
|
|
// anything built from it: they keep no record of it themselves, and a replayed old merge should
|
|
// not rebuild them either.
|
|
if isHistory(m.MergedAt, lastLookAt(entries, m)) {
|
|
packaging = nil
|
|
}
|
|
touched := whatTheMergeTouched(from, entries, m)
|
|
for _, e := range touched {
|
|
if err := inv.SourceMoved(ctx, e.Manifest.Module, m.Commit); err != nil {
|
|
return notNow(err)
|
|
}
|
|
}
|
|
moved := append(append([]inventory.Entry{}, touched...), packaging...)
|
|
if len(moved) == 0 {
|
|
fmt.Printf("%s/%s merged into %s (%.8s); it changed nothing any module the mesh holds is "+
|
|
"built from\n", m.Owner, m.Repo, m.Base, m.Commit)
|
|
return nil
|
|
}
|
|
against, err := inv.BuiltAgainst(ctx)
|
|
if err != nil {
|
|
return notNow(err)
|
|
}
|
|
ordered := orderByBases(moved, against)
|
|
names := make([]string, 0, len(ordered))
|
|
for _, e := range ordered {
|
|
names = append(names, e.Manifest.Module)
|
|
}
|
|
fmt.Printf("%s/%s merged into %s (%.8s); building %s\n",
|
|
m.Owner, m.Repo, m.Base, m.Commit, strings.Join(names, ", "))
|
|
if len(packaging) > 0 {
|
|
var also []string
|
|
for _, e := range packaging {
|
|
also = append(also, e.Manifest.Module)
|
|
}
|
|
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 len(failed) > 0 {
|
|
fmt.Printf("%d of %d not built: %s\n", len(failed), len(ordered), strings.Join(failed, ", "))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// sourceIs is whether a recorded source is the repository and branch a merge announced. A source on
|
|
// the git seat is recorded as its path on the forge; one elsewhere as the URL it was cloned from.
|
|
// An empty recorded ref is the repository's default branch, which is what a merge into the base
|
|
// branch of the forge's default means.
|
|
func sourceIs(s inventory.Source, m link.SourceMoved) bool {
|
|
if !sameRepository(s.Repository, m) {
|
|
return false
|
|
}
|
|
return s.Ref == "" || s.Ref == m.Base
|
|
}
|
|
|
|
// sameRepository is whether a recorded repository is the one a merge names, in either spelling it
|
|
// may have been recorded in: a path on the git seat, or the URL it was cloned from.
|
|
func sameRepository(repository string, m link.SourceMoved) bool {
|
|
want := strings.ToLower(m.Owner + "/" + m.Repo)
|
|
repo := strings.ToLower(strings.TrimSuffix(repository, ".git"))
|
|
return repo == want || strings.HasSuffix(repo, "/"+want) ||
|
|
(m.CloneURL != "" && repo == strings.ToLower(strings.TrimSuffix(m.CloneURL, ".git")))
|
|
}
|
|
|
|
// readsFrom is whether a module's build read the repository a merge names: the second repository its
|
|
// recipe packages source from. Its ref must be the branch that moved, or unset — the same rule a
|
|
// module's own source follows.
|
|
func readsFrom(read []inventory.ReadRepository, m link.SourceMoved) bool {
|
|
for _, r := range read {
|
|
if sameRepository(r.Repository, m) && (r.Ref == "" || r.Ref == m.Base) {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// lastLookAt is the most recent look at this repository by anything built from it.
|
|
func lastLookAt(entries []inventory.Entry, m link.SourceMoved) time.Time {
|
|
var newest time.Time
|
|
for _, e := range entries {
|
|
if sameRepository(e.Source.Repository, m) && e.Source.Seen.After(newest) {
|
|
newest = e.Source.Seen
|
|
}
|
|
}
|
|
return newest
|
|
}
|
|
|
|
// whatTheMergeTouched narrows the modules built from a repository to the ones the merge changed.
|
|
//
|
|
// **A change inside no module's own directory is a change to what they share.** The forge lists the
|
|
// files a merge changed; a module is affected when one of them is inside its own directory, when it
|
|
// is built from the repository's root — everything there is its source — or when some changed file
|
|
// belongs to no module's directory at all, which is how a shared file, a build recipe or a
|
|
// dependency at the root rebuilds everything built from that repository.
|
|
//
|
|
// A change inside *another* module's directory is that module's business and not this one's, even
|
|
// when the mesh does not hold that module: `known` is every module this repository is known to hold,
|
|
// whatever branch it was registered from. That is also the limit of this — a repository whose shared
|
|
// code sits inside a directory the mesh has never seen a module in reads as shared, and everything
|
|
// is rebuilt. Rebuilding too much is the safe direction: the fault this whole path exists for is a
|
|
// mesh that believes it is current and is not (novox/hq 04-ISSUES/131).
|
|
func whatTheMergeTouched(candidates, known []inventory.Entry, m link.SourceMoved) []inventory.Entry {
|
|
// Nothing said about the files, or not all of them said: everything built from it is affected.
|
|
if len(m.Paths) == 0 || m.PathsTruncated {
|
|
return candidates
|
|
}
|
|
var dirs []string
|
|
for _, e := range known {
|
|
if e.Source.Path != "" && sameRepository(e.Source.Repository, m) {
|
|
dirs = append(dirs, e.Source.Path)
|
|
}
|
|
}
|
|
for _, p := range m.Paths {
|
|
if !insideAny(p, dirs) {
|
|
return candidates
|
|
}
|
|
}
|
|
var out []inventory.Entry
|
|
for _, e := range candidates {
|
|
if e.Source.Path == "" || anyInside(m.Paths, e.Source.Path) {
|
|
out = append(out, e)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// inside is whether a changed file is in a directory: that directory itself, or under it.
|
|
func inside(path, dir string) bool {
|
|
dir = strings.Trim(dir, "/")
|
|
path = strings.TrimPrefix(path, "/")
|
|
return path == dir || strings.HasPrefix(path, dir+"/")
|
|
}
|
|
|
|
// insideAny is whether a changed file is in any of these directories.
|
|
func insideAny(path string, dirs []string) bool {
|
|
for _, dir := range dirs {
|
|
if inside(path, dir) {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// anyInside is whether any of these changed files is in a directory.
|
|
func anyInside(paths []string, dir string) bool {
|
|
for _, p := range paths {
|
|
if inside(p, dir) {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// orderByBases is the entries with every base before what stands on it: a module whose build stood
|
|
// on another's artifact comes after that module. Entries outside the set are not waited for — they
|
|
// are not being rebuilt. Stable for what has no order between it.
|
|
//
|
|
// `against` is what each module's newest build stood on (inventory.BuiltAgainst): the edges are
|
|
// derived from builds, not declared, because a recorded manifest no longer carries `build.on`.
|
|
func orderByBases(entries []inventory.Entry, against map[string][]string) []inventory.Entry {
|
|
inSet := map[string]bool{}
|
|
for _, e := range entries {
|
|
inSet[e.Manifest.Module] = true
|
|
}
|
|
var out []inventory.Entry
|
|
placed := map[string]bool{}
|
|
var place func(e inventory.Entry, seen map[string]bool)
|
|
place = func(e inventory.Entry, seen map[string]bool) {
|
|
name := e.Manifest.Module
|
|
if placed[name] || seen[name] {
|
|
return
|
|
}
|
|
seen[name] = true
|
|
for _, base := range entries {
|
|
if base.Manifest.Module != name && inSet[base.Manifest.Module] && standsOnModule(e, base.Manifest.Module, against) {
|
|
place(base, seen)
|
|
}
|
|
}
|
|
placed[name] = true
|
|
out = append(out, e)
|
|
}
|
|
for _, e := range entries {
|
|
place(e, map[string]bool{})
|
|
}
|
|
return out
|
|
}
|
|
|
|
// standsOn is whether anything in the set is built on the named module's artifacts.
|
|
func standsOn(entries []inventory.Entry, module string, against map[string][]string) bool {
|
|
for _, e := range entries {
|
|
if standsOnModule(e, module, against) {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// standsOnModule is whether an entry's build stood on the named module: by what its newest build
|
|
// recorded it was handed (`artifact-store://<module>/<artifact>@…`, the module's own artifact), or
|
|
// — for a module registered from a manifest and not yet built — by the base its manifest names.
|
|
func standsOnModule(e inventory.Entry, module string, against map[string][]string) bool {
|
|
if e.Manifest.Module == module {
|
|
return false
|
|
}
|
|
if e.Manifest.Build != nil {
|
|
for _, on := range e.Manifest.Build.On {
|
|
if on.Module == module {
|
|
return true
|
|
}
|
|
}
|
|
}
|
|
prefix := catalogue.ArtifactStoreScheme + module + "/"
|
|
for _, ref := range against[e.Manifest.Module] {
|
|
if strings.HasPrefix(ref, prefix) {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// isHistory is whether a merge made at mergedAt predates the last time the source was seen. A merge
|
|
// with no time on it is taken as news: refusing it would silence a forge that says less.
|
|
func isHistory(mergedAt string, seen time.Time) bool {
|
|
if mergedAt == "" || seen.IsZero() {
|
|
return false
|
|
}
|
|
at, err := time.Parse(time.RFC3339, mergedAt)
|
|
if err != nil {
|
|
return false
|
|
}
|
|
return at.Before(seen)
|
|
}
|