// 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 "settings": return settingsCommand(ctx, args[1:]) case "plan": return planCommand(ctx, args[1:]) case "push": return pushCommand(ctx, args[1:]) case "status": return statusCommand(ctx) 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 create a node record node list the nodes this mesh knows about token issue --node a one-time right to join, for an existing record token issue --new 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 send a node a signed declaration overlay place [flags] say where a node is and how it is reached overlay show the private network, as the mesh computes it module add register a module from its manifest module list what modules this mesh knows about module moved the source has a newer commit than the mesh built module forget remove one, unless a node is running it status what the mesh is behind on, and which nodes assign put a module on a node unassign take it off settings set what a module's config should say, for the whole mesh settings set --node ...or for one machine settings clear [--node ] take a layer away plan what that node would run, and why push [] 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("")+`. 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 , 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 ") } 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 , or token issue --new ") } 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 or --new : 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, 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 [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 [--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 , module list, or module forget ") } inv, err := openInventory(ctx) if err != nil { return err } defer inv.Close() switch args[0] { case "add": set := flag.NewFlagSet("module add", flag.ContinueOnError) repo := set.String("source", "", "where this module comes from") ref := set.String("ref", "", "the branch followed there") commit := set.String("commit", "", "the commit this manifest was read at") positionals, err := parseAround(set, args[1:]) if err != nil { return err } if len(positionals) != 1 { return errors.New("module add [--source --ref --commit ]") } raw, err := os.ReadFile(positionals[0]) if err != nil { return err } m, err := catalogue.ParseManifest(raw) if err != nil { return err } // Provenance together or not at all. A source with no commit cannot be compared against // anything, so it would record where the module came from and still never be able to say // the mesh is behind it — which is the one thing recording it is for. if (*repo == "") != (*commit == "") { return errors.New("--source and --commit go together: a source with no commit " + "cannot be compared against anything, and a commit with no source has nothing " + "to be compared with") } if err := inv.RegisterModule(ctx, m, inventory.Source{ Repository: *repo, Ref: *ref, BuiltFrom: *commit, }); err != nil { return err } fmt.Printf("%s registered", m.Module) if *commit != "" { fmt.Printf(" from %s", short(*commit)) } 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 "moved": if len(args) != 3 { return errors.New("module moved — the source has a newer commit") } if err := inv.SourceMoved(ctx, args[1], args[2]); err != nil { return err } from, err := inv.SourceOf(ctx, args[1]) if err != nil { return err } if from.Current() { fmt.Printf("%s is current at %s\n", args[1], short(from.Head)) return nil } fmt.Printf("%s is behind: the mesh holds %s and the source has %s\n", args[1], short(from.BuiltFrom), short(from.Head)) fmt.Println(" build it and `module add` the result to catch up") return nil case "forget": if len(args) != 2 { return errors.New("module forget ") } 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, moved and forget", args[0]) } } func assignCommand(ctx context.Context, verb string, args []string) error { if len(args) != 2 { return fmt.Errorf("%s ", 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, catalogue.SettingsBy, error) { shelf, err := inv.Catalogue(ctx) if err != nil { return catalogue.Resolution{}, nil, err } assigned, err := inv.Assigned(ctx, nodeName) if err != nil { return catalogue.Resolution{}, nil, err } capabilities, err := inv.ProfileOf(ctx, nodeName) if err != nil { return catalogue.Resolution{}, nil, err } places, err := inv.Overlays(ctx) if err != nil { return catalogue.Resolution{}, nil, 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...) } resolved, err := catalogue.Resolve(shelf, assigned, catalogue.Node{Name: nodeName, Site: site, Capabilities: capabilities}, elsewhere) if err != nil { return catalogue.Resolution{}, nil, err } // Settings for everything that resolved, including modules nobody assigned directly: a // requirement pulled in by something else is still configurable, and finding out that it is // not only when you try would be an arbitrary line nobody could predict. settings := catalogue.SettingsBy{} var stray []string for _, m := range resolved.Modules { layers, err := inv.SettingsFor(ctx, nodeName, m.Module) if err != nil { return catalogue.Resolution{}, nil, err } if len(layers) == 0 { continue } settings[m.Module] = layers stray = append(stray, catalogue.UnusedSettings(m, layers)...) } if len(stray) > 0 { // Somebody set something that reaches no file. Said here rather than discovered by the // machine not behaving differently, which is the slowest way there is. return catalogue.Resolution{}, nil, fmt.Errorf( "these settings reach nothing:\n - %s", strings.Join(stray, "\n - ")) } return resolved, settings, nil } func planCommand(ctx context.Context, args []string) error { set := flag.NewFlagSet("plan", flag.ContinueOnError) // Because "one resource" does not tell you whether the settings landed. Being able to read // the file before it is sent is the difference between believing a merge worked and knowing. show := set.Bool("files", false, "print the files this node would be given") positionals, err := parseAround(set, args) if err != nil { return err } if len(positionals) != 1 { return errors.New("plan [--files]") } args = positionals inv, err := openInventory(ctx) if err != nil { return err } defer inv.Close() plan, settings, 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) } resources, err := plan.Declaration(settings) if err != nil { return err } for module, layers := range settings { for _, layer := range layers { fmt.Printf(" %-20s settings from %s\n", module, layer.From) } } fmt.Printf("\n%d resource(s)\n", len(resources)) if *show { for _, r := range resources { content, ok := r["content"].(string) if !ok { continue } fmt.Printf("\n--- %v %v ---\n%s", r["id"], r["path"], content) } } 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 [] — 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, settings, err := planFor(ctx, inv, n.Name) if err != nil { refusals = append(refusals, fmt.Sprintf("%s:\n%v", n.Name, err)) continue } fromModules, err := plan.Declaration(settings) if err != nil { refusals = append(refusals, fmt.Sprintf("%s:\n%v", n.Name, err)) continue } sending = append(sending, ready{n, append(resources.Resources, fromModules...)}) } 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 } // short is a commit as a person refers to it. func short(commit string) string { if len(commit) > 8 { return commit[:8] } return commit } // statusCommand answers "did my change go out?". // // novox/hq ADR 0010 names losing that question as the real risk of replacing a pipeline with a // comparison: it is answerable today by opening a pipeline, and something has to replace that or // this is worse to live with whatever its other properties. // // The answer is not "a job succeeded". It is which modules the mesh has not built from what their // source now has, and which machines are running the old one. func statusCommand(ctx context.Context) error { inv, err := openInventory(ctx) if err != nil { return err } defer inv.Close() behind, err := inv.Behind(ctx) if err != nil { return err } if len(behind) == 0 { fmt.Println("every module with a source is built from what that source has") return nil } var names []string for m := range behind { names = append(names, m) } sort.Strings(names) fmt.Printf("%d module(s) behind their source:\n\n", len(behind)) for _, m := range names { from, err := inv.SourceOf(ctx, m) if err != nil { return err } fmt.Printf(" %-20s holds %s, source has %s\n", m, short(from.BuiltFrom), short(from.Head)) if nodes := behind[m]; len(nodes) > 0 { // The part somebody actually wants. A module being out of date is a fact about the // catalogue; machines running the old one is the thing with consequences. fmt.Printf(" %-20s running on %s\n", "", strings.Join(nodes, ", ")) } else { fmt.Printf(" %-20s assigned to nothing\n", "") } } return nil } // parseAround reads flags that may sit before, after or between positional arguments. // // The standard library stops at the first non-flag argument, so `module add thing.json --source x` // parses no flags at all and silently ignores every one of them. The host learned this the same // way and says so in its own parser: a flag that is quietly dropped is the fault this project // keeps naming, and it looks exactly like success. func parseAround(set *flag.FlagSet, args []string) ([]string, error) { var positionals []string rest := args for { if err := set.Parse(rest); err != nil { return nil, err } rest = set.Args() if len(rest) == 0 { return positionals, nil } positionals = append(positionals, rest[0]) rest = rest[1:] } } func settingsCommand(ctx context.Context, args []string) error { if len(args) == 0 { return errors.New("settings set [--node ], or settings clear [--node ]") } inv, err := openInventory(ctx) if err != nil { return err } defer inv.Close() set := flag.NewFlagSet("settings", flag.ContinueOnError) node := set.String("node", "", "one machine, rather than the whole mesh") positionals, err := parseAround(set, args[1:]) if err != nil { return err } where := "the whole mesh" if *node != "" { where = *node } switch args[0] { case "set": if len(positionals) != 2 { return errors.New("settings set [--node ]") } raw, err := os.ReadFile(positionals[1]) if err != nil { return err } var values map[string]any if err := json.Unmarshal(raw, &values); err != nil { return fmt.Errorf("%s is not a settings file: %w", positionals[1], err) } if err := inv.SetSettings(ctx, *node, positionals[0], values); err != nil { return err } var keys []string for k := range values { keys = append(keys, k) } sort.Strings(keys) fmt.Printf("%s on %s: %s\n", positionals[0], where, strings.Join(keys, ", ")) fmt.Println(" run `push` to send it") return nil case "clear": if len(positionals) != 1 { return errors.New("settings clear [--node ]") } if err := inv.ClearSettings(ctx, *node, positionals[0]); err != nil { return err } fmt.Printf("%s on %s is back to what the module says\n", positionals[0], where) return nil default: return fmt.Errorf("settings has no %q; it has set and clear", args[0]) } }