Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
81e76fa485 | ||
|
|
12ed35e87d | ||
|
|
525f10b858 | ||
|
|
0547316cf2 | ||
|
|
9fe9b5349c | ||
|
|
d1e488efaf |
@@ -31,6 +31,53 @@ import (
|
|||||||
// control plane may send a machine is bounded by the declaration language. This is the shape the
|
// control plane may send a machine is bounded by the declaration language. This is the shape the
|
||||||
// builder module will take when it is given work over the broker; today a person runs it, and the
|
// builder module will take when it is given work over the broker; today a person runs it, and the
|
||||||
// mesh records the result the same way either way.
|
// mesh records the result the same way either way.
|
||||||
|
// buildOn rebuilds every module the mesh holds that stands on the named module's artifacts — the
|
||||||
|
// rebuild a changed base needs, which nothing else asks for: their sources did not move, and
|
||||||
|
// "behind" does not see a base that did (novox/hq 04-ISSUES/131). Bases first among them too.
|
||||||
|
func buildOn(ctx context.Context, base string, wait time.Duration) error {
|
||||||
|
open, err := openStores(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
defer open.Close()
|
||||||
|
held, err := open.inventory.Catalogued(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
var on []inventory.Entry
|
||||||
|
for _, e := range held {
|
||||||
|
if e.Manifest.Build == nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, b := range e.Manifest.Build.On {
|
||||||
|
if b.Module == base {
|
||||||
|
on = append(on, e)
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(on) == 0 {
|
||||||
|
fmt.Printf("nothing the mesh holds stands on %s\n", base)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
on = orderByBases(on)
|
||||||
|
fmt.Printf("%d module(s) stand on %s:\n", len(on), base)
|
||||||
|
var failed []string
|
||||||
|
for _, e := range on {
|
||||||
|
fmt.Printf("--- %s\n", e.Manifest.Module)
|
||||||
|
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
|
||||||
|
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, wait); err != nil {
|
||||||
|
fmt.Printf(" %v\n", err)
|
||||||
|
failed = append(failed, e.Manifest.Module)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(failed) > 0 {
|
||||||
|
return fmt.Errorf("%d of %d could not be built: %s", len(failed), len(on), strings.Join(failed, ", "))
|
||||||
|
}
|
||||||
|
fmt.Printf("\n%d module(s) rebuilt on %s. `push --behind` sends them on\n", len(on), base)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func buildCommand(ctx context.Context, args []string) error {
|
func buildCommand(ctx context.Context, args []string) error {
|
||||||
set := flag.NewFlagSet("build", flag.ContinueOnError)
|
set := flag.NewFlagSet("build", flag.ContinueOnError)
|
||||||
ref := set.String("ref", "", "the branch, tag or commit to build")
|
ref := set.String("ref", "", "the branch, tag or commit to build")
|
||||||
@@ -46,6 +93,7 @@ func buildCommand(ctx context.Context, args []string) error {
|
|||||||
// retype each repository is asking them to be the loop. Naming a repository and asking which
|
// retype each repository is asking them to be the loop. Naming a repository and asking which
|
||||||
// ones need building are different requests, so they are not combined.
|
// ones need building are different requests, so they are not combined.
|
||||||
behind := set.Bool("behind", false, "every module the mesh holds older than its source has")
|
behind := set.Bool("behind", false, "every module the mesh holds older than its source has")
|
||||||
|
on := set.String("on", "", "rebuild every module that stands on this module's artifacts — the rebuild a changed base needs")
|
||||||
// A repository on the mesh's own forge, named by its path there (novox/hq ADR 0111). Without it
|
// A repository on the mesh's own forge, named by its path there (novox/hq ADR 0111). Without it
|
||||||
// the repository is external, cloned exactly as given — see source.go.
|
// the repository is external, cloned exactly as given — see source.go.
|
||||||
self := set.Bool("self", false, "the repository is a path on the forge holding the git seat")
|
self := set.Bool("self", false, "the repository is a path on the forge holding the git seat")
|
||||||
@@ -53,6 +101,13 @@ func buildCommand(ctx context.Context, args []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
if *on != "" {
|
||||||
|
if len(positionals) != 0 || *behind || *self {
|
||||||
|
return errors.New("build --on <module> names a base and nothing else")
|
||||||
|
}
|
||||||
|
return buildOn(ctx, *on, *wait)
|
||||||
|
}
|
||||||
|
|
||||||
if *behind {
|
if *behind {
|
||||||
if len(positionals) != 0 || *self {
|
if len(positionals) != 0 || *self {
|
||||||
return errors.New("build <repository> or build --behind, not both: one names a " +
|
return errors.New("build <repository> or build --behind, not both: one names a " +
|
||||||
@@ -61,7 +116,7 @@ func buildCommand(ctx context.Context, args []string) error {
|
|||||||
return buildBehind(ctx, *wait)
|
return buildBehind(ctx, *wait)
|
||||||
}
|
}
|
||||||
if len(positionals) != 1 {
|
if len(positionals) != 1 {
|
||||||
return errors.New("build <repository> [--self] [--path P] [--ref R] [--wait D] [--dry-run]")
|
return errors.New("build <repository> [--self] [--path P] [--ref R] [--wait D] [--dry-run] | build --behind | build --on <module>")
|
||||||
}
|
}
|
||||||
source := buildSource{Repository: positionals[0]}
|
source := buildSource{Repository: positionals[0]}
|
||||||
if *self {
|
if *self {
|
||||||
@@ -329,6 +384,10 @@ func buildBehind(ctx context.Context, wait time.Duration) error {
|
|||||||
}
|
}
|
||||||
fmt.Println()
|
fmt.Println()
|
||||||
|
|
||||||
|
// Bases first: a module built before the module it stands on is built against the old one
|
||||||
|
// and reports success (novox/hq 04-ISSUES/131).
|
||||||
|
stale = orderByBases(stale)
|
||||||
|
|
||||||
var failed []string
|
var failed []string
|
||||||
for _, e := range stale {
|
for _, e := range stale {
|
||||||
fmt.Printf("--- %s\n", e.Manifest.Module)
|
fmt.Printf("--- %s\n", e.Manifest.Module)
|
||||||
|
|||||||
@@ -184,6 +184,7 @@ func usage() {
|
|||||||
operator key show the operator key, and what it can recover
|
operator key show the operator key, and what it can recover
|
||||||
build <repository> [--ref R] have a build machine build it, and record what came out
|
build <repository> [--ref R] have a build machine build it, and record what came out
|
||||||
build --behind build every module the mesh holds older than its source
|
build --behind build every module the mesh holds older than its source
|
||||||
|
build --on <module> rebuild every module that stands on this module's artifacts, bases first
|
||||||
builds [<module>] what has been built lately, and what came of it
|
builds [<module>] what has been built lately, and what came of it
|
||||||
builder issue <name> a broker account for a build machine, scoped to build work,
|
builder issue <name> a broker account for a build machine, scoped to build work,
|
||||||
delivered as the builder module's broker secret (module add it first)
|
delivered as the builder module's broker secret (module add it first)
|
||||||
|
|||||||
@@ -0,0 +1,63 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-controller/internal/catalogue"
|
||||||
|
"github.com/novox/mesh-controller/internal/inventory"
|
||||||
|
"github.com/novox/mesh-controller/internal/link"
|
||||||
|
)
|
||||||
|
|
||||||
|
func entry(module string, on ...string) inventory.Entry {
|
||||||
|
b := &catalogue.Build{}
|
||||||
|
for _, o := range on {
|
||||||
|
b.On = append(b.On, catalogue.BuildsOn{Arg: "X", Module: o, Artifact: "runtime"})
|
||||||
|
}
|
||||||
|
return inventory.Entry{Manifest: catalogue.Manifest{Module: module, Build: b}}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A module built before the module it stands on is built against the old one and reports success
|
||||||
|
// (novox/hq 04-ISSUES/131). So bases come first, however the set arrived.
|
||||||
|
func TestBasesAreBuiltBeforeWhatStandsOnThem(t *testing.T) {
|
||||||
|
in := []inventory.Entry{entry("app", "runtime"), entry("runtime", "base"), entry("other"), entry("base")}
|
||||||
|
got := orderByBases(in)
|
||||||
|
pos := map[string]int{}
|
||||||
|
for i, e := range got {
|
||||||
|
pos[e.Manifest.Module] = i
|
||||||
|
}
|
||||||
|
if !(pos["base"] < pos["runtime"] && pos["runtime"] < pos["app"]) {
|
||||||
|
t.Fatalf("bases not first: %v", pos)
|
||||||
|
}
|
||||||
|
if len(got) != 4 {
|
||||||
|
t.Fatalf("an entry was lost or doubled: %d", len(got))
|
||||||
|
}
|
||||||
|
// A base outside the set is not waited for: it is not being rebuilt.
|
||||||
|
got = orderByBases([]inventory.Entry{entry("app", "elsewhere")})
|
||||||
|
if len(got) != 1 {
|
||||||
|
t.Fatalf("a dependency outside the set changed the set: %v", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A merge names a repository the way the forge does; a source is recorded the way a build was
|
||||||
|
// asked for. The two meet on owner/repo and branch, whichever form the record took.
|
||||||
|
func TestAMergeMatchesTheSourcesBuiltFromIt(t *testing.T) {
|
||||||
|
m := link.SourceMoved{Owner: "novox", Repo: "mesh-controller", Base: "main",
|
||||||
|
CloneURL: "http://forge.internal:20000/novox/mesh-controller.git"}
|
||||||
|
for _, s := range []inventory.Source{
|
||||||
|
{Repository: "http://forge.internal:20000/novox/mesh-controller.git", Ref: "main"},
|
||||||
|
{Repository: "novox/mesh-controller", Seat: "git", Ref: ""},
|
||||||
|
{Repository: "https://elsewhere.example/novox/mesh-controller", Ref: "main"},
|
||||||
|
} {
|
||||||
|
if !sourceIs(s, m) {
|
||||||
|
t.Errorf("%+v was not matched by the merge", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, s := range []inventory.Source{
|
||||||
|
{Repository: "novox/mesh-host", Seat: "git"},
|
||||||
|
{Repository: "http://forge.internal:20000/novox/mesh-controller.git", Ref: "release"},
|
||||||
|
} {
|
||||||
|
if sourceIs(s, m) {
|
||||||
|
t.Errorf("%+v was matched by a merge that is not its", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -5,6 +5,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"os"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -33,7 +34,7 @@ import (
|
|||||||
// ability to change things, not the services its modules are serving — measured on 2026-09-27, when
|
// ability to change things, not the services its modules are serving — measured on 2026-09-27, when
|
||||||
// a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
|
// a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
|
||||||
|
|
||||||
const rolloutUsage = "rollout check | rollout mint [--again] | rollout --confirm"
|
const rolloutUsage = "rollout check | rollout mint [--again] | rollout hand <node> | rollout --confirm"
|
||||||
|
|
||||||
func rolloutCommand(ctx context.Context, args []string) error {
|
func rolloutCommand(ctx context.Context, args []string) error {
|
||||||
switch {
|
switch {
|
||||||
@@ -41,6 +42,8 @@ func rolloutCommand(ctx context.Context, args []string) error {
|
|||||||
return rolloutCheck(ctx)
|
return rolloutCheck(ctx)
|
||||||
case len(args) == 1 && args[0] == "mint":
|
case len(args) == 1 && args[0] == "mint":
|
||||||
return rolloutMint(ctx, false)
|
return rolloutMint(ctx, false)
|
||||||
|
case len(args) == 2 && args[0] == "hand":
|
||||||
|
return rolloutHand(ctx, args[1])
|
||||||
case len(args) == 2 && args[0] == "mint" && args[1] == "--again":
|
case len(args) == 2 && args[0] == "mint" && args[1] == "--again":
|
||||||
// Every credential minted afresh, whether or not one exists — for a mint that was wrong
|
// Every credential minted afresh, whether or not one exists — for a mint that was wrong
|
||||||
// before anything was pushed. Afterwards nothing that received the old one still works,
|
// before anything was pushed. Afterwards nothing that received the old one still works,
|
||||||
@@ -387,3 +390,62 @@ func providesBus(m catalogue.Manifest) bool {
|
|||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// rolloutHand mints a machine its credential for the new bus afresh and prints its membership
|
||||||
|
// once, for an operator to carry by hand — the rescue for a machine that cannot be reached over
|
||||||
|
// any bus: rotated while it still held the old password, or reachable only by ssh. The plaintext
|
||||||
|
// exists on this terminal and then only where it is written; the store keeps the hash, and the
|
||||||
|
// sealed copy in the machine's declaration is replaced too, so the next push says the same.
|
||||||
|
func rolloutHand(ctx context.Context, node string) error {
|
||||||
|
open, err := openStores(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
defer open.Close()
|
||||||
|
inv := open.inventory
|
||||||
|
|
||||||
|
known, err := broker.FromEnvironment()
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("the bus's certificate is not known to this process: %w", err)
|
||||||
|
}
|
||||||
|
busAddress, _, err := broker.OnNATS()
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if busAddress == "" {
|
||||||
|
return errors.New("this control plane is not on the new bus, so there is no membership to hand out")
|
||||||
|
}
|
||||||
|
_, _, bare := broker.CredentialIn(busAddress)
|
||||||
|
if _, after, has := strings.Cut(bare, "://"); has {
|
||||||
|
bare = after
|
||||||
|
}
|
||||||
|
if _, err := inv.NodeByName(ctx, node); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
p := broker.Principal{Kind: broker.KindNode, Node: node}
|
||||||
|
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusNode, Node: node})
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
membership, _ := json.Marshal(map[string]string{
|
||||||
|
"broker": bare, "fingerprint": known.Fingerprint, "password": password, "transport": "nats",
|
||||||
|
})
|
||||||
|
key, err := inv.SealingKeyOf(ctx, node)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
sealed, err := secrets.Seal(key, membership)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := inv.PutBusMembership(ctx, node, sealed); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
// The one line of output is the membership itself, so it can be piped to the machine without
|
||||||
|
// being read on the way. Everything else goes to stderr.
|
||||||
|
fmt.Fprintf(os.Stderr, "%s's credential is minted afresh. Write this to %s on it and restart its host; "+
|
||||||
|
"then push the machine running the bus so the user list carries the new hash.\n",
|
||||||
|
node, catalogue.BusMembershipPath)
|
||||||
|
fmt.Println(string(membership))
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import (
|
|||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
"strings"
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/novox/mesh-controller/internal/inventory"
|
"github.com/novox/mesh-controller/internal/inventory"
|
||||||
"github.com/novox/mesh-controller/internal/link"
|
"github.com/novox/mesh-controller/internal/link"
|
||||||
@@ -218,3 +219,127 @@ func notNow(err error) error {
|
|||||||
}
|
}
|
||||||
return 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)
|
||||||
|
}
|
||||||
|
var moved []inventory.Entry
|
||||||
|
for _, e := range entries {
|
||||||
|
if !sourceIs(e.Source, m) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if e.Source.BuiltFrom == m.Commit {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if err := inv.SourceMoved(ctx, e.Manifest.Module, m.Commit); err != nil {
|
||||||
|
return notNow(err)
|
||||||
|
}
|
||||||
|
moved = append(moved, e)
|
||||||
|
}
|
||||||
|
if len(moved) == 0 {
|
||||||
|
fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds is built from it\n",
|
||||||
|
m.Owner, m.Repo, m.Base, m.Commit)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
ordered := orderByBases(moved)
|
||||||
|
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, ", "))
|
||||||
|
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) {
|
||||||
|
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 {
|
||||||
|
want := strings.ToLower(m.Owner + "/" + m.Repo)
|
||||||
|
repo := strings.ToLower(strings.TrimSuffix(s.Repository, ".git"))
|
||||||
|
matches := repo == want || strings.HasSuffix(repo, "/"+want) ||
|
||||||
|
(m.CloneURL != "" && strings.EqualFold(strings.TrimSuffix(s.Repository, ".git"), strings.TrimSuffix(m.CloneURL, ".git")))
|
||||||
|
if !matches {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
return s.Ref == "" || s.Ref == m.Base
|
||||||
|
}
|
||||||
|
|
||||||
|
// orderByBases is the entries with every base before what stands on it: a module whose build names
|
||||||
|
// another's artifact under build.on 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.
|
||||||
|
func orderByBases(entries []inventory.Entry) []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
|
||||||
|
if e.Manifest.Build != nil {
|
||||||
|
for _, on := range e.Manifest.Build.On {
|
||||||
|
if on.Module == "" || on.Module == name || !inSet[on.Module] {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, base := range entries {
|
||||||
|
if base.Manifest.Module == on.Module {
|
||||||
|
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) bool {
|
||||||
|
for _, e := range entries {
|
||||||
|
if e.Manifest.Build == nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, on := range e.Manifest.Build.On {
|
||||||
|
if on.Module == module {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|||||||
@@ -175,6 +175,9 @@ var ControllerFollows = []string{
|
|||||||
// message on the control branch. Same three audiences, one publish: whoever asked, this, and the
|
// message on the control branch. Same three audiences, one publish: whoever asked, this, and the
|
||||||
// catalogue.
|
// catalogue.
|
||||||
seatEventSubject("mesh-build-machine", "built"),
|
seatEventSubject("mesh-build-machine", "built"),
|
||||||
|
// The forge's merges: what moved a source, so the mesh builds what that source produces
|
||||||
|
// without anybody telling it (novox/hq 04-ISSUES/131). Appended, because the index is a name.
|
||||||
|
moduleEventSubject("gitea", "pull.merged"),
|
||||||
}
|
}
|
||||||
|
|
||||||
// moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so
|
// moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so
|
||||||
|
|||||||
+1
-1
@@ -25,7 +25,7 @@ accounts {
|
|||||||
users = [
|
users = [
|
||||||
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
|
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
|
||||||
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] }
|
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] }
|
||||||
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] }
|
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] }
|
||||||
allow_responses: { max: 1, ttl: "1m" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
|
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
|
||||||
|
|||||||
@@ -117,6 +117,19 @@ type Announcement struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Upgraded is what the catalogue says when a module's current version moves.
|
// Upgraded is what the catalogue says when a module's current version moves.
|
||||||
|
// SourceMoved is what the forge announces when a pull request is merged: which repository, into
|
||||||
|
// which branch, producing which commit. The mesh matches it against every module's recorded
|
||||||
|
// source and builds what moved, bases first.
|
||||||
|
type SourceMoved struct {
|
||||||
|
Owner string `json:"owner"`
|
||||||
|
Repo string `json:"repo"`
|
||||||
|
Base string `json:"base"`
|
||||||
|
Head string `json:"head"`
|
||||||
|
Commit string `json:"merge_commit_sha"`
|
||||||
|
CloneURL string `json:"clone_url"`
|
||||||
|
HTMLURL string `json:"html_url"`
|
||||||
|
}
|
||||||
|
|
||||||
type Upgraded struct {
|
type Upgraded struct {
|
||||||
Module string `json:"module"`
|
Module string `json:"module"`
|
||||||
Commit string `json:"commit"`
|
Commit string `json:"commit"`
|
||||||
|
|||||||
@@ -30,6 +30,9 @@ const (
|
|||||||
KindHeartbeat = "heartbeat"
|
KindHeartbeat = "heartbeat"
|
||||||
KindBuilt = "built"
|
KindBuilt = "built"
|
||||||
KindModuleMoved = "module-moved"
|
KindModuleMoved = "module-moved"
|
||||||
|
// KindSourceMoved is the forge announcing a merge: a source moved, and what it produces is
|
||||||
|
// built without anybody telling the mesh (novox/hq 04-ISSUES/131).
|
||||||
|
KindSourceMoved = "source-moved"
|
||||||
KindCatchUp = "catch-up"
|
KindCatchUp = "catch-up"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -67,6 +67,10 @@ func Current(conn *amqp.Connection, channel *amqp.Channel) Inbound {
|
|||||||
// are events get their own queue each, and only when something is listening.
|
// are events get their own queue each, and only when something is listening.
|
||||||
func (c *currentInbound) Also(kind string) error {
|
func (c *currentInbound) Also(kind string) error {
|
||||||
switch kind {
|
switch kind {
|
||||||
|
case KindSourceMoved:
|
||||||
|
// Not followed on the bus the mesh is leaving: the forge's merges are announced on the
|
||||||
|
// new one, and this transport goes with the move (design 28, task 5.5).
|
||||||
|
return nil
|
||||||
case KindModuleMoved:
|
case KindModuleMoved:
|
||||||
if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil {
|
if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil {
|
||||||
return err
|
return err
|
||||||
|
|||||||
@@ -48,7 +48,7 @@ func Nats(js *broker.JetStream) Inbound {
|
|||||||
// whatever was asked for — and not at all when nothing was.
|
// whatever was asked for — and not at all when nothing was.
|
||||||
func (n *natsInbound) Also(kind string) error {
|
func (n *natsInbound) Also(kind string) error {
|
||||||
switch kind {
|
switch kind {
|
||||||
case KindModuleMoved, KindCatchUp:
|
case KindModuleMoved, KindCatchUp, KindSourceMoved:
|
||||||
n.follows[kind] = true
|
n.follows[kind] = true
|
||||||
return nil
|
return nil
|
||||||
default:
|
default:
|
||||||
@@ -180,6 +180,8 @@ func kindOfSubject(subject string) (string, bool) {
|
|||||||
return KindModuleMoved, true
|
return KindModuleMoved, true
|
||||||
case broker.ControllerFollows[1]:
|
case broker.ControllerFollows[1]:
|
||||||
return KindCatchUp, true
|
return KindCatchUp, true
|
||||||
|
case broker.ControllerFollows[3]:
|
||||||
|
return KindSourceMoved, true
|
||||||
case BuildOutcome():
|
case BuildOutcome():
|
||||||
// A build's outcome is the role's event now, so it arrives on the events stream rather than
|
// A build's outcome is the role's event now, so it arrives on the events stream rather than
|
||||||
// the control branch — and is acted on by the same handler, because what the controller does
|
// the control branch — and is acted on by the same handler, because what the controller does
|
||||||
|
|||||||
@@ -374,3 +374,20 @@ func TestNatsAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
|
|||||||
t.Fatalf("a report was recorded by a store that never came back")
|
t.Fatalf("a report was recorded by a store that never came back")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// A merge announcement is not what these tests are about; taken and forgotten.
|
||||||
|
func (t *toldAbout) SourceMoved(context.Context, SourceMoved) error { return nil }
|
||||||
|
|
||||||
|
// Every subject the controller follows decodes to a kind, and decoding never reaches past the
|
||||||
|
// list: the day the list was three long and the decoder named a fourth, every message panicked
|
||||||
|
// the control plane (2026-09-28).
|
||||||
|
func TestEverySubjectTheControllerFollowsDecodesToAKind(t *testing.T) {
|
||||||
|
for _, subject := range broker.ControllerFollows {
|
||||||
|
if _, ok := kindOfSubject(subject); !ok {
|
||||||
|
t.Errorf("%s is followed and decodes to nothing", subject)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if _, ok := kindOfSubject("mesh.mod.nobody.event.nothing"); ok {
|
||||||
|
t.Error("a subject nobody follows decoded to a kind")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -58,6 +58,8 @@ type Upgrader interface {
|
|||||||
// stop every upgrade behind it — except the store unreachable for the moment, which is asked
|
// stop every upgrade behind it — except the store unreachable for the moment, which is asked
|
||||||
// again for a bounded time (novox/hq issue 083).
|
// again for a bounded time (novox/hq issue 083).
|
||||||
Upgraded(ctx context.Context, u Upgraded) error
|
Upgraded(ctx context.Context, u Upgraded) error
|
||||||
|
// SourceMoved is a merge on the forge: build what that source produces, bases first.
|
||||||
|
SourceMoved(ctx context.Context, m SourceMoved) error
|
||||||
}
|
}
|
||||||
|
|
||||||
// Server acts on what nodes and modules say.
|
// Server acts on what nodes and modules say.
|
||||||
@@ -101,6 +103,9 @@ func (s *Server) Records(r Recorder) { s.recorder = r }
|
|||||||
|
|
||||||
// Follows says what to do about upgrades, and asks for them to be delivered.
|
// Follows says what to do about upgrades, and asks for them to be delivered.
|
||||||
func (s *Server) Follows(u Upgrader) error {
|
func (s *Server) Follows(u Upgrader) error {
|
||||||
|
if err := s.inbound.Also(KindSourceMoved); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
if err := s.inbound.Also(KindModuleMoved); err != nil {
|
if err := s.inbound.Also(KindModuleMoved); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -232,6 +237,8 @@ func (s *Server) act(ctx context.Context, m Control) {
|
|||||||
s.wasBuilt(ctx, m)
|
s.wasBuilt(ctx, m)
|
||||||
case KindModuleMoved:
|
case KindModuleMoved:
|
||||||
s.moved(ctx, m)
|
s.moved(ctx, m)
|
||||||
|
case KindSourceMoved:
|
||||||
|
s.sourceMoved(ctx, m)
|
||||||
case KindCatchUp:
|
case KindCatchUp:
|
||||||
s.catchingUp(ctx, m)
|
s.catchingUp(ctx, m)
|
||||||
default:
|
default:
|
||||||
@@ -603,3 +610,24 @@ func short(commit string) string {
|
|||||||
}
|
}
|
||||||
return commit
|
return commit
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// sourceMoved acts on the forge's announcement of a merge. Taken whatever happens: a build that
|
||||||
|
// fails is reported by the build itself, and re-delivering the merge would only re-fail it.
|
||||||
|
func (s *Server) sourceMoved(ctx context.Context, m Control) {
|
||||||
|
var moved SourceMoved
|
||||||
|
if err := json.Unmarshal(m.Body(), &moved); err != nil {
|
||||||
|
s.log.Printf("a merge announcement could not be read: %v", err)
|
||||||
|
_ = m.Took()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
m.About("merge " + moved.Owner + "/" + moved.Repo + " into " + moved.Base)
|
||||||
|
if moved.Commit == "" || moved.Repo == "" {
|
||||||
|
s.log.Printf("a merge announcement named no repository or no commit; ignored")
|
||||||
|
_ = m.Took()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if err := s.upgrader.SourceMoved(ctx, moved); err != nil {
|
||||||
|
s.log.Printf("%s/%s moved to %.8s and the mesh could not act on it: %v", moved.Owner, moved.Repo, moved.Commit, err)
|
||||||
|
}
|
||||||
|
_ = m.Took()
|
||||||
|
}
|
||||||
|
|||||||
@@ -152,3 +152,5 @@ func TestAnUpgradeHandledDuringShutdownIsLeftForTheBus(t *testing.T) {
|
|||||||
t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to)
|
t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (u upgradesWith) SourceMoved(context.Context, SourceMoved) error { return nil }
|
||||||
|
|||||||
Reference in New Issue
Block a user