Files
mesh-controller/cmd/mesh-controller/push.go
T
jschoubben d6ee77f17c The consistency cascade sends tolerantly and stops loudly (issue 057 review)
Two robustness fixes to the ADR 0083 cascade, from an adversarial review:

- It routed swept machines through sendTo, which is all-or-nothing — so
  one swept machine's compose error failed the operator's named push and
  skipped its --wait, the intolerance the main path exists to avoid
  (ADR 0066). It now composes them through composeEach, exactly as the
  named send does: a machine that cannot be worked out is a refusal in
  the final report, and the rest are still sent. composeEach's tolerance
  is already covered by TestOneUnresolvableNodeStillLetsTheRestBeSent.
- The fixed 4-round cap could stop a real cascade short in silence. The
  loop is now bounded by the node count (a node is flushed once and never
  revisited, so it cannot run longer) and says so if the guard is ever
  hit, rather than passing over an unfinished cascade quietly.

Scope is unchanged: a named push still flushes every machine left behind,
per ADR 0083 as accepted.
2026-09-20 13:27:12 +02:00

625 lines
22 KiB
Go

package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"flag"
"fmt"
"os"
"sort"
"strings"
"time"
"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/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})
// And what the catalogue decided a build meant. The builder's own result is already handled
// above; this is the other half — the control plane is the only one of the three that knows
// which machines run the thing, so it is the one that acts (novox/hq ADR 0072).
if err := server.Follows(following{open}); err != nil {
return err
}
// And a catalogue that has just started, asking for what it missed. The same type answers
// both: what a build meant and what the builds were are two questions about one record.
if err := server.Answers(following{open}); err != nil {
return err
}
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")
// For a named node, wait until it reports applying exactly what it was sent, so `push <node>`
// means "this node is now what it was told" — a command right after does not race the apply
// (novox/hq ADR 0010). 0 waits for nothing, which is the old fire-and-forget.
wait := set.Duration("wait", 0,
"for a named node, how long to wait for it to report applying what it was sent (0: do not wait)")
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)
})
sentDigest := map[string]string{}
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
}
digest := digestOf(body)
if err := inv.RecordSent(ctx, record.ID, digest); err != nil {
return err
}
sentDigest[s.node] = digest
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.resources))
}
fmt.Printf("\n%d node(s) told\n", len(sending))
// **A named push leaves the mesh consistent, not just the machine it named** (novox/hq
// issue 057, ADR 0083). Assigning a cross-node consumer mints a provision, and the PROVIDER's
// grant list is a pure read of secrets already issued — so after the named node is current,
// other machines can be behind *as a consequence*: their declaration now differs from what
// they were last sent. Those are flushed too, by name, in the push's own output.
//
// Compared against what each machine was last SENT, not against a before/after of this push:
// the mint usually happened at `assign` or `module issue`, before this command ran, so the
// only durable signal is "what it should be" versus "what it last received". A machine behind
// for an unrelated reason is caught here too, which is not a cost — a named push that knew a
// machine was behind and left it so would be the very silence this removes. Bounded: a
// flushed send may itself mint, so this converges over a few rounds.
if len(args) == 1 {
flushed := map[string]bool{args[0]: true}
// Bounded by the node count: a node is marked flushed the round it is handled and is
// never handled twice, so the loop cannot run more than len(nodes) rounds. The bound is
// a guard against a logic error, not a real limit — if it were ever hit, that is a bug
// rather than a cascade legitimately still converging, so it is said rather than passed
// over in silence, unlike the earlier fixed cap that could stop a real cascade short.
rounds := 0
for {
would, err := wouldSend(ctx, open, nodes)
if err != nil {
return err
}
behind, err := inv.Waiting(ctx, would)
if err != nil {
return err
}
var also []string
for _, m := range behind {
if !flushed[m.Node] {
also = append(also, m.Node)
}
}
if len(also) == 0 {
break
}
if rounds++; rounds > len(nodes) {
fmt.Printf("\nstopped cascading after %d rounds with %s still behind — this "+
"should not happen; run `push --behind` to finish\n",
rounds-1, strings.Join(also, ", "))
break
}
sort.Strings(also)
fmt.Printf("\nthis push left %s behind — a provision granted from there, or a "+
"declaration since changed; sending it too\n", strings.Join(also, ", "))
// Tolerantly, exactly as the named send above: a machine that cannot be composed is
// collected as a refusal and reported at the end, and the others are still sent
// (novox/hq ADR 0066). The earlier cut routed these through sendTo, which is
// all-or-nothing — so one swept machine's compose error failed the operator's named
// push and skipped its --wait, the very intolerance the main path exists to avoid.
sending, refused := composeEach(also, func(node string) ([]map[string]any, error) {
plan, settings, err := planFor(ctx, open, node)
if err != nil {
return nil, err
}
reportUnhostable(node, plan)
return declarationWith(ctx, open, node, plan, settings, gens, Allocating)
})
refusals = append(refusals, refused...)
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))
}
// Every candidate this round is marked handled — the sent ones so they are not
// re-listed, and the refused ones so a machine that cannot be composed does not make
// the loop spin on it for ever. Its refusal is already in the report.
for _, name := range also {
flushed[name] = true
}
}
}
// A named node is a request to make THAT node current now, so it waits for the node to say it
// applied exactly this. A whole-mesh or --behind push does not wait: it is a sweep, and blocking
// on the slowest machine would hold back the report on all the others.
if *wait > 0 && len(args) == 1 {
if err := waitForApplied(ctx, inv, args[0], sentDigest[args[0]], *wait); err != nil {
return err
}
}
return couldNotBeResolved(refusals, len(sending))
}
// waitForApplied blocks until the node reports it applied exactly the declaration just sent, or the
// wait runs out. A report of failure or refusal for that same declaration ends the wait at once —
// there is nothing to wait for, and the reason is the node's own.
func waitForApplied(ctx context.Context, inv *inventory.Inventory, node, digest string, wait time.Duration) error {
if digest == "" {
return nil // nothing was sent to this node
}
deadline := time.Now().Add(wait)
for {
doing, said, err := inv.DoingOf(ctx, node)
if err != nil {
return err
}
if said && doing.Declared == digest {
switch doing.Outcome {
case inventory.OutcomeApplied:
fmt.Printf("%s applied it\n", node)
return nil
case inventory.OutcomeFailed:
return fmt.Errorf("%s applied what it was sent but %d resource(s) failed", node, len(doing.Failed))
case inventory.OutcomeRefused:
return fmt.Errorf("%s refused what it was sent: %s", node, doing.Refused)
}
}
if time.Now().After(deadline) {
return fmt.Errorf("%s did not report applying what it was sent within %s "+
"(it may still be converging; check `status`)", node, wait)
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(500 * time.Millisecond):
}
}
}
// 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
}