Task 3.9's other half and 1.4's missing client. The derivation is pure and unit-tested; only "does the server accept this" needs one running, behind MESH_TEST_NATS so the ordinary suite stays offline. A seat's work queue is created at registration, not assignment, so work queues until a holder appears — a stream created at assignment would make "the holder is not here yet" mean "your messages are gone". Named after the seat, because the holder can change and the queued work must not care. A holder's worker uses a queue group even though the seat guarantees one holder: the seat is authority, the queue group is delivery, and tying them together means the day somebody allows two holders every message is processed twice with nothing reporting it. One consumer per module carrying every filter, because its ack permission is derived from its name. And a real bug the live server caught: a durable name may not contain a dot, but an ack subject is $JS.ACK.<stream>.<consumer>, so the single string that read correctly inside the permission was rejected as a consumer name. Split in two, beside the permission that has to match. Unfixed, the symptom would have been every message redelivered forever with a permission list that looks right — which is the failure design 25 §4 warns about.
147 lines
4.1 KiB
Go
147 lines
4.1 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,
|
|
"BUILDS": RetentionWorkQueue,
|
|
"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)
|
|
}
|
|
}
|
|
}
|