The half of the module system that was built and never proved. Unassigning i3 removed i3's file AND xorg's, because xorg was only there to satisfy i3 -- the node's own record agrees, and the resolution the mesh sends no longer mentions either. That works because a declaration removes what the mesh previously declared and nothing else, which is 04-ISSUES/010's fix carrying its weight here: the substrate the machine raised for itself is untouched by any of it. Tests for the storage layer, which had none. The ones worth naming: A module a machine is running cannot be forgotten -- not a fault, it means the mesh would lose the ability to describe what is on that machine. Removing a node DOES take its assignments, and the asymmetry is deliberate: a node that is gone cannot be running anything. A node that has never reported has NO capabilities rather than all of them. That refuses anything needing one, which is wrong but visible -- where assuming it can do everything would assign work it cannot do and find out on the machine. And a capability the node reported as ABSENT is not counted: reading the list without the verdict would let a module onto a machine that said no. `overlay push` is gone, replaced by `push`, which sends a node its network and its modules as one declaration. Two commands that overlap is how a mesh ends up half-configured by whichever was run. The old name answers with where to go, and answers before opening a database -- needing one would turn a redirect into a connection error.
937 lines
27 KiB
Go
937 lines
27 KiB
Go
// Command mesh-control is the control plane: everything that needs to know about more than one
|
|
// node (novox/hq ADR 0006).
|
|
//
|
|
// It runs as one process holding several contexts, each owning its own store. Today it holds one,
|
|
// `inventory`, and does one thing with it — brings its schema up to date, which is step 3 of the
|
|
// bootstrap in novox/hq 07-the-substrate and the step the first node cannot get past without.
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"os"
|
|
"os/signal"
|
|
"sort"
|
|
"strings"
|
|
"syscall"
|
|
"time"
|
|
|
|
"github.com/novox/mesh-control/internal/broker"
|
|
"github.com/novox/mesh-control/internal/catalogue"
|
|
"github.com/novox/mesh-control/internal/identity"
|
|
"github.com/novox/mesh-control/internal/inventory"
|
|
"github.com/novox/mesh-control/internal/link"
|
|
"github.com/novox/mesh-control/internal/overlay"
|
|
"github.com/novox/mesh-control/internal/store"
|
|
"github.com/novox/mesh-control/internal/token"
|
|
)
|
|
|
|
// version is stamped at link time. Unset in a development build, and it says so rather than
|
|
// claiming a number.
|
|
var version = "development build"
|
|
|
|
// held is a context this process was granted, and the schema it carries.
|
|
//
|
|
// novox/hq ADR 0006 names seven. One is built. The list is short because the others do not exist
|
|
// yet, not because they are optional.
|
|
var held = []struct {
|
|
name string
|
|
migrations func() ([]store.Migration, error)
|
|
}{
|
|
{inventory.Name, inventory.Migrations},
|
|
{identity.Name, identity.Migrations},
|
|
}
|
|
|
|
func main() {
|
|
if err := run(); err != nil {
|
|
fmt.Fprintf(os.Stderr, "mesh-control: %v\n", err)
|
|
os.Exit(1)
|
|
}
|
|
}
|
|
|
|
func run() error {
|
|
args := os.Args[1:]
|
|
if len(args) == 0 {
|
|
usage()
|
|
return fmt.Errorf("no command given")
|
|
}
|
|
|
|
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
|
defer stop()
|
|
|
|
switch args[0] {
|
|
case "migrate":
|
|
return migrate(ctx)
|
|
case "node":
|
|
return nodeCommand(ctx, args[1:])
|
|
case "token":
|
|
return tokenCommand(ctx, args[1:])
|
|
case "identity":
|
|
return identityCommand(ctx, args[1:])
|
|
case "broker":
|
|
return brokerCommand(args[1:])
|
|
case "serve":
|
|
return serve(ctx)
|
|
case "declare":
|
|
return declare(ctx, args[1:])
|
|
case "overlay":
|
|
return overlayCommand(ctx, args[1:])
|
|
case "module":
|
|
return moduleCommand(ctx, args[1:])
|
|
case "assign", "unassign":
|
|
return assignCommand(ctx, args[0], args[1:])
|
|
case "plan":
|
|
return planCommand(ctx, args[1:])
|
|
case "push":
|
|
return pushCommand(ctx, args[1:])
|
|
case "version":
|
|
fmt.Println(version)
|
|
return nil
|
|
case "help", "-h", "--help":
|
|
usage()
|
|
return nil
|
|
default:
|
|
usage()
|
|
return fmt.Errorf("%q is not a command", args[0])
|
|
}
|
|
}
|
|
|
|
func usage() {
|
|
fmt.Fprint(os.Stderr, `mesh-control — the control plane
|
|
|
|
migrate bring each context's schema up to date
|
|
node add <name> create a node record
|
|
node list the nodes this mesh knows about
|
|
token issue --node <name> a one-time right to join, for an existing record
|
|
token issue --new <name> create the record and issue for it
|
|
identity show this control plane's signing key
|
|
broker show where the broker is, and what to expect there
|
|
serve consume what nodes say, and answer
|
|
declare <node> <file> send a node a signed declaration
|
|
overlay place <node> [flags] say where a node is and how it is reached
|
|
overlay show the private network, as the mesh computes it
|
|
module add <file> register a module from its manifest
|
|
module list what modules this mesh knows about
|
|
module forget <name> remove one, unless a node is running it
|
|
assign <node> <module> put a module on a node
|
|
unassign <node> <module> take it off
|
|
plan <node> what that node would run, and why
|
|
push [<node>] send a node everything it should be
|
|
version what this binary is
|
|
|
|
Each context reaches its own store through its own credential (novox/hq ADR 0008), named
|
|
`+store.Variable("<context>")+`. This process holds:
|
|
|
|
`)
|
|
for _, c := range held {
|
|
fmt.Fprintf(os.Stderr, " %-12s database %-12s from %s\n",
|
|
c.name, store.Database(c.name), store.Variable(c.name))
|
|
}
|
|
fmt.Fprintln(os.Stderr)
|
|
}
|
|
|
|
// migrate brings every held context's schema up to date.
|
|
//
|
|
// Reported per context and per migration, because this runs during a bootstrap on a machine with
|
|
// nothing else on it — the output is the only account of what happened, and "migrated" is not one.
|
|
func migrate(ctx context.Context) error {
|
|
for _, c := range held {
|
|
migrations, err := c.migrations()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
s, err := store.Open(ctx, c.name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer s.Close()
|
|
|
|
// The bootstrap raises PostgreSQL moments before this runs, and a container that is
|
|
// running is not a database that will answer — a distinction this project has already
|
|
// paid for once, when a crash-looping database reported itself as up between restarts.
|
|
if err := s.Ready(ctx, 60*time.Second); err != nil {
|
|
return err
|
|
}
|
|
|
|
done, err := s.Migrate(ctx, migrations)
|
|
for _, m := range done {
|
|
fmt.Printf("%s: applied %04d-%s\n", c.name, m.Number, m.Name)
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(done) == 0 {
|
|
applied, err := s.AppliedMigrations(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("%s: already up to date — %d migration(s)\n", c.name, len(applied))
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// openInventory connects and waits, the way every command that touches it needs to.
|
|
func openInventory(ctx context.Context) (*inventory.Inventory, error) {
|
|
inv, err := inventory.Open(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if err := inv.Ready(ctx, 30*time.Second); err != nil {
|
|
inv.Close()
|
|
return nil, err
|
|
}
|
|
return inv, nil
|
|
}
|
|
|
|
func nodeCommand(ctx context.Context, args []string) error {
|
|
if len(args) == 0 {
|
|
return errors.New("node add <name>, or node list")
|
|
}
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
|
|
switch args[0] {
|
|
case "add":
|
|
if len(args) != 2 {
|
|
return errors.New("node add <name>")
|
|
}
|
|
node, err := inv.AddNode(ctx, args[1])
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("added %s (%s)\n", node.Name, node.ID)
|
|
return nil
|
|
|
|
case "list":
|
|
nodes, err := inv.Nodes(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(nodes) == 0 {
|
|
// Said rather than printed as nothing: an empty list and a failed read must never
|
|
// look the same, and this command answering "none" is only honest because getting
|
|
// here means the store answered.
|
|
fmt.Println("this mesh has no node records yet")
|
|
return nil
|
|
}
|
|
for _, n := range nodes {
|
|
fmt.Printf("%-20s %-14s %s\n", n.Name, heardFrom(n), n.ID)
|
|
}
|
|
return nil
|
|
|
|
default:
|
|
return fmt.Errorf("node has no %q; it has add and list", args[0])
|
|
}
|
|
}
|
|
|
|
func tokenCommand(ctx context.Context, args []string) error {
|
|
if len(args) == 0 || args[0] != "issue" {
|
|
return errors.New("token issue --node <name>, or token issue --new <name>")
|
|
}
|
|
|
|
set := flag.NewFlagSet("token issue", flag.ContinueOnError)
|
|
existing := set.String("node", "", "issue for a node record that already exists")
|
|
fresh := set.String("new", "", "create the node record, then issue for it")
|
|
validFor := set.Duration("for", time.Hour, "how long the token may be used")
|
|
if err := set.Parse(args[1:]); err != nil {
|
|
return err
|
|
}
|
|
|
|
// Exactly one, because the difference is what the token binds to. A command that guessed
|
|
// would sometimes create a second record for a machine that already has one.
|
|
if (*existing == "") == (*fresh == "") {
|
|
return errors.New("give exactly one of --node <name> or --new <name>: the first is a " +
|
|
"machine the mesh already has a record for, the second is one it has never seen")
|
|
}
|
|
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
|
|
name := *existing
|
|
if *fresh != "" {
|
|
node, err := inv.AddNode(ctx, *fresh)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
name = node.Name
|
|
}
|
|
|
|
issued, err := inv.IssueToken(ctx, name, *validFor)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Assembled from two contexts by the process that holds both grants. Neither reads the
|
|
// other's store (novox/hq ADR 0008) — each is asked for its own part.
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
key, err := ident.Establish(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// The account is created before the token is handed over, which is what removes the
|
|
// chicken-and-egg entirely: the mesh runs the broker, so a joining node's credentials can
|
|
// exist before it does. The one-time secret IS the password, so a node's first connection is
|
|
// already authenticated and enrolment is what happens over it.
|
|
if management, err := broker.ManagementFromEnvironment(); err == nil {
|
|
if err := management.CreateNodeAccount(ctx, issued.Node.Name, issued.Secret); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("broker account %s created, scoped to %s and the %s exchange\n\n",
|
|
issued.Node.Name, link.QueueFor(issued.Node.Name), link.Exchange)
|
|
} else if !errors.Is(err, broker.ErrNotConfigured) {
|
|
return err
|
|
}
|
|
|
|
made := token.Token{Signer: key.Public, Secret: issued.Secret}
|
|
|
|
// Absent is a state, not a failure: a control plane can hold records and a key before it has
|
|
// a broker. What it cannot do is issue a token anybody could use, and Missing() says so.
|
|
known, err := broker.FromEnvironment()
|
|
switch {
|
|
case err == nil:
|
|
made.Broker, made.Fingerprint = known.Address, known.Fingerprint
|
|
case errors.Is(err, broker.ErrNotConfigured):
|
|
default:
|
|
return err
|
|
}
|
|
encoded, err := made.Encode()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
fmt.Printf("token for %s, usable once, until %s\n\n %s\n\n",
|
|
issued.Node.Name, issued.Expires.Format(time.RFC3339), encoded)
|
|
fmt.Println("This is the only time it is shown. What is stored is a hash of the secret.")
|
|
|
|
if missing := made.Missing(); len(missing) > 0 {
|
|
fmt.Printf("\nINCOMPLETE — this token cannot be used to join anything yet. Missing:\n")
|
|
for _, m := range missing {
|
|
fmt.Printf(" - %s\n", m)
|
|
}
|
|
fmt.Printf("\nSet %s and %s once the broker is raised.\n",
|
|
broker.AddressVar, broker.CertificateVar)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func openIdentity(ctx context.Context) (*identity.Identity, error) {
|
|
ident, err := identity.Open(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if err := ident.Ready(ctx, 30*time.Second); err != nil {
|
|
ident.Close()
|
|
return nil, err
|
|
}
|
|
return ident, nil
|
|
}
|
|
|
|
func identityCommand(ctx context.Context, args []string) error {
|
|
if len(args) == 0 || args[0] != "show" {
|
|
return errors.New("identity show")
|
|
}
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
// Establish rather than read: a control plane asked for its identity before it has one should
|
|
// get one, not an error. Generating it is idempotent, so this is safe to run at any time.
|
|
key, err := ident.Establish(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("signing key %s\n", key.ID)
|
|
fmt.Printf("fingerprint %s\n", key.Fingerprint())
|
|
fmt.Printf("created %s\n", key.Created.Format(time.RFC3339))
|
|
fmt.Printf("\nThe public half of this travels in every enrolment token. A node believes a\n" +
|
|
"declaration because it carries a signature this key made (novox/hq ADR 0004).\n")
|
|
return nil
|
|
}
|
|
|
|
func brokerCommand(args []string) error {
|
|
if len(args) == 0 || args[0] != "show" {
|
|
return errors.New("broker show")
|
|
}
|
|
known, err := broker.FromEnvironment()
|
|
if errors.Is(err, broker.ErrNotConfigured) {
|
|
fmt.Printf("no broker configured. Set %s and %s.\n\n"+
|
|
"Until then tokens carry the signing key and the one-time secret, and say what they\n"+
|
|
"are missing. They cannot be used to join.\n",
|
|
broker.AddressVar, broker.CertificateVar)
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("address %s\n", known.Address)
|
|
fmt.Printf("fingerprint %s\n", known.Fingerprint)
|
|
fmt.Print("\nThe fingerprint is computed from the certificate on disk, never configured. A\n" +
|
|
"node checks it before sending anything (novox/hq ADR 0004).\n")
|
|
return nil
|
|
}
|
|
|
|
// serve is the control plane running: one connection to the broker, one queue, one consumer.
|
|
func serve(ctx context.Context) error {
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
|
|
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()
|
|
|
|
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.
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
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"
|
|
|
|
func overlayCIDR() string {
|
|
if v := strings.TrimSpace(os.Getenv(OverlayCIDRVar)); v != "" {
|
|
return v
|
|
}
|
|
return "10.42.0.0/16"
|
|
}
|
|
|
|
func overlayCommand(ctx context.Context, args []string) error {
|
|
if len(args) == 0 {
|
|
return errors.New("overlay place <node> [flags], or overlay show")
|
|
}
|
|
// Answered before anything is opened. A message about which command to use should not need a
|
|
// database to say so, and needing one turns a redirect into a connection error.
|
|
if args[0] == "push" {
|
|
return errors.New("`overlay push` is now `push`, which sends a node its network AND " +
|
|
"what its assignments resolve to — the two are computed from one picture of the " +
|
|
"mesh, and sending them separately would let them disagree")
|
|
}
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
|
|
switch args[0] {
|
|
case "place":
|
|
return overlayPlace(ctx, inv, args[1:])
|
|
case "show":
|
|
return overlayShow(ctx, inv)
|
|
|
|
default:
|
|
return fmt.Errorf("overlay has no %q; it has place and show", args[0])
|
|
}
|
|
}
|
|
|
|
func overlayPlace(ctx context.Context, inv *inventory.Inventory, args []string) error {
|
|
if len(args) == 0 {
|
|
return errors.New("overlay place <node> [--endpoint host:port] [--site name] [--hub]")
|
|
}
|
|
node := args[0]
|
|
|
|
set := flag.NewFlagSet("overlay place", flag.ContinueOnError)
|
|
endpoint := set.String("endpoint", "", "where this node can be dialled, or empty for nowhere")
|
|
site := set.String("site", "", "where this machine physically is, or empty if it roams")
|
|
hub := set.Bool("hub", false, "this node is the hub every other routes through")
|
|
if err := set.Parse(args[1:]); err != nil {
|
|
return err
|
|
}
|
|
|
|
// Declared, all three. The address is evidence of reachability and is not the fact, and hub
|
|
// election by address prefix fails silently (novox/hq ADR 0007).
|
|
if err := inv.SetPlace(ctx, node, *endpoint, *site, *hub, ""); err != nil {
|
|
return err
|
|
}
|
|
found, err := inv.NodeByName(ctx, node)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
address, err := inv.AssignAddress(ctx, found.ID, overlayCIDR())
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
fmt.Printf("%s is at %s on the overlay\n", node, address)
|
|
switch {
|
|
case *hub:
|
|
fmt.Println(" the hub — every node not sharing a site routes through it")
|
|
case *endpoint == "":
|
|
fmt.Println(" not dialable — it opens every path itself")
|
|
}
|
|
if *site != "" {
|
|
fmt.Printf(" at %s, so it peers directly with anything else there\n", *site)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// graph reads every node's place and computes the network. Every node at once, which is the whole
|
|
// reason this is the control plane's work.
|
|
func graph(ctx context.Context, inv *inventory.Inventory) ([]overlay.Node, overlay.Graph, error) {
|
|
places, err := inv.Overlays(ctx)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
nodes := make([]overlay.Node, 0, len(places))
|
|
for _, p := range places {
|
|
nodes = append(nodes, overlay.Node{
|
|
Name: p.Name, Key: p.Key, Endpoint: p.Endpoint,
|
|
Site: p.Site, Hub: p.Hub, Address: p.Address,
|
|
})
|
|
}
|
|
computed, err := overlay.Compute(nodes, overlayCIDR())
|
|
return nodes, computed, err
|
|
}
|
|
|
|
func overlayShow(ctx context.Context, inv *inventory.Inventory) error {
|
|
nodes, computed, err := graph(ctx, inv)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(nodes) == 0 {
|
|
fmt.Println("this mesh has no nodes")
|
|
return nil
|
|
}
|
|
|
|
for _, n := range nodes {
|
|
place := n.Address
|
|
if place == "" {
|
|
// Said, not skipped. A node with no place is a node with no network, and it should
|
|
// be visible here rather than quietly absent from a list of who is on it.
|
|
place = "no address — run `overlay place`"
|
|
}
|
|
fmt.Printf("%-16s %-14s", n.Name, place)
|
|
switch {
|
|
case n.Hub:
|
|
fmt.Print(" hub")
|
|
case !n.Reachable():
|
|
fmt.Print(" not dialable")
|
|
}
|
|
if n.Site != "" {
|
|
fmt.Printf(" at %s", n.Site)
|
|
}
|
|
fmt.Println()
|
|
for _, p := range computed[n.Name] {
|
|
fmt.Printf(" → %-14s %-18s %s\n", p.Name, p.Allowed, p.Why)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SilentFor is how long a node may be quiet before the mesh says so.
|
|
//
|
|
// A node speaks every minute, so three of them missed is a gap rather than a slow one. The number
|
|
// is not the point — being able to say "out of touch" at all is, and nothing could before.
|
|
const SilentFor = 3 * time.Minute
|
|
|
|
// heardFrom says when a node was last heard from, in a form somebody can act on.
|
|
//
|
|
// "never" and "an hour ago" are different answers and are kept different. A node that has never
|
|
// spoken did not finish joining; a node last heard from an hour ago is running an hour-old
|
|
// picture of the mesh.
|
|
func heardFrom(n inventory.Node) string {
|
|
silent, ever := n.Silent()
|
|
switch {
|
|
case !ever:
|
|
return "never spoken"
|
|
case silent > SilentFor:
|
|
return "out of touch " + roughly(silent)
|
|
default:
|
|
return "here"
|
|
}
|
|
}
|
|
|
|
// roughly is a duration a person reads rather than parses.
|
|
func roughly(d time.Duration) string {
|
|
switch {
|
|
case d < time.Hour:
|
|
return fmt.Sprintf("%dm", int(d.Minutes()))
|
|
case d < 48*time.Hour:
|
|
return fmt.Sprintf("%dh", int(d.Hours()))
|
|
default:
|
|
return fmt.Sprintf("%dd", int(d.Hours()/24))
|
|
}
|
|
}
|
|
|
|
func moduleCommand(ctx context.Context, args []string) error {
|
|
if len(args) == 0 {
|
|
return errors.New("module add <file>, module list, or module forget <name>")
|
|
}
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
|
|
switch args[0] {
|
|
case "add":
|
|
if len(args) != 2 {
|
|
return errors.New("module add <manifest.json>")
|
|
}
|
|
raw, err := os.ReadFile(args[1])
|
|
if err != nil {
|
|
return err
|
|
}
|
|
m, err := catalogue.ParseManifest(raw)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := inv.RegisterModule(ctx, m); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("%s registered", m.Module)
|
|
if len(m.Provides) > 0 {
|
|
fmt.Printf(", providing %s", strings.Join(m.Provides, ", "))
|
|
}
|
|
fmt.Println()
|
|
for _, c := range m.Claims {
|
|
fmt.Printf(" claims %s, one per %s\n", c.Name, c.At())
|
|
}
|
|
return nil
|
|
|
|
case "list":
|
|
shelf, err := inv.Catalogue(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(shelf) == 0 {
|
|
fmt.Println("this mesh knows about no modules yet")
|
|
return nil
|
|
}
|
|
var names []string
|
|
for n := range shelf {
|
|
names = append(names, n)
|
|
}
|
|
sort.Strings(names)
|
|
for _, n := range names {
|
|
m := shelf[n]
|
|
fmt.Printf("%-20s", m.Module)
|
|
if len(m.Provides) > 0 {
|
|
fmt.Printf(" provides %s", strings.Join(m.Provides, ", "))
|
|
}
|
|
for _, c := range m.Claims {
|
|
fmt.Printf(" claims %s/%s", c.At(), c.Name)
|
|
}
|
|
fmt.Println()
|
|
}
|
|
return nil
|
|
|
|
case "forget":
|
|
if len(args) != 2 {
|
|
return errors.New("module forget <name>")
|
|
}
|
|
if err := inv.ForgetModule(ctx, args[1]); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("%s forgotten\n", args[1])
|
|
return nil
|
|
|
|
default:
|
|
return fmt.Errorf("module has no %q; it has add, list and forget", args[0])
|
|
}
|
|
}
|
|
|
|
func assignCommand(ctx context.Context, verb string, args []string) error {
|
|
if len(args) != 2 {
|
|
return fmt.Errorf("%s <node> <module>", verb)
|
|
}
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
|
|
if verb == "unassign" {
|
|
if err := inv.Unassign(ctx, args[0], args[1]); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("%s no longer runs %s — run `push %s` to make it so\n", args[0], args[1], args[0])
|
|
return nil
|
|
}
|
|
if err := inv.Assign(ctx, args[0], args[1]); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("%s is assigned %s\n", args[0], args[1])
|
|
|
|
// Resolved immediately, because an assignment that cannot be applied should be said now
|
|
// rather than at the next push. The assignment is kept either way: it is what a person meant,
|
|
// and the refusal is about the set rather than about this one.
|
|
if _, err := planFor(ctx, inv, args[0]); err != nil {
|
|
fmt.Println()
|
|
return err
|
|
}
|
|
fmt.Printf(" run `push %s` to send it\n", args[0])
|
|
return nil
|
|
}
|
|
|
|
// planFor works out everything a node should run, from what was assigned to it.
|
|
func planFor(ctx context.Context, inv *inventory.Inventory, nodeName string) (catalogue.Resolution, error) {
|
|
shelf, err := inv.Catalogue(ctx)
|
|
if err != nil {
|
|
return catalogue.Resolution{}, err
|
|
}
|
|
assigned, err := inv.Assigned(ctx, nodeName)
|
|
if err != nil {
|
|
return catalogue.Resolution{}, err
|
|
}
|
|
capabilities, err := inv.ProfileOf(ctx, nodeName)
|
|
if err != nil {
|
|
return catalogue.Resolution{}, err
|
|
}
|
|
|
|
places, err := inv.Overlays(ctx)
|
|
if err != nil {
|
|
return catalogue.Resolution{}, err
|
|
}
|
|
var site string
|
|
for _, p := range places {
|
|
if p.Name == nodeName {
|
|
site = p.Site
|
|
}
|
|
}
|
|
|
|
// What every other node already holds, so the claims wider than one machine can be checked.
|
|
// Resolved rather than read from a table: a claim is held by whatever a node actually runs,
|
|
// and a record of it would be a second answer that could disagree with the first.
|
|
var elsewhere []catalogue.Held
|
|
for _, p := range places {
|
|
if p.Name == nodeName {
|
|
continue
|
|
}
|
|
theirs, err := inv.Assigned(ctx, p.Name)
|
|
if err != nil || len(theirs) == 0 {
|
|
continue
|
|
}
|
|
theirCaps, _ := inv.ProfileOf(ctx, p.Name)
|
|
got, err := catalogue.Resolve(shelf, theirs,
|
|
catalogue.Node{Name: p.Name, Site: p.Site, Capabilities: theirCaps}, nil)
|
|
if err != nil {
|
|
// Their set does not resolve either. Not this node's problem to report, and their
|
|
// claims cannot be counted because nothing of theirs is running.
|
|
continue
|
|
}
|
|
elsewhere = append(elsewhere, got.Claims...)
|
|
}
|
|
|
|
return catalogue.Resolve(shelf, assigned,
|
|
catalogue.Node{Name: nodeName, Site: site, Capabilities: capabilities}, elsewhere)
|
|
}
|
|
|
|
func planCommand(ctx context.Context, args []string) error {
|
|
if len(args) != 1 {
|
|
return errors.New("plan <node>")
|
|
}
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
|
|
plan, err := planFor(ctx, inv, args[0])
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(plan.Modules) == 0 {
|
|
fmt.Printf("%s is assigned nothing\n", args[0])
|
|
return nil
|
|
}
|
|
fmt.Printf("%s would run:\n", args[0])
|
|
for _, m := range plan.Modules {
|
|
fmt.Printf(" %-20s %s\n", m.Module, plan.Because[m.Module])
|
|
}
|
|
for _, c := range plan.Claims {
|
|
fmt.Printf(" holds %s, one per %s\n", c.Claim, c.Scope)
|
|
}
|
|
fmt.Printf("\n%d resource(s)\n", len(plan.Declaration()))
|
|
return nil
|
|
}
|
|
|
|
// 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 {
|
|
if len(args) > 1 {
|
|
return errors.New("push [<node>] — one node, or all of them")
|
|
}
|
|
inv, err := openInventory(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer inv.Close()
|
|
|
|
ident, err := openIdentity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer ident.Close()
|
|
|
|
nodes, computed, err := graph(ctx, inv)
|
|
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 overlay.Node
|
|
resources []map[string]any
|
|
}
|
|
var sending []ready
|
|
var refusals []string
|
|
|
|
for _, n := range nodes {
|
|
if len(args) == 1 && n.Name != args[0] {
|
|
continue
|
|
}
|
|
peers, onOverlay := computed[n.Name]
|
|
if !onOverlay {
|
|
fmt.Printf("%s is not on the overlay yet — skipped\n", n.Name)
|
|
continue
|
|
}
|
|
|
|
declaration, err := overlay.Declaration(n, peers, nodes, "")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var resources struct {
|
|
Resources []map[string]any `json:"resources"`
|
|
}
|
|
if err := json.Unmarshal(declaration, &resources); err != nil {
|
|
return err
|
|
}
|
|
|
|
plan, err := planFor(ctx, inv, n.Name)
|
|
if err != nil {
|
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", n.Name, err))
|
|
continue
|
|
}
|
|
sending = append(sending, ready{n, append(resources.Resources, plan.Declaration()...)})
|
|
}
|
|
|
|
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.Name, body, 15*time.Second); err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("sent %s %d resource(s)\n", s.node.Name, len(s.resources))
|
|
}
|
|
fmt.Printf("\n%d node(s) told\n", len(sending))
|
|
return nil
|
|
}
|