Take the foundation's ports as genesis inputs, check them free, and hand them to the controller as the node's settings (hq ADR 0100)

This commit is contained in:
2026-09-22 17:28:36 +02:00
parent 770f589401
commit 3964d9da0a
10 changed files with 902 additions and 14 deletions
+60 -1
View File
@@ -25,6 +25,7 @@ import (
"net/http" "net/http"
"os" "os"
"os/signal" "os/signal"
"strconv"
"syscall" "syscall"
"time" "time"
@@ -126,6 +127,14 @@ const usage = `mesh-bootstrap — make a bare machine into a mesh
--packet-filter which packet filter to run (nftables) --packet-filter which packet filter to run (nftables)
--extras catalogue modules beyond the floor, comma-separated --extras catalogue modules beyond the floor, comma-separated
The foundation's ports are this machine's, each checked free before anything is
raised and kept as the node's setting for the module that binds it:
--store-port 5432 --bus-port 5671 --amqp-port 5672 --management-port 15672
--registry-port 5000 (follows --registry, and must agree with it)
--packages-port 3000 --hub-port 51820/udp
--overlay-range the private network's range (default 10.42.0.0/16); refused
if it overlaps an interface or route the machine already has
The installer carries a builder, not a control plane. What raises a mesh is therefore The installer carries a builder, not a control plane. What raises a mesh is therefore
the same thing that will maintain it, and the control plane a mesh ends up running is the same thing that will maintain it, and the control plane a mesh ends up running is
one it built itself, from a repository and a commit it can name and build again. one it built itself, from a repository and a commit it can name and build again.
@@ -176,7 +185,9 @@ func parseArgs(args []string) (string, bootstrap.Options, bool, error) {
HostService: defaultService, HostService: defaultService,
// Longer than the host's 10s: these probes reach a container runtime that may be busy // Longer than the host's 10s: these probes reach a container runtime that may be busy
// pulling, and a probe that times out on a working machine is a false refusal. // pulling, and a probe that times out on a working machine is a false refusal.
Timeout: 30 * time.Second, Ports: bootstrap.DefaultPorts(),
OverlayRange: bootstrap.DefaultOverlayRange,
Timeout: 30 * time.Second,
// A socket-activated runtime queued behind the network, and a control plane running its // A socket-activated runtime queued behind the network, and a control plane running its
// first `initdb`-shaped wait, are both minutes rather than seconds. // first `initdb`-shaped wait, are both minutes rather than seconds.
Wait: 3 * time.Minute, Wait: 3 * time.Minute,
@@ -210,9 +221,38 @@ func parseArgs(args []string) (string, bootstrap.Options, bool, error) {
return "", opts, false, fmt.Errorf( return "", opts, false, fmt.Errorf(
"unexpected argument %q — try `mesh-bootstrap help`", positionals[0]) "unexpected argument %q — try `mesh-bootstrap help`", positionals[0])
} }
if err := registryAgrees(set, &opts); err != nil {
return "", opts, false, err
}
return command, opts, jsonOut, nil return command, opts, jsonOut, nil
} }
// registryAgrees makes --registry and --registry-port say one port (novox/hq ADR 0100): the
// registry is raised on the port the node gives it, and every node pulls from the address given.
// Either may be said alone and the other follows; said both ways, they must agree.
func registryAgrees(set *flag.FlagSet, opts *bootstrap.Options) error {
said := map[string]bool{}
set.Visit(func(f *flag.Flag) { said[f.Name] = true })
host, portText, err := net.SplitHostPort(opts.Registry)
if err != nil {
return fmt.Errorf("--registry %q is not host:port: %w", opts.Registry, err)
}
port, err := strconv.Atoi(portText)
if err != nil {
return fmt.Errorf("--registry %q does not end in a port", opts.Registry)
}
switch {
case said["registry-port"] && said["registry"] && port != opts.Ports.Registry:
return fmt.Errorf("--registry %s and --registry-port %d name two ports for one registry",
opts.Registry, opts.Ports.Registry)
case said["registry-port"]:
opts.Registry = net.JoinHostPort(host, strconv.Itoa(opts.Ports.Registry))
case said["registry"]:
opts.Ports.Registry = port
}
return nil
}
func newFlagSet(opts *bootstrap.Options, jsonOut *bool) *flag.FlagSet { func newFlagSet(opts *bootstrap.Options, jsonOut *bool) *flag.FlagSet {
set := flag.NewFlagSet("mesh-bootstrap", flag.ContinueOnError) set := flag.NewFlagSet("mesh-bootstrap", flag.ContinueOnError)
set.SetOutput(os.Stderr) set.SetOutput(os.Stderr)
@@ -255,6 +295,25 @@ func newFlagSet(opts *bootstrap.Options, jsonOut *bool) *flag.FlagSet {
set.StringVar(&opts.SDKSource.Ref, "sdk-ref", opts.SDKSource.Ref, set.StringVar(&opts.SDKSource.Ref, "sdk-ref", opts.SDKSource.Ref,
"what of it to build (default main)") "what of it to build (default main)")
set.StringVar(&opts.Site, "site", "main", "where this machine sits, for the private network") set.StringVar(&opts.Site, "site", "main", "where this machine sits, for the private network")
// The foundation's ports are this node's (novox/hq ADR 0100): each is checked free before
// anything is raised, and becomes the node's setting for the module that binds it.
for _, p := range []struct {
name, what string
into *int
}{
{"store-port", "the store", &opts.Ports.Store},
{"bus-port", "the bus (amqps)", &opts.Ports.Bus},
{"amqp-port", "the broker's AMQP", &opts.Ports.AMQP},
{"management-port", "the broker's management, on loopback", &opts.Ports.Management},
{"registry-port", "the registry", &opts.Ports.Registry},
{"packages-port", "the package registry", &opts.Ports.Packages},
{"hub-port", "the private network's hub (udp)", &opts.Ports.Hub},
} {
set.IntVar(p.into, p.name, *p.into, "the machine's port for "+p.what)
}
set.StringVar(&opts.OverlayRange, "overlay-range", opts.OverlayRange,
"the private network's address range; must not overlap a tunnel the machine already runs")
if opts.Answers == nil { if opts.Answers == nil {
opts.Answers = map[string]string{} opts.Answers = map[string]string{}
} }
+26
View File
@@ -164,3 +164,29 @@ func TestTheNodeNameCanBeSaid(t *testing.T) {
t.Errorf("--catalog parsed as %q", opts.Catalogue) t.Errorf("--catalog parsed as %q", opts.Catalogue)
} }
} }
// Defends novox/hq ADR 0100: the foundation's ports are inputs to genesis, and the registry's port
// and the address nodes pull from say one port.
func TestTheFoundationsPortsAreGiven(t *testing.T) {
_, opts, _, err := parseArgs([]string{"--store-port", "5433", "--hub-port", "51821", "--overlay-range", "10.77.0.0/16"})
if err != nil {
t.Fatal(err)
}
if opts.Ports.Store != 5433 || opts.Ports.Hub != 51821 || opts.Ports.Bus != 5671 || opts.OverlayRange != "10.77.0.0/16" {
t.Errorf("ports read as %+v, range %s", opts.Ports, opts.OverlayRange)
}
}
func TestTheRegistrysPortAndAddressAgree(t *testing.T) {
_, opts, _, err := parseArgs([]string{"--registry-port", "5100"})
if err != nil || opts.Registry != "127.0.0.1:5100" {
t.Errorf("--registry-port alone: %s %v", opts.Registry, err)
}
_, opts, _, err = parseArgs([]string{"--registry", "192.0.2.10:5100"})
if err != nil || opts.Ports.Registry != 5100 {
t.Errorf("--registry alone: %d %v", opts.Ports.Registry, err)
}
if _, _, _, err := parseArgs([]string{"--registry", "192.0.2.10:5000", "--registry-port", "5100"}); err == nil {
t.Error("two ports for one registry were accepted")
}
}
+40 -1
View File
@@ -185,6 +185,20 @@ type Options struct {
Prompt func(Choice) (string, error) Prompt func(Choice) (string, error)
// Extras are catalogue modules beyond the floor, asked for by name. // Extras are catalogue modules beyond the floor, asked for by name.
Extras []string Extras []string
// Ports are the ports the foundation binds on this machine (novox/hq ADR 0100). Inputs to
// genesis, each checked free before anything is raised, and then the node's settings for the
// foundation's modules — so adopting the foundation as modules leaves it where it was raised.
// Zero means the catalogue's defaults.
Ports FoundationPorts
// OverlayRange is the private network's address range, checked against every interface and
// route the machine already has. Empty means the mesh's default.
OverlayRange string
// Adopted raises this machine as an adopted node (novox/hq ADR 0100): what is on it is kept
// until each module is taken, its firewall stays in force, and the mesh guards its own ports
// in a table that only refuses. Without it, a machine in use is refused.
Adopted bool
} }
// pivots reports whether this run goes past the foundation. // pivots reports whether this run goes past the foundation.
@@ -287,6 +301,14 @@ type Result struct {
// Stopped names why a run went no further. Empty on a run that pivoted. // Stopped names why a run went no further. Empty on a run that pivoted.
Stopped string `json:"stopped,omitempty"` Stopped string `json:"stopped,omitempty"`
// Adopted, the firewall found, and the ports the foundation was raised on (novox/hq ADR 0100).
Adopted bool `json:"adopted,omitempty"`
Firewall string `json:"firewall,omitempty"`
Ports FoundationPorts `json:"ports"`
// Filter is the packet filter chosen for when the node converges; an adopted genesis loads
// none, and the flip assigns this one.
Filter string `json:"filter-on-converge,omitempty"`
} }
// Run performs the bootstrap, saying what it is doing as it goes. // Run performs the bootstrap, saying what it is doing as it goes.
@@ -339,7 +361,12 @@ func Run(ctx context.Context, o Options, d Deps, say func(string)) (Result, erro
if say == nil { if say == nil {
say = func(string) {} say = func(string) {}
} }
result := Result{DryRun: o.DryRun} result := Result{DryRun: o.DryRun, Adopted: o.Adopted}
o.Ports = o.Ports.orDefaults()
result.Ports = o.Ports
if err := o.Ports.Check(); err != nil {
return result, failed(StepPreflight, err)
}
// ---- 1. preflight ------------------------------------------------------------------- // ---- 1. preflight -------------------------------------------------------------------
say("preflight — what has to be true before anything is changed") say("preflight — what has to be true before anything is changed")
@@ -399,12 +426,24 @@ func Run(ctx context.Context, o Options, d Deps, say func(string)) (Result, erro
if err := RefuseExistingServers(ctx, d.Run, creds); err != nil { if err := RefuseExistingServers(ctx, d.Run, creds); err != nil {
return result, failed(StepBundle, err) return result, failed(StepBundle, err)
} }
// The foundation's ports, its private network's range and its containers' names are checked
// free before anything is raised (novox/hq ADR 0100), each refusal naming what holds it.
if err := CheckTheMachine(ctx, o, d.Run, rewritten.Declaration, say); err != nil {
return result, failed(StepBundle, err)
}
// From here on nothing this installer says contains the values it just made. // From here on nothing this installer says contains the values it just made.
say = Masking(say, creds) say = Masking(say, creds)
root, err := RewriteRoot(&rewritten, creds) root, err := RewriteRoot(&rewritten, creds)
if err != nil { if err != nil {
return result, failed(StepBundle, err) return result, failed(StepBundle, err)
} }
moved, err := RewritePorts(&rewritten, o.Ports, o.OverlayRange)
if err != nil {
return result, failed(StepBundle, err)
}
if moved.Places > 0 {
say(fmt.Sprintf(" ports %d place(s) rewritten to this node's foundation ports", moved.Places))
}
for _, c := range []struct { for _, c := range []struct {
what, path string what, path string
made bool made bool
+10
View File
@@ -124,6 +124,9 @@ func registerAndAssign(ctx context.Context, o Options, control controlPlane, mod
// the only place the reason appears. // the only place the reason appears.
say(indent(refusal)) say(indent(refusal))
} }
if err := prepareModule(ctx, o, control, module, say); err != nil {
return out, err
}
return out, nil return out, nil
} }
@@ -209,3 +212,10 @@ func pinPlaceholder(manifest []byte, reference, module string) ([]byte, int, err
} }
return pinned, places, nil return pinned, places, nil
} }
// prepareModule is what genesis tells the controller about a module on this node once it is
// assigned and before it is pushed: the ports this node gave it (novox/hq ADR 0100).
func prepareModule(ctx context.Context, o Options, control controlPlane, module string,
say func(string)) error {
return setFoundationSettings(ctx, o, control, module, say)
}
+7 -3
View File
@@ -5,6 +5,7 @@ import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"net" "net"
"strconv"
"strings" "strings"
"time" "time"
) )
@@ -87,6 +88,9 @@ func InstallFromCatalogue(ctx context.Context, o Options, control controlPlane,
if _, err := control.tell(ctx, "assign", o.Node, module); err != nil { if _, err := control.tell(ctx, "assign", o.Node, module); err != nil {
return err return err
} }
if err := prepareModule(ctx, o, control, module, say); err != nil {
return err
}
if _, err := pushNode(ctx, o, control, say); err != nil { if _, err := pushNode(ctx, o, control, say); err != nil {
return err return err
} }
@@ -120,7 +124,7 @@ func PlaceOnTheNetwork(ctx context.Context, o Options, control controlPlane,
Name: "endpoint", Name: "endpoint",
Question: "Where do other machines reach this one for the private network? " + Question: "Where do other machines reach this one for the private network? " +
"(host:port; the host other machines dial)", "(host:port; the host other machines dial)",
Default: derivedEndpoint(brokerAddress), Default: derivedEndpoint(brokerAddress, o.Ports.orDefaults().Hub),
}, o.Answers["endpoint"], o.Prompt, say) }, o.Answers["endpoint"], o.Prompt, say)
if err != nil { if err != nil {
return err return err
@@ -202,12 +206,12 @@ func builds(manifest []byte) bool {
// derivedEndpoint is the default place other machines dial for the private network: the same host // derivedEndpoint is the default place other machines dial for the private network: the same host
// they already dial for the broker, on WireGuard's ordinary port. One fact, not two. // they already dial for the broker, on WireGuard's ordinary port. One fact, not two.
func derivedEndpoint(brokerAddress string) string { func derivedEndpoint(brokerAddress string, hub int) string {
host, _, err := net.SplitHostPort(brokerAddress) host, _, err := net.SplitHostPort(brokerAddress)
if err != nil || host == "" { if err != nil || host == "" {
return "" return ""
} }
return net.JoinHostPort(host, "51820") return net.JoinHostPort(host, strconv.Itoa(hub))
} }
func refOr(ref string) string { func refOr(ref string) string {
+2 -2
View File
@@ -5,10 +5,10 @@ import "testing"
// The endpoint other machines dial defaults to the host they already dial — the broker's — on // The endpoint other machines dial defaults to the host they already dial — the broker's — on
// WireGuard's port. One fact, not two that drift. // WireGuard's port. One fact, not two that drift.
func TestTheEndpointDerivesFromTheBrokerAddress(t *testing.T) { func TestTheEndpointDerivesFromTheBrokerAddress(t *testing.T) {
if got := derivedEndpoint("192.0.2.10:5671"); got != "192.0.2.10:51820" { if got := derivedEndpoint("192.0.2.10:5671", 51820); got != "192.0.2.10:51820" {
t.Fatalf("derived %q", got) t.Fatalf("derived %q", got)
} }
if got := derivedEndpoint(""); got != "" { if got := derivedEndpoint("", 51820); got != "" {
t.Fatalf("an endpoint was invented from nothing: %q", got) t.Fatalf("an endpoint was invented from nothing: %q", got)
} }
} }
+3
View File
@@ -122,6 +122,9 @@ func installProvider(ctx context.Context, o Options, control controlPlane, modul
if _, err := control.tell(ctx, "assign", o.Node, module); err != nil { if _, err := control.tell(ctx, "assign", o.Node, module); err != nil {
return err return err
} }
if err := prepareModule(ctx, o, control, module, say); err != nil {
return err
}
if beforePush != nil { if beforePush != nil {
if err := beforePush(); err != nil { if err := beforePush(); err != nil {
return err return err
+14 -7
View File
@@ -42,8 +42,9 @@ const (
// giteaDBRole/giteaDBName is gitea's own database in the foundation store. // giteaDBRole/giteaDBName is gitea's own database in the foundation store.
giteaDBRole = "mesh_gitea" giteaDBRole = "mesh_gitea"
giteaDBName = "mesh_gitea" giteaDBName = "mesh_gitea"
// giteaPort is where the raised server answers on the machine. // defaultGiteaPort is where the raised server answers on the machine unless the node gave the
giteaPort = 3000 // package registry another port (novox/hq ADR 0100).
defaultGiteaPort = 3000
) )
// RaisePackageRegistry puts a working npm registry in front of the base build. It is idempotent: // RaisePackageRegistry puts a working npm registry in front of the base build. It is idempotent:
@@ -60,17 +61,18 @@ func RaisePackageRegistry(ctx context.Context, o Options, d Deps, control contro
} }
say(" seeding gitea's database in the foundation store") say(" seeding gitea's database in the foundation store")
ports := o.Ports.orDefaults()
if err := seedGiteaDatabase(ctx, run, o.Timeout, dbPassword, say); err != nil { if err := seedGiteaDatabase(ctx, run, o.Timeout, dbPassword, say); err != nil {
return err return err
} }
say(" raising the gitea server on that database") say(" raising the gitea server on that database")
if err := raiseGiteaServer(ctx, run, o.Timeout, dbPassword, say); err != nil { if err := raiseGiteaServer(ctx, run, o.Timeout, dbPassword, ports, say); err != nil {
return err return err
} }
say(" waiting for gitea to answer") say(" waiting for gitea to answer")
base := fmt.Sprintf("http://127.0.0.1:%d", giteaPort) base := fmt.Sprintf("http://127.0.0.1:%d", ports.Packages)
if err := waitForGitea(ctx, d, o, base, say); err != nil { if err := waitForGitea(ctx, d, o, base, say); err != nil {
return err return err
} }
@@ -150,7 +152,7 @@ func seedGiteaDatabase(ctx context.Context, run Runner, timeout time.Duration, p
// store's network namespace so `127.0.0.1:5432` reaches postgres, and publishes its own port on the // store's network namespace so `127.0.0.1:5432` reaches postgres, and publishes its own port on the
// machine so the builder and this installer can reach it. Started if absent, left alone if present. // machine so the builder and this installer can reach it. Started if absent, left alone if present.
func raiseGiteaServer(ctx context.Context, run Runner, timeout time.Duration, dbPassword string, func raiseGiteaServer(ctx context.Context, run Runner, timeout time.Duration, dbPassword string,
say func(string)) error { ports FoundationPorts, say func(string)) error {
asking, cancel := context.WithTimeout(ctx, timeout) asking, cancel := context.WithTimeout(ctx, timeout)
defer cancel() defer cancel()
@@ -164,7 +166,7 @@ func raiseGiteaServer(ctx context.Context, run Runner, timeout time.Duration, db
env := []string{ env := []string{
"-e", "GITEA__database__DB_TYPE=postgres", "-e", "GITEA__database__DB_TYPE=postgres",
// The store is reached on the shared network namespace's loopback. // The store is reached on the shared network namespace's loopback.
"-e", "GITEA__database__HOST=127.0.0.1:5432", "-e", fmt.Sprintf("GITEA__database__HOST=127.0.0.1:%d", ports.Store),
"-e", "GITEA__database__NAME=" + giteaDBName, "-e", "GITEA__database__NAME=" + giteaDBName,
"-e", "GITEA__database__USER=" + giteaDBRole, "-e", "GITEA__database__USER=" + giteaDBRole,
"-e", "GITEA__database__PASSWD=" + dbPassword, "-e", "GITEA__database__PASSWD=" + dbPassword,
@@ -174,9 +176,14 @@ func raiseGiteaServer(ctx context.Context, run Runner, timeout time.Duration, db
// package metadata hands npm a tarball URL built from ROOT_URL, and a client only sends its // package metadata hands npm a tarball URL built from ROOT_URL, and a client only sends its
// stored credential to the host it was stored for. A default ROOT_URL of localhost is a // stored credential to the host it was stored for. A default ROOT_URL of localhost is a
// different host than the binding's 127.0.0.1, so the credential would not be sent. // different host than the binding's 127.0.0.1, so the credential would not be sent.
"-e", fmt.Sprintf("GITEA__server__ROOT_URL=http://127.0.0.1:%d/", giteaPort), "-e", fmt.Sprintf("GITEA__server__ROOT_URL=http://127.0.0.1:%d/", ports.Packages),
"-e", "USER_UID=1000", "-e", "USER_GID=1000", "-e", "USER_UID=1000", "-e", "USER_GID=1000",
} }
if ports.Packages != defaultGiteaPort {
// On the machine's network the server binds its own port, so a port given for it is
// the one it is told to listen on.
env = append(env, "-e", fmt.Sprintf("GITEA__server__HTTP_PORT=%d", ports.Packages))
}
args := append([]string{ args := append([]string{
"run", "-d", "--name", giteaBootstrap, "run", "-d", "--name", giteaBootstrap,
// Host network, like the control plane: it reaches the foundation store on the machine's // Host network, like the control plane: it reaches the foundation store on the machine's
+454
View File
@@ -0,0 +1,454 @@
package bootstrap
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net"
"sort"
"strconv"
"strings"
"github.com/novox/mesh-host/internal/declaration"
"github.com/novox/mesh-host/internal/reachable"
"github.com/novox/mesh-host/internal/store"
)
// FoundationPorts are the machine's ports the foundation binds (novox/hq ADR 0100).
//
// **The node's, not the catalogue's.** A machine in use may already hold one — a predecessor's
// registry on 5000, its broker's management port — and a port fixed in the bundle and the manifests
// surfaces as a container that fails to bind, and one changed at genesis would be changed back when
// the foundation is adopted as modules. So each is an input here, checked free, rewritten into the
// bundle, and handed to the controller as that node's setting for the module that binds it.
type FoundationPorts struct {
Store int `json:"store"`
Bus int `json:"bus"`
AMQP int `json:"amqp"`
Management int `json:"management"`
Registry int `json:"registry"`
Packages int `json:"packages"`
Hub int `json:"hub"`
}
// DefaultPorts are the catalogue's numbers.
func DefaultPorts() FoundationPorts {
return FoundationPorts{Store: 5432, Bus: 5671, AMQP: 5672, Management: 15672, Registry: 5000,
Packages: 3000, Hub: 51820}
}
// DefaultOverlayRange is the controller's default private-network range.
const DefaultOverlayRange = "10.42.0.0/16"
// orDefaults fills every port left unsaid.
func (p FoundationPorts) orDefaults() FoundationPorts {
d := DefaultPorts()
for _, f := range []struct{ got, def *int }{
{&p.Store, &d.Store}, {&p.Bus, &d.Bus}, {&p.AMQP, &d.AMQP}, {&p.Management, &d.Management},
{&p.Registry, &d.Registry}, {&p.Packages, &d.Packages}, {&p.Hub, &d.Hub},
} {
if *f.got == 0 {
*f.got = *f.def
}
}
return p
}
// named is each port with what it is and its protocol, in a fixed order.
func (p FoundationPorts) named() []namedPort {
return []namedPort{
{"the store", "tcp", p.Store}, {"the bus", "tcp", p.Bus}, {"the broker's AMQP", "tcp", p.AMQP},
{"the broker's management", "tcp", p.Management}, {"the registry", "tcp", p.Registry},
{"the package registry", "tcp", p.Packages}, {"the private network's hub", "udp", p.Hub},
}
}
type namedPort struct {
what, protocol string
port int
}
// Check refuses a port out of range, or one port given for two things.
func (p FoundationPorts) Check() error {
seen := map[string]string{}
for _, n := range p.named() {
if n.port < 1 || n.port > 65535 {
return fmt.Errorf("%s's port is %d, and a port is 1-65535", n.what, n.port)
}
key := n.protocol + "/" + strconv.Itoa(n.port)
if other, twice := seen[key]; twice {
return fmt.Errorf("%s and %s were both given %s", other, n.what, key)
}
seen[key] = n.what
}
return nil
}
// moduleSettings is what each foundation module is told about its ports on this node: the port it
// declares, to the machine's port it is given. Only what differs from the catalogue — a converged
// genesis on the defaults sets nothing, and so changes nothing it did before.
func (p FoundationPorts) moduleSettings() map[string]map[string]int {
d := DefaultPorts()
out := map[string]map[string]int{}
add := func(module string, declared, given int) {
if given == declared {
return
}
if out[module] == nil {
out[module] = map[string]int{}
}
out[module][strconv.Itoa(declared)] = given
}
add("postgres", d.Store, p.Store)
add("lavinmq", d.Bus, p.Bus)
add("lavinmq", d.AMQP, p.AMQP)
add("lavinmq", d.Management, p.Management)
add(RegistryModule, d.Registry, p.Registry)
return out
}
// PortsSetting is the controller's settings key for a module's given ports.
const PortsSetting = "ports"
// setFoundationSettings tells the controller the ports this node gave a foundation module — and,
// on an adopted node, that the registry is reached from anywhere, as a node pulls from it before it
// has a private-network address (novox/hq ADR 0100). Done after the module is registered and before
// the push that raises it, so the first declaration already names the node's ports.
func setFoundationSettings(ctx context.Context, o Options, control controlPlane, module string,
say func(string)) error {
values := map[string]any{}
if ports := o.Ports.orDefaults().moduleSettings()[module]; len(ports) > 0 {
values[PortsSetting] = ports
}
if o.Adopted && module == RegistryModule {
values["expose"] = map[string]string{strconv.Itoa(DefaultPorts().Registry): "anywhere"}
}
if len(values) == 0 {
return nil
}
raw, err := json.Marshal(values)
if err != nil {
return err
}
remote := "/" + module + "-settings.json"
if err := control.carrying(ctx, module+"-settings.json", raw, remote); err != nil {
return err
}
if _, err := control.tell(ctx, "settings", "set", module, remote, "--node", o.Node); err != nil {
return err
}
say(" settings " + module + " on " + o.Node + ": " + string(raw))
return nil
}
// PortsRewrite says what RewritePorts changed.
type PortsRewrite struct {
Places int
}
// RewritePorts puts the node's foundation ports into the produced bundle, in place of the
// template's, byte for byte like every other rewrite — so the file keeps its comments and a person
// can read what was applied. A port left at its default is not touched, so a genesis on the
// defaults produces exactly the bundle it did before.
//
// Only the machine's side moves: the outer port of each mapping, the addresses the control plane
// dials on the machine's loopback, and the address nodes are told to dial. What a container listens
// on inside itself, and what an action reaches inside the store's own network, stay as they are.
func RewritePorts(r *Rewritten, p FoundationPorts, overlayRange string) (PortsRewrite, error) {
var out PortsRewrite
p = p.orDefaults()
d := DefaultPorts()
bundle := r.Bundle
var err error
replace := func(from, to, what string) {
if err != nil || from == to {
return
}
bundle, err = replaceOnce(bundle, from, to, what)
out.Places++
}
if p.Store != d.Store {
replace(`"ports": ["5432:5432"]`, fmt.Sprintf(`"ports": ["%d:5432"]`, p.Store), "the store's published port")
}
if p.Bus != d.Bus || p.AMQP != d.AMQP || p.Management != d.Management {
replace(`"ports": ["5671:5671", "5672:5672", "127.0.0.1:15672:15672"]`,
fmt.Sprintf(`"ports": ["%d:5671", "%d:5672", "127.0.0.1:%d:15672"]`, p.Bus, p.AMQP, p.Management),
"the broker's published ports")
}
if err != nil {
return out, err
}
// The control plane runs on the machine's network and dials the store and the broker on its
// loopback, so its connection strings name the machine's ports. The schema step reaches the
// store inside the store's own network and keeps the container's port — so these are found by
// the control plane's environment, not by searching for the text.
control, cerr := controlPlaneIn(r.Declaration)
if cerr != nil {
return out, cerr
}
for _, key := range sortedKeys(control.Env) {
value := control.Env[key]
now := value
now = strings.ReplaceAll(now, "@127.0.0.1:5432/", fmt.Sprintf("@127.0.0.1:%d/", p.Store))
now = strings.ReplaceAll(now, "@127.0.0.1:5672/", fmt.Sprintf("@127.0.0.1:%d/", p.AMQP))
if strings.HasSuffix(now, "@127.0.0.1:15672") {
now = strings.TrimSuffix(now, "15672") + strconv.Itoa(p.Management)
}
if key == brokerAddressVar {
if host, port, splitErr := net.SplitHostPort(value); splitErr == nil && port == "5671" {
now = net.JoinHostPort(host, strconv.Itoa(p.Bus))
}
}
if now != value {
replace(`"`+key+`": "`+value+`"`, `"`+key+`": "`+now+`"`, "the control plane's "+key)
}
}
if err != nil {
return out, err
}
// The foundation's own filter, where the template carries one: it admits the bus and the
// registry from anywhere, on whatever port they are.
for _, f := range []struct{ def, now int }{{d.Bus, p.Bus}, {d.Registry, p.Registry}} {
if f.def == f.now {
continue
}
for _, form := range []string{"tcp dport %d accept", "ct original proto-dst %d accept"} {
from, to := fmt.Sprintf(form, f.def), fmt.Sprintf(form, f.now)
if n := bytes.Count(bundle, []byte(from)); n > 0 {
bundle = bytes.ReplaceAll(bundle, []byte(from), []byte(to))
out.Places += n
}
}
}
// The private network's range, when it is not the default, is the controller's to know.
if overlayRange != "" && overlayRange != DefaultOverlayRange {
replace(`"`+brokerAddressVar+`": `,
`"MESH_OVERLAY_CIDR": "`+overlayRange+`",
"`+brokerAddressVar+`": `, "where the control plane is told the private network's range")
if err != nil {
return out, err
}
}
if out.Places == 0 {
return out, nil
}
parsed, perr := declaration.ParseFileTrusted(bundle)
if perr != nil {
return out, fmt.Errorf("the bundle stopped being a declaration after its ports were rewritten, which is this installer's fault: %w", perr)
}
r.Bundle, r.Declaration, r.Resources = bundle, parsed, len(parsed.Resources)
if c, cerr := controlPlaneIn(parsed); cerr == nil {
r.BrokerAddress = c.Env[brokerAddressVar]
}
return out, nil
}
// PortsFree refuses a foundation port something else already holds, naming what holds it. What the
// mesh itself raised on an earlier run of genesis is not counted: ours says which holders are.
func PortsFree(ctx context.Context, run Runner, p FoundationPorts, ours func(reachable.Reach) bool) error {
out, err := run(ctx, "ss", "-Hltunp")
if err != nil {
return fmt.Errorf("cannot read which ports this machine holds, so the foundation's cannot be checked free: %w", err)
}
sockets := reachable.Sockets(out)
var published []reachable.Reach
if ps, err := run(ctx, "docker", "ps", "--format", "{{.Names}}\t{{.Ports}}"); err == nil {
published = reachable.Published(ps)
}
held := reachable.Merge(sockets, published)
var problems []string
for _, n := range p.orDefaults().named() {
var by []string
for _, r := range held {
if r.Protocol != n.protocol || r.Port != n.port || ours(r) {
continue
}
holder := r.By
if holder == "" {
holder = "something ss does not name"
}
if r.Published {
holder = "the container " + r.By
}
if !contains(by, holder) {
by = append(by, holder)
}
}
if len(by) > 0 {
problems = append(problems, fmt.Sprintf("%s's port %s/%d is held by %s",
n.what, n.protocol, n.port, strings.Join(by, ", ")))
}
}
if len(problems) > 0 {
return fmt.Errorf("the foundation's ports must be free before anything is raised:\n - %s\n"+
"Give it another with the matching flag (--store-port, --bus-port, --amqp-port, "+
"--management-port, --registry-port, --packages-port, --hub-port); nothing was changed",
strings.Join(problems, "\n - "))
}
return nil
}
// OverlayClear refuses a private-network range that overlaps an address or a route the machine
// already has — a predecessor's tunnel still running — naming the interface. The mesh's own
// interface is not counted.
func OverlayClear(ctx context.Context, run Runner, overlayRange string) error {
if overlayRange == "" {
overlayRange = DefaultOverlayRange
}
_, mine, err := net.ParseCIDR(overlayRange)
if err != nil {
return fmt.Errorf("the private network's range %q is not a range: %w", overlayRange, err)
}
var clashes []string
if out, err := run(ctx, "ip", "-o", "addr", "show"); err == nil {
for _, line := range strings.Split(out, "\n") {
f := strings.Fields(line)
// 3: wg0 inet 10.42.0.1/24 scope global wg0
if len(f) < 4 || (f[2] != "inet" && f[2] != "inet6") {
continue
}
iface := strings.TrimSuffix(f[1], ":")
if clash(mine, f[3]) && iface != meshInterface {
clashes = append(clashes, fmt.Sprintf("%s holds %s", iface, f[3]))
}
}
} else {
return fmt.Errorf("cannot read this machine's addresses to check the private network's range: %w", err)
}
if out, err := run(ctx, "ip", "-o", "route", "show"); err == nil {
for _, line := range strings.Split(out, "\n") {
f := strings.Fields(line)
// 10.42.0.0/16 dev wg0 proto kernel scope link src 10.42.0.1
if len(f) < 3 || f[0] == "default" {
continue
}
iface := ""
for i := range f {
if f[i] == "dev" && i+1 < len(f) {
iface = f[i+1]
}
}
if iface != meshInterface && clash(mine, f[0]) {
clashes = append(clashes, fmt.Sprintf("%s routes %s", iface, f[0]))
}
}
}
if len(clashes) > 0 {
return fmt.Errorf("the private network's range %s overlaps what this machine already has: %s.\n"+
"A tunnel a predecessor still runs would take the mesh's traffic. Give another range with "+
"--overlay-range; nothing was changed", overlayRange, strings.Join(clashes, "; "))
}
return nil
}
// meshInterface is the private network's own interface, which a re-run finds holding its range.
const meshInterface = "mesh0"
func clash(mine *net.IPNet, other string) bool {
if !strings.Contains(other, "/") {
if ip := net.ParseIP(other); ip != nil {
return mine.Contains(ip)
}
return false
}
ip, theirs, err := net.ParseCIDR(other)
if err != nil {
return false
}
return mine.Contains(theirs.IP) || theirs.Contains(mine.IP) || mine.Contains(ip)
}
// NamesFree refuses a foundation or bundle container name that a container already has, when no
// host made that container and this node has no record of it — a predecessor's container under the
// mesh's name, which raising the foundation would replace.
func NamesFree(ctx context.Context, run Runner, names []string, known store.State) error {
var taken []string
sorted := append([]string{}, names...)
sort.Strings(sorted)
for _, name := range sorted {
out, err := run(ctx, "docker", "inspect", "--format",
"{{index .Config.Labels \"mesh-host.spec\"}}", name)
if err != nil {
continue // no such container
}
label := strings.TrimSpace(out)
if label != "" && label != "<no value>" {
continue
}
if known.Recorded(string(declaration.TypeContainer), name) {
continue
}
taken = append(taken, name)
}
if len(taken) > 0 {
return fmt.Errorf("this machine already runs a container under the name the foundation uses, "+
"and nothing of the mesh's made it: %s.\nRaising the foundation would replace it. Rename or "+
"stop it first; nothing was changed", strings.Join(taken, ", "))
}
return nil
}
// CheckTheMachine is every check genesis makes before raising anything that this machine does not
// already hold what the foundation needs: its ports, its private network's range, its containers'
// names (novox/hq ADR 0100). A re-run of genesis finds the foundation it raised and does not count
// it.
func CheckTheMachine(ctx context.Context, o Options, run Runner, bundle *declaration.Declaration,
say func(string)) error {
known, err := store.Load(o.State)
if err != nil {
return err
}
rerun := len(known.Resources) > 0
names := foundationNames(bundle)
mine := map[string]bool{}
for _, n := range names {
mine[n] = true
}
p := o.Ports.orDefaults()
ours := func(r reachable.Reach) bool {
switch {
case mine[r.By]:
return true
case !rerun:
return false
case r.By == "gitea" && r.Port == p.Packages:
// The package registry runs on the machine's network, so ss names its process.
return true
case r.By == "" && r.Protocol == "udp" && r.Port == p.Hub:
// The private network's hub is a kernel interface and has no process.
return true
}
return false
}
if err := PortsFree(ctx, run, p, ours); err != nil {
return err
}
say(fmt.Sprintf(" ports free store %d, bus %d, amqp %d, management %d, registry %d, packages %d, hub %d/udp",
p.Store, p.Bus, p.AMQP, p.Management, p.Registry, p.Packages, p.Hub))
if err := OverlayClear(ctx, run, o.OverlayRange); err != nil {
return err
}
if err := NamesFree(ctx, run, names, known); err != nil {
return err
}
return nil
}
// foundationNames are the containers the foundation and genesis raise under fixed names.
func foundationNames(bundle *declaration.Declaration) []string {
names := []string{ControlPlaneModule, giteaBootstrap, "mesh-registry"}
for _, n := range containerNames(bundle) {
if !contains(names, n) {
names = append(names, n)
}
}
sort.Strings(names)
return names
}
+286
View File
@@ -0,0 +1,286 @@
package bootstrap
import (
"context"
"errors"
"os"
"path/filepath"
"strings"
"testing"
"time"
"github.com/novox/mesh-host/internal/declaration"
"github.com/novox/mesh-host/internal/reachable"
"github.com/novox/mesh-host/internal/store"
)
// Defends novox/hq ADR 0100: the foundation's ports are the node's — inputs to genesis, checked free,
// rewritten into the bundle, and handed to the controller as the node's settings.
func producedBundle(t *testing.T) Rewritten {
t.Helper()
template, err := os.ReadFile("../../examples/foundation-first-node.lock")
if err != nil {
t.Skip("no example bundle beside this checkout")
}
r, err := Rewrite(template, "sha256:"+strings.Repeat("ab", 32))
if err != nil {
t.Fatal(err)
}
if _, err := RewriteRoot(&r, RootCredentials{Store: "s", Broker: "b"}); err != nil {
t.Fatal(err)
}
return r
}
func containerNamed(d *declaration.Declaration, name string) *declaration.Container {
for _, r := range d.Resources {
if c, ok := r.(*declaration.Container); ok && c.Name == name {
return c
}
}
return nil
}
func TestTheDefaultPortsLeaveTheBundleAsItWas(t *testing.T) {
r := producedBundle(t)
before := string(r.Bundle)
got, err := RewritePorts(&r, DefaultPorts(), DefaultOverlayRange)
if err != nil {
t.Fatal(err)
}
if got.Places != 0 || string(r.Bundle) != before {
t.Errorf("the default ports rewrote %d place(s)", got.Places)
}
}
func TestGivenPortsMoveOnlyTheMachinesSide(t *testing.T) {
r := producedBundle(t)
p := FoundationPorts{Store: 5433, Bus: 5771, AMQP: 5772, Management: 15673, Registry: 5100}
got, err := RewritePorts(&r, p, "10.77.0.0/16")
if err != nil {
t.Fatal(err)
}
if got.Places == 0 {
t.Fatal("nothing was rewritten")
}
storeC := containerNamed(r.Declaration, "mesh-store")
if len(storeC.Ports) != 1 || storeC.Ports[0] != "5433:5432" {
t.Errorf("the store publishes %v", storeC.Ports)
}
broker := containerNamed(r.Declaration, "mesh-broker")
if strings.Join(broker.Ports, " ") != "5771:5671 5772:5672 127.0.0.1:15673:15672" {
t.Errorf("the broker publishes %v", broker.Ports)
}
control, err := controlPlaneIn(r.Declaration)
if err != nil {
t.Fatal(err)
}
for key, value := range control.Env {
if strings.HasPrefix(key, "MESH_STORE_") && !strings.Contains(value, "@127.0.0.1:5433/") {
t.Errorf("%s still dials %s", key, value)
}
}
if !strings.Contains(control.Env["MESH_BROKER_AMQP"], "@127.0.0.1:5772/") ||
!strings.HasSuffix(control.Env["MESH_BROKER_MANAGEMENT"], "@127.0.0.1:15673") {
t.Errorf("the broker is dialled at %s and %s", control.Env["MESH_BROKER_AMQP"], control.Env["MESH_BROKER_MANAGEMENT"])
}
if control.Env["MESH_BROKER_ADDRESS"] != "192.0.2.10:5771" || r.BrokerAddress != "192.0.2.10:5771" {
t.Errorf("nodes are told to dial %s (%s)", control.Env["MESH_BROKER_ADDRESS"], r.BrokerAddress)
}
if control.Env["MESH_OVERLAY_CIDR"] != "10.77.0.0/16" {
t.Errorf("the control plane is told the range %q", control.Env["MESH_OVERLAY_CIDR"])
}
// The schema step reaches the store inside its own network, on the container's port.
text := string(r.Bundle)
if !strings.Contains(text, `MESH_STORE_INVENTORY=postgres://postgres:s@127.0.0.1:5432/inventory`) {
t.Error("the schema step's connection, inside the store's network, was moved off the container's port")
}
for _, want := range []string{"tcp dport 5771 accept", "ct original proto-dst 5771 accept",
"tcp dport 5100 accept", "ct original proto-dst 5100 accept"} {
if !strings.Contains(text, want) {
t.Errorf("the base filter does not say %q", want)
}
}
if strings.Contains(text, "dport 5671 accept") || strings.Contains(text, "dport 5000 accept") {
t.Error("the base filter still admits a default port")
}
}
func TestATemplateThatDoesNotSayItsPortsAsExpectedIsRefused(t *testing.T) {
r := producedBundle(t)
r.Bundle = []byte(strings.Replace(string(r.Bundle), `"ports": ["5432:5432"]`, `"ports": [ "5432:5432" ]`, 1))
if _, err := RewritePorts(&r, FoundationPorts{Store: 5433}, ""); err == nil {
t.Error("a store port the installer could not find was silently left")
}
}
func TestTwoThingsOnOnePortAreRefused(t *testing.T) {
p := DefaultPorts()
p.Registry = p.Store
if err := p.Check(); err == nil {
t.Error("the registry and the store were both given one port")
}
p = DefaultPorts()
p.Hub = 5432 // udp, beside the store's tcp: two different ports
if err := p.Check(); err != nil {
t.Errorf("a udp port beside a tcp one of the same number was refused: %v", err)
}
}
// machineRunner answers ss, docker ps, docker inspect and ip from fixtures.
type machineRunner struct {
ss, ps, addrs, routes string
unlabelled map[string]bool
labelled map[string]bool
}
func (m machineRunner) run(_ context.Context, name string, args ...string) (string, error) {
switch {
case name == "ss":
return m.ss, nil
case name == "docker" && args[0] == "ps":
return m.ps, nil
case name == "docker" && args[0] == "inspect":
n := args[len(args)-1]
if m.labelled[n] {
return "abc\n", nil
}
if m.unlabelled[n] {
return "\n", nil
}
return "", errors.New("no such container")
case name == "ip" && args[1] == "addr":
return m.addrs, nil
case name == "ip" && args[1] == "route":
return m.routes, nil
}
return "", nil
}
func noneOurs(reachable.Reach) bool { return false }
func TestABusyPortIsRefusedNamingItsHolder(t *testing.T) {
m := machineRunner{
ss: "tcp LISTEN 0 4096 0.0.0.0:5000 0.0.0.0:* users:((\"docker-proxy\",pid=1,fd=7))\n" +
"tcp LISTEN 0 4096 127.0.0.1:15672 0.0.0.0:* users:((\"beam.smp\",pid=2,fd=7))\n",
ps: "predecessor-registry\t0.0.0.0:5000->5000/tcp\n",
}
err := PortsFree(context.Background(), m.run, DefaultPorts(), noneOurs)
if err == nil {
t.Fatal("held ports were not refused")
}
for _, want := range []string{"predecessor-registry", "beam.smp", "tcp/5000", "tcp/15672", "--registry-port"} {
if !strings.Contains(err.Error(), want) {
t.Errorf("the refusal does not say %q: %v", want, err)
}
}
p := DefaultPorts()
p.Registry, p.Management = 5100, 15673
if err := PortsFree(context.Background(), m.run, p, noneOurs); err != nil {
t.Errorf("other ports given and still refused: %v", err)
}
}
func TestTheFoundationsOwnContainersAreNotCountedOnARerun(t *testing.T) {
m := machineRunner{
ss: "tcp LISTEN 0 4096 0.0.0.0:5432 0.0.0.0:* users:((\"docker-proxy\",pid=1,fd=7))\n",
ps: "mesh-store\t0.0.0.0:5432->5432/tcp\n",
}
ours := func(r reachable.Reach) bool { return r.By == "mesh-store" }
if err := PortsFree(context.Background(), m.run, DefaultPorts(), ours); err != nil {
t.Errorf("the foundation's own store was counted as holding its port: %v", err)
}
}
func TestAnOverlappingTunnelIsRefusedNamingItsInterface(t *testing.T) {
m := machineRunner{
addrs: "1: lo inet 127.0.0.1/8 scope host lo\n5: wg0 inet 10.42.3.1/24 scope global wg0\n7: mesh0 inet 10.42.0.1/16 scope global mesh0\n",
routes: "default via 192.0.2.1 dev eth0\n10.42.3.0/24 dev wg0 proto kernel scope link src 10.42.3.1\n",
}
err := OverlayClear(context.Background(), m.run, "")
if err == nil || !strings.Contains(err.Error(), "wg0") || strings.Contains(err.Error(), "mesh0") {
t.Fatalf("the overlap was not named by its interface alone: %v", err)
}
if err := OverlayClear(context.Background(), m.run, "10.77.0.0/16"); err != nil {
t.Errorf("a clear range was refused: %v", err)
}
}
func TestAPredecessorsContainerUnderTheMeshsNameIsRefused(t *testing.T) {
m := machineRunner{unlabelled: map[string]bool{"mesh-registry": true}, labelled: map[string]bool{"mesh-store": true}}
err := NamesFree(context.Background(), m.run, []string{"mesh-store", "mesh-registry", "mesh-broker"}, store.State{})
if err == nil || !strings.Contains(err.Error(), "mesh-registry") || strings.Contains(err.Error(), "mesh-store") {
t.Fatalf("names: %v", err)
}
known := store.State{Resources: []store.Applied{{ID: "x", Type: "container", Target: "mesh-registry"}}}
if err := NamesFree(context.Background(), m.run, []string{"mesh-registry"}, known); err != nil {
t.Errorf("a container this node has a record of was refused: %v", err)
}
}
// controlRecorder is a control plane that answers everything and writes down what it was told,
// with the content of every settings file carried to it.
type controlRecorder struct {
told []string
settings map[string]string
}
func (c *controlRecorder) run(_ context.Context, name string, args ...string) (string, error) {
if name == "docker" && args[0] == "cp" {
raw, _ := os.ReadFile(args[1])
if strings.HasSuffix(args[2], "-settings.json") {
c.settings[filepath.Base(args[2])] = string(raw)
}
return "", nil
}
if name == "docker" && args[0] == "exec" {
c.told = append(c.told, strings.Join(args[3:], " "))
}
return "", nil
}
func (c *controlRecorder) index(prefix string) int {
for i, t := range c.told {
if strings.HasPrefix(t, prefix) {
return i
}
}
return -1
}
func TestTheNodesPortsAreSetBeforeTheModuleIsPushed(t *testing.T) {
t.Setenv("TMPDIR", t.TempDir())
c := &controlRecorder{settings: map[string]string{}}
control := controlPlane{container: "temp-mesh-controller", run: c.run, timeout: time.Second}
o := Options{Node: "anchor", Ports: FoundationPorts{Registry: 5100}, Wait: time.Second}
if _, err := installModule(context.Background(), o, control, RegistryModule, []byte(`{}`), quietly); err != nil {
t.Fatal(err)
}
set, push := c.index("settings set distribution"), c.index("push anchor")
if set < 0 || push < 0 || set > push {
t.Fatalf("settings were not set before the push: %v", c.told)
}
if add := c.index("module add"); add > set {
t.Errorf("settings were set before the module existed: %v", c.told)
}
if !strings.Contains(c.told[set], "--node anchor") {
t.Errorf("the settings are not the node's: %s", c.told[set])
}
if got := c.settings["distribution-settings.json"]; got != `{"ports":{"5000":5100}}` {
t.Errorf("the registry was told %s", got)
}
}
func TestAGenesisOnTheDefaultsSetsNoSettings(t *testing.T) {
t.Setenv("TMPDIR", t.TempDir())
c := &controlRecorder{settings: map[string]string{}}
control := controlPlane{container: "temp-mesh-controller", run: c.run, timeout: time.Second}
o := Options{Node: "anchor", Wait: time.Second}
if _, err := installModule(context.Background(), o, control, RegistryModule, []byte(`{}`), quietly); err != nil {
t.Fatal(err)
}
if c.index("settings") >= 0 {
t.Errorf("a converged genesis on the default ports set settings: %v", c.told)
}
}