package main import ( "context" "encoding/json" "errors" "net" "slices" "strings" "sync/atomic" "testing" "time" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" "golang.org/x/net/dns/dnsmessage" "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/conditions" "github.com/novox/mesh-controller/internal/link" "github.com/novox/mesh-controller/internal/testbus" ) // The self-check (novox/hq to-be 45 §4): a probe that fails raises its condition, one that cannot run // raises probe-failed for itself and is never a pass, and every run ends with its heartbeat. // withProbes runs the test with a registry of its own. func withProbes(t *testing.T, probes ...probe) { t.Helper() before := probeRegistry probeRegistry = probes t.Cleanup(func() { probeRegistry = before }) } func TestARunKeepsWhatEachProbeFoundAndSaysItsHeartbeat(t *testing.T) { failing := conditions.Observation{Scope: conditions.ScopeMachine, ID: "anchor", Token: "refused", Severity: conditions.Urgent, Summary: "anchor's node-engine would refuse its declaration"} var broken atomic.Bool broken.Store(true) withProbes(t, probe{ID: "P1", Asserts: "passes", Kind: "never", Phase: 1, run: func(context.Context, *doctor) ([]conditions.Observation, error) { return nil, nil }}, probe{ID: "P2", Asserts: "finds a fault", Kind: "declaration-refused", Phase: 1, run: func(context.Context, *doctor) ([]conditions.Observation, error) { if broken.Load() { return []conditions.Observation{failing}, nil } return nil, nil }}, probe{ID: "P3", Asserts: "cannot run", Kind: "x", Phase: 1, run: func(context.Context, *doctor) ([]conditions.Observation, error) { if broken.Load() { return nil, errors.New("the store is away") } return nil, nil }}, probe{ID: "P4", Asserts: "hangs", Kind: "x", Phase: 1, run: func(ctx context.Context, _ *doctor) ([]conditions.Observation, error) { if broken.Load() { <-ctx.Done() time.Sleep(50 * time.Millisecond) } return nil, nil }}, probe{ID: "P5", Asserts: "later", Kind: "x", Phase: 2, Deferred: "not yet"}, ) before := probeWithin probeWithin = 200 * time.Millisecond t.Cleanup(func() { probeWithin = before }) store := conditions.NewInMemory() told := &conditions.Told{} k := conditions.NewKeeper(t.Context(), conditions.Options{Store: store, History: store}) defer k.Close(context.Background()) d := &doctor{keeper: k, teller: told, host: "anchor"} // A probe that cannot run is said in the verdict at once, and as a condition when the next run cannot // run it either (novox/hq issue 277): one timed-out question is not a probe gone blind. first := d.runOnce(t.Context(), "a test") if first.Counts != (doctorCounts{Passed: 1, Failed: 1, FailedToRun: 2, Deferred: 1}) { t.Fatalf("counted %+v", first.Counts) } if open, _ := k.Open(t.Context()); len(open) != 1 || open[0].Key != "machine.anchor.refused" { t.Fatalf("after one run, open: %+v", open) } run := d.runOnce(t.Context(), "a test") if run.Counts != (doctorCounts{Passed: 1, Failed: 1, FailedToRun: 2, Deferred: 1}) { t.Fatalf("counted %+v", run.Counts) } open, err := k.Open(t.Context()) if err != nil { t.Fatal(err) } var keys []string for _, c := range open { keys = append(keys, c.Key+"="+c.Kind) } for _, want := range []string{"machine.anchor.refused=declaration-refused", "probe.P3.failed=probe-failed", "probe.P4.failed=probe-failed"} { if !slices.Contains(keys, want) { t.Errorf("%s is not open: %v", want, keys) } } // The heartbeat, in the shape mesh-watcher reads (the contract with the operator's channel). if len(told.Names) != 2 || told.Names[1] != conditions.HeartbeatEvent { t.Fatalf("said %v", told.Names) } if d.lastRunEnded().IsZero() || d.lastRun().Run != run.Run { t.Fatal("the run is not the last verdict") } body, _ := json.Marshal(run) var shape map[string]any _ = json.Unmarshal(body, &shape) for _, field := range []string{"run", "at", "interval-seconds", "counts", "probes", "controller"} { if _, ok := shape[field]; !ok { t.Errorf("the heartbeat carries no %q: %s", field, body) } } // Mended: the next run clears every one of them. broken.Store(false) d.runOnce(t.Context(), "a test") if open, _ := k.Open(t.Context()); len(open) != 0 { t.Fatalf("a passing run left open %+v", open) } } // **The registry says what each probe asserts**, and a probe not built says why and when. func TestTheRegistryIsTheDesignsLiveForm(t *testing.T) { seen := map[string]bool{} for _, p := range probeRegistry { if seen[p.ID] { t.Errorf("%s twice", p.ID) } seen[p.ID] = true if p.Asserts == "" || p.From == "" || p.Kind == "" { t.Errorf("%s does not say what it asserts, where from, or what it raises", p.ID) } if (p.run == nil) != (p.Deferred != "") || (p.Deferred != "" && p.Phase <= 1) { t.Errorf("%s is run and deferred, or neither, or deferred out of Phase 1: %+v", p.ID, p) } } for _, id := range []string{"D1", "D2", "D3", "D4", "D5", "D6", "D7", "D8", "D9", "D10"} { if !seen[id] { t.Errorf("to-be 45 §4 has %s and the registry does not", id) } } } // **The doctor and the watchdogs watch each other**: watchdogs that stopped are DW; a self-check that // stopped is S10 (signals_test.go). func TestWatchdogsThatStoppedAreSaid(t *testing.T) { w := &watchdogs{started: time.Now().Add(-time.Hour)} d := &doctor{watchdogs: w, host: "anchor"} got, err := probeWatchdogs(t.Context(), d) if err != nil || len(got) != 1 || got[0].Severity != conditions.Urgent { t.Fatalf("%+v %v", got, err) } w.ticked = time.Now() if got, _ := probeWatchdogs(t.Context(), d); len(got) != 0 { t.Fatalf("%+v", got) } } // **D1 composes every machine of a healthy mesh and the host's own validator takes each.** func TestEveryMachineOfAHealthyMeshComposesAndValidates(t *testing.T) { open := aMesh(t) got, err := probeDeclarations(t.Context(), &doctor{open: open}) if err != nil { t.Fatal(err) } if len(got) != 0 { t.Fatalf("a healthy mesh failed D1: %+v", got) } } // **D1 names a machine nothing can be sent to**, and the network that cannot be computed for it. func TestAMachineWhoseDeclarationDoesNotComposeIsSaid(t *testing.T) { open := aMesh(t) ctx := t.Context() one, two := rivals() register(t, open, one) register(t, open, two) for _, m := range []string{"rival-one", "rival-two"} { if _, err := assign(ctx, open, "laptop", m); err != nil && m == "rival-one" { t.Fatal(err) } } got, err := probeDeclarations(ctx, &doctor{open: open}) linted(got) if err != nil { t.Fatal(err) } if len(got) != 1 || got[0].Key() != "machine.laptop.uncomposable" || !strings.Contains(got[0].Summary, "the-seat") { t.Fatalf("%+v", got) } } // **D2: a resolver answering NXDOMAIN for IPv6 is wrong** — musl takes it as no such name (issue 262). func TestAResolverAnsweringNoSuchNameForIPv6IsWrong(t *testing.T) { answerAs := func(rcode dnsmessage.RCode) string { conn, err := net.ListenPacket("udp", "127.0.0.1:0") if err != nil { t.Fatal(err) } t.Cleanup(func() { _ = conn.Close() }) go func() { buf := make([]byte, 1500) for { n, from, err := conn.ReadFrom(buf) if err != nil { return } var q dnsmessage.Message if q.Unpack(buf[:n]) != nil { continue } reply := dnsmessage.Message{Header: dnsmessage.Header{ID: q.ID, Response: true}, Questions: q.Questions} if q.Questions[0].Type == dnsmessage.TypeA { reply.Answers = []dnsmessage.Resource{{Header: dnsmessage.ResourceHeader{Name: q.Questions[0].Name, Type: dnsmessage.TypeA, Class: dnsmessage.ClassINET}, Body: &dnsmessage.AResource{A: [4]byte{10, 77, 0, 1}}}} } else { reply.RCode = rcode } packed, _ := reply.Pack() _, _ = conn.WriteTo(packed, from) } }() _, port, _ := net.SplitHostPort(conn.LocalAddr().String()) return port } before := resolverPort t.Cleanup(func() { resolverPort = before }) resolverPort = answerAs(dnsmessage.RCodeSuccess) v4, rcode, err := askResolver(t.Context(), "127.0.0.1", "anchor.internal", dnsmessage.TypeA) if err != nil || rcode != dnsmessage.RCodeSuccess || !slices.Equal(v4, []string{"10.77.0.1"}) { t.Fatalf("%v %v %v", v4, rcode, err) } v6, rcode, err := askResolver(t.Context(), "127.0.0.1", "anchor.internal", dnsmessage.TypeAAAA) if err != nil || rcode != dnsmessage.RCodeSuccess || len(v6) != 0 { t.Fatalf("NODATA read as %v %v %v", v6, rcode, err) } resolverPort = answerAs(dnsmessage.RCodeNameError) if _, rcode, _ := askResolver(t.Context(), "127.0.0.1", "anchor.internal", dnsmessage.TypeAAAA); rcode != dnsmessage.RCodeNameError { t.Fatalf("NXDOMAIN read as %v", rcode) } } // **D6, D7: what the controller defines is what it finds**, and a consumer deleted or a stream // redefined is said — against a real bus, raised by the same derivation the controller starts with. func TestNatsTheBusIsWhatTheControllerDefines(t *testing.T) { url := testbus.URL(t) open := aMesh(t) js, err := broker.Dial(url) if err != nil { t.Fatal(err) } t.Cleanup(js.Close) for _, s := range []string{"CONTROL", "NODES", "ASSIGNMENTS", "EVENTS"} { _ = js.Context().DeleteStream(s) } if _, err := assertBusObjects(t.Context(), open.inventory, js); err != nil { t.Fatal(err) } if err := js.EnsureControllerBuckets(); err != nil { t.Fatal(err) } d := &doctor{open: open, js: js} for _, p := range []func(context.Context, *doctor) ([]conditions.Observation, error){probeConsumers, probeStreams} { got, err := p(t.Context(), d) if err != nil || len(got) != 0 { t.Fatalf("a bus just raised fails: %+v %v", got, err) } } if err := js.Context().DeleteConsumer("NODES", "laptop"); err != nil { t.Fatal(err) } info, err := js.Context().StreamInfo("EVENTS") if err != nil { t.Fatal(err) } cfg := info.Config cfg.MaxMsgsPerSubject = 3 if _, err := js.Context().UpdateStream(&cfg); err != nil { t.Fatal(err) } consumers, err := probeConsumers(t.Context(), d) if err != nil || len(consumers) != 1 || consumers[0].Key() != "bus.NODES.laptop.missing" { t.Fatalf("the deleted consumer: %+v %v", consumers, err) } streams, err := probeStreams(t.Context(), d) if err != nil || len(streams) != 1 || !strings.Contains(streams[0].Summary, "per subject") { t.Fatalf("the redefined stream: %+v %v", streams, err) } } // **S9 hears the bus**: a consumer that gives up on a message, and one deleted, as the server says. func TestNatsTheBusSaysAConsumerGaveUpAndOneWasDeleted(t *testing.T) { url := testbus.URL(t) conn, err := nats.Connect(url) if err != nil { t.Fatal(err) } defer conn.Close() heard := make(chan *nats.Msg, 16) for _, subject := range broker.BusAdvisories { if _, err := conn.ChanSubscribe(subject, heard); err != nil { t.Fatal(err) } } api, _ := jetstream.New(conn) _ = api.DeleteStream(t.Context(), "SEAT_ADVISED") stream, err := api.CreateStream(t.Context(), jetstream.StreamConfig{Name: "SEAT_ADVISED", Subjects: []string{"advised.>"}}) if err != nil { t.Fatal(err) } defer func() { _ = api.DeleteStream(context.Background(), "SEAT_ADVISED") }() consumer, err := stream.CreateConsumer(t.Context(), jetstream.ConsumerConfig{Durable: "SEAT_ADVISED_worker", AckPolicy: jetstream.AckExplicitPolicy, MaxDeliver: 1, AckWait: 100 * time.Millisecond}) if err != nil { t.Fatal(err) } if _, err := api.Publish(t.Context(), "advised.x", []byte("x")); err != nil { t.Fatal(err) } if _, err := consumer.Fetch(1, jetstream.FetchMaxWait(time.Second)); err != nil { t.Fatal(err) } // A seat's worker, so the deletion is of a consumer the mesh names (link.MeshNamed). // Not acknowledged: after its one delivery the consumer gives up on it — on the next fetch. time.Sleep(300 * time.Millisecond) _, _ = consumer.Fetch(1, jetstream.FetchMaxWait(300*time.Millisecond)) if err := stream.DeleteConsumer(t.Context(), "SEAT_ADVISED_worker"); err != nil { t.Fatal(err) } kinds := map[string]string{} deadline := time.After(5 * time.Second) for len(kinds) < 2 { select { case m := <-heard: if a, ok := link.ReadAdvisory(m.Subject, m.Data); ok { kinds[a.Kind] = a.Said } case <-deadline: t.Fatalf("the bus said only %v", kinds) } } if !strings.Contains(kinds["max-deliveries"], "gave up") || !strings.Contains(kinds["consumer-lost"], "was deleted") { t.Fatalf("%v", kinds) } } // **D10 compares a node-engine with what it is delivered as, not with its commit** (2026-10-06: every // machine read as behind right after a push sent it the current build — it says the digest-named // directory it runs from, and the mesh holds a commit). func TestANodeEngineIsJudgedByTheVersionItIsDeliveredAs(t *testing.T) { m := catalogue.Manifest{Module: "mesh-host", Resources: []map[string]any{ {"id": "launcher", "type": "file", "path": "/usr/lib/nox-mesh-host/launch"}, {"id": "host", "type": "archive", "path": "/usr/lib/nox-mesh-host/versions/31045596c83a"}, {"id": "unfilled", "type": "archive", "path": "/usr/lib/x/versions/${version}"}, }} delivered := deliveredVersions(m) if !slices.Equal(delivered, []string{"31045596c83a"}) { t.Fatalf("%v", delivered) } commit := "1545b00a9f0c" for _, c := range []struct { reported string behind bool }{ {"31045596c83a", false}, // the live case: current, and was called behind {"0123456789ab", true}, // another delivery {"1545b00a", false}, // placed by hand, stamped with the commit {"", false}, // not said } { if got := engineBehind(c.reported, delivered, commit); got != c.behind { t.Errorf("%q behind = %v, want %v", c.reported, got, c.behind) } } if engineBehind("31045596c83a", nil, commit) { t.Error("behind a mesh that holds no delivered build") } } // **Every seat verb a probe calls is one it declares, and one the controller is granted** — derived // from the registry, so a probe added with a question the bus would refuse fails here, not live. func TestEverySeatVerbAProbeAsksIsGranted(t *testing.T) { granted, err := broker.PermissionsFor(broker.Principal{Kind: broker.KindController, PasswordHash: "x"}) if err != nil { t.Fatal(err) } asked := 0 for _, p := range probeRegistry { for _, v := range p.Asks { asked++ subject := link.NodeSeatToolSubject(v.Seat, v.Verb, "anchor") if !slices.ContainsFunc(granted.Publish, func(pattern string) bool { return subjectMatches(pattern, subject) }) { t.Errorf("%s asks %s.%s and the controller may not publish %s", p.ID, v.Seat, v.Verb, subject) } if !slices.Contains(broker.VerbsTheSelfCheckAsks, v) { t.Errorf("%s asks %s.%s, which broker.VerbsTheSelfCheckAsks does not name", p.ID, v.Seat, v.Verb) } } } if asked == 0 { t.Fatal("no probe asks a seat verb: D8 lost its declaration") } // And a probe asking what it did not declare is refused before anything is sent. ctx := context.WithValue(t.Context(), probeAsksKey{}, probe{ID: "DX"}) if _, err := askSeatTool(ctx, nil, "node-intrusion-prevention", "banned", "anchor"); err == nil || !strings.Contains(err.Error(), "does not declare") { t.Fatalf("an undeclared question was asked: %v", err) } } // subjectMatches is the bus's matching of a permission pattern against a subject. func subjectMatches(pattern, subject string) bool { p, s := strings.Split(pattern, "."), strings.Split(subject, ".") for i, tok := range p { if tok == ">" { return len(s) > i } if i >= len(s) || (tok != "*" && tok != s[i]) { return false } } return len(p) == len(s) }