`status` hung. It composes a declaration for every node to answer *is this machine running what I would send it*, and composing one assigns each module a machine port — so the question wrote to the database, and wrote to the same rows as the machine it was asking about. `port_assignment` is unique on (node, machine). Two transactions inserting the same port do not race, they queue: the second waits on the index until the first commits. A status polled every two seconds while a node applies is two writers on those rows, and the poll stopped returning rather than returning something wrong — which is the better failure of the two, and still a failure. The latent version of this was there before anything polled: two compositions running at once could both allocate. So allocation belongs to the send path alone. The mesh chooses a port when it commits to sending one; every other caller reads what was chosen. A module with nothing assigned has never been sent, which is precisely what "waiting" means — the read needs no number to be right about that, and inventing one would make the answer worse. Named rather than passed as a bare bool: at three call sites, `true` and `false` say nothing about which of these two things is meant. Checked by the lab, which now polls status throughout an apply.
414 lines
13 KiB
Go
414 lines
13 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"os"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/novox/mesh-control/internal/broker"
|
|
"github.com/novox/mesh-control/internal/inventory"
|
|
"github.com/novox/mesh-control/internal/link"
|
|
)
|
|
|
|
// sending it, and holding the link that carries it.
|
|
//
|
|
// Split out of main.go, which had reached 2,769 lines because appending was always the
|
|
// cheapest next step. That is how novox/hq ADR 0001 records `hal/sdk` reaching 34,636:
|
|
// nothing in it was wrong, and no one edit was the one that should have been a new file.
|
|
|
|
// serve is the control plane running: one connection to the broker, one queue, one consumer.
|
|
func serve(ctx context.Context) error {
|
|
open, err := openStores(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer open.Close()
|
|
inv := open.inventory
|
|
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
// Established at start rather than on first use. A control plane that cannot sign is one
|
|
// whose declarations every node correctly refuses, and that should be a startup failure
|
|
// rather than something discovered at the first declaration.
|
|
key, err := ident.Establish(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("signing as %s\n", key.Fingerprint()[:16])
|
|
|
|
management, err := broker.ManagementFromEnvironment()
|
|
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
|
|
return err
|
|
}
|
|
|
|
// Where the broker is and what to expect there, so a node can be told how to come back
|
|
// without a person and a new token.
|
|
known, err := broker.FromEnvironment()
|
|
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
|
|
return err
|
|
}
|
|
if errors.Is(err, broker.ErrNotConfigured) {
|
|
fmt.Printf("no broker address configured, so enrolled nodes will not be told how to "+
|
|
"reconnect. Set %s and %s.\n", broker.AddressVar, broker.CertificateVar)
|
|
}
|
|
|
|
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known}
|
|
server, err := link.Connect(work, work)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
// And build results nobody was waiting for. A build triggered any other way than `build`
|
|
// would otherwise be reported into the void, which is the same as not reporting it.
|
|
server.Records(builds{inv})
|
|
|
|
return server.Serve(ctx)
|
|
}
|
|
|
|
// declare sends one node a declaration, signed.
|
|
//
|
|
// Signed here rather than trusted from the broker: a node connects to the broker and takes
|
|
// instruction from the control plane behind it, and those are two identities. If a node believed
|
|
// whatever arrived on its queue, a compromised broker could forge declarations — and since the
|
|
// host applies whatever the link delivers, that is the whole machine (novox/hq ADR 0004).
|
|
func declare(ctx context.Context, args []string) error {
|
|
if len(args) != 2 {
|
|
return errors.New("declare <node> <declaration.json>")
|
|
}
|
|
node, path := args[0], args[1]
|
|
|
|
raw, err := os.ReadFile(path)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
// The node has to exist before it can be told anything. Publishing to a queue nobody consumes
|
|
// would sit there looking like success.
|
|
open, err := openStores(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer open.Close()
|
|
inv := open.inventory
|
|
if _, err := inv.NodeByName(ctx, node); err != nil {
|
|
return err
|
|
}
|
|
|
|
server, err := link.Connect(nil, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
if err := link.Declare(ctx, server.Channel(), ident, node, raw, 15*time.Second); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
|
return nil
|
|
}
|
|
|
|
// OverlayCIDRVar is the range the mesh allocates node addresses from.
|
|
const OverlayCIDRVar = "MESH_OVERLAY_CIDR"
|
|
|
|
// pushCommand sends nodes everything they should be: their place on the network, and what their
|
|
// assignments resolve to.
|
|
//
|
|
// One declaration, not two. A node holding its network and not its modules, or the reverse, is
|
|
// half-configured for as long as that lasts — and the two are computed from the same picture of
|
|
// the mesh, so sending them apart would let them disagree.
|
|
func pushCommand(ctx context.Context, args []string) error {
|
|
set := flag.NewFlagSet("push", flag.ContinueOnError)
|
|
// Only the machines that need it.
|
|
//
|
|
// **A command rather than a timer, to begin with.** Something that re-pushes on a schedule is
|
|
// a scheduler over this, and building the scheduler first would mean two paths to one act
|
|
// with nothing to compare them against. A person can run this; so can cron; so can whatever
|
|
// eventually watches.
|
|
behind := set.Bool("behind", false,
|
|
"only machines whose last declaration was refused or partly failed")
|
|
positionals, err := parseAround(set, args)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
args = positionals
|
|
if len(args) > 1 {
|
|
return errors.New("push [<node>] [--behind] — one node, or all of them")
|
|
}
|
|
if len(args) == 1 && *behind {
|
|
// Naming a machine and asking for the ones that need it are two different requests, and
|
|
// guessing which was meant would sometimes push to a machine somebody did not name.
|
|
return errors.New("push <node> or push --behind, not both: one names a machine and the " +
|
|
"other asks which machines need one")
|
|
}
|
|
open, err := openStores(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer open.Close()
|
|
inv := open.inventory
|
|
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
// Every node, not only the ones on the private network. A machine that was never given the
|
|
// network module still takes modules, and iterating the network here is what used to make
|
|
// "on the network" and "managed" the same thing.
|
|
nodes, err := inv.Nodes(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Which machines are not in the state they were sent, when that is what was asked for.
|
|
var needsOne map[string]inventory.Doing
|
|
if *behind {
|
|
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
needsOne = map[string]inventory.Doing{}
|
|
for _, d := range wrong {
|
|
needsOne[d.Node] = d
|
|
}
|
|
// **And every machine not running what the mesh would send it.** "Behind" used to mean
|
|
// only "failed or refused", so a machine that applied cleanly and whose declaration has
|
|
// since changed was not behind — and novox/hq ADR 0010's question, *did my change go
|
|
// out?*, was answerable only for the machines that broke.
|
|
would, err := wouldSend(ctx, open, nodes)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
waiting, err := inv.Waiting(ctx, would)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, m := range waiting {
|
|
if _, already := needsOne[m.Node]; already {
|
|
continue
|
|
}
|
|
needsOne[m.Node] = inventory.Doing{Node: m.Node, Outcome: "waiting"}
|
|
}
|
|
if len(needsOne) == 0 {
|
|
// Said rather than doing nothing quietly. "Nothing needed one" and "this did not run"
|
|
// must never look the same.
|
|
fmt.Println("every machine is doing what it was told")
|
|
return nil
|
|
}
|
|
}
|
|
gens, err := generators(ctx, open)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
server, err := link.Connect(nil, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
// Every node is resolved before anything is sent. A push that configured three nodes and then
|
|
// refused on the fourth would leave the mesh in a state nobody asked for, and the fourth is
|
|
// exactly where a claim collision shows up.
|
|
type ready struct {
|
|
node string
|
|
resources []map[string]any
|
|
}
|
|
var sending []ready
|
|
var refusals []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"))
|
|
}
|
|
}
|
|
plan, settings, err := planFor(ctx, open, n.Name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", n.Name, err))
|
|
continue
|
|
}
|
|
// 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.
|
|
resources, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Allocating)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", n.Name, err))
|
|
continue
|
|
}
|
|
if len(resources) == 0 {
|
|
fmt.Printf("%s is assigned nothing — skipped\n", n.Name)
|
|
continue
|
|
}
|
|
sending = append(sending, ready{n.Name, resources})
|
|
}
|
|
|
|
if len(refusals) > 0 {
|
|
return fmt.Errorf("nothing was sent. %d node(s) could not be resolved:\n\n%s",
|
|
len(refusals), strings.Join(refusals, "\n\n"))
|
|
}
|
|
|
|
for _, s := range sending {
|
|
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": s.resources})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil {
|
|
return err
|
|
}
|
|
// After it is away, not before. A digest recorded for something that failed to send would
|
|
// make the machine look current for a declaration it never received.
|
|
record, err := inv.NodeByName(ctx, s.node)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.resources))
|
|
}
|
|
fmt.Printf("\n%d node(s) told\n", len(sending))
|
|
return nil
|
|
}
|
|
|
|
// sendTo resolves and sends to exactly the machines named, or refuses without sending anything.
|
|
//
|
|
// The same all-or-nothing rule push follows, and for the same reason: a rotation that reached the
|
|
// consumer and refused on the provider would leave one end holding a credential the other has
|
|
// never heard of — which is the state this whole mechanism exists to make impossible.
|
|
func sendTo(ctx context.Context, open *stores, names []string) error {
|
|
inv := open.inventory
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
gens, err := generators(ctx, open)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
type ready struct {
|
|
node string
|
|
resources []map[string]any
|
|
}
|
|
var sending []ready
|
|
var refusals []string
|
|
for _, name := range names {
|
|
plan, settings, err := planFor(ctx, open, name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
resources, err := declarationWith(ctx, open, name, plan, settings, gens, Allocating)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
continue
|
|
}
|
|
sending = append(sending, ready{name, resources})
|
|
}
|
|
if len(refusals) > 0 {
|
|
return fmt.Errorf("nothing was sent. %d machine(s) could not be resolved:\n\n%s",
|
|
len(refusals), strings.Join(refusals, "\n\n"))
|
|
}
|
|
|
|
server, err := link.Connect(nil, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer server.Close()
|
|
|
|
for _, s := range sending {
|
|
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": s.resources})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil {
|
|
return err
|
|
}
|
|
record, err := inv.NodeByName(ctx, s.node)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.resources))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// digestOf is what the mesh compares to answer "has this machine been sent what it should be".
|
|
//
|
|
// Over the same bytes that are sent, so the comparison is of the thing itself rather than of
|
|
// something derived beside it that could drift from it.
|
|
func digestOf(body []byte) string {
|
|
sum := sha256.Sum256(body)
|
|
return hex.EncodeToString(sum[:])
|
|
}
|
|
|
|
// wouldSend is the digest of what each machine should be right now.
|
|
//
|
|
// Machines that do not resolve are left out rather than reported as waiting: "this machine cannot
|
|
// be worked out" is a different problem with a different remedy, and `plan` is where it is said.
|
|
func wouldSend(ctx context.Context, open *stores,
|
|
nodes []inventory.Node) (map[string]string, error) {
|
|
|
|
gens, err := generators(ctx, open)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out := map[string]string{}
|
|
for _, n := range nodes {
|
|
plan, settings, err := planFor(ctx, open, n.Name)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
resources, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": resources})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out[n.Name] = digestOf(body)
|
|
}
|
|
return out, nil
|
|
}
|