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, } 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 }