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