Files
mesh-controller/cmd/mesh-control/main.go
T
jschoubben 8b974deb42 A working private network, and four reasons it did not work
Three machines across two sites, two of them behind no reachable address, all
nine paths open. The mesh computes the graph, delivers it as a declaration, and
the nodes bring it up.

Every fault below looked like success from inside the mesh: the graph was
right, the files were right, the services were up, every node reported it had
applied. None was reachable by reasoning.

A running interface does not re-read its configuration. A node joins, every
existing node's peer list changes, the file is replaced -- and the service is
already running, so nothing reloads it. Fixed as declared state rather than a
command: the service must reflect the file. A command to restart would be an
action, and the link may not carry one. The host refused exactly that, which is
how this shape was arrived at.

A hub sharing a site with a spoke appeared twice in that spoke's peer list --
once as a direct peer, once as the route of last resort. WireGuard takes one
entry per key and refuses the file. The ordinary shape of a small mesh, and in
none of the tests written before it ran.

Two nodes at one site that neither can be dialled were peered directly. Nobody
opens the path, and the direct route is more specific than the hub's, so it
wins and blackholes -- this design's own warning arriving in its
implementation. They now route through the hub unless one end can be dialled.

And Docker sets the FORWARD policy to DROP, so a hub with ip_forward enabled
carried nothing between its spokes. The substrate at tier 1 silently breaks the
network at tier 2, and nothing in either tier's state says so. The hub inserts
its own rule above those chains and removes it on the way down.

Two weak tests found by injection along the way: one asserted the keepalive
rule only against the hub, whose peer entries happen not to set that field at
all, so it tested an absence; the other checked the firewall rules by looking
for FORWARD anywhere, which the PostDown line satisfies on its own.
2026-08-29 18:04:15 +02:00

636 lines
19 KiB
Go

