package main import ( "context" "crypto/sha256" "encoding/hex" "encoding/json" "errors" "flag" "fmt" "github.com/novox/mesh-controller/internal/conditions" "io" "log" "os" "slices" "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) (err error) { // Nothing this process does, or starts, is the operator at the terminal (novox/hq ADR 0266). markServed() // 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. // **The lease, before anything that acts** (novox/hq to-be 45 §6): asserting the bus's objects is the // controller's to do, and so is everything after. A controller starting while another holds it waits // here, said; one that loses it stops: every act's gate closes at once, and ctx ends so the process // exits and is started again as a candidate. busAddress, err := broker.BusAddress() if err != nil { return err } lost, err := theLease.serveUnderTheLease(ctx, inv, busAddress) if err != nil { return err } // Given back before the store closes, so the epoch is recorded as given back rather than found // expired by the next holder. defer theLease.release() ctx, stopActing := context.WithCancel(ctx) defer stopActing() go func() { select { case <-lost: stopActing() case <-ctx.Done(): } // Stopping, whichever way: a verb arriving from here is refused as a handover, so its caller // asks the controller after this one rather than have a command started and killed with this // process (novox/hq issue 289). handingOver.Store(true) }() defer func() { select { case <-lost: if err == nil { err = errors.New("the controller lease was lost; this controller stopped acting and exits, to " + "be started again as a candidate") } default: } }() work := link.Enrolment{Inventory: inv, Identity: ident, Broker: known, OnNATS: true} // `status` from a summary kept current here (novox/hq to-be 45 Phase 0): a machine saying // something new is one thing that moves it, so the listener nudges it. statusFrom = newStatusSummary(composeStatus(open)) server, err := connectLink(ctx, inv, work, nudgingListener{Enrolment: work, summary: statusFrom, open: open}) 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 the pending assignments settled on a tick of their own (novox/hq ADR 0261): made once their module // is registered, ended with why when its build will not register it, raised and cleared as conditions. // Never by a read. go settlingPending(ctx, open) // The durations the core's bounds are set from are kept a month (novox/hq to-be 45 Phase 0). go forgettingOldDurations(ctx, 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 the merges the bus announced and never handed over, read back on a timer and acted on late // rather than never (novox/hq issue 266). go catchingUpOnMerges(ctx, open, server) // 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{keeper: func() *conditions.Keeper { return conditionsFrom }}); err != nil { return err } // And what each machine says of what it runs between its reports, kept and raised from (novox/hq ADR // 0240): the gate, `node show` and the conditions read it. server.Hears(moduleHealth{inv: inv, keeper: func() *conditions.Keeper { return conditionsFrom }}) // And what they say about consumers the mesh stopped asking for: one waiting for a person is an // urgent condition, and an act asked of a provider some other way is recorded by hand (ADR 0230). if err := server.KeepsRetirements(retirements{keeper: func() *conditions.Keeper { return conditionsFrom }, record: recordHandActOnTheBus}); 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") } // The hand-act log is counted for `status` on this connection rather than a new one a minute. handActConn = bus.Conn // And everything else this controller does on the bus for a moment (novox/hq issue 327). servingBus.Store(server.JetStream()) // No longer serving: nothing is lent, and dead-letters says it is not read here (novox/hq issue 330). defer servingBus.Store(nil) // And says when it replaced a value given by hand (novox/hq ADR 0228). givenEvents = bus // And a pull request's merge check, asked when the forge announces its head and said when judged // (novox/hq to-be 45 §9). checkEvents = bus // And every walk this controller keeps in a new state, for the delivery it walks (novox/hq ADR 0239). inventory.PlanSaved = func(p inventory.Plan) { sayPlanMoved(ctx, bus, p) } if err := server.Checks(following{open}); err != nil { return err } // Composed now and kept current, before the verb that answers from it is served. go statusFrom.keep(ctx) // The facts snapshot a merge check is fed (novox/hq to-be 45 §9): kept current while this // controller holds the lease, read by the build seat from the artifact store. go exportingFacts(ctx, open, bus.Conn.ConnectedServerVersion) // Every call carries the lease's epoch, and its record is written only under the lease (novox/hq // to-be 45 §6). link.Calls.UnderLease(func() (uint64, error) { return theLease.epoch(ctx) }) // A call that outlasts its caller's patience is followed by `calls` (novox/hq issue 265). link.Calls.Follow = catalogue.ControllerSeatName + ".calls" // And every call is kept on the bus, so a restart of this process keeps what came of each // (novox/hq to-be 45 §6). A bus without the bucket is said and served from memory, as before: // answering no calls at all would be worse than answering them without the record. said := log.New(os.Stdout, "", log.LstdFlags) if keeper, err := link.CallsOnTheBus(ctx, bus.Conn); err != nil { fmt.Printf("calls are kept in memory only, and lost when this controller stops: %v\n", err) } else if err := link.Calls.Durably(ctx, keeper, instance, said); err != nil { fmt.Printf("calls are kept on the bus from now on; the ones kept before could not be read: %v\n", err) } // What is wrong, kept and said (novox/hq to-be 45 §2): the condition store, the watchdogs of the // signals table, what the bus says about itself, and the self-check. A store that cannot be opened // is said and the controller serves on: status then says the conditions cannot be read, and is // not well — louder than not serving at all, and the push that repairs the bus still runs. if stopWatching := watchTheMesh(ctx, open, server, bus); stopWatching != nil { defer stopWatching() } stopServing, err := bus.ServeSeatTools(catalogue.ControllerSeatName, handlers, log.New(os.Stdout, "", log.LstdFlags)) if err != nil { return err } defer stopServing() // And the operator's command on every node, asked through each node's engine (novox/hq ADR 0272). stopCLI, err := bus.ServeCLI(answerMeshCLI(open.inventory), log.New(os.Stdout, "", log.LstdFlags)) if err != nil { return err } defer stopCLI() // 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, 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. // The epoch it carried, if a person wrote one in, is what the machine heard. var carried struct { Epoch uint64 `json:"epoch"` } _ = json.Unmarshal(raw, &carried) if _, err := recordSent(ctx, inv, node, raw, nil, carried.Epoch); 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 ` // 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)") // A push by hand is a repair, and says why (novox/hq to-be 45 §7): required through the seat, // recorded when given at a shell — see handacts.go for why a shell is not refused. why := addHandActFlags(set) // Every machine at once is said, not defaulted to (novox/hq ADR 0217): a bare `push` sent the // whole mesh on 2026-10-05 when its usage was wanted. all := set.Bool("all", false, "every machine") // A running module's data moving is sent only when said, per module (novox/hq ADR 0217). move := set.String("move", "", "modules whose data may move with this push, comma-separated") positionals, err := parseAround(set, args) if err != nil { return err } args = positionals if len(args) > 1 { return errors.New("push | push --behind | push --all — one machine, the ones behind, or all of them") } if len(args) == 0 && !*behind && !*all { return errors.New("push names a machine, or says --behind (the ones not running what they " + "should) or --all (every machine) — a push to the whole mesh is not the default (novox/hq ADR 0217)") } if *all && (*behind || len(args) == 1) { return errors.New("push --all is every machine; it takes no machine and no --behind") } movable := map[string]bool{} for _, m := range splitModules(*move) { movable[m] = true } var heldBack []string 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 or push --behind, not both: one names a machine and the " + "other asks which machines need one") } recorded := append([]string(nil), args...) if *behind { recorded = append(recorded, "--behind") } if *all { recorded = append(recorded, "--all") } if *move != "" { recorded = append(recorded, "--move", *move) } why.record(ctx, "push", recorded) 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("no machine named: a push of the whole mesh, and every machine is doing what it was told — nothing sent") 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) } // **A push that named no machine says so, first.** Through the console a machine the caller // meant to name could be lost on the way (novox/hq issue 244): the call arrived empty, ran as // `push --behind`, and every machine behind was pushed by someone who thought they had pushed one. // Its whole-mesh reach is the first line of the answer, with the machines it is about to send. if len(args) == 0 { which := "every machine" if *behind { which = "every machine that is behind" } fmt.Printf("no machine named: this is a push of the WHOLE mesh — %s (%d): %s\n", which, len(asked), strings.Join(asked, ", ")) } // **A whole-mesh push sends no build a gate has not seen** (novox/hq ADR 0236): a machine where one // waits is left, named, for the release plan — `push ` still sends one machine by a person's word. if len(args) == 0 { f, err := readMoveFacts(ctx, inv) if err != nil { return err } var kept []string for _, n := range asked { moves, err := machineMoves(ctx, open, f, n, false) if err != nil { return err } if len(moves) > 0 { fmt.Printf(" %s is left: %d build(s) wait there for a gate, which a release plan sends one machine "+ "at a time (`upgrade backlog` lists them; `push %s` sends it by name)\n", n, len(moves), n) continue } kept = append(kept, n) } asked = kept } // **The bus is replaced only as a planned step** (novox/hq ADR 0236, to-be 45 §8): a machine whose bus // would move is not sent by a push — named, it is refused; otherwise it is left and said. busKept, err := busHeld(ctx, inv, asked) if err != nil { return err } if len(busKept) > 0 { if len(args) == 1 { return fmt.Errorf("%s. Nothing was sent", busKept[args[0]]) } var kept []string for _, n := range asked { if why, h := busKept[n]; h { fmt.Printf(" %s is left: %s\n", n, why) continue } kept = append(kept, n) } asked = kept } // **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) reportKept(held, 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() // A machine whose declaration would move a running module's data is held, and only it (novox/hq // ADR 0217): said, with the command that sends it, and left out of everything below. toSend, err := holdingBack(ctx, inv, sending, movable, &heldBack) if err != nil { return err } // **What it recreates, said before it is sent** (novox/hq ADR 0245): a person's push moves every // module whose build changed, and recreates what changed in it — a gated send said so, a push did not. sayWhatAPushRecreates(ctx, open, toSend, os.Stdout) // 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, toSend) if err != nil { return err } release() told := make([]string, 0, len(toSend)) for _, r := range toSend { told = append(told, r.node) } fmt.Printf("\n%d node(s) told: %s\n", len(toSend), strings.Join(told, ", ")) 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 } // The cascade is a push too, and a machine in it whose data would move is held the same way // (novox/hq ADR 0217). keep := keeping(func(held context.Context, sending []readyNode) ([]readyNode, error) { return holdingBack(held, inv, sending, movable, &heldBack) }) refused, err := flushBehind(ctx, open, nodes, handled, composeForPush(open, gens), bus, holder, keep, 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 } } if len(heldBack) > 0 { sort.Strings(heldBack) return fmt.Errorf("held %s: a running module's data would move — see above; nothing was sent "+ "there (novox/hq ADR 0217)", strings.Join(heldBack, ", ")) } 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) (order, 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. numbered, 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, declared.Epoch = numbered.sequence, numbered.epoch 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, keep keeping) ([]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 keep != nil { if sending, err = keep(held, sending); err != nil { return refused, err } } if _, err := deliver(held, d, holder, sending); err != nil { return refused, err } return refused, nil } // keeping is what a round of sends keeps of the machines it composed, before any is delivered: nil // keeps them all. A push holds back the machines whose running modules' data would move (holdingBack). type keeping func(held context.Context, sending []readyNode) ([]readyNode, error) // holdingBack leaves out each machine whose declaration would move a running module's data, says so, // and names it among those held back (novox/hq ADR 0217); the rest go on to be delivered, grants first // (issue 249). A declaration's body is the same bytes however often it is read. func holdingBack(ctx context.Context, inv *inventory.Inventory, sending []readyNode, movable map[string]bool, heldBack *[]string) ([]readyNode, error) { var out []readyNode for _, s := range sending { body, err := s.declared.Body() if err != nil { return nil, err } moves, err := heldMoves(ctx, inv, s.node, body, movable) if err != nil { return nil, err } if len(moves) > 0 { sayDataHeld(os.Stdout, s.node, moves) *heldBack = append(*heldBack, s.node) continue } out = append(out, s) } return out, 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 { // The bus's objects first, which every declaration implies (novox/hq issue 208): said and raised // when they cannot be, and never what holds the send back (assertOnSend). if bus, ok := b.server.Bus().(link.OverNATS); ok { js := broker.OnConn(bus.Conn) js.Note = func(format string, args ...any) { fmt.Printf(b.indent+" "+format+"\n", args...) } _ = assertOnSend(ctx, b.open.inventory, js, b.indent) } 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, s.declared.Epoch) if err != nil { return "", err } if len(s.declared.Bindings) > 0 { // And where it bound each consumer of a provision that keeps its data (novox/hq ADR 0232): // what the next resolution keeps it at. On the same outliving context as the send's record. kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second) err := b.open.inventory.RecordBindings(kept, s.node, s.declared.Bindings) cancel() if err != nil { return "", fmt.Errorf("%s was sent its declaration, and where its consumers are bound to their "+ "data could not be recorded: %w", s.node, 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 // **No send replaces the bus but its planned step** (novox/hq ADR 0236). A plan's send to the bus's // machine for some other module carried the bus's new build with it on 2026-10-06, and the bus // restarted under every machine with nobody having asked. Refused whole — the send cannot leave the // bus behind and carry the rest (ADR 0221 option 2) — and said with the remedy; the plan tries again. if !busStepSending(ctx) { held, err := busHeld(ctx, inv, names) if err != nil { return nil, err } for _, n := range names { if why, h := held[n]; h { return nil, fmt.Errorf("%w: %s", errBusWaits, why) } } } 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 } added := "" if behind && !slices.Contains(names, holder) { added = holder } names = brokerFirst(names, holder, behind) // **No build moves on a machine without a gate** (novox/hq ADR 0236): a send outside its scope that // would carry one is refused, said with what waits; the bus's machine added for its user list is left. if names, err = ungatedIn(ctx, open, names, added); err != nil { return nil, err } // **A recorded build moves only by a person's push** (novox/hq issue 295, ADR 0245): this send — a // plan's, a release plan's, a rollback's, a healer's, a rotation's — composes every recorded module at // the build its machine runs. The bus step is a person's word for the bus alone. if ctx, err = sendKeeps(ctx, inv); err != nil { return nil, 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 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). numbered, 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) reportKept(ctx, 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, declared.Epoch = numbered.sequence, numbered.epoch 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 } bus, ok := server.Bus().(link.OverNATS) if !ok { return nil } // Which machines are root-free now (novox/hq ADR 0259 §8): a channel's verified sender is composed for the // router only from one, beside a router on one. Judged once per push, by the root-free verb's judgement. records.RootFree = rootFreeNow(ctx, newRootReader(ctx, open.inventory, bus.Conn), records.Nodes, time.Now()) where := broker.PlacementsOf(records, records.Interchangeable) // **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) { return wouldSendFrom(ctx, open, nodes, nil) } // planned is one machine's plan as planFor answered it, for a caller that already asked. type planned struct { plan catalogue.Resolution settings catalogue.SettingsBy } // wouldSendFrom is wouldSend reusing the plans a caller worked out a moment before: resolving a // machine is most of what `status` costs, and it used to resolve every machine twice (novox/hq // to-be 45 Phase 0). A machine absent from plans is worked out here. func wouldSendFrom(ctx context.Context, open *stores, nodes []inventory.Node, plans map[string]planned) (map[string]string, error) { gens, err := generators(ctx, open) if err != nil { return nil, err } out := map[string]string{} for _, n := range nodes { known, have := plans[n.Name] plan, settings := known.plan, known.settings if !have { if plan, settings, err = planFor(ctx, open, n.Name); 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 } // And the epoch it was last sent under, for the same reason: a new holder of the lease is not a // change of the machine (novox/hq to-be 45 §6). if declared.Epoch, err = open.inventory.SentEpoch(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) } } // The mesh's own streams and consumers, the seats' work queues and their workers, and how every // machine hears its declaration — one derivation, which the self-check reads as well (D6, D7). names, err := assertBusObjects(ctx, inv, js) if 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 } // And the controller's own buckets (novox/hq to-be 45 §1): the calls it serves and the acts done // by hand, kept where a restart of this process does not take them. if err := js.EnsureControllerBuckets(); 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: asserted with the rest above, counted here. hearing, err := moduleConsumerCount(ctx, inv) if err != nil { return err } 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) (order, error) { return func(node string) (order, error) { return allot(ctx, inv, node) } } // epochForActs is the lease's gate as a composition asks it; a variable so a test can act under an epoch // without a bus. var epochForActs = func(ctx context.Context) (uint64, error) { return theLease.epoch(ctx) } // order is what a declaration carries of its writer's order (link/order.go): its sequence, and the // epoch of the lease it is composed under — zero for a machine that has not said it reads one. type order struct { sequence int64 epoch uint64 } // allot takes the next sequence for a machine — the number its next declaration carries — under the // lease: a process that may not act takes none, and composes nothing (novox/hq to-be 45 §6). func allot(ctx context.Context, inv *inventory.Inventory, node string) (order, error) { epoch, err := epochForActs(ctx) if err != nil { return order{}, fmt.Errorf("nothing was composed for %s: %w", node, err) } record, err := inv.NodeByName(ctx, node) if err != nil { return order{}, err } if epoch > 0 { reads, err := inv.ReadsEpoch(ctx, record.ID) if err != nil { return order{}, err } if !reads { epoch = 0 } } seq, err := inv.NextSequence(ctx, record.ID) if err != nil { return order{}, err } return order{sequence: seq, epoch: epoch}, nil } // 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, epoch uint64) (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.RecordSentUnder(kept, record.ID, digest, builds, epoch); err != nil { return "", err } // And what it was, summarised, so the next push can be compared with it (novox/hq ADR 0217). // Never a reason to fail a send that is already away: a summary that cannot be kept means the // next comparison has nothing to hold against, which is how every machine starts. if summary, err := summarize(body); err == nil { if raw, err := json.Marshal(summary); err == nil { if err := inv.RecordSentSummary(kept, record.ID, raw); err != nil { fmt.Fprintf(os.Stderr, "%s: what it was sent is not kept for comparison: %v\n", node, 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)) } }