The recurring twin of run-once, one modifier over: a container marked schedule: "<cron>" is run to completion on its cadence, not started as a service and not run once as a gate. The gating rule is deliberately reversed. Installing a schedule records it as present state and reports the node current at once (applySchedule) -- it never runs the container and does not gate what follows. A Scheduler, held for the life of the daemon and re-established from each applied declaration (the declaration is the source of truth, ADR 0018), fires the container off an injected clock. A run that exits non-zero is logged and never fails the apply or flips the node's state, because it happens outside the apply and the store entirely. Runs never stack: a run still going when the next is due is skipped, not started as a second copy. No new host shape and no new action -- schedule is a string on the container the host already has, and the host process runs the container itself rather than installing a system timer (the rejected option 1). A minimal five-field cron (declaration/cron.go) validates on arrival and computes the next due minute; time is injected so the scheduler is tested without the wall clock. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
858 lines
31 KiB
Go
858 lines
31 KiB
Go
// Command mesh-host is tier 0 of the Novox Mesh: the one thing installed by hand, and the
|
|
// only thing that changes a machine.
|
|
//
|
|
// Stage 1 (novox/hq 03-DESIGN/01-to-be/05-the-node-host.md) is profile and inventory only —
|
|
// the host reads what this machine can do and what it is, and reports it. It applies nothing,
|
|
// connects to nothing, and listens on nothing.
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/base64"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"os"
|
|
"os/signal"
|
|
"sort"
|
|
"strings"
|
|
"syscall"
|
|
"text/tabwriter"
|
|
"time"
|
|
|
|
"github.com/novox/mesh-host/internal/apply"
|
|
"github.com/novox/mesh-host/internal/bundle"
|
|
"github.com/novox/mesh-host/internal/declaration"
|
|
"github.com/novox/mesh-host/internal/identity"
|
|
"github.com/novox/mesh-host/internal/inventory"
|
|
"github.com/novox/mesh-host/internal/link"
|
|
"github.com/novox/mesh-host/internal/profile"
|
|
"github.com/novox/mesh-host/internal/store"
|
|
"github.com/novox/mesh-host/internal/system"
|
|
"github.com/novox/mesh-host/internal/upgrade"
|
|
)
|
|
|
|
// version is stamped at build time. Unset in a development build, and said so rather than
|
|
// defaulted to something that looks like a release.
|
|
// builtFor names the operating system this host was built for, set at link time
|
|
// (novox/hq ADR 0005). A host built without one refuses to do anything that touches the
|
|
// machine, rather than guessing and calling a package manager that is not there.
|
|
var builtFor = ""
|
|
|
|
var version = "development build"
|
|
|
|
const usage = `mesh-host — the node host
|
|
|
|
profile what this machine can be asked to do
|
|
inventory what this machine is, and what it holds
|
|
apply FILE make this machine match a declaration from a file
|
|
reconcile make this machine match the declaration this host carries
|
|
bundle show what this host carries
|
|
owned what this host has applied and still owns
|
|
version
|
|
|
|
--json machine-readable output
|
|
--timeout how long any single probe may take (default 10s)
|
|
--state where this node keeps what it knows (default /var/lib/mesh-host/state.json)
|
|
--dry-run read the declaration and refuse it if wrong, but change nothing
|
|
|
|
It connects to nothing and listens on nothing. What it applies comes from a file.
|
|
`
|
|
|
|
func main() {
|
|
// A probe runs a command on a real machine. Ctrl-C must stop the host, not be swallowed by
|
|
// whatever it is waiting for.
|
|
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
|
defer stop()
|
|
|
|
command, opts, err := parseArgs(os.Args[1:])
|
|
if err == nil {
|
|
err = run(ctx, command, opts)
|
|
}
|
|
if err != nil {
|
|
fmt.Fprintf(os.Stderr, "mesh-host: %v\n", err)
|
|
os.Exit(1)
|
|
}
|
|
}
|
|
|
|
type options struct {
|
|
json bool
|
|
timeout time.Duration
|
|
state string
|
|
token string
|
|
nodeName string
|
|
dryRun bool
|
|
file string
|
|
}
|
|
|
|
// parseArgs takes the subcommand first, then its flags.
|
|
//
|
|
// The standard library stops parsing at the first non-flag argument, so `mesh-host inventory
|
|
// --json` left `--json` sitting in the positional arguments and printed text — a flag the user
|
|
// passed, silently ignored, with a successful exit. That is the fault this whole project keeps
|
|
// naming, so the parser takes the subcommand off the front and parses what follows.
|
|
func parseArgs(args []string) (string, options, error) {
|
|
opts := options{timeout: 10 * time.Second, state: store.DefaultPath}
|
|
|
|
command := ""
|
|
if len(args) > 0 {
|
|
command = args[0]
|
|
args = args[1:]
|
|
}
|
|
|
|
set := flag.NewFlagSet("mesh-host", flag.ContinueOnError)
|
|
set.SetOutput(os.Stderr)
|
|
set.Usage = func() { fmt.Fprint(os.Stderr, usage) }
|
|
set.BoolVar(&opts.json, "json", false, "machine-readable output")
|
|
set.DurationVar(&opts.timeout, "timeout", opts.timeout, "how long any single probe may take")
|
|
set.StringVar(&opts.state, "state", opts.state, "where this node keeps what it knows")
|
|
set.BoolVar(&opts.dryRun, "dry-run", false, "read and check the declaration, change nothing")
|
|
set.StringVar(&opts.token, "token", "", "enrol: the one-time token, carried here by a person")
|
|
set.StringVar(&opts.nodeName, "name", "", "enrol: override the name the token carries")
|
|
|
|
// Parsed in a loop, because the standard library stops at the FIRST non-flag argument.
|
|
// `mesh-host inventory --json` hit that once, and taking the subcommand off the front
|
|
// fixed only half of it: `mesh-host apply decl.json --dry-run` left --dry-run unread in
|
|
// exactly the same way. A flag may sit before, after or between positionals, and one that
|
|
// is silently dropped is the fault this whole project keeps naming.
|
|
var positionals []string
|
|
rest := args
|
|
for {
|
|
if err := set.Parse(rest); err != nil {
|
|
return "", opts, err
|
|
}
|
|
rest = set.Args()
|
|
if len(rest) == 0 {
|
|
break
|
|
}
|
|
positionals = append(positionals, rest[0])
|
|
rest = rest[1:]
|
|
}
|
|
|
|
if command == "apply" {
|
|
if len(positionals) != 1 {
|
|
return "", opts, errors.New("apply needs exactly one declaration file")
|
|
}
|
|
opts.file = positionals[0]
|
|
return command, opts, nil
|
|
}
|
|
// Anything left over was neither the command nor a flag. Refused rather than ignored: a
|
|
// mistyped argument that changes nothing and reports success is worse than an error.
|
|
if len(positionals) > 0 {
|
|
return "", opts, fmt.Errorf("unexpected argument %q — try `mesh-host help`", positionals[0])
|
|
}
|
|
return command, opts, nil
|
|
}
|
|
|
|
func run(ctx context.Context, command string, opts options) error {
|
|
jsonOut, timeout := opts.json, opts.timeout
|
|
switch command {
|
|
case "profile":
|
|
p := profile.Detect(ctx, profile.Default(nil), timeout)
|
|
if jsonOut {
|
|
return writeJSON(p)
|
|
}
|
|
writeProfile(p)
|
|
return nil
|
|
|
|
case "inventory":
|
|
inv := inventory.Collect(ctx, nil, profile.Default(nil), timeout)
|
|
if jsonOut {
|
|
return writeJSON(inv)
|
|
}
|
|
writeInventory(inv)
|
|
return nil
|
|
|
|
case "apply":
|
|
raw, err := os.ReadFile(opts.file)
|
|
if err != nil {
|
|
return fmt.Errorf("reading the declaration: %w", err)
|
|
}
|
|
// ParseTrusted: a file handed to the host by someone already running it as root is
|
|
// not the link. novox/hq ADR 0005 bounds what a REMOTE party may push; someone who
|
|
// can write this file and run this binary can do anything the binary can, so refusing
|
|
// them an action would buy nothing and would make an action untestable except by
|
|
// rebuilding the bundle.
|
|
d, err := declaration.ParseFileTrusted(raw)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return runApply(ctx, opts, d, opts.file)
|
|
|
|
case "reconcile":
|
|
// The first node's path. novox/hq ADR 0004: no mesh reachable means the declaration
|
|
// comes from the bundle the host carries. There is no link yet, so this is currently
|
|
// the only source — which is a stage, not a design, and saying so beats implying the
|
|
// other source exists.
|
|
d, err := bundle.Load(builtFor)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return runApply(ctx, opts, d, "the carried bundle")
|
|
|
|
case "bundle":
|
|
if bundle.IsEmpty(builtFor) {
|
|
fmt.Println("this host carries no bundle")
|
|
return nil
|
|
}
|
|
_, err := bundle.Load(builtFor)
|
|
if err != nil {
|
|
// Asked before it matters, rather than discovered on a first node.
|
|
return fmt.Errorf("this host carries a bundle it cannot itself read: %w", err)
|
|
}
|
|
os.Stdout.Write(bundle.Raw(builtFor))
|
|
return nil
|
|
|
|
case "owned":
|
|
known, err := store.Load(opts.state)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if jsonOut {
|
|
return writeJSON(known)
|
|
}
|
|
if len(known.Resources) == 0 {
|
|
fmt.Println("this host has applied nothing on this machine")
|
|
return nil
|
|
}
|
|
w := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0)
|
|
for _, r := range known.Resources {
|
|
fmt.Fprintf(w, " %s\t%s\t%s\n", r.Type, r.ID, r.Target)
|
|
}
|
|
return w.Flush()
|
|
|
|
case "enrol", "enroll":
|
|
return enrol(ctx, opts)
|
|
|
|
case "run":
|
|
return runLink(ctx, opts)
|
|
|
|
case "version":
|
|
fmt.Println(version)
|
|
return nil
|
|
|
|
case "", "help", "-h", "--help":
|
|
fmt.Fprint(os.Stderr, usage)
|
|
return nil
|
|
|
|
default:
|
|
return fmt.Errorf("unknown command %q — try `mesh-host help`", command)
|
|
}
|
|
}
|
|
|
|
func writeJSON(v any) error {
|
|
enc := json.NewEncoder(os.Stdout)
|
|
enc.SetIndent("", " ")
|
|
return enc.Encode(v)
|
|
}
|
|
|
|
// writeProfile prints every verdict WITH its reason.
|
|
//
|
|
// The reason is not decoration: a capability reported absent with no reason is something
|
|
// nobody can act on, and this is the surface where a person meets that.
|
|
func writeProfile(p profile.Profile) {
|
|
fmt.Printf("%s/%s\n\n", p.Kernel, p.Architecture)
|
|
|
|
w := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0)
|
|
for _, v := range p.Capabilities {
|
|
mark := "no "
|
|
if v.Present {
|
|
mark = "yes"
|
|
}
|
|
fmt.Fprintf(w, " %s\t%s\t%s\n", mark, v.Name, v.Detail)
|
|
}
|
|
w.Flush()
|
|
|
|
if missing := p.Missing(); len(missing) > 0 {
|
|
fmt.Printf("\ncannot be asked to: %v\n", missing)
|
|
}
|
|
}
|
|
|
|
func writeInventory(inv inventory.Inventory) {
|
|
w := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0)
|
|
fmt.Fprintf(w, "machine\t%s\n", inv.Machine)
|
|
if inv.Distribution != "" {
|
|
fmt.Fprintf(w, "distribution\t%s\n", inv.Distribution)
|
|
}
|
|
if inv.Kernel != "" {
|
|
fmt.Fprintf(w, "kernel\t%s\n", inv.Kernel)
|
|
}
|
|
fmt.Fprintf(w, "architecture\t%s/%s\n", inv.OS, inv.Architecture)
|
|
fmt.Fprintf(w, "cpus\t%d\n", inv.CPUs)
|
|
if inv.MemoryKB > 0 {
|
|
fmt.Fprintf(w, "memory\t%d MB\n", inv.MemoryKB/1024)
|
|
}
|
|
fmt.Fprintf(w, "observed\t%s\n", inv.ObservedAt.Format(time.RFC3339))
|
|
w.Flush()
|
|
|
|
fmt.Println()
|
|
writeProfile(inv.Profile)
|
|
|
|
// Printed last and never hidden. An inventory that quietly omits what it could not read
|
|
// is the same fault as a report assembled from intent (novox/hq ADR 0018).
|
|
if len(inv.Unreadable) > 0 {
|
|
fmt.Println("\ncould not read:")
|
|
for _, u := range inv.Unreadable {
|
|
fmt.Printf(" %s\n", u)
|
|
}
|
|
}
|
|
}
|
|
|
|
// runApply reads a declaration and makes the machine match it.
|
|
//
|
|
// The state is loaded before anything is touched and saved after, including when the apply
|
|
// fails part-way: what was applied before the failure is on the machine, and a host that did
|
|
// not record it would believe it owns less than it does and leave that behind forever.
|
|
func runApply(ctx context.Context, opts options, d *declaration.Declaration, source string) error {
|
|
known, err := store.Load(opts.state)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if opts.dryRun {
|
|
fmt.Printf("%s: %d resource(s), version %d — accepted, nothing applied\n",
|
|
source, len(d.Resources), d.Version)
|
|
return nil
|
|
}
|
|
|
|
sys, err := system.For(builtFor)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// Refuse a declaration naming a shape this host cannot apply, before anything is applied.
|
|
// An android host has no `package` applier, and finding that out half way through is the
|
|
// half-configured machine this host exists to prevent (novox/hq ADR 0005).
|
|
if err := system.Check(sys, d); err != nil {
|
|
return err
|
|
}
|
|
// And prove this is the machine the host was built for. Installing the arch host on Alpine
|
|
// must say so once, at the start, rather than failing later inside pacman.
|
|
if err := sys.Confirm(ctx, apply.ExecRunner); err != nil {
|
|
return err
|
|
}
|
|
|
|
report, updated, applyErr := apply.Apply(ctx, sys, d, known, store.OriginCarried,
|
|
apply.ExecRunner, func(line string) {
|
|
if !opts.json {
|
|
fmt.Println(line)
|
|
}
|
|
}, sealOpener(opts.state))
|
|
|
|
// Saved whichever way it went. Recording only on success would lose the footprint of a
|
|
// failed apply, and that footprint is on the machine either way.
|
|
if saveErr := store.Save(opts.state, updated); saveErr != nil {
|
|
if applyErr != nil {
|
|
return fmt.Errorf("%w\n\nand the node's state could not be saved: %v", applyErr, saveErr)
|
|
}
|
|
return saveErr
|
|
}
|
|
if applyErr != nil {
|
|
return applyErr
|
|
}
|
|
|
|
// Only now, and only after a clean apply: this version got as far as a completed
|
|
// reconcile, which is the whole of what "known good" claims (novox/hq ADR 0005). Not
|
|
// health — a disconnected node is ordinary, and a resource that fails is the machine's
|
|
// problem rather than the binary's.
|
|
//
|
|
// A failure to record is reported and does not fail the apply. The apply worked; what is
|
|
// lost is a rollback's ability to come back here, which is worse to hide than to say.
|
|
if version != "" {
|
|
if err := upgrade.RecordKnownGood(upgrade.KnownGoodPath(opts.state), version); err != nil {
|
|
fmt.Fprintf(os.Stderr,
|
|
"mesh-host: applied, but could not record %s as known-good: %v\n"+
|
|
" a rollback would have nothing to return to.\n", version, err)
|
|
}
|
|
}
|
|
// And tell the launcher this start worked. Without it the counter only climbs, and a node
|
|
// that has been up for months rolls itself back on its third ordinary restart.
|
|
if err := upgrade.ClearAttempts(upgrade.AttemptsPath(opts.state)); err != nil {
|
|
fmt.Fprintf(os.Stderr,
|
|
"mesh-host: applied, but could not clear the start counter: %v\n"+
|
|
" this node may roll itself back after a few more restarts.\n", err)
|
|
}
|
|
|
|
if opts.json {
|
|
return writeJSON(report)
|
|
}
|
|
if !report.Changed() {
|
|
fmt.Printf("%s: already matches — %d resource(s) checked\n", source, len(report.Outcomes))
|
|
return nil
|
|
}
|
|
fmt.Printf("%s: applied — %d resource(s)\n", source, len(report.Outcomes))
|
|
return nil
|
|
}
|
|
|
|
// enrol joins this machine to a mesh.
|
|
//
|
|
// novox/hq 09-the-node-lifecycle: the token carries four things, the node dials the broker over
|
|
// the underlay, checks the certificate against the pin *before sending anything*, and presents
|
|
// the one-time secret together with a public key it generated itself.
|
|
//
|
|
// The mesh issues no identity. This machine arrives holding one; what it receives is being known.
|
|
func enrol(ctx context.Context, opts options) error {
|
|
tokenText, name := &opts.token, &opts.nodeName
|
|
if strings.TrimSpace(*tokenText) == "" {
|
|
return errors.New("enrol --token <token>: the token is carried to this machine by a " +
|
|
"person, and is the only thing it needs")
|
|
}
|
|
|
|
// Refused whole if incomplete. A token without the fingerprint would have this machine
|
|
// connect to whatever answers; without the signing key it could not tell a declaration from
|
|
// a forgery, and it applies whatever the link delivers.
|
|
token, err := identity.ParseToken(*tokenText)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// The name comes from the token, because the node cannot work it out: the broker account it
|
|
// authenticates as is named after it, and that account exists before this machine has been
|
|
// told anything. --name remains for a token issued before the name travelled in one, and
|
|
// saying so beats a connection refused with an empty username — which is what this was.
|
|
if strings.TrimSpace(*name) == "" {
|
|
*name = token.Node
|
|
}
|
|
if strings.TrimSpace(*name) == "" {
|
|
return errors.New(
|
|
"this token does not say what the mesh calls this machine, and no --name was given. " +
|
|
"A token issued by a current control plane carries the name")
|
|
}
|
|
|
|
// Before anything else: an already-enrolled machine must not quietly acquire a second
|
|
// identity. The mesh believes the first one, and re-enrolling is a deliberate act that
|
|
// starts with a person issuing a new token for that node record.
|
|
identityPath := identity.Path(opts.state)
|
|
switch existing, err := identity.Load(identityPath); {
|
|
case err == nil:
|
|
return fmt.Errorf(
|
|
"this machine is already node %q. Re-enrolling replaces the identity the mesh "+
|
|
"believes, so it is done deliberately: remove %s first",
|
|
existing.Node, identityPath)
|
|
case errors.Is(err, identity.ErrNoIdentity):
|
|
default:
|
|
return err
|
|
}
|
|
|
|
fmt.Printf("token for broker %s\n", token.Broker)
|
|
fmt.Printf(" pinned certificate %s\n", token.Fingerprint)
|
|
fmt.Printf(" signing key %s\n",
|
|
base64.StdEncoding.EncodeToString(token.Signer)[:16]+"...")
|
|
|
|
// The check that has to happen before this machine says anything.
|
|
conn, err := link.Dial(token.Broker, token.Fingerprint, opts.timeout)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer conn.Close()
|
|
fmt.Println("\nthe broker presented the certificate this token pins")
|
|
|
|
conn.Close()
|
|
|
|
mine, err := identity.Generate(*name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("generated this node's identity: %s\n", mine.PublicBase64())
|
|
|
|
// Its key on the private network, generated here and now for the same reason: the private
|
|
// half must never have been anywhere else. The mesh receives only the public half and uses it
|
|
// to compute a graph it cannot impersonate.
|
|
mine.Overlay, err = identity.GenerateOverlayKey()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("generated this node's overlay key: %s\n", mine.Overlay.Public)
|
|
|
|
// And the key secrets are sealed to. Here, with the others, because the mesh cannot seal
|
|
// anything to a key it has not been told about — a key made later would leave a node that
|
|
// looks enrolled and can receive no credential.
|
|
sealing, err := identity.GenerateSealingKey()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("generated this node's sealing key: %s\n", sealing.Public)
|
|
|
|
// And the key it serves TLS with on its name inside the mesh. Generated here for the same
|
|
// reason as the others: the private half must never have been anywhere else, and the mesh
|
|
// only ever certifies the public one.
|
|
serving, err := identity.GenerateServingKey()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
fmt.Printf("generated this node's serving key: %s\n", serving.Public)
|
|
|
|
// What this machine can be asked to do, gathered before joining rather than after. The
|
|
// control plane cannot decide what a node should run without it, so it travels with the
|
|
// request instead of being asked for in a second round trip.
|
|
detected := profile.Detect(ctx, profile.Default(nil), opts.timeout)
|
|
reported := map[string]any{}
|
|
if raw, err := json.Marshal(detected); err == nil {
|
|
_ = json.Unmarshal(raw, &reported)
|
|
}
|
|
|
|
reply, err := link.Enrol(ctx, token.Broker, token.Fingerprint, *name, token.Secret,
|
|
mine.Public, mine.Overlay.Public, sealing.Public, serving.Public, reported, opts.timeout)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// The mesh's name for this node wins over what the machine called itself: the token was
|
|
// issued for a node record, and that record is what the identity binds to.
|
|
mine.Node = reply.Node
|
|
mine.Membership = identity.Membership{
|
|
Broker: firstNonEmpty(reply.Broker, token.Broker),
|
|
Fingerprint: firstNonEmpty(reply.Fingerprint, token.Fingerprint),
|
|
Signer: firstNonEmpty2(reply.Signer, token.Signer),
|
|
Password: reply.Password,
|
|
}
|
|
if mine.Membership.Password == "" {
|
|
// The mesh did not replace the token's secret, so it is still this node's broker
|
|
// password. Said rather than silently kept: a one-time secret living on as a credential
|
|
// is worth knowing about.
|
|
mine.Membership.Password = token.Secret
|
|
fmt.Println("\nnote: the mesh issued no separate broker password, so the token's secret " +
|
|
"remains this node's credential")
|
|
}
|
|
|
|
// Saved only now, and only once the mesh has said it knows this node. A node holding an
|
|
// identity the mesh has never recorded would believe it had joined and be believed by
|
|
// nobody — worse than not having joined, because nothing would look wrong.
|
|
if err := identity.Save(identityPath, mine); err != nil {
|
|
return fmt.Errorf(
|
|
"the mesh accepted this node as %q and its identity could not be saved: %w\n"+
|
|
"That token is spent, so getting back needs a new one", reply.Node, err)
|
|
}
|
|
|
|
if !mine.Membership.Joined() {
|
|
return fmt.Errorf(
|
|
"the mesh accepted this node as %q but did not say how to reach it again, so this "+
|
|
"identity could not be used after a restart. Nothing was saved", reply.Node)
|
|
}
|
|
|
|
// Written before the identity, so a node that dies between the two has a key file with no
|
|
// identity — which enrols again cleanly — rather than an identity naming a key that is not
|
|
// there, which looks joined and cannot come up.
|
|
if err := os.WriteFile(identity.OverlayKeyPath(opts.state),
|
|
[]byte(mine.Overlay.Private+"\n"), 0o600); err != nil {
|
|
return fmt.Errorf("cannot write this node's overlay key: %w", err)
|
|
}
|
|
if err := os.WriteFile(identity.SealingKeyPath(opts.state),
|
|
[]byte(sealing.Private+"\n"), 0o600); err != nil {
|
|
return fmt.Errorf("cannot write this node's sealing key: %w", err)
|
|
}
|
|
if err := identity.WriteServingKey(identity.ServingKeyPath(opts.state), serving); err != nil {
|
|
return fmt.Errorf("cannot write this node's serving key: %w", err)
|
|
}
|
|
|
|
fmt.Printf("\nenrolled as %s\n", reply.Node)
|
|
fmt.Printf(" identity %s\n", identityPath)
|
|
fmt.Printf(" queue %s\n", reply.Queue)
|
|
return nil
|
|
}
|
|
|
|
func firstNonEmpty(values ...string) string {
|
|
for _, v := range values {
|
|
if strings.TrimSpace(v) != "" {
|
|
return v
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func firstNonEmpty2(values ...[]byte) []byte {
|
|
for _, v := range values {
|
|
if len(v) > 0 {
|
|
return v
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// runLink holds this node's link to the mesh open, applying what arrives.
|
|
//
|
|
// One outbound connection and nothing listening. While it is up this node is enrolled; while it
|
|
// is down it is disconnected, which is an ordinary situation rather than a failure — the machine
|
|
// keeps running whatever it was last told, from its own store.
|
|
func runLink(ctx context.Context, opts options) error {
|
|
// A machine that has not enrolled waits here rather than failing. It is *hosted*: the host is
|
|
// running, it has no identity, and there is nobody to link to — an ordinary state, and the
|
|
// one every machine passes through (novox/hq 09-the-node-lifecycle).
|
|
//
|
|
// Exiting instead would be worse than untidy. The launcher counts a failed start, and three
|
|
// of them roll the binary back — so a freshly installed host, waiting to be enrolled exactly
|
|
// as intended, would undo its own installation.
|
|
mine, err := waitForEnrolment(ctx, opts)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if mine.Node == "" {
|
|
return nil // asked to stop while waiting
|
|
}
|
|
|
|
fmt.Printf("node %s, linking to %s\n", mine.Node, mine.Membership.Broker)
|
|
|
|
say := func(line string) { fmt.Println(line) }
|
|
|
|
// One scheduler for the life of the process, re-established from each applied declaration
|
|
// (novox/hq ADR 0053). It fires scheduled steps on their cadence, surviving across applies and
|
|
// across the reconcile loop; a host restart rebuilds it from the declaration the node kept, the
|
|
// first time either path applies. Its own loop is the thing on the clock — no system timer.
|
|
sched := apply.NewScheduler(apply.SystemClock(), apply.ExecRunner, say)
|
|
go sched.Run(ctx)
|
|
|
|
applier := func(ctx context.Context, raw, signature []byte) link.Report {
|
|
return applyAndKeep(ctx, opts, raw, &store.Declared{Declaration: raw, Signature: signature}, sched)
|
|
}
|
|
|
|
// Two things at once, and the second is what makes disconnection ordinary. The link brings
|
|
// new declarations; this holds the machine in the last one whether the link is up or not. A
|
|
// laptop shut for a week comes back and reconciles — it does not come back and ask what it is
|
|
// (novox/hq ADR 0004).
|
|
go holdTheMachine(ctx, opts, mine, say, sched)
|
|
|
|
return link.HoldRoused(ctx, link.Membership{
|
|
Node: mine.Node,
|
|
Broker: mine.Membership.Broker,
|
|
Fingerprint: mine.Membership.Fingerprint,
|
|
Password: mine.Membership.Password,
|
|
Signer: mine.Membership.Signer,
|
|
}, applier, say, opts.timeout, rousedBySignal(ctx))
|
|
}
|
|
|
|
// rousedBySignal is the machine telling this process that its link is probably stale.
|
|
//
|
|
// **A signal, because nothing may listen on a node** (novox/hq ADR 0004). A socket for this would
|
|
// be a control surface on every machine, reachable by anything that can reach the machine, in
|
|
// exchange for saving twenty seconds — and the whole security argument rests on there not being
|
|
// one. A signal is delivered by the service manager to a process it already supervises.
|
|
//
|
|
// SIGHUP, because that is the signal a long-running program conventionally reads as *look again*,
|
|
// and nothing here is being reloaded from a file that a different signal would suit better.
|
|
//
|
|
// Dropped rather than queued when one arrives while another is unread: two wakes in the same
|
|
// instant are one wake, and a machine that suspends and resumes repeatedly must not build a
|
|
// backlog of reconnections to work through.
|
|
func rousedBySignal(ctx context.Context) link.Roused {
|
|
woken := make(chan os.Signal, 1)
|
|
signal.Notify(woken, syscall.SIGHUP)
|
|
|
|
out := make(chan struct{}, 1)
|
|
go func() {
|
|
defer signal.Stop(woken)
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-woken:
|
|
select {
|
|
case out <- struct{}{}:
|
|
default:
|
|
// One is already waiting to be read. Two wakes in the same instant are one.
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
return out
|
|
}
|
|
|
|
// ReconcileEvery is how often a node re-applies what it was last told.
|
|
//
|
|
// Not driven by the link. Changes are pushed, so this is not polling for them — it is the answer
|
|
// to a machine drifting: a file edited by hand, a container stopped by somebody, a service that
|
|
// died. A node that only acted when told would hold its state exactly until something else
|
|
// changed it, and then for ever.
|
|
const ReconcileEvery = 5 * time.Minute
|
|
|
|
func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, say link.Announce,
|
|
sched *apply.Scheduler) {
|
|
ticker := time.NewTicker(ReconcileEvery)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
}
|
|
|
|
declared, err := store.LoadDeclared(store.DeclaredPath(opts.state), mine.Membership.Signer)
|
|
if errors.Is(err, store.ErrNothingDeclared) {
|
|
// Nothing to hold this machine to yet. Ordinary on a node that has enrolled and not
|
|
// been assigned anything.
|
|
continue
|
|
}
|
|
if err != nil {
|
|
say("cannot re-apply what this node was told: " + err.Error())
|
|
continue
|
|
}
|
|
|
|
report := applyDeclared(ctx, opts, declared, sched)
|
|
switch {
|
|
case report.Refused != "":
|
|
say("what this node was last told no longer applies: " + report.Refused)
|
|
case len(report.Failed) > 0:
|
|
say(fmt.Sprintf("holding this machine: %d applied, and %v", len(report.Applied), report.Failed))
|
|
}
|
|
}
|
|
}
|
|
|
|
// applyDeclared applies a declaration that has already been proved to come from the mesh.
|
|
//
|
|
// Signature checking happens before this is called, in the link. By the time anything here runs,
|
|
// the question "is this from the mesh I joined" is settled — which is why this can treat the
|
|
// bytes as instructions.
|
|
func applyDeclared(ctx context.Context, opts options, raw []byte, sched *apply.Scheduler) link.Report {
|
|
return applyAndKeep(ctx, opts, raw, nil, sched)
|
|
}
|
|
|
|
// applyAndKeep applies a declaration and, when it came from the mesh, keeps it so this node can
|
|
// go on obeying it while disconnected.
|
|
func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.Declared,
|
|
sched *apply.Scheduler) link.Report {
|
|
declared, err := declaration.Parse(raw)
|
|
if err != nil {
|
|
return link.Report{Refused: err.Error()}
|
|
}
|
|
|
|
built, err := system.For(builtFor)
|
|
if err != nil {
|
|
return link.Report{Refused: err.Error()}
|
|
}
|
|
if err := system.Check(built, declared); err != nil {
|
|
return link.Report{Refused: err.Error()}
|
|
}
|
|
|
|
known, err := store.Load(opts.state)
|
|
if err != nil {
|
|
return link.Report{Refused: err.Error()}
|
|
}
|
|
|
|
if err := built.Confirm(ctx, apply.ExecRunner); err != nil {
|
|
return link.Report{Refused: err.Error()}
|
|
}
|
|
|
|
// Declared, not carried. A declaration from the mesh removes only what the mesh previously
|
|
// declared — never what this machine raised for itself from its bundle (04-ISSUES/010).
|
|
outcome, updated, applyErr := apply.Apply(ctx, built, declared, known, store.OriginDeclared,
|
|
apply.ExecRunner, nil, sealOpener(opts.state))
|
|
|
|
// Saved whichever way it went. Recording only on success would lose the footprint of a
|
|
// failed apply, and that footprint is on the machine either way.
|
|
if saveErr := store.Save(opts.state, updated); saveErr != nil {
|
|
return link.Report{Refused: "applied, and the node's state could not be saved: " +
|
|
saveErr.Error()}
|
|
}
|
|
|
|
// Re-establish the scheduled steps from the declaration just applied (novox/hq ADR 0053). Done
|
|
// from the parsed declaration, which is the source of truth (ADR 0018): a schedule newly declared
|
|
// is armed, one whose image, environment or cadence changed is re-armed, and one no longer
|
|
// declared is forgotten — and after a host restart the first apply rebuilds them all. A nil
|
|
// scheduler is the one-shot CLI path, which exits rather than staying up to fire anything.
|
|
if sched != nil {
|
|
sched.Sync(declared)
|
|
}
|
|
|
|
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw)}
|
|
for _, change := range outcome.Outcomes {
|
|
report.Applied = append(report.Applied, change.ID)
|
|
}
|
|
// Kept whichever way it went, so a node that is disconnected next minute still knows what it
|
|
// was told. Saved after applying rather than before: what is kept is what this node acted on.
|
|
if signed != nil {
|
|
if err := store.SaveDeclared(store.DeclaredPath(opts.state), *signed); err != nil {
|
|
report.Failed = map[string]string{"keeping the declaration": err.Error()}
|
|
}
|
|
}
|
|
if applyErr != nil {
|
|
report.Failed = map[string]string{"apply": applyErr.Error()}
|
|
}
|
|
return report
|
|
}
|
|
|
|
// waitForEnrolment returns this node's identity, waiting for one if it has none.
|
|
//
|
|
// It applies the carried bundle first, if there is one, because that is what a first node does
|
|
// before there is a mesh at all — and a machine that has been installed and not yet enrolled
|
|
// should still be whatever its bundle says it is.
|
|
func waitForEnrolment(ctx context.Context, opts options) (identity.Identity, error) {
|
|
const look = 5 * time.Second
|
|
|
|
said := false
|
|
for {
|
|
mine, err := identity.Load(identity.Path(opts.state))
|
|
if err == nil {
|
|
return mine, nil
|
|
}
|
|
if !errors.Is(err, identity.ErrNoIdentity) {
|
|
// An identity that exists and cannot be read is a fault, not a wait. Treating it as
|
|
// "not enrolled yet" would leave a node sitting quietly for ever while the mesh
|
|
// believes it is a member.
|
|
return identity.Identity{}, err
|
|
}
|
|
|
|
if !said {
|
|
fmt.Println("this machine has not joined a mesh, and is waiting to be told which one.")
|
|
fmt.Println(" enrol it with: mesh-host enrol --token <token> --name <name>")
|
|
said = true
|
|
}
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
return identity.Identity{}, nil
|
|
case <-time.After(look):
|
|
}
|
|
}
|
|
}
|
|
|
|
// sealOpener is how a sealed file is opened.
|
|
//
|
|
// Looked up per file rather than held, because most declarations contain no sealed file at all
|
|
// and a node with no key must fail on the one that needs it rather than on every apply.
|
|
func sealOpener(statePath string) apply.Unseal {
|
|
return func(sealed string) ([]byte, error) {
|
|
key, err := identity.LoadSealingKey(identity.SealingKeyPath(statePath))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return key.Unseal(sealed)
|
|
}
|
|
}
|
|
|
|
// carriedPorts is every machine port held by what this host raised from its own bundle.
|
|
//
|
|
// **What the mesh must assign around** (novox/hq ADR 0038). The substrate is not a module: a node
|
|
// raises it before any mesh exists, so the control plane has never heard of the store or the
|
|
// broker. Told this, it can put a module somewhere else; not told, it hands out a port one of them
|
|
// holds and finds out from a container runtime.
|
|
//
|
|
// Only what was carried. What the mesh itself put here it already knows about, and reporting it
|
|
// back would make the machine an authority on the mesh's own bookkeeping.
|
|
// digestOf names a declaration by its bytes, exactly as the mesh names what it sends. The two
|
|
// sides never exchange the digest of different things: this hashes the same raw bytes the mesh
|
|
// hashed when it recorded the send.
|
|
func digestOf(body []byte) string {
|
|
sum := sha256.Sum256(body)
|
|
return hex.EncodeToString(sum[:])
|
|
}
|
|
|
|
func carriedPorts(state store.State) []int {
|
|
seen := map[int]bool{}
|
|
var out []int
|
|
for _, applied := range state.Resources {
|
|
if applied.Origin == store.OriginDeclared {
|
|
continue
|
|
}
|
|
for _, port := range applied.Holds {
|
|
if !seen[port] {
|
|
seen[port] = true
|
|
out = append(out, port)
|
|
}
|
|
}
|
|
}
|
|
sort.Ints(out)
|
|
return out
|
|
}
|