// Command mesh-control is the control plane: everything that needs to know about more than one
// node (novox/hq ADR 0006).
//
// It runs as one process holding several contexts, each owning its own store. Today it holds one,
// `inventory`, and does one thing with it — brings its schema up to date, which is step 3 of the
// bootstrap in novox/hq 07-the-substrate and the step the first node cannot get past without.
package main
import (
"context"
"errors"
"flag"
"fmt"
"os"
"os/signal"
"strings"
"syscall"
"time"
"github.com/novox/mesh-control/internal/broker"
"github.com/novox/mesh-control/internal/identity"
"github.com/novox/mesh-control/internal/inventory"
"github.com/novox/mesh-control/internal/link"
"github.com/novox/mesh-control/internal/overlay"
"github.com/novox/mesh-control/internal/store"
"github.com/novox/mesh-control/internal/token"
)
// version is stamped at link time. Unset in a development build, and it says so rather than
// claiming a number.
var version = "development build"
// held is a context this process was granted, and the schema it carries.
//
// novox/hq ADR 0006 names seven. One is built. The list is short because the others do not exist
// yet, not because they are optional.
var held = []struct {
name string
migrations func() ([]store.Migration, error)
}{
{inventory.Name, inventory.Migrations},
{identity.Name, identity.Migrations},
}
func main() {
if err := run(); err != nil {
fmt.Fprintf(os.Stderr, "mesh-control: %v\n", err)
os.Exit(1)
}
}
func run() error {
args := os.Args[1:]
if len(args) == 0 {
usage()
return fmt.Errorf("no command given")
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
switch args[0] {
case "migrate":
return migrate(ctx)
case "node":
return nodeCommand(ctx, args[1:])
case "token":
return tokenCommand(ctx, args[1:])
case "identity":
return identityCommand(ctx, args[1:])
case "broker":
return brokerCommand(args[1:])
case "serve":
return serve(ctx)
case "declare":
return declare(ctx, args[1:])
case "overlay":
return overlayCommand(ctx, args[1:])
case "version":
fmt.Println(version)
return nil
case "help", "-h", "--help":
usage()
return nil
default:
usage()
return fmt.Errorf("%q is not a command", args[0])
}
}
func usage() {
fmt.Fprint(os.Stderr, `mesh-control — the control plane
migrate bring each context's schema up to date
node add <name> create a node record
node list the nodes this mesh knows about
token issue --node <name> a one-time right to join, for an existing record
token issue --new <name> create the record and issue for it
identity show this control plane's signing key
broker show where the broker is, and what to expect there
serve consume what nodes say, and answer
declare <node> <file> send a node a signed declaration
overlay place <node> [flags] say where a node is and how it is reached
overlay show the private network, as the mesh computes it
overlay push send every node its part of the private network
version what this binary is
Each context reaches its own store through its own credential (novox/hq ADR 0008), named
`+store.Variable("<context>")+`. This process holds:
`)
for _, c := range held {
fmt.Fprintf(os.Stderr, " %-12s database %-12s from %s\n",
c.name, store.Database(c.name), store.Variable(c.name))
}
fmt.Fprintln(os.Stderr)
}
// migrate brings every held context's schema up to date.
//
// Reported per context and per migration, because this runs during a bootstrap on a machine with
// nothing else on it — the output is the only account of what happened, and "migrated" is not one.
func migrate(ctx context.Context) error {
for _, c := range held {
migrations, err := c.migrations()
if err != nil {
return err
}
s, err := store.Open(ctx, c.name)
if err != nil {
return err
}
defer s.Close()
// The bootstrap raises PostgreSQL moments before this runs, and a container that is
// running is not a database that will answer — a distinction this project has already
// paid for once, when a crash-looping database reported itself as up between restarts.
if err := s.Ready(ctx, 60*time.Second); err != nil {
return err
}
done, err := s.Migrate(ctx, migrations)
for _, m := range done {
fmt.Printf("%s: applied %04d-%s\n", c.name, m.Number, m.Name)
}
if err != nil {
return err
}
if len(done) == 0 {
applied, err := s.AppliedMigrations(ctx)
if err != nil {
return err
}
fmt.Printf("%s: already up to date — %d migration(s)\n", c.name, len(applied))
}
}
return nil
}
// openInventory connects and waits, the way every command that touches it needs to.
func openInventory(ctx context.Context) (*inventory.Inventory, error) {
inv, err := inventory.Open(ctx)
if err != nil {
return nil, err
}
if err := inv.Ready(ctx, 30*time.Second); err != nil {
inv.Close()
return nil, err
}
return inv, nil
}
func nodeCommand(ctx context.Context, args []string) error {
if len(args) == 0 {
return errors.New("node add <name>, or node list")
}
inv, err := openInventory(ctx)
if err != nil {
return err
}
defer inv.Close()
switch args[0] {
case "add":
if len(args) != 2 {
return errors.New("node add <name>")
}
node, err := inv.AddNode(ctx, args[1])
if err != nil {
return err
}
fmt.Printf("added %s (%s)\n", node.Name, node.ID)
return nil
case "list":
nodes, err := inv.Nodes(ctx)
if err != nil {
return err
}
if len(nodes) == 0 {
// Said rather than printed as nothing: an empty list and a failed read must never
// look the same, and this command answering "none" is only honest because getting
// here means the store answered.
fmt.Println("this mesh has no node records yet")
return nil
}
for _, n := range nodes {
fmt.Printf("%-20s %s added %s\n", n.Name, n.ID, n.Created.Format(time.RFC3339))
}
return nil
default:
return fmt.Errorf("node has no %q; it has add and list", args[0])
}
}
func tokenCommand(ctx context.Context, args []string) error {
if len(args) == 0 || args[0] != "issue" {
return errors.New("token issue --node <name>, or token issue --new <name>")
}
set := flag.NewFlagSet("token issue", flag.ContinueOnError)
existing := set.String("node", "", "issue for a node record that already exists")
fresh := set.String("new", "", "create the node record, then issue for it")
validFor := set.Duration("for", time.Hour, "how long the token may be used")
if err := set.Parse(args[1:]); err != nil {
return err
}
// Exactly one, because the difference is what the token binds to. A command that guessed
// would sometimes create a second record for a machine that already has one.
if (*existing == "") == (*fresh == "") {
return errors.New("give exactly one of --node <name> or --new <name>: the first is a " +
"machine the mesh already has a record for, the second is one it has never seen")
}
inv, err := openInventory(ctx)
if err != nil {
return err
}
defer inv.Close()
name := *existing
if *fresh != "" {
node, err := inv.AddNode(ctx, *fresh)
if err != nil {
return err
}
name = node.Name
}
issued, err := inv.IssueToken(ctx, name, *validFor)
if err != nil {
return err
}
// Assembled from two contexts by the process that holds both grants. Neither reads the
// other's store (novox/hq ADR 0008) — each is asked for its own part.
ident, err := openIdentity(ctx)
if err != nil {
return err
}
defer ident.Close()
key, err := ident.Establish(ctx)
if err != nil {
return err
}
// The account is created before the token is handed over, which is what removes the
// chicken-and-egg entirely: the mesh runs the broker, so a joining node's credentials can
// exist before it does. The one-time secret IS the password, so a node's first connection is
// already authenticated and enrolment is what happens over it.
if management, err := broker.ManagementFromEnvironment(); err == nil {
if err := management.CreateNodeAccount(ctx, issued.Node.Name, issued.Secret); err != nil {
return err
}
fmt.Printf("broker account %s created, scoped to %s and the %s exchange\n\n",
issued.Node.Name, link.QueueFor(issued.Node.Name), link.Exchange)
} else if !errors.Is(err, broker.ErrNotConfigured) {
return err
}
made := token.Token{Signer: key.Public, Secret: issued.Secret}
// Absent is a state, not a failure: a control plane can hold records and a key before it has
// a broker. What it cannot do is issue a token anybody could use, and Missing() says so.
known, err := broker.FromEnvironment()
switch {
case err == nil:
made.Broker, made.Fingerprint = known.Address, known.Fingerprint
case errors.Is(err, broker.ErrNotConfigured):
default:
return err
}
encoded, err := made.Encode()
if err != nil {
return err
}
fmt.Printf("token for %s, usable once, until %s\n\n %s\n\n",
issued.Node.Name, issued.Expires.Format(time.RFC3339), encoded)
fmt.Println("This is the only time it is shown. What is stored is a hash of the secret.")
if missing := made.Missing(); len(missing) > 0 {
fmt.Printf("\nINCOMPLETE — this token cannot be used to join anything yet. Missing:\n")
for _, m := range missing {
fmt.Printf(" - %s\n", m)
}
fmt.Printf("\nSet %s and %s once the broker is raised.\n",
broker.AddressVar, broker.CertificateVar)
}
return nil
}
func openIdentity(ctx context.Context) (*identity.Identity, error) {
ident, err := identity.Open(ctx)
if err != nil {
return nil, err
}
if err := ident.Ready(ctx, 30*time.Second); err != nil {
ident.Close()
return nil, err
}
return ident, nil
}
func identityCommand(ctx context.Context, args []string) error {
if len(args) == 0 || args[0] != "show" {
return errors.New("identity show")
}
ident, err := openIdentity(ctx)
if err != nil {
return err
}
defer ident.Close()
// Establish rather than read: a control plane asked for its identity before it has one should
// get one, not an error. Generating it is idempotent, so this is safe to run at any time.
key, err := ident.Establish(ctx)
if err != nil {
return err
}
fmt.Printf("signing key %s\n", key.ID)
fmt.Printf("fingerprint %s\n", key.Fingerprint())
fmt.Printf("created %s\n", key.Created.Format(time.RFC3339))
fmt.Printf("\nThe public half of this travels in every enrolment token. A node believes a\n" +
"declaration because it carries a signature this key made (novox/hq ADR 0004).\n")
return nil
}
func brokerCommand(args []string) error {
if len(args) == 0 || args[0] != "show" {
return errors.New("broker show")
}
known, err := broker.FromEnvironment()
if errors.Is(err, broker.ErrNotConfigured) {
fmt.Printf("no broker configured. Set %s and %s.\n\n"+
"Until then tokens carry the signing key and the one-time secret, and say what they\n"+
"are missing. They cannot be used to join.\n",
broker.AddressVar, broker.CertificateVar)
return nil
}
if err != nil {
return err
}
fmt.Printf("address %s\n", known.Address)
fmt.Printf("fingerprint %s\n", known.Fingerprint)
fmt.Print("\nThe fingerprint is computed from the certificate on disk, never configured. A\n" +
"node checks it before sending anything (novox/hq ADR 0004).\n")
return nil
}
// serve is the control plane running: one connection to the broker, one queue, one consumer.
func serve(ctx context.Context) error {
inv, err := openInventory(ctx)
if err != nil {
return err
}
defer inv.Close()
ident, err := openIdentity(ctx)
if err != nil {
return err
}
defer ident.Close()
// Established at start rather than on first use. A control plane that cannot sign is one
// whose declarations every node correctly refuses, and that should be a startup failure
// rather than something discovered at the first declaration.
key, err := ident.Establish(ctx)
if err != nil {
return err
}
fmt.Printf("signing as %s\n", key.Fingerprint()[:16])
management, err := broker.ManagementFromEnvironment()
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
return err
}
// Where the broker is and what to expect there, so a node can be told how to come back
// without a person and a new token.
known, err := broker.FromEnvironment()
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
return err
}
if errors.Is(err, broker.ErrNotConfigured) {
fmt.Printf("no broker address configured, so enrolled nodes will not be told how to "+
"reconnect. Set %s and %s.\n", broker.AddressVar, broker.CertificateVar)
}
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known}
server, err := link.Connect(work, work)
if err != nil {
return err
}
defer server.Close()
return server.Serve(ctx)
}
// declare sends one node a declaration, signed.
//
// Signed here rather than trusted from the broker: a node connects to the broker and takes
// instruction from the control plane behind it, and those are two identities. If a node believed
// whatever arrived on its queue, a compromised broker could forge declarations — and since the
// host applies whatever the link delivers, that is the whole machine (novox/hq ADR 0004).
func declare(ctx context.Context, args []string) error {
if len(args) != 2 {
return errors.New("declare <node> <declaration.json>")
}
node, path := args[0], args[1]
raw, err := os.ReadFile(path)
if err != nil {
return err
}
ident, err := openIdentity(ctx)
if err != nil {
return err
}
defer ident.Close()
// The node has to exist before it can be told anything. Publishing to a queue nobody consumes
// would sit there looking like success.
inv, err := openInventory(ctx)
if err != nil {
return err
}
defer inv.Close()
if _, err := inv.NodeByName(ctx, node); err != nil {
return err
}
server, err := link.Connect(nil, nil)
if err != nil {
return err
}
defer server.Close()
if err := link.Declare(ctx, server.Channel(), ident, node, raw, 15*time.Second); err != nil {
return err
}
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
return nil
}
// OverlayCIDRVar is the range the mesh allocates node addresses from.
const OverlayCIDRVar = "MESH_OVERLAY_CIDR"
func overlayCIDR() string {
if v := strings.TrimSpace(os.Getenv(OverlayCIDRVar)); v != "" {
return v
}
return "10.42.0.0/16"
}
func overlayCommand(ctx context.Context, args []string) error {
if len(args) == 0 {
return errors.New("overlay place <node> [flags], overlay show, or overlay push")
}
inv, err := openInventory(ctx)
if err != nil {
return err
}
defer inv.Close()
switch args[0] {
case "place":
return overlayPlace(ctx, inv, args[1:])
case "show":
return overlayShow(ctx, inv)
case "push":
return overlayPush(ctx, inv)
default:
return fmt.Errorf("overlay has no %q; it has place, show and push", args[0])
}
}
func overlayPlace(ctx context.Context, inv *inventory.Inventory, args []string) error {
if len(args) == 0 {
return errors.New("overlay place <node> [--endpoint host:port] [--site name] [--hub]")
}
node := args[0]
set := flag.NewFlagSet("overlay place", flag.ContinueOnError)
endpoint := set.String("endpoint", "", "where this node can be dialled, or empty for nowhere")
site := set.String("site", "", "where this machine physically is, or empty if it roams")
hub := set.Bool("hub", false, "this node is the hub every other routes through")
if err := set.Parse(args[1:]); err != nil {
return err
}
// Declared, all three. The address is evidence of reachability and is not the fact, and hub
// election by address prefix fails silently (novox/hq ADR 0007).
if err := inv.SetPlace(ctx, node, *endpoint, *site, *hub, ""); err != nil {
return err
}
found, err := inv.NodeByName(ctx, node)
if err != nil {
return err
}
address, err := inv.AssignAddress(ctx, found.ID, overlayCIDR())
if err != nil {
return err
}
fmt.Printf("%s is at %s on the overlay\n", node, address)
switch {
case *hub:
fmt.Println(" the hub — every node not sharing a site routes through it")
case *endpoint == "":
fmt.Println(" not dialable — it opens every path itself")
}
if *site != "" {
fmt.Printf(" at %s, so it peers directly with anything else there\n", *site)
}
return nil
}
// graph reads every node's place and computes the network. Every node at once, which is the whole
// reason this is the control plane's work.
func graph(ctx context.Context, inv *inventory.Inventory) ([]overlay.Node, overlay.Graph, error) {
places, err := inv.Overlays(ctx)
if err != nil {
return nil, nil, err
}
nodes := make([]overlay.Node, 0, len(places))
for _, p := range places {
nodes = append(nodes, overlay.Node{
Name: p.Name, Key: p.Key, Endpoint: p.Endpoint,
Site: p.Site, Hub: p.Hub, Address: p.Address,
})
}
computed, err := overlay.Compute(nodes, overlayCIDR())
return nodes, computed, err
}
func overlayShow(ctx context.Context, inv *inventory.Inventory) error {
nodes, computed, err := graph(ctx, inv)
if err != nil {
return err
}
if len(nodes) == 0 {
fmt.Println("this mesh has no nodes")
return nil
}
for _, n := range nodes {
place := n.Address
if place == "" {
// Said, not skipped. A node with no place is a node with no network, and it should
// be visible here rather than quietly absent from a list of who is on it.
place = "no address — run `overlay place`"
}
fmt.Printf("%-16s %-14s", n.Name, place)
switch {
case n.Hub:
fmt.Print(" hub")
case !n.Reachable():
fmt.Print(" not dialable")
}
if n.Site != "" {
fmt.Printf(" at %s", n.Site)
}
fmt.Println()
for _, p := range computed[n.Name] {
fmt.Printf(" → %-14s %-18s %s\n", p.Name, p.Allowed, p.Why)
}
}
return nil
}
func overlayPush(ctx context.Context, inv *inventory.Inventory) error {
nodes, computed, err := graph(ctx, inv)
if err != nil {
return err
}
ident, err := openIdentity(ctx)
if err != nil {
return err
}
defer ident.Close()
server, err := link.Connect(nil, nil)
if err != nil {
return err
}
defer server.Close()
sent := 0
for _, n := range nodes {
peers, ok := computed[n.Name]
if !ok {
// Skipped, and said. A node with no key or address is not on the network yet, and
// sending it an empty configuration would take down the one it may already have.
fmt.Printf("%s is not on the overlay yet — skipped\n", n.Name)
continue
}
declaration, err := overlay.Declaration(n, peers, "")
if err != nil {
return err
}
if err := link.Declare(ctx, server.Channel(), ident, n.Name, declaration, 15*time.Second); err != nil {
return err
}
fmt.Printf("sent %s its place on the overlay — %d peer(s)\n", n.Name, len(peers))
sent++
}
fmt.Printf("\n%d of %d node(s) told\n", sent, len(nodes))
return nil
}