Files
mesh-controller/cmd/mesh-controller/adoption.go
T
jschoubben 4566c5c9aa Adopt the tunnel as a mesh fact, refuse a mismatched takeover, and rekey after enrolment
Review of the ADR 0105 build (hq ADR 0105). Four things it got wrong and one
path it lacked:

- A predecessor spoke's tunnel names one peer, the hub, routed the whole
  range; recording refused it and the whole enrolment failed. Range-routed
  peers are skipped now — only the hub's peers are ever carried.
- The range and the carried peers were conditions on the node being adopted,
  so converging the hub would have renumbered the mesh and dropped the peers
  still reaching it. They are facts of the tunnel record now, mode aside; the
  takeover alone is declared to an adopted node. Converging the hub is refused
  while a carried peer has not enrolled, naming it.
- A push composed a takeover for a hub whose address or endpoint disagreed
  with the tunnel, which would have the host stop the found interface and
  raise the mesh's where no peer listens. The graph refuses to compose it,
  naming both and the placement that fixes it.
- The host's account said taken or not; "found down and the mesh's not up"
  read as not taken. Three states now, and an account on every takeover.
- A hub that enrolled before this feature holds a key of its own, and
  re-enrolling would rotate every key the mesh sealed credentials to. A node
  now rekeys in a report, signed with its identity key over the key it
  leaves, the key it takes and the tunnel; the mesh verifies against the live
  key, refuses a stale or foreign proof, records key and tunnel, and moves a
  hub to the tunnel's address. `overlay show` names the path for a hub that
  found no tunnel.

Also: a carried IPv6 peer is routed /128, and identity.ForTest exists so the
link can be tested against a real identity store.
2026-09-24 00:02:07 +02:00

622 lines
22 KiB
Go

