Files
mesh-controller/cmd/mesh-control/push.go
T
jschoubben c4030947b0 The routing record is 0066, not 0056
0056 is 'the authority is the control plane, not a database'. A citation
pointing at the wrong decision is worse than none: it reads as corroboration.

Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
2026-09-11 00:09:56 +02:00

473 lines
16 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/catalogue"
"github.com/novox/mesh-control/internal/inventory"
"github.com/novox/mesh-control/internal/link"
)
// reportUnhostable says which of a node's assigned modules the machine cannot run, once per push.
//
// A module whose declared capability has no detector on the machine is on the wrong machine. It is
// kept out of what the node is sent — the healthy modules beside it still converge — and named here
// so it is neither silently dropped nor a reason the whole node fails to push.
func reportUnhostable(node string, plan catalogue.Resolution) {
for _, u := range plan.Unhostable {
for _, c := range u.Missing {
fmt.Printf("%s not applied — %s\n", node, catalogue.WrongMachine(u.Module, c, node))
}
}
}
// 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 {
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
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.
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
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")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
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, open, 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, open)
if err != nil {
return err
}
server, err := link.Connect(nil, nil)
if err != nil {
return err
}
defer server.Close()
// Which machines this push is about, before any of them is worked out.
var asked []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"))
}
}
asked = append(asked, n.Name)
}
sending, refusals := composeEach(asked, func(node string) ([]map[string]any, error) {
plan, settings, err := planFor(ctx, open, node)
if err != nil {
return nil, err
}
// A module assigned here that this machine cannot host is said and left out, not fatal: the
// healthy modules beside it are still resolved and sent. Reported so it is not silently
// dropped — the remedy is to move it, and until then the rest of the node converges.
reportUnhostable(node, plan)
// 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.
return declarationWith(ctx, open, node, plan, settings, gens, Allocating)
})
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 couldNotBeResolved(refusals, len(sending))
}
// readyNode is one machine and the declaration it would be sent.
type readyNode struct {
node string
resources []map[string]any
}
// composeEach works out what each named machine should be, and never lets one machine's answer
// decide another's.
//
// **A machine whose set cannot be worked out is that machine's problem** (novox/hq ADR 0066). A
// whole-mesh push used to refuse outright when any one node failed to resolve, so a single
// unanswerable requirement on a single machine — one module requiring a provision nobody had
// assigned a provider for — left every other machine in the mesh unconverged, including machines
// with no relation to it at all. Nothing was sent anywhere, and the machines that could not be sent
// were the ones with nothing wrong with them.
//
// It is the same rule a92c11b established one level down, where an un-hostable module stopped
// taking down the healthy modules beside it, applied one level up: **the blast radius of a fault is
// the thing that has it.** What could not be worked out is named and returned, so a push still ends
// with a non-zero outcome and nobody mistakes a partial convergence for a whole one.
//
// The all-or-nothing rule is kept where it means something — sendTo, which rotates a credential
// across two machines that must agree — and dropped here, where it never did.
func composeEach(names []string,
compose func(node string) ([]map[string]any, error)) ([]readyNode, []string) {
var sending []readyNode
var refusals []string
for _, name := range names {
resources, err := compose(name)
if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
continue
}
if len(resources) == 0 {
fmt.Printf("%s is assigned nothing — skipped\n", name)
continue
}
sending = append(sending, readyNode{name, resources})
}
return sending, refusals
}
// couldNotBeResolved is what a push ends with when some machines could not be worked out.
//
// **After the rest have been sent, never instead of sending them.** It is still an error, because
// the mesh is not in the state somebody asked for and a command that exits cleanly having skipped a
// machine is a command that lies. What it must not do is decide anything about the machines beside
// it, which is why it says how many were sent.
func couldNotBeResolved(refusals []string, sent int) error {
if len(refusals) == 0 {
return nil
}
return fmt.Errorf(
"%d node(s) could not be resolved and were not sent. %d other node(s) were:\n\n%s",
len(refusals), sent, strings.Join(refusals, "\n\n"))
}
// sendTo resolves and sends to exactly the machines named, or refuses without sending anything.
//
// The all-or-nothing rule push deliberately does NOT follow, and for a reason that holds here and
// not there: 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, open *stores, names []string) error {
inv := open.inventory
ident, err := openIdentity(ctx)
if err != nil {
return err
}
defer ident.Close()
gens, err := generators(ctx, open)
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, open, name)
if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
continue
}
reportUnhostable(name, plan)
resources, err := declarationWith(ctx, open, name, plan, settings, gens, Allocating)
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, open *stores,
nodes []inventory.Node) (map[string]string, error) {
gens, err := generators(ctx, open)
if err != nil {
return nil, err
}
out := map[string]string{}
for _, n := range nodes {
plan, settings, err := planFor(ctx, open, n.Name)
if err != nil {
continue
}
resources, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading)
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
}