Files
mesh-controller/internal/broker/streams_test.go
T
jschoubben 603ad61142 The mesh issues an assignment's subjects: a membership per module on a machine (hq ADR 0160)
For every module on every machine the controller composes what that instance serves — its machine's
address always, the module's plain address in a queue when it is alone or its definition says its
instances are interchangeable — the verbs of the seats it holds at the seats' subjects, where its
events land, and what it may reach, resolved the same way for the modules it invokes. Published
beside the node's declaration on `mesh.assignment.<node>.<module>`, last per subject in a stream
that allows direct reads, and the account may read exactly its own. Composed from the same records
the bus's accounts are, so what a runtime serves and what its account may are one composition.
`instances: interchangeable` is the one fact a definition states for it.

The shape issued is the shape the mesh already had, so nothing moves when the membership arrives;
the runtime that reads it instead of deriving it is the next piece.
2026-10-01 14:43:26 +02:00

275 lines
8.9 KiB
Go

package broker
import (
"errors"
"strings"
"testing"
)
type recorder struct {
seen []Stream
fail string
}
func (r *recorder) EnsureStream(s Stream) error {
if s.Name == r.fail {
return errors.New("refused")
}
r.seen = append(r.seen, s)
return nil
}
// The controller asserts on every start, not only at genesis: a stream that was deleted, or a mesh
// raised from a backup, must converge rather than run without the guarantee its messages assume.
func TestAssertingTwiceIsTheSameAsOnce(t *testing.T) {
a, b := &recorder{}, &recorder{}
if err := AssertMeshStreams(a); err != nil {
t.Fatal(err)
}
if err := AssertMeshStreams(a); err != nil {
t.Fatal(err)
}
if err := AssertMeshStreams(b); err != nil {
t.Fatal(err)
}
if len(a.seen) != 2*len(b.seen) {
t.Fatalf("asserted %d then %d; assertion is not repeatable", len(a.seen), len(b.seen))
}
}
func TestAFailedAssertionNamesItsStream(t *testing.T) {
err := AssertMeshStreams(&recorder{fail: "NODES"})
if err == nil || !strings.Contains(err.Error(), "NODES") {
t.Fatalf("got %v, which does not say which stream failed", err)
}
}
// Two streams matching one subject is accepted by NATS and stores the message twice under two
// retentions. Nothing reports that, so it is refused where the set is written.
func TestNoTwoStreamsClaimTheSameSubject(t *testing.T) {
if clashes := Overlaps(); len(clashes) != 0 {
t.Fatalf("overlapping subject filters: %v", clashes)
}
}
// A heartbeat under mesh.control.> must not be persisted: a lost one is the next one, and a
// stream of them competes for retention with the messages that matter.
func TestHeartbeatsAreNotInTheControlStream(t *testing.T) {
for _, s := range MeshStreams() {
for _, subject := range s.Subjects {
if subject == "mesh.control.>" || strings.Contains(subject, "alive") {
t.Fatalf("stream %s claims %q, which captures heartbeats", s.Name, subject)
}
}
}
}
// The reason the kind token exists: a filter over a module's whole namespace would persist every
// tool call in the mesh.
func TestTheEventsStreamDoesNotCaptureToolCalls(t *testing.T) {
var events Stream
for _, s := range MeshStreams() {
if s.Name == "EVENTS" {
events = s
}
}
// Nothing a tool call rides may match any of the filters — a module's or a seat's.
for _, tool := range []string{
"mesh.mod.billing.tool.status",
"mesh.seat.telegram-sender.tool.status",
"mesh.seat.telegram-sender.accept.send", // work, not an event: its own stream
} {
for _, f := range events.Subjects {
if subjectMatches(f, tool) {
t.Fatalf("%q matches the events filter %q, so it would be persisted here", tool, f)
}
}
}
// And both kinds of event do match.
for _, event := range []string{
"mesh.mod.billing.event.order.placed",
"mesh.seat.telegram-sender.event.delivered",
} {
matched := false
for _, f := range events.Subjects {
if subjectMatches(f, event) {
matched = true
}
}
if !matched {
t.Fatalf("%q matches no events filter, so nothing would keep it", event)
}
}
}
// subjectMatches is NATS subject matching, enough for these filters: `*` is one token, `>` is the
// rest.
func subjectMatches(filter, subject string) bool {
f, s := strings.Split(filter, "."), strings.Split(subject, ".")
for i, tok := range f {
if tok == ">" {
return i <= len(s)
}
if i >= len(s) {
return false
}
if tok != "*" && tok != s[i] {
return false
}
}
return len(f) == len(s)
}
// Each relationship's retention is the thing that makes it what it is (design 29 §4).
func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) {
want := map[string]Retention{
"CONTROL": RetentionWorkQueue,
"NODES": RetentionLastPerSubject,
"EVENTS": RetentionLimits,
"ASSIGNMENTS": RetentionLastPerSubject,
}
got := map[string]Retention{}
for _, s := range MeshStreams() {
got[s.Name] = s.Retention
if s.Why == "" {
t.Errorf("stream %s says no reason it exists", s.Name)
}
}
if len(got) != len(want) {
t.Fatalf("the foundation set is %v", got)
}
for name, r := range want {
if got[name] != r {
t.Errorf("%s retains as %q, expected %q", name, got[name], r)
}
}
}
// The order the bus's objects are asserted in, because getting it wrong is a refusal that names the
// wrong thing: a consumer on a stream that does not exist is refused naming the *stream*, so
// somebody reading it goes looking for a deletion instead of a reversed pair of lines.
func TestTheBusesObjectsAreAssertedStreamsBeforeConsumers(t *testing.T) {
r := &recording{}
if err := Raise(r, []string{"anchor", "laptop"}); err != nil {
t.Fatal(err)
}
// Every stream before every consumer.
firstConsumer := -1
for i, step := range r.steps {
if strings.HasPrefix(step, "consumer ") && firstConsumer < 0 {
firstConsumer = i
}
if strings.HasPrefix(step, "stream ") && firstConsumer >= 0 {
t.Fatalf("a stream was asserted after a consumer: %v", r.steps)
}
}
if firstConsumer < 0 {
t.Fatalf("no consumer was asserted: %v", r.steps)
}
// And every node got one, named after it — without which that node hears nothing while
// everything else about it looks correct.
for _, node := range []string{"anchor", "laptop"} {
if !containsStep(r.steps, "consumer NODES/"+node) {
t.Errorf("%s was given no way to hear its declaration: %v", node, r.steps)
}
}
// And the controller its own, on both streams it reads.
for _, want := range []string{"consumer CONTROL/controller", "consumer EVENTS/controller"} {
if !containsStep(r.steps, want) {
t.Errorf("the controller is missing %s: %v", want, r.steps)
}
}
}
// A seat's work queue is asserted whether or not anybody holds it; the holder's worker only when
// somebody does. **The stream without the consumer is the point**: work queues until a holder
// appears, so installing the module later flushes the backlog instead of having lost it.
func TestASeatsQueueExistsBeforeItsHolderDoes(t *testing.T) {
seats := []DeclaredSeat{{Name: "telegram-sender", Accepts: []string{"send"}}}
unheld := &recording{}
if err := RaiseSeats(unheld, seats, nil); err != nil {
t.Fatal(err)
}
if !containsStep(unheld.steps, "stream SEAT_TELEGRAM_SENDER") {
t.Fatalf("a declared seat got no work queue: %v", unheld.steps)
}
for _, step := range unheld.steps {
if strings.HasPrefix(step, "consumer ") {
t.Fatalf("a seat nobody holds got a worker: %v", unheld.steps)
}
}
held := &recording{}
if err := RaiseSeats(held, seats, map[string]Holder{
"telegram-sender": {Node: "anchor", Module: "telegram"},
}); err != nil {
t.Fatal(err)
}
if !containsStep(held.steps, "consumer SEAT_TELEGRAM_SENDER/SEAT_TELEGRAM_SENDER_worker") {
t.Fatalf("the seat's holder got no worker: %v", held.steps)
}
}
// recording is a connection to the bus that writes down what it was asked for.
type recording struct{ steps []string }
func (r *recording) EnsureStream(s Stream) error {
r.steps = append(r.steps, "stream "+s.Name)
return nil
}
func (r *recording) EnsureConsumer(c Consumer) error {
r.steps = append(r.steps, "consumer "+c.Stream+"/"+c.Name)
return nil
}
func containsStep(steps []string, want string) bool {
for _, s := range steps {
if s == want {
return true
}
}
return false
}
// **Two consumers may share a name, and must not share a delivery subject** (novox/hq
// 04-ISSUES/146).
//
// A push consumer delivers onto an ordinary subject and everything subscribed to it gets a copy.
// The controller holds a consumer called `controller` on CONTROL and another called `controller` on
// EVENTS; while both were given `_DELIVER.controller`, the one process holding both subscriptions
// acted on every message twice — a joining machine enrolled twice from one request, with the second
// enrolment minting a credential that replaced the one the machine had just been handed.
//
// Checked here rather than against a server because it is a property of what the mesh asks for, and
// because the failure it produces is silent: every count is right, nothing is redelivered, and the
// work simply happens twice.
func TestNoTwoConsumersDeliverOntoTheSameSubject(t *testing.T) {
seen := map[string]string{}
for _, c := range MeshConsumers() {
if !c.Push && c.Queue == "" {
continue
}
subject := DeliverSubjectFor(c)
if other, taken := seen[subject]; taken {
t.Errorf("%s on %s and %s deliver onto %s, so whoever holds both acts on every "+
"message twice", c.Name, c.Stream, other, subject)
}
seen[subject] = c.Name + " on " + c.Stream
}
}
// The controller's events consumer is handed one announcement at a time (novox/hq issue 175): a
// merge's handler builds for minutes, and what is queued behind it must wait on the server rather
// than time out on the client and be acted on twice.
func TestTheControllerTakesOneAnnouncementAtATime(t *testing.T) {
for _, c := range MeshConsumers() {
if c.Stream == "EVENTS" && c.Name == ControllerName && c.MaxAckPending != 1 {
t.Fatalf("the events consumer may have %d outstanding; one announcement at a time", c.MaxAckPending)
}
}
}