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.
120 lines
4.5 KiB
Go
120 lines
4.5 KiB
Go
package broker
|
|
|
|
import (
|
|
"strings"
|
|
"testing"
|
|
)
|
|
|
|
func telegramSeat() DeclaredSeat {
|
|
return DeclaredSeat{Name: "telegram-sender", Accepts: []string{"send"}}
|
|
}
|
|
|
|
// The stream exists from registration, not assignment: work queues until a holder appears, so
|
|
// installing the module a week later flushes the backlog rather than having lost it.
|
|
func TestASeatGetsAWorkQueueOfItsOwn(t *testing.T) {
|
|
got := SeatStreams([]DeclaredSeat{telegramSeat()})
|
|
if len(got) != 1 {
|
|
t.Fatalf("expected one stream, got %d", len(got))
|
|
}
|
|
s := got[0]
|
|
if s.Retention != RetentionWorkQueue {
|
|
t.Fatalf("a seat's inbound queue retains as %q; one holder must take each message once", s.Retention)
|
|
}
|
|
if s.Subjects[0] != "mesh.seat.telegram-sender.accept.>" {
|
|
t.Fatalf("filters on %v", s.Subjects)
|
|
}
|
|
}
|
|
|
|
// A seat that only emits and serves needs no stream: its events ride EVENTS and its tools are
|
|
// core request/reply, which is never persisted.
|
|
func TestASeatThatAcceptsNothingGetsNoStream(t *testing.T) {
|
|
if got := SeatStreams([]DeclaredSeat{{Name: "announcer"}}); len(got) != 0 {
|
|
t.Fatalf("a seat with no inbound work got %d stream(s)", len(got))
|
|
}
|
|
}
|
|
|
|
// Retention belongs to whoever owns the namespace, and a seat owns its own.
|
|
func TestASeatsRetentionIsItsOwn(t *testing.T) {
|
|
s := SeatStreams([]DeclaredSeat{{Name: "slow", Accepts: []string{"work"}, RetainSeconds: 30 * 24 * 60 * 60}})
|
|
if s[0].MaxAge != 30*24*60*60 {
|
|
t.Fatalf("the seat's declared retention was not used: %d", s[0].MaxAge)
|
|
}
|
|
d := SeatStreams([]DeclaredSeat{telegramSeat()})
|
|
if d[0].MaxAge == 0 {
|
|
t.Fatal("a seat that declares no retention got an unbounded queue")
|
|
}
|
|
}
|
|
|
|
// NATS refuses an overlapping stream outright, so a clash here is a registration that fails.
|
|
func TestNoDerivedStreamOverlapsTheMeshsOwn(t *testing.T) {
|
|
seats := []DeclaredSeat{telegramSeat(), {Name: "licensing-master", Accepts: []string{"report"}}}
|
|
if c := AllOverlaps(seats); len(c) != 0 {
|
|
t.Fatalf("overlapping filters: %v", c)
|
|
}
|
|
}
|
|
|
|
// One consumer per module, with every consumed subject as a filter — because its ack permission
|
|
// is derived from its name, and a consumer per event would need an ack permission per consumer.
|
|
func TestAModuleGetsOneConsumerCarryingEveryFilter(t *testing.T) {
|
|
c, ok := ConsumerFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
|
|
Consumes: []string{"shop.order.placed", "billing.invoice.sent"}, PasswordHash: "x"})
|
|
if !ok {
|
|
t.Fatal("a module that consumes got no consumer")
|
|
}
|
|
if len(c.Filters) != 2 {
|
|
t.Fatalf("expected both subjects as filters, got %v", c.Filters)
|
|
}
|
|
perms, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
|
|
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
|
|
ack := "$JS.ACK." + c.Stream + "." + c.Name + ".>"
|
|
found := false
|
|
for _, p := range perms.Publish {
|
|
if p == ack {
|
|
found = true
|
|
}
|
|
}
|
|
if !found {
|
|
t.Fatalf("the consumer is named %q but the ack permission is %v; a module could not ack "+
|
|
"its own deliveries", c.Name, perms.Publish)
|
|
}
|
|
}
|
|
|
|
// A module that subscribes to nothing needs no consumer, and creating one would leave an object
|
|
// nothing reads and everything has to maintain.
|
|
func TestAModuleThatConsumesNothingGetsNoConsumer(t *testing.T) {
|
|
if _, ok := ConsumerFor(Principal{Kind: KindModule, Node: "one", Module: "shop",
|
|
Emits: []string{"order.placed"}, PasswordHash: "x"}); ok {
|
|
t.Fatal("a pure emitter got a consumer")
|
|
}
|
|
}
|
|
|
|
// The seat is authority and the queue group is delivery. Tie them together and the day somebody
|
|
// allows two holders, every message is processed twice with nothing reporting it.
|
|
func TestAHoldersWorkerUsesAQueueGroupAnyway(t *testing.T) {
|
|
c, ok := HolderConsumerFor("one", "telegram", telegramSeat())
|
|
if !ok {
|
|
t.Fatal("the holder of a seat with inbound work got no worker")
|
|
}
|
|
if c.Queue == "" {
|
|
t.Fatal("the worker is not in a queue group, so a second holder would double-process")
|
|
}
|
|
if c.Stream != "SEAT_TELEGRAM_SENDER" {
|
|
t.Fatalf("the worker reads %q, not the seat's own stream", c.Stream)
|
|
}
|
|
if c.MaxDeliver == 0 {
|
|
t.Fatal("a failing worker would redeliver forever rather than dead-letter")
|
|
}
|
|
}
|
|
|
|
// The stream is named after the seat, not its holder: the holder can change and the queued work
|
|
// must not care.
|
|
func TestASeatsStreamIsNamedAfterTheSeat(t *testing.T) {
|
|
name := seatStreamName("telegram-sender")
|
|
if strings.Contains(name, "telegram-sender") {
|
|
t.Fatalf("%q keeps characters a stream name may not hold", name)
|
|
}
|
|
if name != "SEAT_TELEGRAM_SENDER" {
|
|
t.Fatalf("unexpected stream name %q", name)
|
|
}
|
|
}
|