A launched bundle's mesh/subscribe binds the module's own durable consumer — EVENTS, <node>_<module>, by name as the module's own runtime bound it, so nothing is lost or replayed in the move — and every event goes to each child of the module that subscribed as mesh/event, acknowledged only when all answered, negatively acknowledged after a short delay when one failed or died, terminated when it is not an event. mesh/ask calls a tool as the module. Every launched bundle is started again when it exits, with backoff, since long-running code waits for no call. Requires SDK 0.1.6.
185 lines
5.3 KiB
Go
185 lines
5.3 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.
|
|
//
|
|
// **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
|
|
}
|