92 lines
3.3 KiB
Go
92 lines
3.3 KiB
Go
package broker
|
|
|
|
import (
|
|
"fmt"
|
|
"os"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// 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.
|
|
//
|
|
// docker run -d --rm --name t -p 14231:4222 nats:2.10-alpine -js
|
|
// MESH_TEST_NATS=nats://127.0.0.1:14231 go test ./internal/broker/ -run TestAConsumerMadeFromNow
|
|
func TestAConsumerMadeFromNowNeverReplaysTheStreamsHistory(t *testing.T) {
|
|
url := os.Getenv("MESH_TEST_NATS")
|
|
if url == "" {
|
|
t.Skip("MESH_TEST_NATS unset")
|
|
}
|
|
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")
|
|
}
|
|
}
|