// Package meshtest raises what the controller would, for tests against a real bus: the ASSIGNMENTS // and EVENTS streams, memberships issued by hand, and a fixture's path. // // **The packages share one bus, so run them one at a time: `go test -p 1 ./...`.** Each test raises // the streams afresh, and the console discovers every runtime that announces itself on the bus // (novox/hq ADR 0197) — a runtime from another package's test is, correctly, found. package meshtest import ( "encoding/json" "os" "path/filepath" "runtime" "strings" "testing" "time" "github.com/nats-io/nats.go" "github.com/novox/mesh-tools/node-tools/internal/bus" ) // URL is the test bus, or the test is skipped. func URL(t *testing.T) string { t.Helper() url := os.Getenv("MESH_TEST_NATS") if url == "" { t.Skip("MESH_TEST_NATS unset") } return url } // Mesh is the controller's job, done by hand. type Mesh struct { nc *nats.Conn js nats.JetStreamContext } // New raises the streams afresh. func New(t *testing.T) *Mesh { t.Helper() nc, err := nats.Connect(URL(t)) if err != nil { t.Fatal(err) } js, _ := nc.JetStream() for _, s := range []string{"ASSIGNMENTS", "EVENTS"} { _ = js.DeleteStream(s) } if _, err := js.AddStream(&nats.StreamConfig{Name: "ASSIGNMENTS", Subjects: []string{"mesh.assignment.>"}, MaxMsgsPerSubject: 1, AllowDirect: true}); err != nil { t.Fatal(err) } if _, err := js.AddStream(&nats.StreamConfig{Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}}); err != nil { t.Fatal(err) } t.Cleanup(nc.Close) return &Mesh{nc: nc, js: js} } // Issue publishes a membership. func (m *Mesh) Issue(t *testing.T, mem bus.Membership) { t.Helper() body, _ := json.Marshal(mem) if _, err := m.js.Publish(bus.MembershipSubject(mem.Node, mem.Module), body); err != nil { t.Fatal(err) } } // NextEvent is the subject the next event under a pattern lands on. func (m *Mesh) NextEvent(t *testing.T, pattern string) <-chan string { t.Helper() ch := make(chan string, 1) sub, err := m.nc.Subscribe(pattern, func(msg *nats.Msg) { select { case ch <- msg.Subject: default: } }) if err != nil { t.Fatal(err) } t.Cleanup(func() { _ = sub.Unsubscribe() }) _ = m.nc.Flush() return ch } // MembershipOf is a membership as the controller issues one on a machine. func MembershipOf(module, node string, plain bool, seats map[string][]string) bus.Membership { own := "mesh.mod." + module serves := []bus.Served{{Subject: own + ".tool.{tool}." + node}} if plain { serves = append(serves, bus.Served{Subject: own + ".tool.{tool}", Queue: "serve." + module}) } var verbs []bus.SeatVerb for seat, vs := range seats { for _, v := range vs { verbs = append(verbs, bus.SeatVerb{Seat: seat, Verb: v, Subject: "mesh.seat." + seat + ".tool." + v + "." + node}) } } return bus.Membership{Node: node, Module: module, Serves: serves, Seats: verbs, Emits: own + ".event.{event}", Tools: own + ".tool.tools"} } // Fixture is a path in node-tools/test/fixtures, which these tests share with the TypeScript ones. func Fixture(name string) string { _, here, _, _ := runtime.Caller(0) return filepath.Join(filepath.Dir(here), "..", "..", "test", "fixtures", name) } // Until retries while the bus answers "no responders" — something not yet served. func Until(t *testing.T, try func() error) { t.Helper() var err error for i := 0; i < 50; i++ { if err = try(); err == nil { return } time.Sleep(100 * time.Millisecond) } t.Fatal(err) } // Logs collects what the runtime says. type Logs struct{ lines []string } // Logf is a logger that keeps the lines. func (l *Logs) Logf(format string, args ...any) { l.lines = append(l.lines, sprintf(format, args...)) } // Has says whether a line contains every fragment. func (l *Logs) Has(fragments ...string) bool { for _, line := range l.lines { all := true for _, f := range fragments { all = all && strings.Contains(line, f) } if all { return true } } return false } // All is every line. func (l *Logs) All() string { return strings.Join(l.lines, "\n") } // Consumer makes a module's durable consumer on a machine as the controller does: pull, on the // EVENTS stream, filtered to what the module consumes, with a short ack wait so a test sees a // redelivery in seconds rather than the mesh's minutes. func (m *Mesh) Consumer(t *testing.T, node, module string, filters []string, ackWait time.Duration) { t.Helper() cfg := &nats.ConsumerConfig{Durable: node + "_" + module, AckPolicy: nats.AckExplicitPolicy, AckWait: ackWait, MaxDeliver: 10, DeliverPolicy: nats.DeliverNewPolicy} if len(filters) == 1 { cfg.FilterSubject = filters[0] } else { cfg.FilterSubjects = filters } if _, err := m.js.AddConsumer("EVENTS", cfg); err != nil { t.Fatal(err) } } // Emit publishes an event as a module would, into the EVENTS stream. func (m *Mesh) Emit(t *testing.T, subject string, body any) { t.Helper() data, _ := json.Marshal(body) if _, err := m.js.Publish(subject, data); err != nil { t.Fatal(err) } } // Pending is what a module's consumer still holds: delivered and not acknowledged, and not yet // delivered. func (m *Mesh) Pending(t *testing.T, node, module string) (ackPending, notDelivered uint64) { t.Helper() info, err := m.js.ConsumerInfo("EVENTS", node+"_"+module) if err != nil { t.Fatal(err) } return uint64(info.NumAckPending), info.NumPending } // Bucket makes a module's state afresh as the controller does from the catalogue (novox/hq ADR 0202), // and answers it for writing what is there before a bundle starts. func (m *Mesh) Bucket(t *testing.T, name string) nats.KeyValue { t.Helper() _ = m.js.DeleteKeyValue(name) kv, err := m.js.CreateKeyValue(&nats.KeyValueConfig{Bucket: name, History: 1, MaxValueSize: 256 * 1024}) if err != nil { t.Fatal(err) } return kv }