The live tests reached one shared bus and assert, read and remove the mesh's own objects by their fixed names, so packages run in parallel deleted what each other read and the suite passed only one package at a time; a red suite read as noise. internal/testbus starts a server per test, linked in at the nats-server release go.mod pins, and a test holds that pin to the catalogue's bus image and to the facts snapshot's bus when there is one, so the tests never run a bus the mesh does not. The waiter test read a timing (the most connections held at one look) and now reads the state it means (the fewest held across the wait). make check runs the packages in parallel under the race detector, with a timeout.
89 lines
3.3 KiB
Go
89 lines
3.3 KiB
Go
package broker
|
|
|
|
import (
|
|
"fmt"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/testbus"
|
|
)
|
|
|
|
// A consumer that reacts to announcements, made on a stream that keeps a week of them, starts from now:
|
|
// made from the start it replays every merge and build of that week (novox/hq issue 248). And a reset
|
|
// re-makes a stuck one from now, its configuration otherwise kept, and refuses a work queue.
|
|
//
|
|
// go test ./internal/broker/ -run TestAConsumerMadeFromNow (each test on a bus of its own: internal/testbus)
|
|
func TestAConsumerMadeFromNowNeverReplaysTheStreamsHistory(t *testing.T) {
|
|
url := testbus.URL(t)
|
|
js, err := Dial(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer js.Close()
|
|
const stream = "HISTORY_TEST"
|
|
_ = js.js.DeleteStream(stream)
|
|
if _, err := js.js.AddStream(&nats.StreamConfig{Name: stream, Subjects: []string{"history.>"}, Storage: nats.MemoryStorage}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = js.js.DeleteStream(stream) }()
|
|
for i := 0; i < 5; i++ {
|
|
if _, err := js.js.Publish("history.merged", []byte(fmt.Sprint(i))); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
fresh := Consumer{Name: "fresh", Stream: stream, Filters: []string{"history.merged"}, Push: true,
|
|
AckWaitSeconds: 30, MaxDeliver: 5, MaxAckPending: 1, FromNow: true}
|
|
if err := js.EnsureConsumer(fresh); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
info, _ := js.js.ConsumerInfo(stream, "fresh")
|
|
if info.NumPending != 0 || info.Config.DeliverPolicy != nats.DeliverNewPolicy {
|
|
t.Fatalf("a consumer made from now holds %d of the stream's past (policy %v)", info.NumPending, info.Config.DeliverPolicy)
|
|
}
|
|
_, _ = js.js.Publish("history.merged", []byte("new"))
|
|
time.Sleep(100 * time.Millisecond)
|
|
if info, _ = js.js.ConsumerInfo(stream, "fresh"); info.NumPending != 1 {
|
|
t.Fatalf("a new announcement is not pending: %d", info.NumPending)
|
|
}
|
|
// Asserted again, it keeps where it is.
|
|
if err := js.EnsureConsumer(fresh); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
// One made from the start, as the server's default makes it, and stuck behind its history.
|
|
old := fresh
|
|
old.Name, old.FromNow = "old", false
|
|
if err := js.EnsureConsumer(old); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if info, _ = js.js.ConsumerInfo(stream, "old"); info.NumPending != 6 {
|
|
t.Fatalf("the default should replay all six: %d", info.NumPending)
|
|
}
|
|
before, after, err := js.ResetConsumer(stream, "old")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
info, _ = js.js.ConsumerInfo(stream, "old")
|
|
if before.Pending != 6 || after.Pending != 0 || info.Config.MaxAckPending != 1 || info.Config.MaxDeliver != 5 ||
|
|
info.Config.AckWait != 30*time.Second || info.Config.DeliverSubject == "" {
|
|
t.Fatalf("before %+v after %+v config %+v", before, after, info.Config)
|
|
}
|
|
|
|
// Never a work queue: what is pending there is work.
|
|
const queue = "QUEUE_TEST"
|
|
_ = js.js.DeleteStream(queue)
|
|
if _, err := js.js.AddStream(&nats.StreamConfig{Name: queue, Subjects: []string{"queue.>"}, Retention: nats.WorkQueuePolicy, Storage: nats.MemoryStorage}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = js.js.DeleteStream(queue) }()
|
|
if _, err := js.js.AddConsumer(queue, &nats.ConsumerConfig{Durable: "w", AckPolicy: nats.AckExplicitPolicy}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, _, err := js.ResetConsumer(queue, "w"); err == nil {
|
|
t.Fatal("a work queue's consumer was reset")
|
|
}
|
|
}
|