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() reportsReaching(t, open, []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}, }, held...) } // reportsReaching has the anchor report, on what it was last sent, what is reachable on it and // holding what is given. func reportsReaching(t *testing.T, open *stores, reachable []link.Reach, 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: reachable, }); 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", // The anchor faces inward and is on the private network: the derived filter admits ssh // from the mesh only. "WILL CLOSE to everything outside the private network — ssh stays open from the mesh", "notes\n replacing the found file /etc/notes.conf (notes.conf), original kept at", "assigns nftables", "the found firewall (ufw) is disabled, never flushed", // What it routes is not a listener: said not to be previewed, and to be dropped. "not previewed: traffic the machine routes that is not a published port", "the derived filter drops it unless a module declares it", } { 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, digestIn(t, preview), ""); 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) } } // novox/hq ADR 0100: the preview says what the derived filter does, rendered as it is rendered. On // a machine that faces inward, ssh is admitted from the private network only, and a port a module // admits from the mesh only closes to everything outside it: both are said to close. func TestThePreviewSaysWhatNarrowsToTheMeshCloses(t *testing.T) { open, _ := anAdoptedAnchor(t) ctx := t.Context() register(t, open, catalogue.Manifest{Module: "store", Version: "1", Listens: []catalogue.Listening{{Port: 5432, From: catalogue.FromMesh}}}) if _, err := assign(ctx, open, "anchor", "store"); err != nil { t.Fatal(err) } reportsReaching(t, open, []link.Reach{ {Protocol: "tcp", Address: "0.0.0.0", Port: 22, By: "sshd"}, {Protocol: "tcp", Address: "0.0.0.0", Port: 5432, By: "postgres"}, {Protocol: "tcp", Address: "10.77.0.1", Port: 5432, By: "postgres"}, {Protocol: "tcp", Address: "0.0.0.0", Port: 8080, By: "hello-web", Published: true, ContainerPort: 80}, }) preview, err := converge(ctx, open, "anchor", false, "", "") if err != nil { t.Fatal(err) } lines := map[string]string{} for _, line := range strings.Split(preview, "\n") { fields := strings.Fields(line) if len(fields) > 1 && strings.HasPrefix(fields[0], "tcp/") { lines[fields[0]+" "+fields[1]] += line + "\n" } } if got := lines["tcp/22 sshd"]; !strings.Contains(got, "WILL CLOSE to everything outside "+ "the private network") { t.Errorf("ssh on an inward machine is not said to close outside the mesh:\n%s", preview) } store := lines["tcp/5432 postgres"] if strings.Count(store, "WILL CLOSE to everything outside the private network") != 1 || !strings.Contains(store, "declared by store (from mesh)") { t.Errorf("the store's narrowing is not said to close, or its mesh address is:\n%s", preview) } if got := lines["tcp/8080 hello-web"]; !strings.Contains(got, "declared by hello-web (from anywhere)") { t.Errorf("a port open to everywhere is not said to stay:\n%s", preview) } } // digestIn is the digest a converge preview printed. func digestIn(t *testing.T, preview string) string { t.Helper() for _, line := range strings.Split(preview, "\n") { if fields := strings.Fields(line); len(fields) == 2 && fields[0] == "preview" { return fields[1] } } t.Fatalf("the preview printed no digest:\n%s", preview) return "" } // The flip acts on the preview the operator saw: it names that preview's digest, and it is refused // when the digest is missing, when anything the preview says has changed since, or when the node's // account of itself is too old to be the machine as it is. func TestTheFlipActsOnlyOnThePreviewTheOperatorSaw(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) } saw := digestIn(t, preview) if !strings.Contains(preview, "converge anchor --yes "+saw) { t.Fatalf("the preview does not say how to act on it:\n%s", preview) } unchanged := func() { t.Helper() if n, _ := open.inventory.NodeByName(ctx, "anchor"); !n.Adopted || len(*sent) != 0 { t.Fatal("a refused flip changed something") } } if _, err := converge(ctx, open, "anchor", true, "", ""); err == nil || !strings.Contains(err.Error(), "--yes "+saw) { t.Fatalf("a flip naming no preview was not refused: %v", err) } unchanged() // Something new is reachable: the preview the operator saw is not what would happen. reportsReaching(t, open, []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: 6000, By: "something-new"}, }, heldFile) _, err = converge(ctx, open, "anchor", true, saw, "") if err == nil || !strings.Contains(err.Error(), "has changed since preview "+saw) { t.Fatalf("a flip on a changed preview was not refused: %v", err) } unchanged() // An account older than the flip trusts is refused, whatever digest is named. again, err := converge(ctx, open, "anchor", false, "", "") if err != nil { t.Fatal(err) } saved := reportFreshFor reportFreshFor = time.Nanosecond _, err = converge(ctx, open, "anchor", true, digestIn(t, again), "") reportFreshFor = saved if err == nil || !strings.Contains(err.Error(), "wait for its next report") { t.Fatalf("a flip on an old account was not refused: %v", err) } unchanged() if _, err := converge(ctx, open, "anchor", true, digestIn(t, again), ""); err != nil { t.Fatalf("the flip on the preview just seen was refused: %v", err) } if n, _ := open.inventory.NodeByName(ctx, "anchor"); n.Adopted { t.Fatal("the flip did not converge the node") } } // The flip holds the node from its checks to its send, so a push composed meanwhile waits and is // composed after it — never sent after it with the node still adopted. And the send inside the // flip is not made to wait on the flip's own hold. func TestTheFlipHoldsTheNodeWhileItSends(t *testing.T) { open, _ := 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) } var heldElsewhere, heldHere error sendNodes = func(inner context.Context, open *stores, names []string) error { // Another caller cannot hold the node while the flip sends it. waiting, cancel := context.WithTimeout(ctx, 300*time.Millisecond) defer cancel() if release, err := open.inventory.HoldNodes(waiting, names); err == nil { release() heldElsewhere = errors.New("another caller held the node while the flip sent it") } // The flip's own send holds it without waiting on itself. _, release, err := holdNodes(inner, open, names) if err != nil { heldHere = err return err } release() return nil } if _, err := converge(ctx, open, "anchor", true, digestIn(t, preview), ""); err != nil { t.Fatal(err) } if heldElsewhere != nil || heldHere != nil { t.Fatalf("%v %v", heldElsewhere, heldHere) } // And given back once it is done. release, err := open.inventory.HoldNodes(ctx, []string{"anchor"}) if err != nil { t.Fatal(err) } release() } // novox/hq ADR 0103: the preview names every kind of thing a module the flip takes holds as found, // not only its files, and the digest changes when any of them does. func TestThePreviewNamesEveryHeldKind(t *testing.T) { open, _ := anAdoptedAnchor(t) ctx := t.Context() if _, err := take(ctx, open, "anchor", "hello-web"); err != nil { t.Fatal(err) } since := time.Now() held := []link.Held{heldFile, {ID: "notes.data", Module: "notes", Kind: "directory", Target: "/var/lib/notes", Since: since}, {ID: "notes.daemon", Module: "notes", Kind: "service", Target: "notes.service", Since: since}, {ID: "notes.seed", Module: "notes", Kind: "archive", Target: "/srv/notes", Since: since}, {ID: "notes.worker", Module: "notes", Kind: "process", Target: "notes-worker", Since: since}, {ID: "notes.account", Module: "notes", Kind: "user", Target: "notes", Since: since}, } reportsHolding(t, open, held...) preview, err := converge(ctx, open, "anchor", false, "", "") if err != nil { t.Fatal(err) } for _, want := range []string{ "replacing the found directory /var/lib/notes (notes.data)", "replacing the found service notes.service (notes.daemon)", "replacing the found archive /srv/notes (notes.seed)", "replacing the found process notes-worker (notes.worker)", "replacing the found user notes (notes.account)", } { if !strings.Contains(preview, want) { t.Errorf("the preview does not say %q:\n%s", want, preview) } } reportsHolding(t, open, held[:len(held)-1]...) fewer, err := converge(ctx, open, "anchor", false, "", "") if err != nil { t.Fatal(err) } if digestIn(t, fewer) == digestIn(t, preview) { t.Fatal("the digest does not change with what is held") } } // novox/hq ADR 0100: the flip acts on what the node said is reachable, so an account naming nothing // is refused. Every machine that is up answers on ssh; nothing reported means the host's collectors // did not, and flipping would close ports the preview never named. func TestTheFlipIsRefusedOnAnAccountNamingNothingReachable(t *testing.T) { open, sent := anAdoptedAnchor(t) ctx := t.Context() if _, err := take(ctx, open, "anchor", "hello-web"); err != nil { t.Fatal(err) } // Only a loopback listener: nothing off the machine, which is the same silence. reportsReaching(t, open, []link.Reach{ {Protocol: "tcp", Address: "127.0.0.1", Port: 15672, By: "mesh-broker"}, }, heldFile) preview, err := converge(ctx, open, "anchor", false, "", "") if err != nil { t.Fatal(err) } if !strings.Contains(preview, "this account looks partial") { t.Errorf("the preview does not mark a partial account:\n%s", preview) } _, err = converge(ctx, open, "anchor", true, digestIn(t, preview), "") if err == nil || !strings.Contains(err.Error(), "says nothing is reachable on it") { t.Fatalf("the flip was not refused on an account naming nothing: %v", err) } if n, _ := open.inventory.NodeByName(ctx, "anchor"); !n.Adopted || len(*sent) != 0 { t.Fatal("a refused flip changed something") } } // An assignment cannot land between a preview and the flip that takes every module: assigning // holds the node, so it waits for whatever is converging it. func TestAssigningWaitsForWhateverIsConvergingTheNode(t *testing.T) { open, _ := 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) } saved, savedPoll := inventory.HoldWaitFor, inventory.HoldPoll inventory.HoldWaitFor, inventory.HoldPoll = time.Second, 50*time.Millisecond defer func() { inventory.HoldWaitFor, inventory.HoldPoll = saved, savedPoll }() var whileFlipping error sendNodes = func(context.Context, *stores, []string) error { _, whileFlipping = assign(ctx, open, "anchor", "notes") return nil } if _, err := converge(ctx, open, "anchor", true, digestIn(t, preview), ""); err != nil { t.Fatal(err) } if !errors.Is(whileFlipping, inventory.ErrNodeBusy) { t.Fatalf("an assignment landed while the node was being converged: %v", whileFlipping) } }