Files
mesh-controller/cmd/mesh-controller/modules.go
T
jschoubben 79a9e17df0 The broker credential resolves the mesh-broker seat, not the hub (issue 059)
An adversarial review of the 055 fix found it encoded the wrong invariants, latent while
every mesh keeps its broker on the hub. Now: the address is the overlay name of the node
ASSIGNED a module claiming the mesh-broker seat (the hub stands in only while nothing holds
the seat — genesis); "on the overlay" is what whereEveryoneIs answers (resolved the
networking module), not "has an address"; a portless genesis address defaults to 5671
instead of silently disabling the path; a second `overlay place --hub` is refused rather
than last-write-wins; and `overlay place` says that earlier credentials keep their old
address. A test now binds the controller's own module.json to its seat, so deleting the
claim fails the suite.

https://claude.ai/code/session_01D6qtiYU3P9jk3pnAXyAFyx
2026-09-18 01:02:26 +02:00

555 lines
18 KiB
Go

package main
import (
"context"
"crypto/rand"
"encoding/base64"
"encoding/json"
"errors"
"flag"
"fmt"
"net"
"os"
"sort"
"strings"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/overlay"
)
// the catalogue: what exists, what is assigned, and how it is configured.
//
// Split out of main.go, which had reached 2,769 lines because appending was always the
// cheapest next step. That is how novox/hq ADR 0001 records `hal/sdk` reaching 34,636:
// nothing in it was wrong, and no one edit was the one that should have been a new file.
// provided is what comes with the control plane rather than from a repository.
//
// WireGuard, the names, and the domain module over both. The first two are here because the code
// that works out their files is here:
// a peer list is derived from every machine at once, so it cannot be written in a manifest, and
// whatever computes it has to live wherever the whole picture is.
//
// **It is a module in every other respect** — assigned, unassigned, resolved, settled, and absent
// from a machine nobody gave it to.
func providedModules() []catalogue.Manifest {
var out []catalogue.Manifest
for _, raw := range []map[string]any{
// **Two used to be here and are gone**: one wrote the mesh's names into a hosts file, the
// other wrote the same machines as wildcards for a resolver to read. Neither ran software
// and neither could be swapped for anything, which is the test of whether a thing is a
// module at all (novox/hq ADR 0040). They existed because computed output needed somewhere
// to live, and now a module says where it wants it — `facts` in its own manifest.
overlay.Manifest(), overlay.DomainManifest(),
} {
var m catalogue.Manifest
b, _ := json.Marshal(raw)
_ = json.Unmarshal(b, &m)
out = append(out, m)
}
return out
}
var provided = providedModules()
func moduleCommand(ctx context.Context, args []string) error {
if len(args) == 0 {
return errors.New("module add <file>, module list, or module forget <name>")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
switch args[0] {
case "add":
set := flag.NewFlagSet("module add", flag.ContinueOnError)
repo := set.String("source", "", "where this module comes from")
ref := set.String("ref", "", "the branch followed there")
commit := set.String("commit", "", "the commit this manifest was read at")
positionals, err := parseAround(set, args[1:])
if err != nil {
return err
}
if len(positionals) != 1 {
return errors.New("module add <manifest.json> [--source <repo> --ref <branch> --commit <sha>]")
}
raw, err := os.ReadFile(positionals[0])
if err != nil {
return err
}
m, err := catalogue.ParseManifest(raw)
if err != nil {
return err
}
// Provenance together or not at all. A source with no commit cannot be compared against
// anything, so it would record where the module came from and still never be able to say
// the mesh is behind it — which is the one thing recording it is for.
if (*repo == "") != (*commit == "") {
return errors.New("--source and --commit go together: a source with no commit " +
"cannot be compared against anything, and a commit with no source has nothing " +
"to be compared with")
}
if err := inv.RegisterModule(ctx, m, inventory.Source{
Repository: *repo, Ref: *ref, BuiltFrom: *commit,
}); err != nil {
return err
}
fmt.Printf("%s registered", m.Module)
if *commit != "" {
fmt.Printf(" from %s", short(*commit))
}
if len(m.Provides) > 0 {
fmt.Printf(", providing %s", describeOffers(m.Provides))
}
fmt.Println()
for _, c := range m.Claims {
fmt.Printf(" claims %s, one per %s\n", c.Name, c.At())
}
return nil
case "list":
// The catalogue: what exists, where it came from, whether it is current, and who runs it.
// The provenance was recorded from the first build and nothing showed it, which made
// "is this current?" a question you could only answer by reading the database.
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
if len(entries) == 0 {
fmt.Println("this mesh knows about no modules yet")
return nil
}
var stale int
for _, e := range entries {
m := e.Manifest
fmt.Printf("%-18s %-8s", m.Module, m.Version)
switch {
case e.Provided:
fmt.Printf(" %-22s", "with the control plane")
case e.Source.Repository == "":
// Handed over by hand. Legitimate — it is how a module is fixed in a hurry — and
// worth saying, because nothing can rebuild it.
fmt.Printf(" %-22s", "handed over")
case !e.Source.Current():
stale++
fmt.Printf(" %-22s", "behind "+short(e.Source.BuiltFrom)+" < "+short(e.Source.Head))
default:
fmt.Printf(" %-22s", "built "+short(e.Source.BuiltFrom))
}
if len(e.On) > 0 {
fmt.Printf(" on %s", strings.Join(e.On, ", "))
} else {
fmt.Printf(" on nothing")
}
fmt.Println()
var says []string
if len(m.Provides) > 0 {
says = append(says, "provides "+describeOffers(m.Provides))
}
if len(m.Requires) > 0 {
says = append(says, "requires "+strings.Join(m.Requires, ", "))
}
for _, c := range m.Claims {
says = append(says, "claims "+c.At()+"/"+c.Name)
}
if len(m.Capabilities) > 0 {
says = append(says, "needs "+strings.Join(m.Capabilities, ", "))
}
if len(says) > 0 {
fmt.Printf(" %s\n", strings.Join(says, " · "))
}
}
if stale > 0 {
fmt.Printf("\n%d module(s) behind their source — `build --behind` to catch up\n", stale)
}
return nil
case "moved":
if len(args) != 3 {
return errors.New("module moved <name> <commit> — the source has a newer commit")
}
if err := inv.SourceMoved(ctx, args[1], args[2]); err != nil {
return err
}
from, err := inv.SourceOf(ctx, args[1])
if err != nil {
return err
}
if from.Current() {
fmt.Printf("%s is current at %s\n", args[1], short(from.Head))
return nil
}
fmt.Printf("%s is behind: the mesh holds %s and the source has %s\n",
args[1], short(from.BuiltFrom), short(from.Head))
fmt.Printf(" run `build %s` to catch up\n", from.Repository)
return nil
case "forget":
// **What goes with it is said before it goes** (novox/hq 04-ISSUES/017). The settings, the
// module's own secrets and the ports the mesh chose all cascade off the module row, so
// `forget` used to destroy them and report "forgotten" — an action succeeding into a state
// its own verify would reject, and a sealed secret is not recoverable afterwards.
set := flag.NewFlagSet("module forget", flag.ContinueOnError)
andHeld := set.Bool("and-what-it-holds", false,
"discard its settings, its own secrets and its ports along with it")
positionals, err := parseAround(set, args[1:])
if err != nil {
return err
}
if len(positionals) != 1 {
return errors.New("module forget <name> [--and-what-it-holds]")
}
if !*andHeld {
if err := inv.ForgetModule(ctx, positionals[0]); err != nil {
return err
}
fmt.Printf("%s forgotten\n", positionals[0])
return nil
}
held, err := inv.DiscardModule(ctx, positionals[0])
if err != nil {
return err
}
fmt.Printf("%s forgotten\n", positionals[0])
for _, line := range held.Lines() {
// Said after the fact as well as before it: this is the only record that these
// existed, and the next person to ask why the module came back empty reads it here.
fmt.Printf(" discarded%s\n", strings.TrimPrefix(line, " "))
}
return nil
case "issue":
// A module's broker account, scoped by its emits and consumes (novox/hq ADR 0043) and
// sealed to the machine that will run it — the generic case the builder was the first of.
set := flag.NewFlagSet("module issue", flag.ContinueOnError)
forNode := set.String("node", "",
"the machine that will run it, so the credential is delivered instead of printed")
positionals, err := parseAround(set, args[1:])
if err != nil {
return err
}
if len(positionals) != 1 {
return errors.New("module issue <module> --node <machine>")
}
module := positionals[0]
if *forNode == "" {
return errors.New("module issue needs --node: a module's account is sealed to the " +
"machine that runs it, and the mesh cannot read it back to print")
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
return err
}
m, ok := shelf[module]
if !ok {
return fmt.Errorf("this mesh knows no module %q; `module add` it first", module)
}
management, err := broker.ManagementFromEnvironment()
if err != nil {
return err
}
// The foundation owns the bus; make sure it exists before a module binds onto it.
if err := management.EnsureEventExchanges(ctx); err != nil {
return err
}
secret := make([]byte, 32)
if _, err := rand.Read(secret); err != nil {
return err
}
password := base64.RawURLEncoding.EncodeToString(secret)
account, err := management.CreateModuleAccount(ctx, *forNode, module, password, m.Emits, m.Consumes)
if err != nil {
return err
}
// A consumer's queue, with its dead-letter, is the foundation's to declare — its own account
// may not (ADR 0043). Made now, so it exists before the module binds onto it.
if len(m.Consumes) > 0 {
if err := management.EnsureModuleQueue(ctx, *forNode, module); err != nil {
return err
}
}
known, err := broker.FromEnvironment()
if err != nil {
return fmt.Errorf("cannot deliver a credential without knowing where the broker is: %w", err)
}
brokerAddr, err := brokerReachableAt(ctx, inv, known, *forNode)
if err != nil {
return err
}
// The URL and what verifies the broker, together — a mesh's broker presents its own
// certificate, in no public trust store, so a URL alone fails at TLS (as `builder issue`).
held, err := json.Marshal(struct {
URL string `json:"url"`
Fingerprint string `json:"fingerprint,omitempty"`
Node string `json:"node"`
Module string `json:"module"`
}{
URL: fmt.Sprintf("amqps://%s:%s@%s/", account, password, brokerAddr),
Fingerprint: known.Fingerprint,
// The node and module the account is for, so the runtime names its queue as the mesh
// scoped it (<node>.<module>.events) without a manifest having to interpolate a node.
Node: *forNode,
Module: module,
})
if err != nil {
return err
}
if err := inv.AcceptSecretForModule(ctx, *forNode, module, "broker", string(held)); err != nil {
return err
}
fmt.Printf("broker account %s created for %s, scoped to what it emits and consumes\n",
account, module)
fmt.Printf(" sealed to %s. It arrives with the next push — `push %s` to send it\n",
*forNode, *forNode)
return nil
default:
return fmt.Errorf("module has no %q; it has add, list, moved, forget and issue", args[0])
}
}
func assignCommand(ctx context.Context, verb string, args []string) error {
if len(args) != 2 {
return fmt.Errorf("%s <node> <module>", verb)
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
// The act itself is in acts.go, so the command API refuses exactly what this refuses
// (novox/hq ADR 0035). What differs between the surfaces is how the answer is printed.
act := assign
if verb == "unassign" {
act = unassign
}
said, err := act(ctx, open, args[0], args[1])
if said != "" {
fmt.Println(said)
}
if err != nil {
fmt.Println()
return err
}
return nil
}
func settingsCommand(ctx context.Context, args []string) error {
if len(args) == 0 {
return errors.New("settings set <module> <file> [--node <node>], or settings clear <module> [--node <node>]")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
set := flag.NewFlagSet("settings", flag.ContinueOnError)
node := set.String("node", "", "one machine, rather than the whole mesh")
positionals, err := parseAround(set, args[1:])
if err != nil {
return err
}
where := "the whole mesh"
if *node != "" {
where = *node
}
switch args[0] {
case "set":
if len(positionals) != 2 {
return errors.New("settings set <module> <settings.json> [--node <node>]")
}
raw, err := os.ReadFile(positionals[1])
if err != nil {
return err
}
var values map[string]any
if err := json.Unmarshal(raw, &values); err != nil {
return fmt.Errorf("%s is not a settings file: %w", positionals[1], err)
}
if err := inv.SetSettings(ctx, *node, positionals[0], values); err != nil {
return err
}
var keys []string
for k := range values {
keys = append(keys, k)
}
sort.Strings(keys)
fmt.Printf("%s on %s: %s\n", positionals[0], where, strings.Join(keys, ", "))
fmt.Println(" run `push` to send it")
return nil
case "clear":
if len(positionals) != 1 {
return errors.New("settings clear <module> [--node <node>]")
}
if err := inv.ClearSettings(ctx, *node, positionals[0]); err != nil {
return err
}
fmt.Printf("%s on %s is back to what the module says\n", positionals[0], where)
return nil
default:
return fmt.Errorf("settings has no %q; it has set and clear", args[0])
}
}
// describeOffers says what a module provides, and marks the ones answered from anywhere in the
// mesh — because "provides a database" and "provides a shell" are read the same way and mean
// entirely different things about where the answer has to be.
func describeOffers(offers []catalogue.Offer) string {
var out []string
for _, o := range offers {
if o.At() == catalogue.ScopeMesh {
out = append(out, o.Name+" (from anywhere in the mesh)")
continue
}
out = append(out, o.Name)
}
return strings.Join(out, ", ")
}
// pinCommand says which node a machine gets a provision from.
//
// Needed only when more than one could answer, and recordable before that -- a mesh with one
// database should not change where an existing machine gets its data the day a second one
// arrives.
func pinCommand(ctx context.Context, args []string, setting bool) error {
if setting && len(args) != 3 {
return errors.New("pin <node> <provision> <from-node>")
}
if !setting && len(args) != 2 {
return errors.New("unpin <node> <provision>")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
if !setting {
if err := inv.UnpinProvision(ctx, args[0], args[1]); err != nil {
return err
}
fmt.Printf("%s is no longer told where to get %s from\n", args[0], args[1])
return nil
}
if args[0] == args[2] {
// Allowed by nothing here, and worth saying rather than resolving into a confusing
// refusal later: a node providing something to itself is a node-scoped provision, and
// this field is for the other kind.
return fmt.Errorf("%s cannot get %s from itself; that would be a provision this machine "+
"provides, which does not need saying", args[0], args[1])
}
if err := inv.PinProvision(ctx, args[0], args[1], args[2]); err != nil {
return err
}
fmt.Printf("%s gets %s from %s\n", args[0], args[1], args[2])
fmt.Printf(" run `push %s` to send it\n", args[0])
return nil
}
// theBrokerSeat is the mesh-scoped seat the broker module claims (novox/hq ADR 0079: a
// foundation seat is named after the server it guards).
const theBrokerSeat = "mesh-broker"
// brokerReachableAt is the broker's address as the given node can reach it.
//
// The genesis address (MESH_BROKER_ADDRESS) is the broker's public endpoint — right for a machine
// that can dial it, wrong across the mesh, where the path is the overlay. A node that is ON the
// overlay is given the overlay name of the node HOLDING the broker — the one assigned a module
// claiming the mesh-broker seat — because that is where the broker is, whatever else the topology
// says. Only when nothing holds the seat yet (genesis raised the broker as plumbing and no module
// has adopted it) does the hub stand in, which is where the foundation is by convention.
//
// "On the overlay" is what `whereEveryoneIs` answers — a machine that RESOLVED the networking
// module — not "has an address", which is true of every placed machine and says nothing about
// whether anything can reach it (novox/hq issue 059). A node not on the overlay — at genesis,
// before any `overlay place`, which is when the builder's account is issued — keeps the genesis
// address, so nothing about bring-up changes. This is issue 055, corrected by 059.
func brokerReachableAt(ctx context.Context, inv *inventory.Inventory, known broker.Broker, node string) (string, error) {
shelf, err := inv.Catalogue(ctx)
if err != nil {
return "", err
}
onNetwork, err := whereEveryoneIs(ctx, inv, shelf)
if err != nil {
// No catalogue yet is genesis, and at genesis the genesis address is the right one.
return known.Address, nil
}
if onNetwork[node] == "" {
return known.Address, nil
}
// Where the broker is: the node whose assigned set includes a module claiming the seat.
holders := map[string]bool{}
for name, m := range shelf {
for _, c := range m.Claims {
if c.Name == theBrokerSeat && c.At() == catalogue.ScopeMesh {
holders[name] = true
}
}
}
brokerAt := ""
if len(holders) > 0 {
overlays, err := inv.Overlays(ctx)
if err != nil {
return "", err
}
for _, o := range overlays {
assigned, err := inv.Assigned(ctx, o.Name)
if err != nil {
continue
}
for _, a := range assigned {
if holders[a] && onNetwork[o.Name] != "" {
brokerAt = onNetwork[o.Name]
}
}
}
}
if brokerAt == "" {
// Nothing holds the seat (or its node is not on the overlay): the hub, where the
// foundation is by convention — and only a hub that is itself on the overlay, or the
// name we hand out routes nowhere.
overlays, err := inv.Overlays(ctx)
if err != nil {
return "", err
}
for _, o := range overlays {
if o.Hub && onNetwork[o.Name] != "" {
brokerAt = onNetwork[o.Name]
}
}
}
if brokerAt == "" {
return known.Address, nil
}
// A portless genesis address is a working configuration (amqps defaults to 5671), and
// silently keeping the public address would disable this whole path — so the port defaults
// rather than the fix dissolving (novox/hq issue 059).
_, port, err := net.SplitHostPort(known.Address)
if err != nil {
port = "5671"
}
return net.JoinHostPort(brokerAt, port), nil
}