Files
mesh-controller/cmd/mesh-controller/modules.go
T
jschoubben aecac5bda2 One bus: the AMQP transport is gone from the controller
The mesh runs on the seat's bus alone (novox/hq ADR 0131, design 28 task 5.5). The old
transport's consume loop, build request, tool ask, management API and account scoping are
deleted, and the bus switch with them; the controller connects to the broker seat and to
nothing else. The store-window tests keep their assertions on a bus-less fake, and the tests
that only made sense for the old transport's in-memory holding go with it.
2026-09-28 03:36:16 +02:00

605 lines
21 KiB
Go

package main
import (
"context"
"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)
}
if err := mayIssue(m); err != nil {
return err
}
// Which bus this mesh is on. A module gets a credential for exactly one, and the two are
// made in entirely different ways: on the bus the mesh runs on today an account is a
// management call, and on the bus being built it is a row the next composition writes into
// the server's user list (novox/hq design 25 §4).
busAddress, err := broker.BusAddress()
if err != nil {
return err
}
return issueOnTheNewBus(ctx, inv, m, *forNode, busAddress)
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
}
// mayIssue says whether a module can be given a broker account. The account is delivered as the
// module's own secret named broker; a module that declares none has nothing to read it with, and
// an account nothing reads is an orphan on the bus (novox/hq 04-ISSUES/078). Said before the
// account is made, not after.
func mayIssue(m catalogue.Manifest) error {
if _, reads := m.OwnSecrets["broker"]; !reads {
return fmt.Errorf("%s declares no own secret named broker, so there is nowhere to deliver "+
"an account; a module that speaks on the bus declares \"own-secrets\": {\"broker\": <path>}", m.Module)
}
return nil
}
// issueOnTheNewBus gives an assigned module its credential on the bus being built.
//
// **Three things differ from a management call, and each is the point of the move.** The credential
// is minted into the mesh's records and becomes usable at the next composition, so there is no
// server to be reachable for this to work. The password travels beside the address rather than inside
// it, because the runtime's contract already separates them and a credential embedded in a URL is one
// that leaks into every log line that prints a connection. And the module's durable consumer is
// derived from what it declared rather than declared by name, so a module cannot ask for delivery of
// something it did not say it consumes.
func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue.Manifest,
node, busAddress string) error {
user := broker.Principal{Kind: broker.KindModule, Node: node, Module: m.Module}.Username()
password, err := inv.MintBusPassword(ctx, inventory.BusUser{
Username: user, Kind: inventory.BusModule, Node: node, Module: m.Module,
})
if err != nil {
return err
}
// Where the module is told to find the bus, and what certificate it must present. The same pair
// a node is told, for the same reason: a mesh's bus presents its own certificate, in no public
// trust store, so an address alone fails at TLS.
known, err := broker.FromEnvironment()
if err != nil {
return fmt.Errorf("cannot deliver a credential without knowing where the bus is: %w", err)
}
reachable, err := brokerReachableAt(ctx, inv, known, node)
if err != nil {
return err
}
return issueWith(ctx, inv, m, node, busAddress, known, reachable, user, password)
}
// issueWith is the delivery half: the minted password sealed to the machine as the module's broker
// secret, and the module's consumer created where the bus can be reached. Split from the minting
// so the move can issue every module against a bus whose address it worked out itself
// (`rollout mint`, design 28 task 5.2) rather than the one in this process's environment.
func issueWith(ctx context.Context, inv *inventory.Inventory, m catalogue.Manifest,
node, busAddress string, known broker.Broker, reachable, user, password string) error {
held, err := json.Marshal(struct {
URL string `json:"url"`
Fingerprint string `json:"fingerprint,omitempty"`
Node string `json:"node"`
Module string `json:"module"`
User string `json:"user"`
Password string `json:"password"`
}{
URL: "nats://" + reachable, Fingerprint: known.Fingerprint,
Node: node, Module: m.Module, User: user, Password: password,
})
if err != nil {
return err
}
if err := inv.AcceptSecretForModule(ctx, node, m.Module, "broker", string(held)); err != nil {
return err
}
// And how it hears what it consumes. Derived from its declaration, and only when it declared
// something: a module that consumes nothing needs no consumer, and creating one would be a
// durable subscription nobody reads.
if consumer, needed := broker.ConsumerFor(broker.Principal{
Kind: broker.KindModule, Node: node, Module: m.Module,
Emits: m.Emits, Consumes: m.Consumes, Serves: m.Tools,
}); needed {
if busAddress == "" {
fmt.Printf(" %s consumes; its consumer is created when the bus is reachable (`push`, then "+
"`rollout mint` again is harmless)\n", m.Module)
} else {
js, err := broker.Dial(busAddress)
if err != nil {
return fmt.Errorf("the credential is minted and the mesh cannot reach the bus to create "+
"how %s hears what it consumes: %w", m.Module, err)
}
defer js.Close()
if err := js.EnsureConsumer(consumer); err != nil {
return err
}
}
}
fmt.Printf("bus user %s minted for %s, scoped to what it emits and consumes\n", user, m.Module)
fmt.Printf(" sealed to %s. It arrives with the next push — `push %s` to send it\n", node, node)
fmt.Printf(" and it works once the bus has been told: the user list is composed into the " +
"machine holding mesh-broker\n")
return nil
}