Files
mesh-host/internal/bootstrap/ports.go
T

461 lines
16 KiB
Go

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
}
if name == giteaBootstrap && len(known.Resources) > 0 {
// Genesis raises the package registry itself, by hand and before the host records
// anything of it, so on a re-run it is found under its own name with no label and no
// record. A machine that carries what an earlier genesis raised made it.
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
}