The identity provider failed every consumer for a day and status called the mesh well (hq issue 179). The controller now follows every provider's provisioner.failing/recovered, keeps the newest failing word per provider, machine and consumer (migration 0065), and status, its JSON and node show name it until it recovers. Every module that receives contributions is granted the two events, so no manifest can forget them.
1253 lines
51 KiB
Go
1253 lines
51 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"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 {
|
|
// The one process whose log is read over time, so the one that says each change to a node's
|
|
// unmet seat dependencies once (novox/hq ADR 0207).
|
|
logUnheldChanges = true
|
|
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, open})
|
|
// Open plans move on a timer as well as on outcomes (novox/hq ADR 0162): a tier waiting for
|
|
// machines to report moves when they have, and a plan left by a replaced controller resumes.
|
|
go planTicker(ctx, open)
|
|
// 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
|
|
}
|
|
// And what providers say about consumers they keep failing, kept for `status` (novox/hq ADR
|
|
// 0224): a provider's journal must not be the only place that says so.
|
|
if err := server.Watches(standings{inv}); err != nil {
|
|
return err
|
|
}
|
|
|
|
// And the mesh's own verbs, as the seat this control plane holds (novox/hq ADR 0154). Served
|
|
// from the store's row, so what the seat declares is what is answered.
|
|
handlers, behind, err := seatToolHandlers()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(behind) > 0 {
|
|
// Said once, loudly, and then served anyway (novox/hq ADR 0185): the mesh keeps answering
|
|
// while whatever put an older control plane here is undone.
|
|
fmt.Printf("this control plane is behind the %s row: it cannot run %s. "+
|
|
"Those answer the reason when called; everything else is served as usual\n",
|
|
catalogue.ControllerSeatName, strings.Join(behind, ", "))
|
|
}
|
|
bus, isNATS := server.Bus().(link.OverNATS)
|
|
if !isNATS {
|
|
return errors.New("the mesh's verbs are served over the bus, and this control plane is not on it")
|
|
}
|
|
stopServing, err := bus.ServeSeatTools(catalogue.ControllerSeatName, handlers, log.New(os.Stdout, "", log.LstdFlags))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer stopServing()
|
|
// And says so on the bus (novox/hq ADR 0197): what it serves, as the NATS services protocol asks.
|
|
stopAnnouncing, err := bus.Announce(seatAnnouncement(handlers), log.New(os.Stdout, "", log.LstdFlags))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer stopAnnouncing()
|
|
|
|
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
|
|
}
|
|
// Written down like every other send (novox/hq issue 204): a declaration a person sent by hand
|
|
// is still what the machine was last told, and status must not read it as current for the one
|
|
// the mesh would compose. Which builds it carried is recorded as not known (novox/hq issue 259):
|
|
// the mesh did not compose it, so a push that does not name this machine treats it as held.
|
|
if _, err := recordSent(ctx, inv, node, raw, nil); 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)
|
|
}
|
|
|
|
// **The machine holding the bus first** (novox/hq issue 249): its declaration carries the bus's
|
|
// user list, and a module's new grants are refused by the bus until that list says them. Among
|
|
// the machines asked it goes first; not among them and behind, it is added — a named push whose
|
|
// module gained a state would otherwise send the code and leave the right to use it for the
|
|
// cascade below, after.
|
|
holder, holderBehind, err := brokerBehind(ctx, open, asked)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// **Not when its own modules are held back** (novox/hq issue 259, ADR 0221): added rather than
|
|
// named, it is sent its whole declaration, and a build its policy records or a plan has not sent
|
|
// it yet would go with the user list. Named and left, with what that costs.
|
|
saidHeld := map[string]bool{}
|
|
if holderBehind && len(args) == 1 {
|
|
held, err := heldMachines(ctx, open, []string{holder})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if why, isHeld := held[holder]; isHeld {
|
|
sayHeld(os.Stdout, holder, why)
|
|
fmt.Printf("%s holds the bus, and the user list it would carry has changed: until it is "+
|
|
"sent, the bus may refuse what this push's machines were newly granted\n", holder)
|
|
holderBehind = false
|
|
saidHeld[holder] = true
|
|
}
|
|
}
|
|
asked = brokerFirst(asked, holder, holderBehind)
|
|
|
|
// 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
|
|
}
|
|
// Each machine's unmet seat dependencies (novox/hq ADR 0207), said after the sends: in full
|
|
// for a machine named, as a count for each of many — the full list is `status`'s.
|
|
unheld := map[string][]catalogue.Unheld{}
|
|
sending, refusals := composeEach(asked, allotting(held, inv), 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)
|
|
unheld[node] = plan.Unheld
|
|
// 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.
|
|
declared, err := declarationWith(held, open, node, plan, settings, gens, Allocating)
|
|
if err == nil {
|
|
reportLeftOut(node, declared)
|
|
}
|
|
return declared, err
|
|
})
|
|
|
|
defer release()
|
|
// Each machine's memberships first, then the declarations (novox/hq issue 249, ADR 0160): a push
|
|
// is the one most operators run, and on 2026-10-01 it was the one path that issued none.
|
|
bus := overTheBus{open: open, server: server, signer: ident}
|
|
sentDigest, err := deliver(ctx, bus, holder, sending)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
release()
|
|
fmt.Printf("\n%d node(s) told\n", len(sending))
|
|
reportUnheldPushed(os.Stdout, len(args) == 1, asked, unheld)
|
|
|
|
// **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". Bounded: a
|
|
// flushed send may itself mint, so this converges over a few rounds.
|
|
//
|
|
// **Except a machine a policy or a plan holds back** (novox/hq issue 259, ADR 0221): one whose
|
|
// modules would move to a build their upgrade policy records rather than rolls out, or that an
|
|
// open plan has not sent it yet. It is named, with why, and left for a push that names it.
|
|
if len(args) == 1 {
|
|
handled := map[string]bool{args[0]: true}
|
|
for n := range saidHeld {
|
|
handled[n] = true
|
|
}
|
|
refused, err := flushBehind(ctx, open, nodes, handled, composeForPush(open, gens), bus, holder, os.Stdout)
|
|
refusals = append(refusals, refused...)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
// 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, allot func(node string) (int64, error),
|
|
compose func(node string) (sendable, error)) ([]readyNode, []string) {
|
|
|
|
var sending []readyNode
|
|
var refusals []string
|
|
for _, name := range names {
|
|
// **Numbered before it is composed, not before it is sent** (novox/hq issue 204). The
|
|
// number says where this declaration stands against every other the mesh composed for the
|
|
// machine, and the host refuses one lower than the last it applied. Taken at send time, as
|
|
// it was, a declaration composed a minute ago — before an assignment changed — went out with
|
|
// a number higher than one composed after the change and sent before it, and the machine
|
|
// took the older content as the newer word: on 2026-10-02 a runtime assigned and applied on
|
|
// two machines was undone two seconds later by exactly that. Taken here, before the first
|
|
// read, what was composed earlier is numbered lower whatever order the sends happen in.
|
|
seq, err := allot(name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
declared, err := compose(name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
declared.Sequence = seq
|
|
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. Their memberships go before their declarations, as every send's do (issue 249).
|
|
func sendRound(ctx context.Context, open *stores, names []string,
|
|
compose func(held context.Context, node string) (sendable, error),
|
|
d delivery, holder string) ([]string, error) {
|
|
held, release, err := holdNodes(ctx, open, names)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer release()
|
|
sending, refused := composeEach(names, allotting(held, open.inventory), func(node string) (sendable, error) {
|
|
return compose(held, node)
|
|
})
|
|
if _, err := deliver(held, d, holder, sending); err != nil {
|
|
return refused, err
|
|
}
|
|
return refused, nil
|
|
}
|
|
|
|
// delivery is the two acts of sending machines what they should be, apart, so the order between
|
|
// them is one function's and can be read and tested there (novox/hq issue 249).
|
|
type delivery interface {
|
|
// grant issues what the machines' modules may do — each module's state raised and its
|
|
// membership issued — for every machine about to be sent.
|
|
grant(ctx context.Context, sending []readyNode) error
|
|
// declare sends one machine its declaration and records it sent, answering the digest.
|
|
declare(ctx context.Context, s readyNode, body []byte) (string, error)
|
|
}
|
|
|
|
// errGrants marks a send that stopped because what the machines' modules may do could not be issued
|
|
// (novox/hq issue 249). Nothing about the machines is wrong; asked again, it is likely to work, so an
|
|
// announcement that hits it is held and asked again.
|
|
var errGrants = errors.New("what the machines' modules may do on the bus could not be issued, and code " +
|
|
"sent before its grants is refused there")
|
|
|
|
// grantsRefused is a grant that failed for some machines and not others: their memberships could not
|
|
// be issued, by machine, and only those machines are held back.
|
|
type grantsRefused struct{ nodes map[string]error }
|
|
|
|
func (g *grantsRefused) Error() string {
|
|
names := make([]string, 0, len(g.nodes))
|
|
for n := range g.nodes {
|
|
names = append(names, n)
|
|
}
|
|
sort.Strings(names)
|
|
return fmt.Sprintf("the memberships of %s could not be issued; the first: %v",
|
|
strings.Join(names, ", "), g.nodes[names[0]])
|
|
}
|
|
|
|
// deliver sends the machines their declarations: **the machine holding the bus, then the grants,
|
|
// then the rest** (novox/hq issue 249).
|
|
//
|
|
// A merge gave a module a new state; its bundle reached every machine within a minute, and the
|
|
// machines' permissions on the bus did not include the state until somebody pushed by hand: the code
|
|
// arrived before the right to use it. A module that read its new state on start failed its start; the
|
|
// one that was there retried for two minutes. The memberships were issued after the declarations —
|
|
// "because the runtime it is for arrives with it" — and a membership is retained last-per-subject on
|
|
// the bus (internal/link/bus.go), so issued first it waits for the runtime that arrives after it. A
|
|
// runtime still on the old code merely holds a grant it does not use yet.
|
|
//
|
|
// **The holder's declaration before the grants, though.** The bus's user list travels in it, and the
|
|
// controller's own right to publish memberships and raise buckets is in that list (the precedent of
|
|
// issue 183): grants first, and a grant the controller is not yet allowed to make would hold the very
|
|
// declaration that allows it — a lock only a hand on the broker could open. A runtime already running
|
|
// on that machine follows a membership issued after its declaration, as it always has.
|
|
//
|
|
// **A grant that cannot be issued holds back what it concerns, and says so as an error.** It used to
|
|
// be said and passed over — "the machines keep what they derive until the next push" — which reported
|
|
// a rollout done that had delivered code its machines could not run. A membership that failed holds
|
|
// back its own machine; a failure that names no machine (the buckets) holds back every machine but the
|
|
// holder, already sent. The error carries errGrants, so the caller's rollout is not marked sent and is
|
|
// tried again.
|
|
//
|
|
// The grants are issued, not waited on: a membership is a retained message the runtime reads when it
|
|
// comes, and the bus answers its publication; nothing here waits for a runtime to have read one.
|
|
func deliver(ctx context.Context, d delivery, holder string, sending []readyNode) (map[string]string, error) {
|
|
digests := map[string]string{}
|
|
if len(sending) == 0 {
|
|
return digests, nil
|
|
}
|
|
send := func(s readyNode) error {
|
|
// The number is inside the signed bytes, so a replayed older declaration cannot borrow a
|
|
// newer one's (novox/hq 04-ISSUES/107); it was taken when the composition began (issue 204).
|
|
body, err := s.declared.Body()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
digest, err := d.declare(ctx, s, body)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
digests[s.node] = digest
|
|
return nil
|
|
}
|
|
var rest []readyNode
|
|
for _, s := range sending {
|
|
if holder != "" && s.node == holder {
|
|
if err := send(s); err != nil {
|
|
return digests, err
|
|
}
|
|
continue
|
|
}
|
|
rest = append(rest, s)
|
|
}
|
|
held := map[string]error{}
|
|
if err := d.grant(ctx, sending); err != nil {
|
|
var some *grantsRefused
|
|
if !errors.As(err, &some) {
|
|
var names []string
|
|
for _, s := range rest {
|
|
names = append(names, s.node)
|
|
}
|
|
if len(names) == 0 {
|
|
return digests, fmt.Errorf("%w: %w", errGrants, err)
|
|
}
|
|
return digests, fmt.Errorf("%w; %s not sent: %w", errGrants, strings.Join(names, ", "), err)
|
|
}
|
|
held = some.nodes
|
|
}
|
|
var notSent []string
|
|
for _, s := range rest {
|
|
if _, refused := held[s.node]; refused {
|
|
notSent = append(notSent, s.node)
|
|
continue
|
|
}
|
|
if err := send(s); err != nil {
|
|
return digests, err
|
|
}
|
|
}
|
|
if len(notSent) > 0 {
|
|
return digests, fmt.Errorf("%w; %s not sent: %w", errGrants, strings.Join(notSent, ", "),
|
|
&grantsRefused{nodes: held})
|
|
}
|
|
if len(held) > 0 {
|
|
// Only the holder's own memberships failed, and it was sent before them.
|
|
return digests, fmt.Errorf("%w: %w", errGrants, &grantsRefused{nodes: held})
|
|
}
|
|
return digests, nil
|
|
}
|
|
|
|
// overTheBus is delivery as the mesh does it: memberships on the bus, declarations signed.
|
|
type overTheBus struct {
|
|
open *stores
|
|
server *link.Server
|
|
signer link.Signer
|
|
// indent is put before each "sent" line, for the callers whose output is nested.
|
|
indent string
|
|
}
|
|
|
|
func (b overTheBus) grant(ctx context.Context, sending []readyNode) error {
|
|
return issueMemberships(ctx, b.open, b.server, sending)
|
|
}
|
|
|
|
func (b overTheBus) declare(ctx context.Context, s readyNode, body []byte) (string, error) {
|
|
if err := link.Declare(ctx, b.server.Bus(), b.signer, 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.
|
|
digest, err := recordSent(ctx, b.open.inventory, s.node, body, s.declared.Builds)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
if s.declared.BusUsers != "" {
|
|
// And the user list it carried, so the next send reads whether it must go first from the
|
|
// list alone (novox/hq issue 249). On the same outliving context as the send's record.
|
|
kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
|
|
err := b.open.inventory.RecordSentBusUsers(kept, s.node, digestOf([]byte(s.declared.BusUsers)))
|
|
cancel()
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
}
|
|
fmt.Printf("%ssent %s %d resource(s)\n", b.indent, s.node, len(s.declared.Resources))
|
|
return digest, nil
|
|
}
|
|
|
|
// brokerFirst is the machines to send in the order a grant needs (novox/hq issue 249): the machine
|
|
// holding the bus first — its declaration carries the bus's user list (composeBusUsers), and a
|
|
// module's new permissions are refused by the bus until that list says them. Among the machines it
|
|
// is moved to the front; not among them, it is added only when it is behind.
|
|
func brokerFirst(names []string, holder string, behind bool) []string {
|
|
if holder == "" {
|
|
return names
|
|
}
|
|
present := false
|
|
rest := make([]string, 0, len(names))
|
|
for _, n := range names {
|
|
if n == holder {
|
|
present = true
|
|
continue
|
|
}
|
|
rest = append(rest, n)
|
|
}
|
|
if !present && !behind {
|
|
return names
|
|
}
|
|
return append([]string{holder}, rest...)
|
|
}
|
|
|
|
// brokerBehind is the machine holding the bus — the one whose declaration carries the user list —
|
|
// and, when it is not among the machines named, whether the user list it would be sent now differs
|
|
// from the one it was last sent (novox/hq issue 249).
|
|
//
|
|
// **The user list alone, not the whole declaration.** Read from the whole declaration, any change
|
|
// pending on that machine — an upgrade its policy records rather than rolls out — went with every
|
|
// send anywhere, and a module running there always put it in its first wave. A digest of the list
|
|
// last sent is kept for this (ADR 0043: the list is composed on each push, never kept itself).
|
|
func brokerBehind(ctx context.Context, open *stores, names []string) (string, bool, error) {
|
|
inv := open.inventory
|
|
holders, err := seatHolders(ctx, inv)
|
|
if err != nil {
|
|
return "", false, err
|
|
}
|
|
h, held := holders[theBrokerSeat]
|
|
if !held || h.Node == "" {
|
|
return "", false, nil
|
|
}
|
|
shelf, err := inv.Catalogue(ctx)
|
|
if err != nil {
|
|
return "", false, err
|
|
}
|
|
if m, known := shelf[h.Module]; !known || m.BusUsers == "" {
|
|
// A holder that is sent no user list carries no grant: nothing to send first.
|
|
return "", false, nil
|
|
}
|
|
for _, n := range names {
|
|
if n == h.Node {
|
|
return h.Node, false, nil
|
|
}
|
|
}
|
|
plan, _, err := planFor(ctx, open, h.Node)
|
|
if err != nil {
|
|
// It cannot be worked out: sending it would refuse the whole send, and `plan` says why.
|
|
return h.Node, false, nil
|
|
}
|
|
list, _, err := busUserList(ctx, inv, plan.Modules)
|
|
if err != nil {
|
|
return h.Node, false, nil
|
|
}
|
|
sent, err := inv.SentBusUsers(ctx, h.Node)
|
|
if err != nil {
|
|
return "", false, err
|
|
}
|
|
return h.Node, userListBehind(list, sent), nil
|
|
}
|
|
|
|
// userListBehind is whether the user list composed now is not the one last sent, by its digest. An
|
|
// empty list composed is never behind: there is nothing for it to carry.
|
|
func userListBehind(now, sentDigest string) bool {
|
|
return now != "" && digestOf([]byte(now)) != sentDigest
|
|
}
|
|
|
|
// 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 {
|
|
_, err := sendToEach(ctx, open, names)
|
|
return err
|
|
}
|
|
|
|
// sendToEach is sendTo, answering the machines it sent: those named, and before them the machine
|
|
// holding the bus when its user list must go first (novox/hq issue 249) — so a caller that waits for
|
|
// the machines it sent waits for that one too.
|
|
func sendToEach(ctx context.Context, open *stores, names []string) ([]string, error) {
|
|
inv := open.inventory
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer ident.Close()
|
|
|
|
gens, err := generators(ctx, open)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
holder, behind, err := brokerBehind(ctx, open, names)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
names = brokerFirst(names, holder, behind)
|
|
|
|
// 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 nil, err
|
|
}
|
|
defer release()
|
|
|
|
var sending []readyNode
|
|
var refusals []string
|
|
for _, name := range names {
|
|
// Numbered before composing, for the reason composeEach gives (novox/hq issue 204).
|
|
seq, err := allot(ctx, inv, name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
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
|
|
}
|
|
declared.Sequence = seq
|
|
reportLeftOut(name, declared)
|
|
sending = append(sending, readyNode{name, declared})
|
|
}
|
|
if len(refusals) > 0 {
|
|
return nil, 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 nil, err
|
|
}
|
|
defer server.Close()
|
|
|
|
// And every assignment on those machines its membership (novox/hq ADR 0160): composed from the
|
|
// same records the bus's accounts are, so what a runtime serves and what its account may are one
|
|
// composition. **Issued before the declarations** (novox/hq issue 249): the runtime the
|
|
// membership is for arrives with the declaration, and a membership waits for it on the bus; the
|
|
// code arriving first was refused its own state until somebody pushed.
|
|
if _, err := deliver(ctx, overTheBus{open: open, server: server, signer: ident, indent: " "}, holder, sending); err != nil {
|
|
return nil, err
|
|
}
|
|
sent := make([]string, 0, len(sending))
|
|
for _, s := range sending {
|
|
sent = append(sent, s.node)
|
|
}
|
|
return sent, nil
|
|
}
|
|
|
|
// issueMemberships publishes the membership of every module on the machines just sent.
|
|
//
|
|
// Each carries what its module receives and the private network's addresses, from the same
|
|
// composition as the declaration it was sent (novox/hq ADR 0167): a provider reads what it is
|
|
// given on the bus, and the file written beside it says the same thing.
|
|
func issueMemberships(ctx context.Context, open *stores, server *link.Server, sent []readyNode) error {
|
|
records, err := open.inventory.BusRecords(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
where := broker.PlacementsOf(records, records.Interchangeable)
|
|
bus, ok := server.Bus().(link.OverNATS)
|
|
if !ok {
|
|
return nil
|
|
}
|
|
// **Every declared state's bucket, before the memberships that name it** (novox/hq ADR 0201). The
|
|
// raise at start asserts them too, but a module registered and assigned since would otherwise have
|
|
// its bucket only after the control plane next restarts — found the first time a module declared
|
|
// state: its bundle asked for a bucket that did not exist. Idempotent and cheap.
|
|
//
|
|
// **A failure here is the send's failure** (novox/hq issue 249). It was said and the push stood,
|
|
// because the declarations were already away; they are sent after this now — all but the bus's
|
|
// own machine, sent before it (deliver) — and a module whose state does not exist is a module
|
|
// that fails its start, so they are not sent and the caller tries again rather than reporting the
|
|
// rollout done.
|
|
buckets, err := open.inventory.DeclaredBuckets(ctx)
|
|
if err != nil {
|
|
return fmt.Errorf("the modules' state could not be read, so no bucket was asserted: %w", err)
|
|
}
|
|
if _, err := broker.RaiseBuckets(broker.OnConn(bus.Conn), buckets); err != nil {
|
|
return fmt.Errorf("the modules' state could not be asserted on the bus: %w", err)
|
|
}
|
|
// Every membership is tried, and the first failure named once.
|
|
issued := 0
|
|
refused := map[string]error{}
|
|
for _, s := range sent {
|
|
node := s.node
|
|
for _, d := range records.Assigned[node] {
|
|
membership := broker.MembershipFor(node, d, where)
|
|
membership.Mesh = s.declared.Mesh
|
|
for requirement, given := range s.declared.Received[d.Module] {
|
|
raw, err := json.Marshal(given)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if membership.Receives == nil {
|
|
membership.Receives = map[string]json.RawMessage{}
|
|
}
|
|
membership.Receives[requirement] = raw
|
|
}
|
|
body, err := json.Marshal(membership)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := bus.PublishMembership(ctx, node, d.Module, body); err != nil {
|
|
if refused[node] == nil {
|
|
refused[node] = fmt.Errorf("%s: %w", d.Module, err)
|
|
}
|
|
continue
|
|
}
|
|
issued++
|
|
}
|
|
}
|
|
if issued > 0 {
|
|
fmt.Printf(" issued %d membership(s)\n", issued)
|
|
}
|
|
if len(refused) > 0 {
|
|
// Returned, never passed over (novox/hq issue 249): the declarations of the machines they are
|
|
// for are not sent, and the rollout that asked is tried again rather than waiting for a push.
|
|
return &grantsRefused{nodes: refused}
|
|
}
|
|
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
|
|
}
|
|
// Composed with the number the machine was LAST sent, so this is byte for byte what it was
|
|
// sent when nothing else changed. A fresh number here would make every machine read as
|
|
// behind for ever (novox/hq 04-ISSUES/107).
|
|
if declared.Sequence, err = open.inventory.Sequence(ctx, n.ID); err != nil {
|
|
return nil, err
|
|
}
|
|
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 each work queue's cancelled set (novox/hq ADR 0219), so a holder taking an ask can ask
|
|
// whether it was cancelled the moment it took it.
|
|
if err := broker.RaiseCancelledSets(js, inventory.MeshSeats()); err != nil {
|
|
return err
|
|
}
|
|
// Every module's state (novox/hq ADR 0201), from the catalogue: a bucket exists from
|
|
// registration, so a module reading one may watch it before its owner runs anywhere. One that
|
|
// nothing declares any more is said and kept — what it holds is data.
|
|
buckets, err := inv.DeclaredBuckets(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
undeclared, err := broker.RaiseBuckets(js, buckets)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(undeclared) > 0 {
|
|
fmt.Printf("the bus holds state nothing declares any more, kept because it is data: %s — "+
|
|
"removing it is a person's act\n", strings.Join(undeclared, ", "))
|
|
}
|
|
// 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 _, c := range broker.ConsumersOf(users) {
|
|
if err := js.EnsureConsumer(c.Consumer); err != nil {
|
|
return fmt.Errorf("how %s on %s hears what it consumes: %w", c.Module, c.Node, err)
|
|
}
|
|
hearing++
|
|
}
|
|
fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, %d module(s) "+
|
|
"can hear what they consume, and %d bucket(s) of state\n", broker.BareAddress(address), len(names), hearing, len(buckets))
|
|
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
|
|
}
|
|
|
|
// number gives one send the next sequence for its node (novox/hq 04-ISSUES/107).
|
|
// allotting is allot over one inventory, in the shape composeEach takes.
|
|
func allotting(ctx context.Context, inv *inventory.Inventory) func(node string) (int64, error) {
|
|
return func(node string) (int64, error) { return allot(ctx, inv, node) }
|
|
}
|
|
|
|
// allot takes the next sequence for a machine — the number its next declaration carries.
|
|
func allot(ctx context.Context, inv *inventory.Inventory, node string) (int64, error) {
|
|
record, err := inv.NodeByName(ctx, node)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return inv.NextSequence(ctx, record.ID)
|
|
}
|
|
|
|
// recordSent writes down what a machine was just sent, and returns the digest.
|
|
//
|
|
// **On a context that outlives the caller's** (novox/hq issue 204). The record is written after the
|
|
// declaration is away, so a send that failed is never recorded as current — and a controller being
|
|
// replaced mid-send had its context cancelled between the two, so the machine was told and the mesh
|
|
// never wrote it down: status read "applied, current" over a machine that had just been sent
|
|
// something else. What was sent was sent; the record of it must not depend on the sender living
|
|
// another second. Bounded, so a store that is away does not hold a dying process open for ever.
|
|
//
|
|
// And the build of each module it carried (novox/hq issue 259, ADR 0221), nil when that is not known:
|
|
// what tells a machine held back by a policy or a plan from one a push left behind.
|
|
func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body []byte,
|
|
builds map[string]string) (string, error) {
|
|
kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
|
|
defer cancel()
|
|
record, err := inv.NodeByName(kept, node)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
digest := digestOf(body)
|
|
if err := inv.RecordSent(kept, record.ID, digest, builds); err != nil {
|
|
return "", err
|
|
}
|
|
return digest, nil
|
|
}
|
|
|
|
// reportUnheldPushed says what a push's machines lack of the seats their modules depend on
|
|
// (novox/hq ADR 0207): every line for a machine the push named, since that is the machine somebody
|
|
// is looking at, and one line per machine otherwise — a list per machine across the mesh is the
|
|
// hundred lines that buried the one that mattered. Nothing for a machine that lacks nothing.
|
|
func reportUnheldPushed(w io.Writer, named bool, asked []string, unheld map[string][]catalogue.Unheld) {
|
|
for _, node := range asked {
|
|
lines := unheld[node]
|
|
if len(lines) == 0 {
|
|
continue
|
|
}
|
|
if named {
|
|
fmt.Fprintf(w, "\n%s has %d unmet seat dependenc(ies) (novox/hq ADR 0207):\n", node, len(lines))
|
|
for _, u := range lines {
|
|
fmt.Fprintf(w, " %s\n", u)
|
|
}
|
|
continue
|
|
}
|
|
fmt.Fprintf(w, "%s: %d unmet seat dependenc(ies) — see `status`\n", node, len(lines))
|
|
}
|
|
}
|