Take a module, converge a node after a preview, and return it to adopted, from the command line and the API (hq ADR 0100)
This commit is contained in:
@@ -2,10 +2,15 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"slices"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
)
|
||||
|
||||
@@ -73,3 +78,340 @@ func adoptedNodes(nodes []inventory.Node) []string {
|
||||
}
|
||||
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, fmt.Sprintf(" %s %s (%s)", h.Kind, h.Target, h.ID))
|
||||
}
|
||||
}
|
||||
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
|
||||
}
|
||||
|
||||
// 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.
|
||||
func converge(ctx context.Context, open *stores, node string, yes bool, filter string) (string, error) {
|
||||
inv := open.inventory
|
||||
if filter == "" {
|
||||
filter = DefaultFilter
|
||||
}
|
||||
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)
|
||||
}
|
||||
|
||||
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"))
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
preview := previewOf(node, reported, rules, with.Foundation, plan, taken, filter, runs[filter])
|
||||
if !yes {
|
||||
return preview + fmt.Sprintf("\n\nNothing has changed. Run `converge %s --yes` to do it.", node), nil
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
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.
|
||||
func previewOf(node string, reported inventory.Adoption, rules []catalogue.Rule, foundation []int,
|
||||
plan catalogue.Resolution, taken []string, filter string, filterAssigned bool) 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)
|
||||
}
|
||||
fmt.Fprintf(&b, " %-44s %s\n", what, fate(r, rules, foundation))
|
||||
}
|
||||
if len(reported.Reachable) == 0 {
|
||||
b.WriteString(" nothing reported\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)
|
||||
for _, h := range reported.Held {
|
||||
if h.Module == m && h.Kind == "file" {
|
||||
fmt.Fprintf(&b, " replacing the found file %s", h.Target)
|
||||
if h.Kept != "" {
|
||||
fmt.Fprintf(&b, " (original kept at %s)", h.Kept)
|
||||
}
|
||||
b.WriteString("\n")
|
||||
}
|
||||
}
|
||||
}
|
||||
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)
|
||||
}
|
||||
return strings.TrimRight(b.String(), "\n")
|
||||
}
|
||||
|
||||
// fate is what the derived filter does to one reachable thing: which module declares it and from
|
||||
// where, or that it will close.
|
||||
func fate(r inventory.Reach, rules []catalogue.Rule, foundation []int) string {
|
||||
if r.Protocol == "tcp" && r.Port == catalogue.SSHPort {
|
||||
return "stays open — ssh is never closed"
|
||||
}
|
||||
for _, port := range foundation {
|
||||
if r.Protocol == "tcp" && r.Port == port {
|
||||
return "stays open — the mesh's own, from anywhere"
|
||||
}
|
||||
}
|
||||
for _, rule := range rules {
|
||||
if rule.Port != r.Port || rule.Protocol != r.Protocol {
|
||||
continue
|
||||
}
|
||||
if rule.From == catalogue.FromMachine {
|
||||
return fmt.Sprintf("WILL CLOSE to the network — declared by %s for this machine only",
|
||||
strings.Join(rule.Because, ", "))
|
||||
}
|
||||
return fmt.Sprintf("declared by %s (from %s)", strings.Join(rule.Because, ", "), rule.From)
|
||||
}
|
||||
return "WILL CLOSE — no module assigned here declares it"
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
if err := inv.SetAdopted(ctx, node, true); err != nil {
|
||||
return "", err
|
||||
}
|
||||
said := fmt.Sprintf("%s is adopted; what was taken on it stays taken", node)
|
||||
if err := sendNodes(ctx, open, []string{node}); err != nil {
|
||||
return said + "\n and it could not be sent: run `push " + node + "`", err
|
||||
}
|
||||
return said + "\n sent: the host unloads the mesh's filter and enables the firewall it found", nil
|
||||
}
|
||||
|
||||
// takeCommand, convergeCommand and adoptCommand are the command line's adapters to the acts above.
|
||||
func takeCommand(ctx context.Context, args []string) error {
|
||||
if len(args) != 2 {
|
||||
return errors.New("take <node> <module>")
|
||||
}
|
||||
return runAct(ctx, func(open *stores) (string, error) { return take(ctx, open, args[0], args[1]) })
|
||||
}
|
||||
|
||||
func convergeCommand(ctx context.Context, args []string) error {
|
||||
set := flag.NewFlagSet("converge", flag.ContinueOnError)
|
||||
yes := set.Bool("yes", false, "do it; without it, only the preview")
|
||||
filter := set.String("filter", DefaultFilter, "the module that loads the mesh's filter")
|
||||
positionals, err := parseAround(set, args)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(positionals) != 1 {
|
||||
return errors.New("converge <node> [--yes] [--filter nftables]")
|
||||
}
|
||||
return runAct(ctx, func(open *stores) (string, error) {
|
||||
return converge(ctx, open, positionals[0], *yes, *filter)
|
||||
})
|
||||
}
|
||||
|
||||
func adoptCommand(ctx context.Context, args []string) error {
|
||||
if len(args) != 1 {
|
||||
return errors.New("adopt <node>")
|
||||
}
|
||||
return runAct(ctx, func(open *stores) (string, error) { return adopt(ctx, open, args[0]) })
|
||||
}
|
||||
|
||||
func runAct(ctx context.Context, act func(*stores) (string, error)) error {
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer open.Close()
|
||||
said, err := act(open)
|
||||
if said != "" {
|
||||
fmt.Println(said)
|
||||
}
|
||||
if err != nil && said != "" {
|
||||
fmt.Println()
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user