Files
mesh-controller/cmd/mesh-control/push.go
T
jschoubben 8623613704 Split main.go along the seams it already had
2,769 lines and 59 functions, holding command parsing, store opening,
resolution, the board, rotation, licences, builds and status rendering.
Nothing in it was wrong. It grew because appending was always the cheapest next
step, and no single edit was the one that should have been a new file.

That is exactly how novox/hq ADR 0001 records `hal/sdk` reaching 155 files and
34,636 lines — "containing code from every context", with each addition
avoiding a cycle and none of them the mistake. This is the same shape at 8% of
the size, which is why it is worth doing now rather than noting.

Eight files, along boundaries that already existed: what a machine is; the
private network; the catalogue; working out what one machine should be; sending
it; builds; the three questions; and reaching each context's store. main.go
keeps what a main is for — parsing arguments and dispatching.

A pure move. No behaviour changed, no test changed, and the gate is green
before and after — which is the only thing that makes a refactor this size
safe to do in one commit.
2026-08-31 13:30:54 +02:00

410 lines
13 KiB
Go

package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"flag"
"fmt"
"os"
"strings"
"time"
"github.com/novox/mesh-control/internal/broker"
"github.com/novox/mesh-control/internal/inventory"
"github.com/novox/mesh-control/internal/link"
)
// sending it, and holding the link that carries it.
//
// 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.
// 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()
// And build results nobody was waiting for. A build triggered any other way than `build`
// would otherwise be reported into the void, which is the same as not reporting it.
server.Records(builds{inv})
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"
// pushCommand sends nodes everything they should be: their place on the network, and what their
// assignments resolve to.
//
// One declaration, not two. A node holding its network and not its modules, or the reverse, is
// half-configured for as long as that lasts — and the two are computed from the same picture of
// the mesh, so sending them apart would let them disagree.
func pushCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("push", flag.ContinueOnError)
// Only the machines that need it.
//
// **A command rather than a timer, to begin with.** Something that re-pushes on a schedule is
// a scheduler over this, and building the scheduler first would mean two paths to one act
// with nothing to compare them against. A person can run this; so can cron; so can whatever
// eventually watches.
behind := set.Bool("behind", false,
"only machines whose last declaration was refused or partly failed")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
args = positionals
if len(args) > 1 {
return errors.New("push [<node>] [--behind] — one node, or all of them")
}
if len(args) == 1 && *behind {
// Naming a machine and asking for the ones that need it are two different requests, and
// guessing which was meant would sometimes push to a machine somebody did not name.
return errors.New("push <node> or push --behind, not both: one names a machine and the " +
"other asks which machines need one")
}
inv, err := openInventory(ctx)
if err != nil {
return err
}
defer inv.Close()
ident, err := openIdentity(ctx)
if err != nil {
return err
}
defer ident.Close()
// Every node, not only the ones on the private network. A machine that was never given the
// network module still takes modules, and iterating the network here is what used to make
// "on the network" and "managed" the same thing.
nodes, err := inv.Nodes(ctx)
if err != nil {
return err
}
// Which machines are not in the state they were sent, when that is what was asked for.
var needsOne map[string]inventory.Doing
if *behind {
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
if err != nil {
return err
}
needsOne = map[string]inventory.Doing{}
for _, d := range wrong {
needsOne[d.Node] = d
}
// **And every machine not running what the mesh would send it.** "Behind" used to mean
// only "failed or refused", so a machine that applied cleanly and whose declaration has
// since changed was not behind — and novox/hq ADR 0010's question, *did my change go
// out?*, was answerable only for the machines that broke.
would, err := wouldSend(ctx, inv, nodes)
if err != nil {
return err
}
waiting, err := inv.Waiting(ctx, would)
if err != nil {
return err
}
for _, m := range waiting {
if _, already := needsOne[m.Node]; already {
continue
}
needsOne[m.Node] = inventory.Doing{Node: m.Node, Outcome: "waiting"}
}
if len(needsOne) == 0 {
// Said rather than doing nothing quietly. "Nothing needed one" and "this did not run"
// must never look the same.
fmt.Println("every machine is doing what it was told")
return nil
}
}
gens, err := generators(ctx, inv)
if err != nil {
return err
}
server, err := link.Connect(nil, nil)
if err != nil {
return err
}
defer server.Close()
// Every node is resolved before anything is sent. A push that configured three nodes and then
// refused on the fourth would leave the mesh in a state nobody asked for, and the fourth is
// exactly where a claim collision shows up.
type ready struct {
node string
resources []map[string]any
}
var sending []ready
var refusals []string
for _, n := range nodes {
if len(args) == 1 && n.Name != args[0] {
continue
}
if *behind {
doing, needs := needsOne[n.Name]
if !needs {
continue
}
// A machine that has been failing the same way for a long time is not going to stop
// because it was asked again. Said, and pushed to anyway — refusing would leave no
// way to retry after fixing the cause, and this is a command somebody ran.
//
// Only for machines that reported something. One that is merely waiting has no report
// to be old, and saying it had been failing since the zero time would be a sentence
// about nothing.
if since := time.Since(doing.At); doing.Outcome != "waiting" && since > 6*time.Hour {
fmt.Printf("%s has been %s since %s; pushing again anyway, but the cause is "+
"unlikely to be timing\n",
n.Name, doing.Outcome, doing.At.Local().Format("2006-01-02 15:04"))
}
}
plan, settings, err := planFor(ctx, inv, n.Name)
if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", n.Name, err))
continue
}
// The private network is in here with everything else. It used to be composed separately
// and prepended, which meant every machine with an address was on it and no machine could
// be kept off. It is a module now, so it arrives the way a module does.
resources, err := declarationWith(ctx, inv, n.Name, plan, settings, gens)
if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", n.Name, err))
continue
}
if len(resources) == 0 {
fmt.Printf("%s is assigned nothing — skipped\n", n.Name)
continue
}
sending = append(sending, ready{n.Name, resources})
}
if len(refusals) > 0 {
return fmt.Errorf("nothing was sent. %d node(s) could not be resolved:\n\n%s",
len(refusals), strings.Join(refusals, "\n\n"))
}
for _, s := range sending {
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": s.resources})
if err != nil {
return err
}
if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil {
return err
}
// After it is away, not before. A digest recorded for something that failed to send would
// make the machine look current for a declaration it never received.
record, err := inv.NodeByName(ctx, s.node)
if err != nil {
return err
}
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
return err
}
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.resources))
}
fmt.Printf("\n%d node(s) told\n", len(sending))
return nil
}
// sendTo resolves and sends to exactly the machines named, or refuses without sending anything.
//
// The same all-or-nothing rule push follows, and for the same reason: a rotation that reached the
// consumer and refused on the provider would leave one end holding a credential the other has
// never heard of — which is the state this whole mechanism exists to make impossible.
func sendTo(ctx context.Context, inv *inventory.Inventory, names []string) error {
ident, err := openIdentity(ctx)
if err != nil {
return err
}
defer ident.Close()
gens, err := generators(ctx, inv)
if err != nil {
return err
}
type ready struct {
node string
resources []map[string]any
}
var sending []ready
var refusals []string
for _, name := range names {
plan, settings, err := planFor(ctx, inv, name)
if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
continue
}
resources, err := declarationWith(ctx, inv, name, plan, settings, gens)
if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
continue
}
sending = append(sending, ready{name, resources})
}
if len(refusals) > 0 {
return fmt.Errorf("nothing was sent. %d machine(s) could not be resolved:\n\n%s",
len(refusals), strings.Join(refusals, "\n\n"))
}
server, err := link.Connect(nil, nil)
if err != nil {
return err
}
defer server.Close()
for _, s := range sending {
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": s.resources})
if err != nil {
return err
}
if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil {
return err
}
record, err := inv.NodeByName(ctx, s.node)
if err != nil {
return err
}
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
return err
}
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.resources))
}
return nil
}
// digestOf is what the mesh compares to answer "has this machine been sent what it should be".
//
// Over the same bytes that are sent, so the comparison is of the thing itself rather than of
// something derived beside it that could drift from it.
func digestOf(body []byte) string {
sum := sha256.Sum256(body)
return hex.EncodeToString(sum[:])
}
// wouldSend is the digest of what each machine should be right now.
//
// Machines that do not resolve are left out rather than reported as waiting: "this machine cannot
// be worked out" is a different problem with a different remedy, and `plan` is where it is said.
func wouldSend(ctx context.Context, inv *inventory.Inventory,
nodes []inventory.Node) (map[string]string, error) {
gens, err := generators(ctx, inv)
if err != nil {
return nil, err
}
out := map[string]string{}
for _, n := range nodes {
plan, settings, err := planFor(ctx, inv, n.Name)
if err != nil {
continue
}
resources, err := declarationWith(ctx, inv, n.Name, plan, settings, gens)
if err != nil {
continue
}
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": resources})
if err != nil {
return nil, err
}
out[n.Name] = digestOf(body)
}
return out, nil
}