The worker a holder bound was a push consumer in a queue group with one ask in flight: right for one holder, and with two it would still be a queue of one — the server hands a pushed ask to whichever subscriber it picks, busy or not, and the in-flight cap is per consumer, not per holder. Now the worker is pulled: every machine holding the seat binds the same durable and fetches one ask when it has finished the last, so an idle machine is the one that takes the next, the asks in flight are bounded by the holders working, and nothing is delivered that nobody asked for — which is also what ended the race issue 186 describes. A holder's grants trade the delivery subject for MSG.NEXT on the worker; the ack grant and the heartbeat that keeps a long build alive stay. Proven against a real bus: the build round trip, a backlog taken by a machine that arrives later, and work handed back by one machine coming round again.
182 lines
7.6 KiB
Go
182 lines
7.6 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 worker is delivery (novox/hq ADR 0190): one worker per seat, shared
|
|
// by every holder and pulled from, so a second holder takes the next ask rather than a copy of the
|
|
// same one — which is what a queue group used to guard, and what pulling one durable gives outright.
|
|
func TestAHoldersWorkerIsOneSharedByItsHolders(t *testing.T) {
|
|
c, ok := HolderConsumerFor("one", "telegram", telegramSeat())
|
|
if !ok {
|
|
t.Fatal("the holder of a seat with inbound work got no worker")
|
|
}
|
|
two, _ := HolderConsumerFor("two", "telegram", telegramSeat())
|
|
if c.Name != two.Name || c.Stream != two.Stream {
|
|
t.Fatal("two holders got two workers, so each would process every ask")
|
|
}
|
|
if c.Push || c.Queue != "" {
|
|
t.Fatal("the worker is pushed, so the server would hand an ask to a busy holder")
|
|
}
|
|
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])
|
|
}
|
|
|
|
// Every holder of a seat shares one worker and pulls from it (novox/hq ADR 0190): no queue group
|
|
// and no delivery subject, because a push consumer hands the next ask to whichever subscriber the
|
|
// server picks, busy or not; and no cap of one in flight, because pulling bounds the asks in flight
|
|
// by the holders that are free — which is what ended the race of issue 186, where asks delivered
|
|
// behind the one being worked expired and were dropped.
|
|
func TestAHoldersWorkerIsPulledByEveryHolder(t *testing.T) {
|
|
c, found := HolderConsumerFor("anchor", "build-agent", DeclaredSeat{Name: "node-build-agent", Accepts: []string{"build"}})
|
|
if !found {
|
|
t.Fatal("a seat that accepts work has no worker")
|
|
}
|
|
if c.Queue != "" || c.Push {
|
|
t.Fatalf("the worker is pushed (queue %q, push %v); a holder pulls when it is free", c.Queue, c.Push)
|
|
}
|
|
if c.MaxAckPending != 0 {
|
|
t.Fatalf("the worker caps asks in flight at %d; pulling bounds them by the holders working", c.MaxAckPending)
|
|
}
|
|
if c.Name != "SEAT_NODE_BUILD_AGENT_worker" || c.Stream != "SEAT_NODE_BUILD_AGENT" {
|
|
t.Fatalf("the worker is %s on %s; one per seat, shared by its holders", c.Name, c.Stream)
|
|
}
|
|
}
|