novox/hq 04-ISSUES/156. Issue 146 put the stream into a push consumer's delivery subject. The server will not move that subject while a subscriber is bound, and answers `consumer name already in use` — a message about the name, for a conflict about the subject. A node is bound to its declaration consumer the whole time it is up: that IS a node listening. So every node consumer in a running mesh became one the assertion could not bring to match, and the control plane crash-looped on the assertion it makes before it serves. A fresh mesh showed nothing, because nothing was bound. Kept rather than deleted and re-made. Re-making moves the subject, and a holder may not be allowed to subscribe to the new one yet: the wider grant travels in the bus's user list, which this same control plane composes and a machine applies minutes later. On the live mesh the nodes are granted `_DELIVER.<node>` and not `_DELIVER.<node>.>`, so re-making would have silenced every machine — worse than the collision it fixes, and harder to undo. Kept rather than fatal, which is what 146's change intended and did not do. The bare subject still delivers, and collides only where one holder has two consumers of one name. That is the controller's own pair, and the controller is not bound to them while it asserts, so those do move. Also: an existing consumer's deliver policy is carried across rather than reasserted, because the server refuses to change it and where a consumer starts is its history. Two tests against a real server: a consumer with a subscriber bound keeps its subject, is reported, and still delivers; one with nothing bound moves, so 146's fix still applies where it matters.
830 lines
32 KiB
Go
830 lines
32 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"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.
|
|
// connectLink opens the controller's link over whichever bus this process is on (design 25: one
|
|
// variable moves it). The streams and this controller's consumers are raised first on the new bus,
|
|
// so nothing served here finds them missing.
|
|
func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) {
|
|
busAddress, err := broker.BusAddress()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if inv != nil {
|
|
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
js, err := broker.Dial(busAddress)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
|
broker.BareAddress(busAddress), err)
|
|
}
|
|
return link.ConnectNats(js, enroller, listener), nil
|
|
}
|
|
|
|
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])
|
|
|
|
// 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)
|
|
}
|
|
|
|
// **Which bus this mesh is on, read once** (novox/hq ADR 0116 step 5). Both clients ship; both
|
|
// being live is refused, because a mesh half on each is one where a declaration goes out on one
|
|
// and the report comes back on the other, and every component logs success while it happens.
|
|
|
|
work := link.Enrolment{Inventory: inv, Identity: ident, Broker: known,
|
|
OnNATS: true}
|
|
server, err := connectLink(ctx, inv, work, work)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
// The bus's own objects, asserted on every start. **Not created once at genesis**: a stream
|
|
// somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was
|
|
// replaced all have records and no objects — and a node whose consumer is missing hears nothing
|
|
// while everything else about it looks correct.
|
|
// 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
|
|
}
|
|
|
|
// **With the inventory, so the bus is raised** (novox/hq ADR 0134, design 30). A module's
|
|
// declaration and how it hears what it consumes move together: its consumer is derived from the
|
|
// same records this declaration is composed from. Raised only when the control plane started
|
|
// serving, a module that gained a `consumes` was sent a declaration it could act on and a
|
|
// consumer that never delivered the event — and nothing anywhere said the two disagreed
|
|
// (found on review, 2026-09-28). Everything the raise does is idempotent.
|
|
server, err := connectLink(ctx, inv, nil, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
if err := link.Declare(ctx, server.Bus(), 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 := connectLink(ctx, nil, 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)
|
|
}
|
|
|
|
// Held from composing to sending, so a converge on one of them cannot send between the two
|
|
// and be overtaken by what was composed before it (novox/hq ADR 0100).
|
|
held, release, err := holdNodes(ctx, open, asked)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
sending, refusals := composeEach(asked, func(node string) (sendable, error) {
|
|
plan, settings, err := planFor(held, open, node)
|
|
if err != nil {
|
|
return sendable{}, 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(held, open, node, plan, settings, gens, Allocating)
|
|
})
|
|
|
|
sentDigest := map[string]string{}
|
|
defer release()
|
|
for _, s := range sending {
|
|
body, err := s.declared.Body()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := link.Declare(ctx, server.Bus(), 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.declared.Resources))
|
|
}
|
|
release()
|
|
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.
|
|
// Held for this round only, and after the last round's were given back, so two pushes
|
|
// cascading into each other's machines never each wait on the other.
|
|
refused, err := sendRound(ctx, open, also,
|
|
func(held context.Context, node string) (sendable, error) {
|
|
plan, settings, err := planFor(held, open, node)
|
|
if err != nil {
|
|
return sendable{}, err
|
|
}
|
|
reportUnhostable(node, plan)
|
|
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
|
},
|
|
func(s readyNode, body []byte) error {
|
|
if err := link.Declare(ctx, server.Bus(), 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.declared.Resources))
|
|
return nil
|
|
})
|
|
refusals = append(refusals, refused...)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// 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
|
|
declared sendable
|
|
}
|
|
|
|
// 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) (sendable, error)) ([]readyNode, []string) {
|
|
|
|
var sending []readyNode
|
|
var refusals []string
|
|
for _, name := range names {
|
|
declared, err := compose(name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
if len(declared.Resources) == 0 {
|
|
// Sent, not skipped (novox/hq issue 127). A node whose declaration composes to
|
|
// nothing may have HELD something before — the broker opening a placement gave it,
|
|
// say — and skipping the empty declaration leaves that last resource in force
|
|
// forever, re-applied by the node's own heartbeat, with no way for the mesh to say
|
|
// it is gone. An empty declaration is the correction: the host drops what the mesh
|
|
// owned and keeps what it found (the adoption envelope still rides along). A node
|
|
// that never held anything applies it as the no-op it is.
|
|
fmt.Printf("%s owns nothing now — sent so it drops what it last held\n", name)
|
|
}
|
|
sending = append(sending, readyNode{name, declared})
|
|
}
|
|
return sending, refusals
|
|
}
|
|
|
|
// sendRound holds the named nodes, composes each and sends each that composed, and gives the hold
|
|
// back on every way out — a body that cannot be marshalled and a send that fails included
|
|
// (novox/hq ADR 0100). A node that cannot be composed is a refusal, not an error: the others are
|
|
// still sent.
|
|
func sendRound(ctx context.Context, open *stores, names []string,
|
|
compose func(held context.Context, node string) (sendable, error),
|
|
send func(s readyNode, body []byte) error) ([]string, error) {
|
|
held, release, err := holdNodes(ctx, open, names)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer release()
|
|
sending, refused := composeEach(names, func(node string) (sendable, error) {
|
|
return compose(held, node)
|
|
})
|
|
for _, s := range sending {
|
|
body, err := s.declared.Body()
|
|
if err != nil {
|
|
return refused, err
|
|
}
|
|
if err := send(s, body); err != nil {
|
|
return refused, err
|
|
}
|
|
}
|
|
return refused, nil
|
|
}
|
|
|
|
// 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
|
|
}
|
|
|
|
// Held from composing to sending (novox/hq ADR 0100); a caller that holds them already —
|
|
// converge, which flips the node and then sends it — is not made to wait on itself.
|
|
ctx, release, err := holdNodes(ctx, open, names)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer release()
|
|
|
|
var sending []readyNode
|
|
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)
|
|
declared, 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, readyNode{name, declared})
|
|
}
|
|
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 := connectLink(ctx, nil, nil, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
for _, s := range sending {
|
|
body, err := s.declared.Body()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := link.Declare(ctx, server.Bus(), 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.declared.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
|
|
}
|
|
declared, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
body, err := declared.Body()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out[n.Name] = digestOf(body)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// raiseTheBus asserts the streams and consumers the mesh's own traffic needs.
|
|
//
|
|
// **Every start, and it says what it did.** The objects are the mesh's, created by nothing else —
|
|
// the controller is their only writer (design 25 §3) — so a mesh that came up without them is one
|
|
// where nodes connect, authenticate, and hear nothing. Said rather than silent for the reason the
|
|
// first line of `serve` is said: a log that is quiet on success and loud on failure reads as broken
|
|
// when it is working.
|
|
func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string) error {
|
|
js, err := broker.Dial(address)
|
|
if err != nil {
|
|
return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
|
broker.BareAddress(address), err)
|
|
}
|
|
defer js.Close()
|
|
// What the raise decided not to fail over. Said, for the reason everything else here is said:
|
|
// a consumer kept as it was is a difference between what the mesh asked for and what the bus
|
|
// holds, and one nobody would find by reading either (novox/hq 04-ISSUES/156).
|
|
js.Note = func(format string, args ...any) {
|
|
fmt.Printf(" "+format+"\n", args...)
|
|
}
|
|
|
|
// **Its own user, before anything else.** The controller's account is created by the installer at
|
|
// a bootstrap password, before there is a controller to mint one — so nothing recorded a hash for
|
|
// it, and the first composition would leave the writer out of the file it was writing. Recorded
|
|
// only if absent: a credential the mesh minted since is the one that counts.
|
|
// **Its own user, before anything else it does here.** The controller's account is created by the
|
|
// installer at a bootstrap password, before there is a controller to mint one — so nothing
|
|
// recorded a hash for it, and the first composition would leave the writer out of the file it was
|
|
// writing: a bus nothing can connect to, produced by the thing connected to it. Recorded only if
|
|
// absent, so a restart cannot put the bootstrap credential back over a rotated one.
|
|
if user, password, _ := broker.CredentialIn(address); user != "" && password != "" {
|
|
if err := inv.SeedBusUser(ctx, inventory.BusUser{
|
|
Username: user, Kind: inventory.BusController,
|
|
}, password); err != nil {
|
|
return fmt.Errorf("cannot record the credential this control plane is using: %w", err)
|
|
}
|
|
}
|
|
|
|
nodes, err := inv.Nodes(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
names := make([]string, 0, len(nodes))
|
|
for _, n := range nodes {
|
|
names = append(names, n.Name)
|
|
}
|
|
if err := broker.Raise(js, names); err != nil {
|
|
return err
|
|
}
|
|
// The work queues of the mesh's own roles (novox/hq ADR 0121). The queue before the holder,
|
|
// deliberately: work queues until somebody arrives to do it, so assigning a build machine a week
|
|
// after something started asking for builds flushes the backlog instead of having lost it.
|
|
// With the seats' holders, so each role's work queue gets the consumer its holder takes
|
|
// work from. Passed as nil until the first live raise, which left the build machine bound to a
|
|
// consumer nothing had created (2026-09-28).
|
|
holders, err := seatHolders(ctx, inv)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
|
|
return err
|
|
}
|
|
// And how every module hears what it consumes. Derived from the same records the user list is
|
|
// composed from, so a module the mesh grants a consumer's subjects has that consumer waiting.
|
|
// Done on every raise, not only when a credential is issued: every module moved onto this bus
|
|
// by the rollout was issued on the old one, and came up with nothing to bind to (2026-09-28).
|
|
records, err := inv.BusRecords(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
users, err := broker.Users(records)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
hearing := 0
|
|
for _, p := range users {
|
|
consumer, needed := broker.ConsumerFor(p)
|
|
if !needed {
|
|
continue
|
|
}
|
|
if err := js.EnsureConsumer(consumer); err != nil {
|
|
return fmt.Errorf("how %s on %s hears what it consumes: %w", p.Module, p.Node, err)
|
|
}
|
|
hearing++
|
|
}
|
|
fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, and %d module(s) "+
|
|
"can hear what they consume\n", broker.BareAddress(address), len(names), hearing)
|
|
return nil
|
|
}
|
|
|
|
// seatHolders is who holds each of the mesh's seats, by seat name: the record where a handover
|
|
// wrote one, and the assigned module claiming the seat otherwise — the same derivation the
|
|
// resolver makes, read from the catalogue rather than re-resolved.
|
|
func seatHolders(ctx context.Context, inv *inventory.Inventory) (map[string]broker.Holder, error) {
|
|
out := map[string]broker.Holder{}
|
|
entries, err := inv.Catalogued(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, e := range entries {
|
|
if len(e.On) == 0 {
|
|
continue
|
|
}
|
|
for _, c := range e.Manifest.Claims {
|
|
seat, known := catalogue.SeatNamed(c.Name)
|
|
if !known {
|
|
continue
|
|
}
|
|
if _, taken := out[seat.Name]; !taken {
|
|
out[seat.Name] = broker.Holder{Node: e.On[0], Module: e.Manifest.Module}
|
|
}
|
|
}
|
|
}
|
|
recorded, err := inv.Holdings(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, h := range recorded {
|
|
if seat, known := catalogue.SeatNamed(h.Claim); known {
|
|
out[seat.Name] = broker.Holder{Node: h.Node, Module: h.Module}
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|