diff --git a/cmd/mesh-controller/adopting_test.go b/cmd/mesh-controller/adopting_test.go new file mode 100644 index 0000000..1936b9a --- /dev/null +++ b/cmd/mesh-controller/adopting_test.go @@ -0,0 +1,255 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "reflect" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" + "github.com/novox/mesh-controller/internal/overlay" +) + +// novox/hq ADR 0100: taking a module is its cutover; converging a node is one act, previewed, and +// refused while a found container is held; returning to adopted keeps what was taken. + +// anAdoptedAnchor is aMesh with the anchor adopted, running a served module the predecessor also +// runs and a module with only a file, and a filter module in the catalogue. +func anAdoptedAnchor(t *testing.T) (*stores, *[]string) { + t.Helper() + open := aMesh(t) + ctx := t.Context() + register(t, open, catalogue.Manifest{Module: "hello-web", Version: "1", + Listens: []catalogue.Listening{{Port: 8080, From: catalogue.FromEverywhere}}, + Resources: []map[string]any{ + {"id": "page", "type": "file", "path": "/var/lib/hello-web/index.html", "content": "hi"}, + {"id": "server", "type": "container", "name": "hello-web", "ports": []any{"8080:80"}, + "image": "registry.example/hello@sha256:" + strings.Repeat("a", 64)}, + }}) + register(t, open, catalogue.Manifest{Module: "notes", Version: "1", + Resources: []map[string]any{ + {"id": "conf", "type": "file", "path": "/etc/notes.conf", "content": "x"}, + }}) + register(t, open, catalogue.Manifest{Module: "nftables", Version: "1", + Filtering: &catalogue.Filtering{Into: "/etc/nftables.conf"}, + Resources: []map[string]any{{"id": "load", "type": "service", "unit": "mesh-filter.service", + "state": "running", "restart-on": []any{"filtering"}}}}) + if err := open.inventory.SetAdopted(ctx, "anchor", true); err != nil { + t.Fatal(err) + } + for _, m := range []string{"hello-web", "notes"} { + if _, err := assign(ctx, open, "anchor", m); err != nil { + t.Fatal(err) + } + } + sent := &[]string{} + saved := sendNodes + sendNodes = func(_ context.Context, _ *stores, names []string) error { + *sent = append(*sent, names...) + return nil + } + t.Cleanup(func() { sendNodes = saved }) + return open, sent +} + +// reportsHolding has the anchor report, on what it was last sent, holding what is given. +func reportsHolding(t *testing.T, open *stores, held ...link.Held) { + t.Helper() + ctx := t.Context() + body, err := composed(t, open, "anchor").Body() + if err != nil { + t.Fatal(err) + } + record, err := open.inventory.NodeByName(ctx, "anchor") + if err != nil { + t.Fatal(err) + } + if err := open.inventory.RecordSent(ctx, record.ID, digestOf(body)); err != nil { + t.Fatal(err) + } + if err := (link.Enrolment{Inventory: open.inventory}).Heard(ctx, link.Report{ + Node: "anchor", Applied: []string{"hello-web.x"}, Declared: digestOf(body), + Firewall: "ufw", Held: held, + Reachable: []link.Reach{ + {Protocol: "tcp", Address: "0.0.0.0", Port: 22, By: "sshd"}, + {Protocol: "tcp", Address: "0.0.0.0", Port: 8080, By: "hello-web", Published: true, + ContainerPort: 80}, + {Protocol: "tcp", Address: "0.0.0.0", Port: 5000, By: "predecessor-registry", + Published: true, ContainerPort: 5000}, + {Protocol: "tcp", Address: "127.0.0.1", Port: 15672, By: "mesh-broker", + Published: true, ContainerPort: 15672}, + }, + }); err != nil { + t.Fatal(err) + } +} + +var ( + heldContainer = link.Held{ID: "hello-web.server", Module: "hello-web", Kind: "container", + Target: "hello-web", Since: time.Now()} + heldFile = link.Held{ID: "notes.conf", Module: "notes", Kind: "file", + Target: "/etc/notes.conf", Since: time.Now(), Kept: "/var/lib/mesh-host/kept/abc-notes.conf"} +) + +func TestTakingAModuleNotOnTheNodeIsRefused(t *testing.T) { + open, _ := anAdoptedAnchor(t) + if _, err := take(t.Context(), open, "anchor", "nftables"); !errors.Is(err, inventory.ErrNotAssigned) { + t.Fatalf("taking an unassigned module gave %v", err) + } + if _, err := take(t.Context(), open, "laptop", "network"); !errors.Is(err, inventory.ErrNotAdopted) { + t.Fatalf("taking on a converged node gave %v", err) + } +} + +func TestConvergingIsRefusedOnAPreviewThatWouldBeStale(t *testing.T) { + open, sent := anAdoptedAnchor(t) + _, err := converge(t.Context(), open, "anchor", false, "") + if err == nil || !strings.Contains(err.Error(), "has not reported") { + t.Fatalf("a node that never reported was previewed: %v", err) + } + if len(*sent) != 0 { + t.Fatal("a refused converge sent something") + } +} + +func TestTheFlipIsRefusedWhileAFoundContainerIsHeld(t *testing.T) { + open, sent := anAdoptedAnchor(t) + reportsHolding(t, open, heldContainer, heldFile) + _, err := converge(t.Context(), open, "anchor", true, "") + if err == nil || !strings.Contains(err.Error(), "take anchor hello-web once its data has moved") { + t.Fatalf("the flip was not refused while hello-web holds its found container: %v", err) + } + if n, _ := open.inventory.NodeByName(t.Context(), "anchor"); !n.Adopted || len(*sent) != 0 { + t.Fatal("a refused flip changed something") + } +} + +func TestTakingNamesWhatItReplaces(t *testing.T) { + open, _ := anAdoptedAnchor(t) + reportsHolding(t, open, heldContainer, heldFile) + said, err := take(t.Context(), open, "anchor", "hello-web") + if err != nil { + t.Fatal(err) + } + if !strings.Contains(said, "container hello-web (hello-web.server)") || + !strings.Contains(said, "push anchor") { + t.Fatalf("taking did not say what it replaces and what to run:\n%s", said) + } +} + +func TestConvergingPreviewsThenChangesAndAdoptingKeepsWhatWasTaken(t *testing.T) { + open, sent := anAdoptedAnchor(t) + ctx := t.Context() + if _, err := take(ctx, open, "anchor", "hello-web"); err != nil { + t.Fatal(err) + } + reportsHolding(t, open, heldFile) + + preview, err := converge(ctx, open, "anchor", false, "") + if err != nil { + t.Fatal(err) + } + for _, want := range []string{ + "tcp/8080 hello-web (published, container port 80)", + "declared by hello-web (from anywhere)", + "WILL CLOSE — no module assigned here declares it", + "ssh is never closed", + "notes\n replacing the found file /etc/notes.conf (original kept at", + "assigns nftables", + "the found firewall (ufw) is disabled, never flushed", + } { + if !strings.Contains(preview, want) { + t.Errorf("the preview does not say %q:\n%s", want, preview) + } + } + if strings.Contains(preview, "15672") { + t.Errorf("a loopback listener is in the preview:\n%s", preview) + } + for _, line := range strings.Split(preview, "\n") { + if strings.Contains(line, "5000") && !strings.Contains(line, "WILL CLOSE") { + t.Errorf("an undeclared published port is not said to close: %s", line) + } + } + if n, _ := open.inventory.NodeByName(ctx, "anchor"); !n.Adopted || len(*sent) != 0 { + t.Fatal("the preview changed something") + } + + if _, err := converge(ctx, open, "anchor", true, ""); err != nil { + t.Fatal(err) + } + n, _ := open.inventory.NodeByName(ctx, "anchor") + if n.Adopted { + t.Fatal("converge --yes left the node adopted") + } + taken, _ := open.inventory.Taken(ctx, "anchor") + if !reflect.DeepEqual(taken, []string{"hello-web", overlay.Name, "nftables", "notes"}) { + t.Fatalf("the flip took %v", taken) + } + if !reflect.DeepEqual(*sent, []string{"anchor"}) { + t.Fatalf("the flip sent %v", *sent) + } + declared := composed(t, open, "anchor") + if declared.Adoption != nil { + t.Fatal("a converged node is still sent an adoption envelope") + } + if !hasID(declared.Resources, "nftables.filtering") { + t.Fatal("the converged node is not declared the mesh's filter") + } + + if _, err := adopt(ctx, open, "anchor"); err != nil { + t.Fatal(err) + } + if _, err := adopt(ctx, open, "anchor"); err == nil { + t.Fatal("adopting an adopted node was not refused") + } + again, _ := open.inventory.Taken(ctx, "anchor") + if !reflect.DeepEqual(again, taken) { + t.Fatalf("returning to adopted lost what was taken: %v", again) + } + declared = composed(t, open, "anchor") + if declared.Adoption == nil || len(declared.Adoption.Untaken) != 0 { + t.Fatalf("returned to adopted, the envelope is %+v", declared.Adoption) + } + if hasID(declared.Resources, "nftables.filtering") || hasID(declared.Resources, "nftables.load") { + t.Fatal("returned to adopted, the mesh's filter is still declared") + } +} + +func hasID(resources []map[string]any, id string) bool { + for _, r := range resources { + if r["id"] == id { + return true + } + } + return false +} + +// The command API refuses a flip exactly as the command line does, in the same words. +func TestTheApiRefusesTheFlipInTheCommandLinesWords(t *testing.T) { + open, _ := anAdoptedAnchor(t) + reportsHolding(t, open, heldContainer) + _, direct := converge(t.Context(), open, "anchor", true, "") + if direct == nil { + t.Fatal("the flip was not refused") + } + got := asking(t, letIn{}, "POST", "/converge", `{"node":"anchor","yes":true}`) + if got.Code != http.StatusConflict { + t.Fatalf("got %d: %s", got.Code, got.Body.String()) + } + var said map[string]any + if err := json.Unmarshal(got.Body.Bytes(), &said); err != nil { + t.Fatal(err) + } + if said["refused"] != direct.Error() { + t.Fatalf("the API said %q and the command line %q", said["refused"], direct.Error()) + } + if got := asking(t, letIn{}, "POST", "/take", `{"node":"anchor"}`); got.Code != http.StatusBadRequest { + t.Fatalf("a take naming no module got %d", got.Code) + } +} diff --git a/cmd/mesh-controller/adoption.go b/cmd/mesh-controller/adoption.go index eec29d4..bd03494 100644 --- a/cmd/mesh-controller/adoption.go +++ b/cmd/mesh-controller/adoption.go @@ -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 ") + } + 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 [--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 ") + } + 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 +} diff --git a/cmd/mesh-controller/api.go b/cmd/mesh-controller/api.go index 1d97030..182c7c9 100644 --- a/cmd/mesh-controller/api.go +++ b/cmd/mesh-controller/api.go @@ -94,19 +94,29 @@ func (n notYet) Who(*http.Request) (string, error) { func commands(who Authenticator) http.Handler { mux := http.NewServeMux() - mux.HandleFunc("POST /assign", acting(who, func(ctx context.Context, open *stores, in request) (string, error) { + mux.HandleFunc("POST /assign", acting(who, true, func(ctx context.Context, open *stores, in request) (string, error) { return assign(ctx, open, in.Node, in.Module) })) - mux.HandleFunc("POST /unassign", acting(who, func(ctx context.Context, open *stores, in request) (string, error) { + mux.HandleFunc("POST /unassign", acting(who, true, func(ctx context.Context, open *stores, in request) (string, error) { return unassign(ctx, open, in.Node, in.Module) })) + // Adoption (novox/hq ADR 0100): the same acts as `take`, `converge` and `adopt`. + mux.HandleFunc("POST /take", acting(who, true, func(ctx context.Context, open *stores, in request) (string, error) { + return take(ctx, open, in.Node, in.Module) + })) + mux.HandleFunc("POST /converge", acting(who, false, func(ctx context.Context, open *stores, in request) (string, error) { + return converge(ctx, open, in.Node, in.Yes, in.Filter) + })) + mux.HandleFunc("POST /adopt", acting(who, false, func(ctx context.Context, open *stores, in request) (string, error) { + return adopt(ctx, open, in.Node) + })) // Anything else is said plainly, because a command surface answering 404 to a verb somebody // expected is indistinguishable from one that is down. mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { refuse(w, http.StatusNotFound, fmt.Errorf( - "%s %s is not something this mesh can be asked; it accepts POST /assign and "+ - "POST /unassign", r.Method, r.URL.Path)) + "%s %s is not something this mesh can be asked; it accepts POST /assign, "+ + "POST /unassign, POST /take, POST /converge and POST /adopt", r.Method, r.URL.Path)) }) return mux } @@ -114,11 +124,16 @@ func commands(who Authenticator) http.Handler { type request struct { Node string `json:"node"` Module string `json:"module"` + // Yes and Filter are converge's: do it rather than preview it, and which module loads the + // mesh's filter. + Yes bool `json:"yes,omitempty"` + Filter string `json:"filter,omitempty"` } // acting is the shape every route shares: authenticate, read, act, answer. func acting( who Authenticator, + needsModule bool, do func(context.Context, *stores, request) (string, error), ) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { @@ -131,8 +146,12 @@ func acting( refuse(w, http.StatusBadRequest, fmt.Errorf("this is not a request this understands: %w", err)) return } - if in.Node == "" || in.Module == "" { - refuse(w, http.StatusBadRequest, errors.New(`both "node" and "module" are needed`)) + if in.Node == "" || (needsModule && in.Module == "") { + if needsModule { + refuse(w, http.StatusBadRequest, errors.New(`both "node" and "module" are needed`)) + } else { + refuse(w, http.StatusBadRequest, errors.New(`"node" is needed`)) + } return } diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index 00cc99b..0fc5b2c 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -98,6 +98,12 @@ func run() error { return moduleCommand(ctx, args[1:]) case "assign", "unassign": return assignCommand(ctx, args[0], args[1:]) + case "take": + return takeCommand(ctx, args[1:]) + case "converge": + return convergeCommand(ctx, args[1:]) + case "adopt": + return adoptCommand(ctx, args[1:]) case "settings": return settingsCommand(ctx, args[1:]) case "secret": @@ -155,6 +161,9 @@ func usage() { api --issuer URL [--listen A] assign and unassign over http, for a surface that is not here assign put a module on a node unassign take it off + take cut a module over on an adopted node, once its data has moved + converge [--yes] [--filter nftables] preview, then make, an adopted node converged + adopt return a converged node to adopted; what was taken stays taken settings set what a module's config should say, for the whole mesh settings set --node ...or for one machine settings clear [--node ] take a layer away diff --git a/cmd/mesh-controller/plan.go b/cmd/mesh-controller/plan.go index 3781267..919817f 100644 --- a/cmd/mesh-controller/plan.go +++ b/cmd/mesh-controller/plan.go @@ -318,10 +318,30 @@ const ( func declarationWith(ctx context.Context, open *stores, node string, plan catalogue.Resolution, settings catalogue.SettingsBy, gens map[string]catalogue.Generator, choosing Choosing) (sendable, error) { + with, record, err := renderingFor(ctx, open, node, plan, settings, gens, choosing) + if err != nil { + return sendable{}, err + } + composed, err := plan.Compose(with) + if err != nil { + return sendable{}, err + } + // And what was taken on it, said in every declaration it is sent from this one place. + adoption, err := adoptionOf(ctx, open.inventory, record, plan, composed) + if err != nil { + return sendable{}, err + } + return sendable{Resources: composed.Resources, Adoption: adoption}, nil +} + +// renderingFor is everything a node's declaration is composed with, and the node's record. +func renderingFor(ctx context.Context, open *stores, node string, + plan catalogue.Resolution, settings catalogue.SettingsBy, + gens map[string]catalogue.Generator, choosing Choosing) (catalogue.Rendering, inventory.Node, error) { inv := open.inventory grants, err := grantsFor(ctx, open, node) if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } // Where this machine puts what each module needs reachable (novox/hq ADR 0038). // @@ -334,7 +354,7 @@ func declarationWith(ctx context.Context, open *stores, node string, if choosing == Reading { held, err := inv.PortsFor(ctx, node) if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } for _, a := range held { if already[a.Module] == nil { @@ -358,7 +378,7 @@ func declarationWith(ctx context.Context, open *stores, node string, case mayAssign && choosing == Allocating: at, err := inv.PortFor(ctx, node, m.Module, l.Port, l.Fixed) if err != nil { - return sendable{}, fmt.Errorf( + return catalogue.Rendering{}, inventory.Node{}, fmt.Errorf( "%s needs %d reachable on %s and it could not be assigned: %w", m.Module, l.Port, node, err) } @@ -399,7 +419,7 @@ func declarationWith(ctx context.Context, open *stores, node string, } } if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } if needed[m.Module] == nil { needed[m.Module] = map[string]string{} @@ -417,7 +437,7 @@ func declarationWith(ctx context.Context, open *stores, node string, } issued, meshCA, err := certificateFor(ctx, open, node) if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } certificate, authority = issued, meshCA break @@ -430,11 +450,11 @@ func declarationWith(ctx context.Context, open *stores, node string, // One reading of the catalogue for the three questions below that resolve the whole mesh. shelf, err := inv.Catalogue(ctx) if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } private, err := onThePrivateNetwork(ctx, inv, shelf) if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } // And every machine's name, so a container can reach one. The same set that writes the @@ -442,7 +462,7 @@ func declarationWith(ctx context.Context, open *stores, node string, // about where another machine is. names, err := namesInTheMesh(ctx, inv, shelf) if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } // And every routed name → the node that serves it (novox/hq ADR 0066). Alongside the @@ -451,7 +471,7 @@ func declarationWith(ctx context.Context, open *stores, node string, // to serve and knows nothing about what they mean. routes, err := routeNamesInTheMesh(ctx, open) if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } for name, at := range routes { names[name] = at @@ -480,7 +500,7 @@ func declarationWith(ctx context.Context, open *stores, node string, continue } if kept, err = inv.OperatorExport(ctx); err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } break } @@ -489,21 +509,13 @@ func declarationWith(ctx context.Context, open *stores, node string, // the declaration carries openings and the mesh's guard in place of a filter. record, err := inv.NodeByName(ctx, node) if err != nil { - return sendable{}, err + return catalogue.Rendering{}, inventory.Node{}, err } - composed, err := plan.Compose(catalogue.Rendering{ + return catalogue.Rendering{ Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports, Certificate: certificate, Authority: authority, Mesh: private, Names: names, - Suffix: overlay.Suffix(), Foundation: foundation, Kept: kept, Adopted: record.Adopted}) - if err != nil { - return sendable{}, err - } - // And what was taken on it, said in every declaration it is sent from this one place. - adoption, err := adoptionOf(ctx, inv, record, plan, composed) - if err != nil { - return sendable{}, err - } - return sendable{Resources: composed.Resources, Adoption: adoption}, nil + Suffix: overlay.Suffix(), Foundation: foundation, Kept: kept, Adopted: record.Adopted, + }, record, nil } // routeNamesInTheMesh is every routed name and the address of the node that serves it (novox/hq diff --git a/internal/catalogue/declaration.go b/internal/catalogue/declaration.go index d4f5ffa..59faa42 100644 --- a/internal/catalogue/declaration.go +++ b/internal/catalogue/declaration.go @@ -244,17 +244,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri // Once, from every module's listens -- not per module. A module receiving only its own ports // would write a rule set that closed every other module on the machine. Each module's per-node // exposure settings override its listens' source first (novox/hq ADR 0046). - exposure := map[string]map[int]string{} - for _, m := range r.Modules { - e, err := Exposure(m, with.Settings[m.Module]) - if err != nil { - return nil, err - } - if e != nil { - exposure[m.Module] = e - } - } - rules, err := r.Filtering(with.Generators, with.Ports, exposure) + rules, err := r.Rules(with) if err != nil { return nil, err } @@ -627,6 +617,23 @@ func (r Resolution) guarded(out []map[string]any, owner map[string]string, with return ports } +// Rules is the rule set this node's filter is derived from: every module's listens, what was +// computed for this machine, and each module's per-node exposure. The same answer whether the node +// is adopted or converged — the one loads it as a filter, the other declares it as openings. +func (r Resolution) Rules(with Rendering) ([]Rule, error) { + exposure := map[string]map[int]string{} + for _, m := range r.Modules { + e, err := Exposure(m, with.Settings[m.Module]) + if err != nil { + return nil, err + } + if e != nil { + exposure[m.Module] = e + } + } + return r.Filtering(with.Generators, with.Ports, exposure) +} + // Contribution is one module telling the answer to a requirement what it needs from it. type Contribution struct { // From is the module that said it, so the provider and a person reading the file can tell