Its event queue is durable, so a running catalogue misses nothing. What it cannot have is what was announced before it first ran — and on a fresh mesh that is never arbitrary: the shared base, the store the catalogue runs on, and the catalogue itself are each necessarily built BEFORE a catalogue exists to hear about them. The graph's foundation is the part it never sees. So it says it is catching up, and the control plane re-announces what it recorded, oldest first, marked as a replay. Oldest first because a graph is built in the order things happened: registering a module that stands on a base before the base would point an edge at a version nothing has seen, and the shape of a fresh mesh guarantees the base is both first and the one that was missed. The replayer hands announcements back rather than publishing them, because the wire belongs to the link package and a replay building its own events could drift from what the builder emits — the one thing it must match exactly, since the catalogue has a single handler for both. Its own queue and its own consumer: two consumers on one queue split its messages, and a catch-up request going to whichever half was not listening is a gap that looks like a working mesh. Toward novox/hq 04-ISSUES/050. Claude-Session: https://claude.ai/code/session_01D6qtiYU3P9jk3pnAXyAFyx
484 lines
17 KiB
Go
484 lines
17 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})
|
|
// 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")
|
|
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
|
|
}
|