Files
mesh-host/cmd/mesh-host/main.go
T

1004 lines
37 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"
"path/filepath"
"sort"
"strings"
"sync"
"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/firewall"
"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/reachable"
"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
}
// An adopted node keeps the firewall it was found with (novox/hq ADR 0100), so a host that
// cannot speak that firewall must say so now — before the mesh records a node it could never
// open anything on.
if token.Adopted {
kind, name, err := firewall.Detect(ctx, apply.ExecRunner)
if err != nil {
return err
}
if kind == firewall.Unsupported {
return fmt.Errorf(
"this token joins this machine adopted, keeping the firewall found on it, and it is "+
"filtered by %s, which no host speaks yet. Nothing was enrolled", name)
}
fmt.Printf("joining adopted: what is on this machine is kept, and its firewall (%s) stays in force\n",
string(kind))
}
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)
}
// Signed with the identity just generated, so the mesh can tell this machine from anyone else
// who knows its public key (novox/hq issue 083).
proof := mine.Sign(link.EnrolProof(token.Secret, mine.Public, mine.Overlay.Public,
sealing.Public, serving.Public))
reply, err := link.Enrol(ctx, token.Broker, token.Fingerprint, *name, token.Secret,
mine.Public, mine.Overlay.Public, sealing.Public, serving.Public, reported, proof, 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).
// Reports a reconcile has to make unasked — what an adopted node holds changed, or its
// firewall did — go out over the link when it is up (novox/hq ADR 0100).
outbox := make(chan link.Unasked, 1)
watch := &adoptionWatch{}
applier = watch.noting(applier)
go holdTheMachine(ctx, opts, mine, say, sched, func(r link.Report) {
if !watch.differs(r) {
return
}
select {
case <-outbox:
// An older one nobody has published yet; this one says everything it did.
default:
}
// Counted as said only once the broker has taken it: queued and lost — the link down, the
// publish refused — the change would never be said again (novox/hq ADR 0100).
outbox <- link.Unasked{Report: r, Done: func(published bool) {
if published {
watch.said(r)
}
}}
})
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), outbox)
}
// adoptionWatch remembers what the node last said about what it holds and its firewall, so a
// reconcile speaks unasked only when that changed.
type adoptionWatch struct {
mu sync.Mutex
last string
}
// fingerprint is what a report says about adoption: each hold and whether it changed, the
// firewall, and what is reachable on the machine — which only an adopted node reports, and which
// is what the controller previews a flip from, so a port that opens or closes between deliveries
// must reach it too (novox/hq ADR 0100).
func adoptionFingerprint(r link.Report) string {
parts := []string{"firewall=" + r.Firewall}
for _, h := range r.Held {
parts = append(parts, "held "+h.ID+"="+h.Changed)
}
for _, reach := range r.Reachable {
parts = append(parts, fmt.Sprintf("reach %s %s:%d %s %v %d", reach.Protocol, reach.Address,
reach.Port, reach.By, reach.Published, reach.ContainerPort))
}
sort.Strings(parts[1:])
return strings.Join(parts, "\n")
}
// differs says whether a report says anything the last one that went out did not. It records
// nothing: what was said is what reached the mesh, not what was written down to send.
func (w *adoptionWatch) differs(r link.Report) bool {
w.mu.Lock()
defer w.mu.Unlock()
return adoptionFingerprint(r) != w.last
}
// said records a report the mesh has actually been told.
func (w *adoptionWatch) said(r link.Report) {
w.mu.Lock()
defer w.mu.Unlock()
w.last = adoptionFingerprint(r)
}
// changed is differs and said together, for a report published as it is made.
func (w *adoptionWatch) changed(r link.Report) bool {
if !w.differs(r) {
return false
}
w.said(r)
return true
}
// noting wraps the applier, so a report the link publishes after a delivery counts as said.
func (w *adoptionWatch) noting(apply link.Applier) link.Applier {
return func(ctx context.Context, raw, signature []byte) link.Report {
r := apply(ctx, raw, signature)
w.changed(r)
return r
}
}
// 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, publish func(link.Report)) {
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)
// A reconcile is otherwise silent. On an adopted node it speaks when what it holds or
// its firewall changed, because that is how a predecessor still writing is caught
// (novox/hq ADR 0100); publish decides whether anything did.
if publish != nil && report.Refused == "" && (len(report.Held) > 0 || report.Firewall != "") {
publish(report)
}
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)
}
// applying serialises applies on this node.
//
// **Two things apply here: the link and the reconcile loop**, and each reads the node's state,
// acts on the machine, and writes the state back. Run at the same time they interleave, and the
// one that saves last writes a state read before the other acted — losing what the first recorded:
// a hold, the firewall found here, a resource just applied. The machine would then be one thing
// and its record another, which is the fault every read-back in this package exists to prevent.
var applying sync.Mutex
// applyAndKeep applies a declaration and, when it came from the mesh, keeps it so this node can
// go on obeying it while disconnected. One at a time, whoever asks.
func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.Declared,
sched *apply.Scheduler) link.Report {
applying.Lock()
defer applying.Unlock()
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.ApplyKeeping(ctx, built, declared, known, store.OriginDeclared,
apply.ExecRunner, nil, sealOpener(opts.state), apply.KeepIn(filepath.Dir(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 {
held := map[string]bool{}
for _, h := range updated.Held {
held[h.ID] = true
}
sched.Sync(declared, held)
}
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw)}
// What this node found and holds, its firewall, and what is reachable on it — so an adopted
// node never reads as converged (novox/hq ADR 0100).
for _, h := range updated.Held {
report.Held = append(report.Held, link.Held{ID: h.ID, Module: h.Module, Kind: h.Kind,
Target: h.Target, Since: h.Since, Changed: h.Changed, Kept: h.Kept})
}
if declared.Adoption != nil {
if updated.Firewall != nil {
report.Firewall = updated.Firewall.Kind
}
reached, err := reachable.Collect(ctx, apply.ExecRunner)
if err != nil {
fmt.Fprintf(os.Stderr, "mesh-host: applied, and could not read what is reachable here: %v\n", err)
}
report.Reachable = reached
}
for _, change := range outcome.Outcomes {
// What is held is not what this machine owns: it was found, and is kept as it was until
// its module is taken (novox/hq ADR 0100).
if change.Action == "held" {
continue
}
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 foundation 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
}