package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"flag"
"fmt"
"slices"
"sort"
"strings"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/overlay"
)
// A node is adopted or converged (novox/hq ADR 0100), and it is said to be adopted wherever the
// mesh reports a node's state: node list, node show, status and the board.
// showMode is the node show lines about a node's mode and what was taken on it.
func showMode(ctx context.Context, inv *inventory.Inventory, node inventory.Node) error {
if !node.Adopted {
fmt.Printf(" mode converged\n")
return nil
}
fmt.Printf(" mode adopted since %s\n",
node.AdoptedSince.Local().Format(time.DateTime))
taken, err := inv.Taken(ctx, node.Name)
if err != nil {
return err
}
if len(taken) == 0 {
fmt.Printf(" taken nothing yet\n")
} else {
fmt.Printf(" taken %s\n", strings.Join(taken, ", "))
}
said, err := inv.AdoptionOf(ctx, node.Name)
if err != nil {
return err
}
if said.At.IsZero() {
fmt.Printf(" it has not yet said what it found\n")
return nil
}
fmt.Printf(" firewall found %s\n", orNone(said.Firewall))
if err := showTunnel(ctx, inv, node.Name); err != nil {
return err
}
if len(said.Held) == 0 {
fmt.Printf(" holding nothing found\n")
}
for _, h := range said.Held {
// A held thing that changed is how a predecessor still writing is caught: said first.
line := fmt.Sprintf(" holds %-11s %s %s, for %s", h.Kind, h.Target, h.ID, h.Module)
if h.Changed != "" {
line += " — " + strings.ToUpper(h.Changed) + " by something other than the mesh"
}
fmt.Println(line)
if h.Kept != "" {
fmt.Printf(" %-17s original kept at %s\n", "", h.Kept)
}
}
fmt.Printf(" as of %s\n", said.At.Local().Format(time.DateTime))
return nil
}
// showTunnel is the node show lines about the tunnel an adopted node found and carried (novox/hq
// ADR 0105): what it presented at enrolment, and what it last said about taking it over.
func showTunnel(ctx context.Context, inv *inventory.Inventory, name string) error {
tunnel, err := inv.TunnelOf(ctx, name)
if errors.Is(err, inventory.ErrNoTunnel) {
return nil
}
if err != nil {
return err
}
fmt.Printf(" tunnel found %s on port %d, %s in %s, %d peer(s)\n",
tunnel.Interface, tunnel.Port, tunnel.Address, tunnel.Range, len(tunnel.Peers))
carried, said, err := inv.CarriedTunnelOf(ctx, name)
if err != nil {
return err
}
switch {
case !said:
fmt.Printf(" %-17s not yet taken over — the node has not said so\n", "")
case carried.State == inventory.CarriedTaken:
fmt.Printf(" %-17s taken over: %s is down and disabled, never flushed; the mesh's interface "+
"runs with its key, port and %d peer(s)\n", "", carried.Interface, carried.Peers)
case carried.State == inventory.CarriedDown:
fmt.Printf(" %-17s TUNNEL DOWN: %s is stopped and the mesh's interface is not up — the peers "+
"reach nothing. On the machine: systemctl start %s\n", "", carried.Interface,
"wg-quick@"+carried.Interface)
default:
fmt.Printf(" %-17s NOT taken over: %s is still the interface the peers reach\n", "", carried.Interface)
}
if said && carried.Note != "" {
fmt.Printf(" %-17s %s\n", "", carried.Note)
}
if said && carried.Kept != "" {
fmt.Printf(" %-17s its configuration's original kept at %s\n", "", carried.Kept)
}
return nil
}
func orNone(s string) string {
if s == "" {
return "none reported"
}
return s
}
// adoptedNodes are the names of every adopted node, in the order given.
func adoptedNodes(nodes []inventory.Node) []string {
var out []string
for _, n := range nodes {
if n.Adopted {
out = append(out, n.Name)
}
}
return out
}
// The operator's acts on an adopted node (novox/hq ADR 0100). Called by the command line and the
// command API alike, so a refusal is the same refusal in the same words at both (ADR 0035).
// sendNodes sends the named machines what they should be now. A variable so a test can see what
// an act would send without a broker.
var sendNodes = sendTo
// DefaultFilter is the module converging a node assigns to load the mesh's derived filter.
const DefaultFilter = "nftables"
// take is a module's cutover on an adopted node: the operator's act, done when that module's data
// has moved. From the next push its resources converge there like any other, replacing what the
// node found and holds for it.
func take(ctx context.Context, open *stores, node, module string) (string, error) {
inv := open.inventory
assigned, err := inv.Assigned(ctx, node)
if err != nil {
return "", err
}
if !slices.Contains(assigned, module) {
if plan, _, err := planFor(ctx, open, node); err == nil {
if why, runs := plan.Because[module]; runs {
return "", fmt.Errorf("%w: %s runs on %s because %s — assign it to %s to take it",
inventory.ErrNotAssigned, module, node, why, node)
}
}
}
if err := inv.Take(ctx, node, module); err != nil {
return "", err
}
said := fmt.Sprintf("%s is taken on %s", module, node)
reported, err := inv.AdoptionOf(ctx, node)
if err != nil {
return "", err
}
var replaces []string
for _, h := range reported.Held {
if h.Module == module {
replaces = append(replaces, " "+heldLine(h))
}
}
if len(replaces) > 0 {
said += "; the next push replaces what the node found and holds for it:\n" +
strings.Join(replaces, "\n")
}
return said + fmt.Sprintf("\n run `push %s` to cut it over", node), nil
}
// reportFreshFor is how old a node's account of itself may be for the flip to act on it. A
// variable so a test can age a report without waiting.
var reportFreshFor = 15 * time.Minute
// converge previews, and with yes makes, the flip of an adopted node to converged: every module it
// runs is taken, the filter module is assigned to load the mesh's derived filter in place of the
// guard, and the found firewall is retired — disabled, never flushed — by the host.
//
// Refused while an assigned module still holds a found container: each service is taken on its
// own, when its data has moved, never by the flip. And refused on a preview that would be stale:
// what is reachable is the node's last account, so that account must be of what it was last sent.
//
// **The flip acts on the preview the operator saw** and on nothing else. The preview ends with a
// short digest of what it said — every reachable thing and its fate, the modules the flip takes and
// the filter — and yes must name that digest: if anything the preview would say has changed since,
// the flip is refused rather than done on a preview nobody read. And it is refused on an account
// older than reportFreshFor: what was reachable then is not evidence of what is reachable now.
func converge(ctx context.Context, open *stores, node string, yes bool, digest string,
filter string) (string, error) {
inv := open.inventory
if filter == "" {
filter = DefaultFilter
}
if yes {
// Held from the checks to the send, so no push composed before the flip is sent after it
// and returns the node to adopted.
held, release, err := holdNodes(ctx, open, []string{node})
if err != nil {
return "", err
}
defer release()
ctx = held
}
record, err := inv.NodeByName(ctx, node)
if err != nil {
return "", err
}
if !record.Adopted {
return "", fmt.Errorf("%s is converged already; there is nothing to flip", node)
}
reports, err := inv.LastReports(ctx)
if err != nil {
return "", err
}
current := false
for _, r := range reports {
if r.Node == node {
current = r.Current
}
}
reported, err := inv.AdoptionOf(ctx, node)
if err != nil {
return "", err
}
if !current || reported.At.IsZero() {
return "", fmt.Errorf("%s has not reported on the declaration it was last sent, so what it "+
"says is reachable may not be the machine as it is: run `push %s --wait 2m` and "+
"converge once it has applied", node, node)
}
assignedWhenPreviewed, err := inv.Assigned(ctx, node)
if err != nil {
return "", err
}
plan, settings, err := planFor(ctx, open, node)
if err != nil {
return "", err
}
runs := map[string]bool{}
for _, m := range plan.Modules {
runs[m.Module] = true
}
var holding []string
for _, h := range reported.Held {
if h.Kind == "container" && runs[h.Module] {
holding = append(holding, fmt.Sprintf(" %s holds the found container %s — take %s %s "+
"once its data has moved", h.Module, h.Target, node, h.Module))
}
}
if len(holding) > 0 {
sort.Strings(holding)
return "", fmt.Errorf("%s still holds what it found, and a service is taken on its own, "+
"never by the flip:\n%s", node, strings.Join(holding, "\n"))
}
// And refused while a peer of the tunnel this hub took over has not enrolled (novox/hq ADR
// 0105): the flip loads the derived filter and retires the found firewall, and a machine the
// mesh has no record of is not one the filter admits — it would go dark.
if _, hubName, adopted, err := inv.AdoptedTunnel(ctx); err != nil {
return "", err
} else if adopted && hubName == node {
carried, err := inv.CarriedPeers(ctx)
if err != nil {
return "", err
}
var waiting []string
for _, c := range carried {
if c.EnrolledAs == "" {
waiting = append(waiting, fmt.Sprintf(" %s at %s", overlay.CarriedName(c.PublicKey), c.Address))
}
}
if len(waiting) > 0 {
return "", fmt.Errorf("%s carries peers of the tunnel it took over that have not enrolled, and "+
"converging would cut them off — enrol each first (`overlay show` says which are enrolled):\n%s",
node, strings.Join(waiting, "\n"))
}
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
return "", err
}
filterModule, known := shelf[filter]
if !known {
return "", fmt.Errorf("%w: %s — converging assigns it to load the mesh's filter; "+
"name another with --filter", inventory.ErrNoSuchModule, filter)
}
if filterModule.Filtering == nil {
return "", fmt.Errorf("%s loads no filter of the mesh's; name a module that does with --filter",
filter)
}
gens, err := generators(ctx, open)
if err != nil {
return "", err
}
with, _, err := renderingFor(ctx, open, node, plan, settings, gens, Reading)
if err != nil {
return "", err
}
rules, err := plan.Rules(with)
if err != nil {
return "", err
}
taken, err := inv.Taken(ctx, node)
if err != nil {
return "", err
}
derived := derivedFilter{rules: rules, foundation: with.Foundation, mesh: with.Mesh,
outward: plan.PublicDomain != ""}
preview, saw := previewOf(node, reported, derived, plan, taken, filter, runs[filter])
preview += "\n\n preview " + saw
if !yes {
return preview + fmt.Sprintf("\n\nNothing has changed. Run `converge %s --yes %s` to do "+
"it.", node, saw), nil
}
// An account naming nothing reachable is not an account of a machine: every machine answers
// on ssh, and the host's collectors failing — `ss` refusing, or the container runtime not
// answering, which drops every published port at once — leaves exactly this. Flipping on it
// would close ports the preview never named.
if yes && countReachable(reported) == 0 {
return preview, fmt.Errorf("%s says nothing is reachable on it, which no machine that is "+
"up ever is: its account looks partial — whatever reads what is listening, or what "+
"the container runtime publishes, did not answer. Fix that on the machine and run "+
"`push %s --wait 2m`, then preview again", node, node)
}
if age := time.Since(reported.At); age > reportFreshFor {
return preview, fmt.Errorf("%s last said what is reachable on it %s ago, and the flip acts "+
"only on an account newer than %s: wait for its next report, or run `push %s --wait 2m`, "+
"then preview again", node, age.Round(time.Second), reportFreshFor, node)
}
if digest == "" {
return preview, fmt.Errorf("converging %s acts on the preview you saw: name its digest, "+
"`converge %s --yes %s`, once you have read it", node, node, saw)
}
if digest != saw {
return preview, fmt.Errorf("what converging %s would do has changed since preview %s "+
"(it is now %s): read the preview above, and run `converge %s --yes %s` if it is "+
"what you want", node, digest, saw, node, saw)
}
// The flip. The filter first, and only kept if the node still resolves with it: a node that
// cannot be worked out would be sent nothing, and would sit with its guard and no filter.
assigned, err := inv.Assigned(ctx, node)
if err != nil {
return "", err
}
// And nothing assigned since the preview was composed: the flip takes every module the node
// runs, and one assigned in between would be taken without ever having been previewed.
if !slices.Equal(assigned, assignedWhenPreviewed) {
return preview, fmt.Errorf("what %s runs changed while this was converging (it is now %s): "+
"the flip takes every module on the node, so read the preview again", node,
strings.Join(assigned, ", "))
}
if !slices.Contains(assigned, filter) {
if err := inv.Assign(ctx, node, filter); err != nil {
return "", err
}
if _, _, err := planFor(ctx, open, node); err != nil {
_ = inv.Unassign(ctx, node, filter)
return "", fmt.Errorf("%s cannot run %s, so it was not converged: %w", node, filter, err)
}
}
took, err := inv.Converge(ctx, node)
if err != nil {
return "", err
}
said := preview + fmt.Sprintf("\n\n%s is converged", node)
if len(took) > 0 {
said += "; took " + strings.Join(took, ", ")
}
if err := sendNodes(ctx, open, []string{node}); err != nil {
return said + "\n and it could not be sent: run `push " + node + "`", err
}
return said + "\n sent: the host loads the mesh's filter and disables the firewall it found", nil
}
// previewOf is what converging a node will change, before it changes it, and a short digest of
// what it said: every reachable thing and its fate, the modules the flip takes and the filter. The
// digest is what the flip is asked to act on, so it changes whenever any of those would.
func previewOf(node string, reported inventory.Adoption, derived derivedFilter,
plan catalogue.Resolution, taken []string, filter string, filterAssigned bool) (string, string) {
var said []string
var b strings.Builder
fmt.Fprintf(&b, "converging %s\n", node)
fmt.Fprintf(&b, "\n reachable on the machine now, as it reported at %s:\n",
reported.At.Local().Format(time.DateTime))
for _, r := range reported.Reachable {
if loopback(r.Address) {
continue
}
what := fmt.Sprintf("%s/%d", r.Protocol, r.Port)
if r.By != "" {
what += " " + r.By
}
if r.Published {
what += fmt.Sprintf(" (published, container port %d)", r.ContainerPort)
}
fate := derived.fate(r)
fmt.Fprintf(&b, " %-44s %s\n", what, fate)
said = append(said, fmt.Sprintf("reach %s %s %s", r.Address, what, fate))
}
if countReachable(reported) == 0 {
// Said as what it is: no machine that is up is reachable on nothing, so this is an
// account that did not come back, not a machine with nothing on it.
b.WriteString(" nothing reported — this account looks partial, and the flip is " +
"refused on it\n")
said = append(said, "reach nothing reported")
}
// What the machine routes for others is not a listener and not a published port, so nothing
// above can show it; the derived filter's forward chain drops it all the same.
b.WriteString(" not previewed: traffic the machine routes that is not a published port " +
"(a tunnel, NAT in the found firewall) — the derived filter drops it unless a module " +
"declares it\n")
isTaken := map[string]bool{}
for _, m := range taken {
isTaken[m] = true
}
var takes []string
for _, m := range plan.Modules {
if !isTaken[m.Module] {
takes = append(takes, m.Module)
}
}
if !filterAssigned && !isTaken[filter] {
takes = append(takes, filter)
}
sort.Strings(takes)
b.WriteString("\n the flip takes:\n")
if len(takes) == 0 {
b.WriteString(" nothing — every module is taken already\n")
}
for _, m := range takes {
fmt.Fprintf(&b, " %s\n", m)
said = append(said, "take "+m)
// Every kind it holds — a directory, a service, an archive, a process, a user as well as a
// file (novox/hq ADR 0103) — each said, and each part of what the flip is asked to act on.
for _, h := range reported.Held {
if h.Module != m {
continue
}
fmt.Fprintf(&b, " replacing the found %s", heldLine(h))
if h.Kept != "" {
fmt.Fprintf(&b, ", original kept at %s", h.Kept)
}
b.WriteString("\n")
said = append(said, "replace "+m+" "+heldLine(h)+" "+h.Kept)
}
}
if !filterAssigned {
fmt.Fprintf(&b, "\n and assigns %s, which loads the mesh's filter in place of its guard\n", filter)
}
fw := reported.Firewall
if fw == "" || fw == "none" {
b.WriteString(" no firewall was found on the machine; the mesh's filter is its first\n")
} else {
fmt.Fprintf(&b, " the found firewall (%s) is disabled, never flushed: its configuration stays on disk\n", fw)
}
said = append(said, fmt.Sprintf("filter %s assigned=%t firewall=%s", filter, filterAssigned, fw))
// Sorted: the same account, reported in another order, is the same preview.
sort.Strings(said)
sum := sha256.Sum256([]byte(strings.Join(said, "\n")))
return strings.TrimRight(b.String(), "\n"), hex.EncodeToString(sum[:])[:12]
}
// derivedFilter is what the filter the flip loads is rendered from, as AsNftables renders it.
type derivedFilter struct {
rules []catalogue.Rule
foundation []int
// mesh is every address on the private network; outward says the machine faces outside.
mesh []string
outward bool
}
// closesOutside is what a narrowing from everywhere to the private network is called: it closes.
const closesOutside = "WILL CLOSE to everything outside the private network"
// fate is what the derived filter does to one reachable thing: which module declares it and from
// where, or that it will close — wholly, or to everything outside the private network. Rendered
// exactly as AsNftables admits it, ssh included.
func (d derivedFilter) fate(r inventory.Reach) string {
// Bound to an address on the private network, it was never reachable from outside it, so
// admitting it from the mesh narrows nothing.
onMesh := slices.Contains(d.mesh, strings.Trim(r.Address, "[]"))
if r.Protocol == "tcp" && r.Port == catalogue.SSHPort {
// From everywhere only when the machine faces outward or the mesh has no addresses to
// narrow it to; otherwise from the private network only.
if d.outward || len(d.mesh) == 0 || onMesh {
return "stays open — ssh is never closed"
}
return closesOutside + " — ssh stays open from the mesh, never closed there"
}
for _, port := range d.foundation {
if r.Protocol == "tcp" && r.Port == port {
return "stays open — the mesh's own, from anywhere"
}
}
for _, rule := range d.rules {
if rule.Port != r.Port || rule.Protocol != r.Protocol {
continue
}
by := strings.Join(rule.Because, ", ")
switch rule.From {
case catalogue.FromMachine:
return fmt.Sprintf("WILL CLOSE to the network — declared by %s for this machine only", by)
case catalogue.FromMesh:
if len(d.mesh) == 0 {
return fmt.Sprintf("WILL CLOSE — declared by %s from the mesh, and this node "+
"knows no mesh addresses", by)
}
if !onMesh {
return fmt.Sprintf("%s — declared by %s from the mesh only", closesOutside, by)
}
}
return fmt.Sprintf("declared by %s (from %s)", by, rule.From)
}
return "WILL CLOSE — no module assigned here declares it"
}
// countReachable is how much of a node's account of itself names something off the machine.
// Loopback is left out for the same reason the preview leaves it out: nothing outside reaches it,
// so a report of loopback alone says nothing about what the filter would close.
func countReachable(reported inventory.Adoption) int {
n := 0
for _, r := range reported.Reachable {
if !loopback(r.Address) {
n++
}
}
return n
}
// heldLine is one thing a node holds as found, as take and the converge preview both say it.
func heldLine(h inventory.Held) string {
return fmt.Sprintf("%s %s (%s)", h.Kind, h.Target, h.ID)
}
// loopback is an address nothing off the machine reaches.
func loopback(address string) bool {
a := strings.Trim(address, "[]")
return strings.HasPrefix(a, "127.") || a == "::1" || a == "localhost"
}
// adopt returns a converged node to adopted: the mesh's filter is unloaded, the guard restored,
// the found firewall enabled again and the openings converged through it once more. What was
// taken stays taken.
func adopt(ctx context.Context, open *stores, node string) (string, error) {
inv := open.inventory
record, err := inv.NodeByName(ctx, node)
if err != nil {
return "", err
}
if record.Adopted {
return "", fmt.Errorf("%s is adopted already", node)
}
held, release, err := holdNodes(ctx, open, []string{node})
if err != nil {
return "", err
}
defer release()
ctx = held
if err := inv.SetAdopted(ctx, node, true); err != nil {
return "", err
}
said := fmt.Sprintf("%s is adopted; what was taken on it stays taken", node)
if err := sendNodes(ctx, open, []string{node}); err != nil {
return said + "\n and it could not be sent: run `push " + node + "`", err
}
return said + "\n sent: the host unloads the mesh's filter and enables the firewall it found", nil
}
// takeCommand, convergeCommand and adoptCommand are the command line's adapters to the acts above.
func takeCommand(ctx context.Context, args []string) error {
if len(args) != 2 {
return errors.New("take <node> <module>")
}
return runAct(ctx, func(open *stores) (string, error) { return take(ctx, open, args[0], args[1]) })
}
func convergeCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("converge", flag.ContinueOnError)
yes := set.String("yes", "", "do it, naming the digest the preview printed; without it, only "+
"the preview")
filter := set.String("filter", DefaultFilter, "the module that loads the mesh's filter")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) != 1 {
return errors.New("converge <node> [--yes <digest>] [--filter nftables]")
}
return runAct(ctx, func(open *stores) (string, error) {
return converge(ctx, open, positionals[0], *yes != "", *yes, *filter)
})
}
func adoptCommand(ctx context.Context, args []string) error {
if len(args) != 1 {
return errors.New("adopt <node>")
}
return runAct(ctx, func(open *stores) (string, error) { return adopt(ctx, open, args[0]) })
}
func runAct(ctx context.Context, act func(*stores) (string, error)) error {
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
said, err := act(open)
if said != "" {
fmt.Println(said)
}
if err != nil && said != "" {
fmt.Println()
}
return err
}