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") } }