A merge-check.sh, the repository's layer of the mesh's merge check (mesh/repo-check): format, vet, and node-tools' suite under the race detector. The tests shared one bus and had to run one package at a time; like the controller's (#97), each now starts a server of its own at the release go.mod pins, held to the catalogue's bus image by a test. What cannot run in the check — bundles that need @novox/mesh-sdk — is said as not tested.
234 lines
7.2 KiB
Go
234 lines
7.2 KiB
Go
// 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.
|
|
//
|
|
// **Every test has a bus of its own** (novox/hq ADR 0237, the controller's internal/testbus): a server
|
|
// linked in at the release go.mod pins — the release the mesh runs, held so by a test beside this one —
|
|
// started for the one test and gone after it. The packages shared one bus and had to run one at a time:
|
|
// each test raised the streams afresh, and the console discovers every runtime that announces itself on
|
|
// the bus (ADR 0197), so a runtime from another package's test was, correctly, found. Now the suite runs
|
|
// as Go runs it, in parallel and under the race detector. A person may still point a run at a bus of
|
|
// their own with MESH_TEST_NATS_EXTERNAL=1 and MESH_TEST_NATS; then the run is theirs to serialise.
|
|
package meshtest
|
|
|
|
import (
|
|
"encoding/json"
|
|
"os"
|
|
"path/filepath"
|
|
"runtime"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats-server/v2/server"
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
|
)
|
|
|
|
// Version is the release of the server the tests run.
|
|
const Version = server.VERSION
|
|
|
|
// URL is the bus of this test alone: started at its first ask, the same one at every ask after — a test
|
|
// that dials twice, or hands the address to the code it tests, reaches one bus — and shut down when the
|
|
// test ends.
|
|
func URL(t *testing.T) string {
|
|
t.Helper()
|
|
if os.Getenv("MESH_TEST_NATS_EXTERNAL") == "1" {
|
|
if url := os.Getenv("MESH_TEST_NATS"); url != "" {
|
|
return url
|
|
}
|
|
t.Fatal("MESH_TEST_NATS_EXTERNAL=1 and MESH_TEST_NATS names no bus")
|
|
}
|
|
if url, ok := buses.Load(t); ok {
|
|
return url.(string)
|
|
}
|
|
opts := &server.Options{Host: "127.0.0.1", Port: server.RANDOM_PORT, JetStream: true, StoreDir: t.TempDir(),
|
|
NoLog: true, NoSigs: true}
|
|
s, err := server.NewServer(opts)
|
|
if err != nil {
|
|
t.Fatalf("a bus for this test could not be made: %v", err)
|
|
}
|
|
go s.Start()
|
|
if !s.ReadyForConnections(30 * time.Second) {
|
|
s.Shutdown()
|
|
t.Fatal("a bus for this test did not come up within 30s")
|
|
}
|
|
url := s.ClientURL()
|
|
buses.Store(t, url)
|
|
t.Cleanup(func() {
|
|
buses.Delete(t)
|
|
s.Shutdown()
|
|
s.WaitForShutdown()
|
|
})
|
|
return url
|
|
}
|
|
|
|
// buses are the running tests' buses, by test.
|
|
var buses sync.Map
|
|
|
|
// 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 0201),
|
|
// 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
|
|
}
|