Both branches changed the seat set from the same starting point, so every number collided and every `mesh-*` name existed twice. The trunk's numbers and names win: this branch's records became 0129/0130 and its migrations 0037/0038, and the hardcoded rename map gave way to the trunk's `seat_alias` table — a rename is a row now (ADR 0122), not a recompile. Three of my checks were wrong and the merge is what showed it: A seat with an empty protocol is a marker, not an incomplete declaration. Most node-scoped seats are markers — which module is this machine's packet filter — and refusing one refused most of the set, the showcase module included. A mistyped field name is already refused by the parser, so an empty protocol was written as one deliberately. A claim on a seat this manifest does not declare is not the parser's to judge. A module may hold a seat another module declared; that is the whole reason ADR 0126 has callers name the seat and not its provider. Whether the seat exists is a fact about the catalogue, so the refusal is at registration, where every declaration is in view. And a seat may share a name with the provision it delivers. `git`, the npm registry and the artifact store still do, because renaming a delivering seat cascades to every consumer requiring it, with a window where a holder stops resolving mid-flight. The trunk deferred exactly those three on purpose. Full suite green against a real NATS and store.
749 lines
28 KiB
Go
749 lines
28 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"os"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/catalogue"
|
|
"github.com/novox/mesh-controller/internal/inventory"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// reportUnhostable says which of a node's assigned modules the machine cannot run, once per push.
|
|
//
|
|
// A module whose declared capability has no detector on the machine is on the wrong machine. It is
|
|
// kept out of what the node is sent — the healthy modules beside it still converge — and named here
|
|
// so it is neither silently dropped nor a reason the whole node fails to push.
|
|
func reportUnhostable(node string, plan catalogue.Resolution) {
|
|
for _, u := range plan.Unhostable {
|
|
for _, c := range u.Missing {
|
|
fmt.Printf("%s not applied — %s\n", node, catalogue.WrongMachine(u.Module, c, node))
|
|
}
|
|
}
|
|
}
|
|
|
|
// sending it, and holding the link that carries it.
|
|
//
|
|
// Split out of main.go, which had reached 2,769 lines because appending was always the
|
|
// cheapest next step. That is how novox/hq ADR 0001 records `hal/sdk` reaching 34,636:
|
|
// nothing in it was wrong, and no one edit was the one that should have been a new file.
|
|
|
|
// serve is the control plane running: one connection to the broker, one queue, one consumer.
|
|
func serve(ctx context.Context) error {
|
|
open, err := openStores(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer open.Close()
|
|
inv := open.inventory
|
|
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
// Established at start rather than on first use. A control plane that cannot sign is one
|
|
// whose declarations every node correctly refuses, and that should be a startup failure
|
|
// rather than something discovered at the first declaration.
|
|
key, err := ident.Establish(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("signing as %s\n", key.Fingerprint()[:16])
|
|
|
|
management, err := broker.ManagementFromEnvironment()
|
|
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
|
|
return err
|
|
}
|
|
|
|
// Where the broker is and what to expect there, so a node can be told how to come back
|
|
// without a person and a new token.
|
|
known, err := broker.FromEnvironment()
|
|
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
|
|
return err
|
|
}
|
|
if errors.Is(err, broker.ErrNotConfigured) {
|
|
fmt.Printf("no broker address configured, so enrolled nodes will not be told how to "+
|
|
"reconnect. Set %s and %s.\n", broker.AddressVar, broker.CertificateVar)
|
|
}
|
|
|
|
// **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.
|
|
busAddress, onNATS, err := broker.OnNATS()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
|
|
return err
|
|
}
|
|
|
|
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
|
|
OnNATS: onNATS}
|
|
server, err := link.Connect(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.
|
|
if onNATS {
|
|
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
// And build results nobody was waiting for. A build triggered any other way than `build`
|
|
// would otherwise be reported into the void, which is the same as not reporting it.
|
|
server.Records(builds{inv})
|
|
// And what the catalogue decided a build meant. The builder's own result is already handled
|
|
// above; this is the other half — the control plane is the only one of the three that knows
|
|
// which machines run the thing, so it is the one that acts (novox/hq ADR 0072).
|
|
if err := server.Follows(following{open}); err != nil {
|
|
return err
|
|
}
|
|
// And a catalogue that has just started, asking for what it missed. The same type answers
|
|
// both: what a build meant and what the builds were are two questions about one record.
|
|
if err := server.Answers(following{open}); err != nil {
|
|
return err
|
|
}
|
|
|
|
return server.Serve(ctx)
|
|
}
|
|
|
|
// declare sends one node a declaration, signed.
|
|
//
|
|
// Signed here rather than trusted from the broker: a node connects to the broker and takes
|
|
// instruction from the control plane behind it, and those are two identities. If a node believed
|
|
// whatever arrived on its queue, a compromised broker could forge declarations — and since the
|
|
// host applies whatever the link delivers, that is the whole machine (novox/hq ADR 0004).
|
|
func declare(ctx context.Context, args []string) error {
|
|
if len(args) != 2 {
|
|
return errors.New("declare <node> <declaration.json>")
|
|
}
|
|
node, path := args[0], args[1]
|
|
|
|
raw, err := os.ReadFile(path)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
// The node has to exist before it can be told anything. Publishing to a queue nobody consumes
|
|
// would sit there looking like success.
|
|
open, err := openStores(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer open.Close()
|
|
inv := open.inventory
|
|
if _, err := inv.NodeByName(ctx, node); err != nil {
|
|
return err
|
|
}
|
|
|
|
server, err := link.Connect(nil, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
|
return nil
|
|
}
|
|
|
|
// OverlayCIDRVar is the range the mesh allocates node addresses from.
|
|
const OverlayCIDRVar = "MESH_OVERLAY_CIDR"
|
|
|
|
// pushCommand sends nodes everything they should be: their place on the network, and what their
|
|
// assignments resolve to.
|
|
//
|
|
// One declaration, not two. A node holding its network and not its modules, or the reverse, is
|
|
// half-configured for as long as that lasts — and the two are computed from the same picture of
|
|
// the mesh, so sending them apart would let them disagree.
|
|
func pushCommand(ctx context.Context, args []string) error {
|
|
set := flag.NewFlagSet("push", flag.ContinueOnError)
|
|
// Only the machines that need it.
|
|
//
|
|
// **A command rather than a timer, to begin with.** Something that re-pushes on a schedule is
|
|
// a scheduler over this, and building the scheduler first would mean two paths to one act
|
|
// with nothing to compare them against. A person can run this; so can cron; so can whatever
|
|
// eventually watches.
|
|
behind := set.Bool("behind", false,
|
|
"only machines whose last declaration was refused or partly failed")
|
|
// For a named node, wait until it reports applying exactly what it was sent, so `push <node>`
|
|
// means "this node is now what it was told" — a command right after does not race the apply
|
|
// (novox/hq ADR 0010). 0 waits for nothing, which is the old fire-and-forget.
|
|
wait := set.Duration("wait", 0,
|
|
"for a named node, how long to wait for it to report applying what it was sent (0: do not wait)")
|
|
positionals, err := parseAround(set, args)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
args = positionals
|
|
if len(args) > 1 {
|
|
return errors.New("push [<node>] [--behind] — one node, or all of them")
|
|
}
|
|
if len(args) == 1 && *behind {
|
|
// Naming a machine and asking for the ones that need it are two different requests, and
|
|
// guessing which was meant would sometimes push to a machine somebody did not name.
|
|
return errors.New("push <node> or push --behind, not both: one names a machine and the " +
|
|
"other asks which machines need one")
|
|
}
|
|
open, err := openStores(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer open.Close()
|
|
inv := open.inventory
|
|
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
// Every node, not only the ones on the private network. A machine that was never given the
|
|
// network module still takes modules, and iterating the network here is what used to make
|
|
// "on the network" and "managed" the same thing.
|
|
nodes, err := inv.Nodes(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Which machines are not in the state they were sent, when that is what was asked for.
|
|
var needsOne map[string]inventory.Doing
|
|
if *behind {
|
|
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
needsOne = map[string]inventory.Doing{}
|
|
for _, d := range wrong {
|
|
needsOne[d.Node] = d
|
|
}
|
|
// **And every machine not running what the mesh would send it.** "Behind" used to mean
|
|
// only "failed or refused", so a machine that applied cleanly and whose declaration has
|
|
// since changed was not behind — and novox/hq ADR 0010's question, *did my change go
|
|
// out?*, was answerable only for the machines that broke.
|
|
would, err := wouldSend(ctx, open, nodes)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
waiting, err := inv.Waiting(ctx, would)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, m := range waiting {
|
|
if _, already := needsOne[m.Node]; already {
|
|
continue
|
|
}
|
|
needsOne[m.Node] = inventory.Doing{Node: m.Node, Outcome: "waiting"}
|
|
}
|
|
if len(needsOne) == 0 {
|
|
// Said rather than doing nothing quietly. "Nothing needed one" and "this did not run"
|
|
// must never look the same.
|
|
fmt.Println("every machine is doing what it was told")
|
|
return nil
|
|
}
|
|
}
|
|
gens, err := generators(ctx, open)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
server, err := link.Connect(nil, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
// Which machines this push is about, before any of them is worked out.
|
|
var asked []string
|
|
for _, n := range nodes {
|
|
if len(args) == 1 && n.Name != args[0] {
|
|
continue
|
|
}
|
|
if *behind {
|
|
doing, needs := needsOne[n.Name]
|
|
if !needs {
|
|
continue
|
|
}
|
|
// A machine that has been failing the same way for a long time is not going to stop
|
|
// because it was asked again. Said, and pushed to anyway — refusing would leave no
|
|
// way to retry after fixing the cause, and this is a command somebody ran.
|
|
//
|
|
// Only for machines that reported something. One that is merely waiting has no report
|
|
// to be old, and saying it had been failing since the zero time would be a sentence
|
|
// about nothing.
|
|
if since := time.Since(doing.At); doing.Outcome != "waiting" && since > 6*time.Hour {
|
|
fmt.Printf("%s has been %s since %s; pushing again anyway, but the cause is "+
|
|
"unlikely to be timing\n",
|
|
n.Name, doing.Outcome, doing.At.Local().Format("2006-01-02 15:04"))
|
|
}
|
|
}
|
|
asked = append(asked, n.Name)
|
|
}
|
|
|
|
// Held from composing to sending, so a converge on one of them cannot send between the two
|
|
// and be overtaken by what was composed before it (novox/hq ADR 0100).
|
|
held, release, err := holdNodes(ctx, open, asked)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
sending, refusals := composeEach(asked, func(node string) (sendable, error) {
|
|
plan, settings, err := planFor(held, open, node)
|
|
if err != nil {
|
|
return sendable{}, err
|
|
}
|
|
// A module assigned here that this machine cannot host is said and left out, not fatal: the
|
|
// healthy modules beside it are still resolved and sent. Reported so it is not silently
|
|
// dropped — the remedy is to move it, and until then the rest of the node converges.
|
|
reportUnhostable(node, plan)
|
|
// The private network is in here with everything else. It used to be composed separately
|
|
// and prepended, which meant every machine with an address was on it and no machine could
|
|
// be kept off. It is a module now, so it arrives the way a module does.
|
|
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
|
})
|
|
|
|
sentDigest := map[string]string{}
|
|
defer release()
|
|
for _, s := range sending {
|
|
body, err := s.declared.Body()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil {
|
|
return err
|
|
}
|
|
// After it is away, not before. A digest recorded for something that failed to send would
|
|
// make the machine look current for a declaration it never received.
|
|
record, err := inv.NodeByName(ctx, s.node)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
digest := digestOf(body)
|
|
if err := inv.RecordSent(ctx, record.ID, digest); err != nil {
|
|
return err
|
|
}
|
|
sentDigest[s.node] = digest
|
|
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
|
}
|
|
release()
|
|
fmt.Printf("\n%d node(s) told\n", len(sending))
|
|
|
|
// **A named push leaves the mesh consistent, not just the machine it named** (novox/hq
|
|
// issue 057, ADR 0083). Assigning a cross-node consumer mints a provision, and the PROVIDER's
|
|
// grant list is a pure read of secrets already issued — so after the named node is current,
|
|
// other machines can be behind *as a consequence*: their declaration now differs from what
|
|
// they were last sent. Those are flushed too, by name, in the push's own output.
|
|
//
|
|
// Compared against what each machine was last SENT, not against a before/after of this push:
|
|
// the mint usually happened at `assign` or `module issue`, before this command ran, so the
|
|
// only durable signal is "what it should be" versus "what it last received". A machine behind
|
|
// for an unrelated reason is caught here too, which is not a cost — a named push that knew a
|
|
// machine was behind and left it so would be the very silence this removes. Bounded: a
|
|
// flushed send may itself mint, so this converges over a few rounds.
|
|
if len(args) == 1 {
|
|
flushed := map[string]bool{args[0]: true}
|
|
// Bounded by the node count: a node is marked flushed the round it is handled and is
|
|
// never handled twice, so the loop cannot run more than len(nodes) rounds. The bound is
|
|
// a guard against a logic error, not a real limit — if it were ever hit, that is a bug
|
|
// rather than a cascade legitimately still converging, so it is said rather than passed
|
|
// over in silence, unlike the earlier fixed cap that could stop a real cascade short.
|
|
rounds := 0
|
|
for {
|
|
would, err := wouldSend(ctx, open, nodes)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
behind, err := inv.Waiting(ctx, would)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var also []string
|
|
for _, m := range behind {
|
|
if !flushed[m.Node] {
|
|
also = append(also, m.Node)
|
|
}
|
|
}
|
|
if len(also) == 0 {
|
|
break
|
|
}
|
|
if rounds++; rounds > len(nodes) {
|
|
fmt.Printf("\nstopped cascading after %d rounds with %s still behind — this "+
|
|
"should not happen; run `push --behind` to finish\n",
|
|
rounds-1, strings.Join(also, ", "))
|
|
break
|
|
}
|
|
sort.Strings(also)
|
|
fmt.Printf("\nthis push left %s behind — a provision granted from there, or a "+
|
|
"declaration since changed; sending it too\n", strings.Join(also, ", "))
|
|
// Tolerantly, exactly as the named send above: a machine that cannot be composed is
|
|
// collected as a refusal and reported at the end, and the others are still sent
|
|
// (novox/hq ADR 0066). The earlier cut routed these through sendTo, which is
|
|
// all-or-nothing — so one swept machine's compose error failed the operator's named
|
|
// push and skipped its --wait, the very intolerance the main path exists to avoid.
|
|
// Held for this round only, and after the last round's were given back, so two pushes
|
|
// cascading into each other's machines never each wait on the other.
|
|
refused, err := sendRound(ctx, open, also,
|
|
func(held context.Context, node string) (sendable, error) {
|
|
plan, settings, err := planFor(held, open, node)
|
|
if err != nil {
|
|
return sendable{}, err
|
|
}
|
|
reportUnhostable(node, plan)
|
|
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
|
},
|
|
func(s readyNode, body []byte) error {
|
|
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body,
|
|
15*time.Second); err != nil {
|
|
return err
|
|
}
|
|
record, err := inv.NodeByName(ctx, s.node)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
|
return nil
|
|
})
|
|
refusals = append(refusals, refused...)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// Every candidate this round is marked handled — the sent ones so they are not
|
|
// re-listed, and the refused ones so a machine that cannot be composed does not make
|
|
// the loop spin on it for ever. Its refusal is already in the report.
|
|
for _, name := range also {
|
|
flushed[name] = true
|
|
}
|
|
}
|
|
}
|
|
|
|
// A named node is a request to make THAT node current now, so it waits for the node to say it
|
|
// applied exactly this. A whole-mesh or --behind push does not wait: it is a sweep, and blocking
|
|
// on the slowest machine would hold back the report on all the others.
|
|
if *wait > 0 && len(args) == 1 {
|
|
if err := waitForApplied(ctx, inv, args[0], sentDigest[args[0]], *wait); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return couldNotBeResolved(refusals, len(sending))
|
|
}
|
|
|
|
// waitForApplied blocks until the node reports it applied exactly the declaration just sent, or the
|
|
// wait runs out. A report of failure or refusal for that same declaration ends the wait at once —
|
|
// there is nothing to wait for, and the reason is the node's own.
|
|
func waitForApplied(ctx context.Context, inv *inventory.Inventory, node, digest string, wait time.Duration) error {
|
|
if digest == "" {
|
|
return nil // nothing was sent to this node
|
|
}
|
|
deadline := time.Now().Add(wait)
|
|
for {
|
|
doing, said, err := inv.DoingOf(ctx, node)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if said && doing.Declared == digest {
|
|
switch doing.Outcome {
|
|
case inventory.OutcomeApplied:
|
|
fmt.Printf("%s applied it\n", node)
|
|
return nil
|
|
case inventory.OutcomeFailed:
|
|
return fmt.Errorf("%s applied what it was sent but %d resource(s) failed", node, len(doing.Failed))
|
|
case inventory.OutcomeRefused:
|
|
return fmt.Errorf("%s refused what it was sent: %s", node, doing.Refused)
|
|
}
|
|
}
|
|
if time.Now().After(deadline) {
|
|
return fmt.Errorf("%s did not report applying what it was sent within %s "+
|
|
"(it may still be converging; check `status`)", node, wait)
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case <-time.After(500 * time.Millisecond):
|
|
}
|
|
}
|
|
}
|
|
|
|
// readyNode is one machine and the declaration it would be sent.
|
|
type readyNode struct {
|
|
node string
|
|
declared sendable
|
|
}
|
|
|
|
// composeEach works out what each named machine should be, and never lets one machine's answer
|
|
// decide another's.
|
|
//
|
|
// **A machine whose set cannot be worked out is that machine's problem** (novox/hq ADR 0066). A
|
|
// whole-mesh push used to refuse outright when any one node failed to resolve, so a single
|
|
// unanswerable requirement on a single machine — one module requiring a provision nobody had
|
|
// assigned a provider for — left every other machine in the mesh unconverged, including machines
|
|
// with no relation to it at all. Nothing was sent anywhere, and the machines that could not be sent
|
|
// were the ones with nothing wrong with them.
|
|
//
|
|
// It is the same rule a92c11b established one level down, where an un-hostable module stopped
|
|
// taking down the healthy modules beside it, applied one level up: **the blast radius of a fault is
|
|
// the thing that has it.** What could not be worked out is named and returned, so a push still ends
|
|
// with a non-zero outcome and nobody mistakes a partial convergence for a whole one.
|
|
//
|
|
// The all-or-nothing rule is kept where it means something — sendTo, which rotates a credential
|
|
// across two machines that must agree — and dropped here, where it never did.
|
|
func composeEach(names []string,
|
|
compose func(node string) (sendable, error)) ([]readyNode, []string) {
|
|
|
|
var sending []readyNode
|
|
var refusals []string
|
|
for _, name := range names {
|
|
declared, err := compose(name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
if len(declared.Resources) == 0 {
|
|
// Sent, not skipped (novox/hq issue 127). A node whose declaration composes to
|
|
// nothing may have HELD something before — the broker opening a placement gave it,
|
|
// say — and skipping the empty declaration leaves that last resource in force
|
|
// forever, re-applied by the node's own heartbeat, with no way for the mesh to say
|
|
// it is gone. An empty declaration is the correction: the host drops what the mesh
|
|
// owned and keeps what it found (the adoption envelope still rides along). A node
|
|
// that never held anything applies it as the no-op it is.
|
|
fmt.Printf("%s owns nothing now — sent so it drops what it last held\n", name)
|
|
}
|
|
sending = append(sending, readyNode{name, declared})
|
|
}
|
|
return sending, refusals
|
|
}
|
|
|
|
// sendRound holds the named nodes, composes each and sends each that composed, and gives the hold
|
|
// back on every way out — a body that cannot be marshalled and a send that fails included
|
|
// (novox/hq ADR 0100). A node that cannot be composed is a refusal, not an error: the others are
|
|
// still sent.
|
|
func sendRound(ctx context.Context, open *stores, names []string,
|
|
compose func(held context.Context, node string) (sendable, error),
|
|
send func(s readyNode, body []byte) error) ([]string, error) {
|
|
held, release, err := holdNodes(ctx, open, names)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer release()
|
|
sending, refused := composeEach(names, func(node string) (sendable, error) {
|
|
return compose(held, node)
|
|
})
|
|
for _, s := range sending {
|
|
body, err := s.declared.Body()
|
|
if err != nil {
|
|
return refused, err
|
|
}
|
|
if err := send(s, body); err != nil {
|
|
return refused, err
|
|
}
|
|
}
|
|
return refused, nil
|
|
}
|
|
|
|
// couldNotBeResolved is what a push ends with when some machines could not be worked out.
|
|
//
|
|
// **After the rest have been sent, never instead of sending them.** It is still an error, because
|
|
// the mesh is not in the state somebody asked for and a command that exits cleanly having skipped a
|
|
// machine is a command that lies. What it must not do is decide anything about the machines beside
|
|
// it, which is why it says how many were sent.
|
|
func couldNotBeResolved(refusals []string, sent int) error {
|
|
if len(refusals) == 0 {
|
|
return nil
|
|
}
|
|
return fmt.Errorf(
|
|
"%d node(s) could not be resolved and were not sent. %d other node(s) were:\n\n%s",
|
|
len(refusals), sent, strings.Join(refusals, "\n\n"))
|
|
}
|
|
|
|
// sendTo resolves and sends to exactly the machines named, or refuses without sending anything.
|
|
//
|
|
// The all-or-nothing rule push deliberately does NOT follow, and for a reason that holds here and
|
|
// not there: a rotation that reached the
|
|
// consumer and refused on the provider would leave one end holding a credential the other has
|
|
// never heard of — which is the state this whole mechanism exists to make impossible.
|
|
func sendTo(ctx context.Context, open *stores, names []string) error {
|
|
inv := open.inventory
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
gens, err := generators(ctx, open)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Held from composing to sending (novox/hq ADR 0100); a caller that holds them already —
|
|
// converge, which flips the node and then sends it — is not made to wait on itself.
|
|
ctx, release, err := holdNodes(ctx, open, names)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer release()
|
|
|
|
var sending []readyNode
|
|
var refusals []string
|
|
for _, name := range names {
|
|
plan, settings, err := planFor(ctx, open, name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
reportUnhostable(name, plan)
|
|
declared, err := declarationWith(ctx, open, name, plan, settings, gens, Allocating)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
sending = append(sending, readyNode{name, declared})
|
|
}
|
|
if len(refusals) > 0 {
|
|
return fmt.Errorf("nothing was sent. %d machine(s) could not be resolved:\n\n%s",
|
|
len(refusals), strings.Join(refusals, "\n\n"))
|
|
}
|
|
|
|
server, err := link.Connect(nil, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
for _, s := range sending {
|
|
body, err := s.declared.Body()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil {
|
|
return err
|
|
}
|
|
record, err := inv.NodeByName(ctx, s.node)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// digestOf is what the mesh compares to answer "has this machine been sent what it should be".
|
|
//
|
|
// Over the same bytes that are sent, so the comparison is of the thing itself rather than of
|
|
// something derived beside it that could drift from it.
|
|
func digestOf(body []byte) string {
|
|
sum := sha256.Sum256(body)
|
|
return hex.EncodeToString(sum[:])
|
|
}
|
|
|
|
// wouldSend is the digest of what each machine should be right now.
|
|
//
|
|
// Machines that do not resolve are left out rather than reported as waiting: "this machine cannot
|
|
// be worked out" is a different problem with a different remedy, and `plan` is where it is said.
|
|
func wouldSend(ctx context.Context, open *stores,
|
|
nodes []inventory.Node) (map[string]string, error) {
|
|
|
|
gens, err := generators(ctx, open)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out := map[string]string{}
|
|
for _, n := range nodes {
|
|
plan, settings, err := planFor(ctx, open, n.Name)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
declared, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
body, err := declared.Body()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out[n.Name] = digestOf(body)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// raiseTheBus asserts the streams and consumers the mesh's own traffic needs.
|
|
//
|
|
// **Every start, and it says what it did.** The objects are the mesh's, created by nothing else —
|
|
// the controller is their only writer (design 25 §3) — so a mesh that came up without them is one
|
|
// where nodes connect, authenticate, and hear nothing. Said rather than silent for the reason the
|
|
// first line of `serve` is said: a log that is quiet on success and loud on failure reads as broken
|
|
// when it is working.
|
|
func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string) error {
|
|
js, err := broker.Dial(address)
|
|
if err != nil {
|
|
return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
|
address, err)
|
|
}
|
|
defer js.Close()
|
|
|
|
// **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.
|
|
if err := broker.RaiseSeats(js, inventory.MeshSeats(), nil); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n",
|
|
address, len(names))
|
|
return nil
|
|
}
|