With the worker consumer's default of many deliveries in flight, every ask behind the one being built was delivered at once, left unacknowledged for the length of the build, redelivered after the ack wait and dropped after the fifth time: on 2026-10-01 twenty-six of forty-three builds asked in two minutes were never built and the queue read as empty (hq issue 186). The holder's worker now has one in flight, and a running build tells the bus it is still working, as the controller's long handlers do, so a build longer than the ack wait is neither redelivered nor counted out.
169 lines
6.8 KiB
Go
169 lines
6.8 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)
|
|
}
|
|
}
|
|
|
|
// A node hears its declaration through a consumer only the controller can make.
|
|
//
|
|
// The three things that would each break it silently: a name other than the node's is one the host
|
|
// cannot acknowledge a delivery from, because its ack grant is derived from the node's name; a
|
|
// filter other than its own declaration subject is a node reading another's; and a pull consumer is
|
|
// one the host cannot bind a channel to without creating something, which it has no authority for.
|
|
func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) {
|
|
c := NodeConsumer("anchor")
|
|
if c.Name != "anchor" {
|
|
t.Fatalf("named %q, so the node cannot ack from it: its grant is $JS.ACK.NODES.anchor.>", c.Name)
|
|
}
|
|
if c.Stream != "NODES" {
|
|
t.Fatalf("on stream %q rather than the one declarations live in", c.Stream)
|
|
}
|
|
if len(c.Filters) != 1 || c.Filters[0] != "mesh.node.anchor.declare" {
|
|
t.Fatalf("filters %v, which is not this node's own declaration and nothing else", c.Filters)
|
|
}
|
|
if !c.Push {
|
|
t.Fatal("pulled, which a host cannot do: pulling needs the JetStream API and a host reaches none of it")
|
|
}
|
|
if c.MaxDeliver != 0 {
|
|
t.Fatalf("max-deliver %d: a declaration a node has not taken yet is not one to dead-letter, "+
|
|
"because the stream holds exactly one per node", c.MaxDeliver)
|
|
}
|
|
|
|
// And the grant the node actually gets has to match, or none of the above matters.
|
|
perms, err := PermissionsFor(Principal{Kind: KindNode, Node: "anchor"})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
// Without the ack grant every declaration a node receives is redelivered for ever; without the
|
|
// subscribe grant its consumer delivers to nobody.
|
|
has(t, perms.Publish, "$JS.ACK.NODES."+c.Name+".>")
|
|
has(t, perms.Subscribe, c.Filters[0])
|
|
}
|
|
|
|
// A holder works one ask at a time, so the server hands it one at a time (novox/hq issue 186):
|
|
// asks queued behind the one being worked wait in the stream rather than being delivered,
|
|
// left to expire and dropped after the fifth redelivery.
|
|
func TestAHoldersWorkerTakesOneAskAtATime(t *testing.T) {
|
|
c, found := HolderConsumerFor("anchor", "builder", DeclaredSeat{Name: "mesh-build-machine", Accepts: []string{"build"}})
|
|
if !found {
|
|
t.Fatal("a seat that accepts work has no worker")
|
|
}
|
|
if c.MaxAckPending != 1 {
|
|
t.Fatalf("the worker may have %d asks in flight; one, so a queue is a queue", c.MaxAckPending)
|
|
}
|
|
}
|