From 3964d9da0ab1ea18e226244abcd23cf3620695fd Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 17:28:36 +0200 Subject: [PATCH] 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) --- cmd/mesh-bootstrap/main.go | 61 +++- cmd/mesh-bootstrap/main_test.go | 26 ++ internal/bootstrap/bootstrap.go | 41 ++- internal/bootstrap/module.go | 10 + internal/bootstrap/phase2.go | 10 +- internal/bootstrap/phase2_test.go | 4 +- internal/bootstrap/phase3.go | 3 + internal/bootstrap/phase_packages.go | 21 +- internal/bootstrap/ports.go | 454 +++++++++++++++++++++++++++ internal/bootstrap/ports_test.go | 286 +++++++++++++++++ 10 files changed, 902 insertions(+), 14 deletions(-) create mode 100644 internal/bootstrap/ports.go create mode 100644 internal/bootstrap/ports_test.go diff --git a/cmd/mesh-bootstrap/main.go b/cmd/mesh-bootstrap/main.go index e48aa1a..7a1fba5 100644 --- a/cmd/mesh-bootstrap/main.go +++ b/cmd/mesh-bootstrap/main.go @@ -25,6 +25,7 @@ import ( "net/http" "os" "os/signal" + "strconv" "syscall" "time" @@ -126,6 +127,14 @@ const usage = `mesh-bootstrap — make a bare machine into a mesh --packet-filter which packet filter to run (nftables) --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 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. @@ -176,7 +185,9 @@ func parseArgs(args []string) (string, bootstrap.Options, bool, error) { HostService: defaultService, // 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. - 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 // first `initdb`-shaped wait, are both minutes rather than seconds. Wait: 3 * time.Minute, @@ -210,9 +221,38 @@ func parseArgs(args []string) (string, bootstrap.Options, bool, error) { return "", opts, false, fmt.Errorf( "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 } +// 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 { set := flag.NewFlagSet("mesh-bootstrap", flag.ContinueOnError) 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, "what of it to build (default main)") 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 { opts.Answers = map[string]string{} } diff --git a/cmd/mesh-bootstrap/main_test.go b/cmd/mesh-bootstrap/main_test.go index 2fadbcb..2a4a4bb 100644 --- a/cmd/mesh-bootstrap/main_test.go +++ b/cmd/mesh-bootstrap/main_test.go @@ -164,3 +164,29 @@ func TestTheNodeNameCanBeSaid(t *testing.T) { 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") + } +} diff --git a/internal/bootstrap/bootstrap.go b/internal/bootstrap/bootstrap.go index 9d3d787..95b86b4 100644 --- a/internal/bootstrap/bootstrap.go +++ b/internal/bootstrap/bootstrap.go @@ -185,6 +185,20 @@ type Options struct { Prompt func(Choice) (string, error) // Extras are catalogue modules beyond the floor, asked for by name. 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. @@ -287,6 +301,14 @@ type Result struct { // Stopped names why a run went no further. Empty on a run that pivoted. 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. @@ -339,7 +361,12 @@ func Run(ctx context.Context, o Options, d Deps, say func(string)) (Result, erro if say == nil { 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 ------------------------------------------------------------------- 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 { 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. say = Masking(say, creds) root, err := RewriteRoot(&rewritten, creds) if err != nil { 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 { what, path string made bool diff --git a/internal/bootstrap/module.go b/internal/bootstrap/module.go index 4ba2729..92c4116 100644 --- a/internal/bootstrap/module.go +++ b/internal/bootstrap/module.go @@ -124,6 +124,9 @@ func registerAndAssign(ctx context.Context, o Options, control controlPlane, mod // the only place the reason appears. say(indent(refusal)) } + if err := prepareModule(ctx, o, control, module, say); err != nil { + return out, err + } return out, nil } @@ -209,3 +212,10 @@ func pinPlaceholder(manifest []byte, reference, module string) ([]byte, int, err } 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) +} diff --git a/internal/bootstrap/phase2.go b/internal/bootstrap/phase2.go index 4eb21a3..431f12c 100644 --- a/internal/bootstrap/phase2.go +++ b/internal/bootstrap/phase2.go @@ -5,6 +5,7 @@ import ( "encoding/json" "fmt" "net" + "strconv" "strings" "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 { return err } + if err := prepareModule(ctx, o, control, module, say); err != nil { + return err + } if _, err := pushNode(ctx, o, control, say); err != nil { return err } @@ -120,7 +124,7 @@ func PlaceOnTheNetwork(ctx context.Context, o Options, control controlPlane, Name: "endpoint", Question: "Where do other machines reach this one for the private network? " + "(host:port; the host other machines dial)", - Default: derivedEndpoint(brokerAddress), + Default: derivedEndpoint(brokerAddress, o.Ports.orDefaults().Hub), }, o.Answers["endpoint"], o.Prompt, say) if err != nil { 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 // 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) if err != nil || host == "" { return "" } - return net.JoinHostPort(host, "51820") + return net.JoinHostPort(host, strconv.Itoa(hub)) } func refOr(ref string) string { diff --git a/internal/bootstrap/phase2_test.go b/internal/bootstrap/phase2_test.go index cd212fc..918a553 100644 --- a/internal/bootstrap/phase2_test.go +++ b/internal/bootstrap/phase2_test.go @@ -5,10 +5,10 @@ import "testing" // 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. 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) } - if got := derivedEndpoint(""); got != "" { + if got := derivedEndpoint("", 51820); got != "" { t.Fatalf("an endpoint was invented from nothing: %q", got) } } diff --git a/internal/bootstrap/phase3.go b/internal/bootstrap/phase3.go index e39b9e3..2d60811 100644 --- a/internal/bootstrap/phase3.go +++ b/internal/bootstrap/phase3.go @@ -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 { return err } + if err := prepareModule(ctx, o, control, module, say); err != nil { + return err + } if beforePush != nil { if err := beforePush(); err != nil { return err diff --git a/internal/bootstrap/phase_packages.go b/internal/bootstrap/phase_packages.go index 1a4ceed..cd82e6e 100644 --- a/internal/bootstrap/phase_packages.go +++ b/internal/bootstrap/phase_packages.go @@ -42,8 +42,9 @@ const ( // giteaDBRole/giteaDBName is gitea's own database in the foundation store. giteaDBRole = "mesh_gitea" giteaDBName = "mesh_gitea" - // giteaPort is where the raised server answers on the machine. - giteaPort = 3000 + // defaultGiteaPort is where the raised server answers on the machine unless the node gave the + // 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: @@ -60,17 +61,18 @@ func RaisePackageRegistry(ctx context.Context, o Options, d Deps, control contro } say(" seeding gitea's database in the foundation store") + ports := o.Ports.orDefaults() if err := seedGiteaDatabase(ctx, run, o.Timeout, dbPassword, say); err != nil { return err } 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 } 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 { 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 // 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, - say func(string)) error { + ports FoundationPorts, say func(string)) error { asking, cancel := context.WithTimeout(ctx, timeout) defer cancel() @@ -164,7 +166,7 @@ func raiseGiteaServer(ctx context.Context, run Runner, timeout time.Duration, db env := []string{ "-e", "GITEA__database__DB_TYPE=postgres", // 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__USER=" + giteaDBRole, "-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 // 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. - "-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", } + 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{ "run", "-d", "--name", giteaBootstrap, // Host network, like the control plane: it reaches the foundation store on the machine's diff --git a/internal/bootstrap/ports.go b/internal/bootstrap/ports.go new file mode 100644 index 0000000..4da7560 --- /dev/null +++ b/internal/bootstrap/ports.go @@ -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 != "" { + 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 +} diff --git a/internal/bootstrap/ports_test.go b/internal/bootstrap/ports_test.go new file mode 100644 index 0000000..49a3e64 --- /dev/null +++ b/internal/bootstrap/ports_test.go @@ -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) + } +}