On 2.10.29 such a consumer was moved past a message now and then without handing it over; the controller's events consumer has seven filters, and a merge on the stream never reached it. The test reproduces the skip on 2.10.29 and keeps the image's release equal to the server it tests.
170 lines
5.7 KiB
Go
170 lines
5.7 KiB
Go
package main
|
|
|
|
import (
|
|
"fmt"
|
|
"math/rand"
|
|
"os"
|
|
"regexp"
|
|
"runtime/debug"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats-server/v2/server"
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// The server these tests run is the server the image runs. The image is pinned by a digest, which
|
|
// names no release, so the Dockerfile says which release it is, and this keeps the two equal: a
|
|
// test of the server below is only worth something about the server the mesh actually runs.
|
|
func TestTheImageIsTheServerTestedHere(t *testing.T) {
|
|
raw, err := os.ReadFile("../../Dockerfile")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
said := regexp.MustCompile(`(?m)^# upstream: nats (\d+\.\d+\.\d+)-alpine$`).FindSubmatch(raw)
|
|
if said == nil {
|
|
t.Fatal("the Dockerfile does not say which nats release its digest is (`# upstream: nats X.Y.Z-alpine`)")
|
|
}
|
|
info, ok := debug.ReadBuildInfo()
|
|
if !ok {
|
|
t.Fatal("no build information to read the tested server's version from")
|
|
}
|
|
for _, dep := range info.Deps {
|
|
if dep.Path == "github.com/nats-io/nats-server/v2" {
|
|
if dep.Version != "v"+string(said[1]) {
|
|
t.Fatalf("the image runs nats %s and these tests run %s: move them together", said[1], dep.Version)
|
|
}
|
|
return
|
|
}
|
|
}
|
|
t.Fatal("these tests run no nats-server")
|
|
}
|
|
|
|
// **A consumer with several filters is handed every message** (novox/hq issue 266).
|
|
//
|
|
// The controller follows the forge's merges, build outcomes and providers' standings on one durable
|
|
// consumer with seven filters, one message at a time. On nats 2.10.29 such a consumer was moved past
|
|
// a message now and then without handing it over: nothing pending, nothing redelivered, the message
|
|
// on the stream and its consumer never told. On 2026-10-06 that message was a merge, and the modules
|
|
// built from that repository were left behind with nobody told. This is that consumer, under traffic
|
|
// shaped like the mesh's — a stream mostly of build logs on ever-new subjects, the followed events
|
|
// few among them — and it fails on 2.10.29 (a few percent of the followed events never arrive) and
|
|
// passes on the release the image pins.
|
|
func TestAConsumerWithSeveralFiltersIsHandedEveryMessage(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip("runs a server under load for seconds")
|
|
}
|
|
s, err := server.NewServer(&server.Options{Port: -1, JetStream: true, StoreDir: t.TempDir(), NoLog: true, NoSigs: true})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
go s.Start()
|
|
if !s.ReadyForConnections(5 * time.Second) {
|
|
t.Fatal("the server did not start")
|
|
}
|
|
defer s.Shutdown()
|
|
|
|
admin, err := nats.Connect(s.ClientURL())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer admin.Close()
|
|
js, _ := admin.JetStream()
|
|
// The events stream and the controller's consumer on it, as the mesh declares them.
|
|
if _, err := js.AddStream(&nats.StreamConfig{Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>", "mesh.seat.*.event.>"},
|
|
MaxAge: 168 * time.Hour, MaxMsgsPerSubject: 10000, Storage: nats.FileStorage}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := js.AddConsumer("EVENTS", &nats.ConsumerConfig{Durable: "controller", DeliverSubject: "_DELIVER.controller.EVENTS",
|
|
FilterSubjects: []string{
|
|
"mesh.mod.mesh-catalog.event.upgraded",
|
|
"mesh.mod.mesh-catalog.event.catching-up",
|
|
"mesh.seat.node-build-agent.event.built",
|
|
"mesh.mod.gitea.event.pull.merged",
|
|
"mesh.seat.mesh-build-machine.event.built",
|
|
"mesh.mod.*.event.provisioner.failing",
|
|
"mesh.mod.*.event.provisioner.recovered",
|
|
},
|
|
AckPolicy: nats.AckExplicitPolicy, AckWait: 30 * time.Second, MaxDeliver: 5, MaxAckPending: 1,
|
|
DeliverPolicy: nats.DeliverNewPolicy}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
controller, err := nats.Connect(s.ClientURL())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer controller.Close()
|
|
cjs, _ := controller.JetStream()
|
|
handed := make(chan *nats.Msg, 64)
|
|
if _, err := cjs.ChanSubscribe("", handed, nats.Bind("EVENTS", "controller")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var got sync.Map
|
|
go func() {
|
|
for m := range handed {
|
|
if meta, err := m.Metadata(); err == nil {
|
|
got.Store(meta.Sequence.Stream, true)
|
|
}
|
|
time.Sleep(time.Duration(rand.Intn(20)) * time.Millisecond)
|
|
_ = m.Ack()
|
|
}
|
|
}()
|
|
|
|
followed := []string{"mesh.mod.gitea.event.pull.merged", "mesh.seat.node-build-agent.event.built",
|
|
"mesh.mod.postgres.event.provisioner.recovered", "mesh.mod.mesh-catalog.event.catching-up"}
|
|
var mu sync.Mutex
|
|
var sent []uint64
|
|
stop := time.Now().Add(5 * time.Second)
|
|
var wg sync.WaitGroup
|
|
for w := 0; w < 4; w++ {
|
|
wg.Add(1)
|
|
go func(w int) {
|
|
defer wg.Done()
|
|
nc, err := nats.Connect(s.ClientURL())
|
|
if err != nil {
|
|
t.Error(err)
|
|
return
|
|
}
|
|
defer nc.Close()
|
|
pjs, _ := nc.JetStream()
|
|
for i := 0; time.Now().Before(stop); i++ {
|
|
if rand.Intn(40) == 0 {
|
|
if ack, err := pjs.Publish(followed[rand.Intn(len(followed))], []byte(`{}`)); err == nil {
|
|
mu.Lock()
|
|
sent = append(sent, ack.Sequence)
|
|
mu.Unlock()
|
|
}
|
|
} else {
|
|
_, _ = pjs.Publish(fmt.Sprintf("mesh.seat.node-build-agent.event.log.build-%d-%d", w, i/50), make([]byte, 200))
|
|
}
|
|
if rand.Intn(100) == 0 {
|
|
time.Sleep(time.Duration(rand.Intn(300)) * time.Millisecond)
|
|
}
|
|
}
|
|
}(w)
|
|
}
|
|
wg.Wait()
|
|
|
|
// Every followed event handed over, given the consumer time to finish.
|
|
deadline := time.Now().Add(20 * time.Second)
|
|
for {
|
|
var missing []uint64
|
|
for _, seq := range sent {
|
|
if _, ok := got.Load(seq); !ok {
|
|
missing = append(missing, seq)
|
|
}
|
|
}
|
|
if len(missing) == 0 {
|
|
return
|
|
}
|
|
info, err := js.ConsumerInfo("EVENTS", "controller")
|
|
if time.Now().After(deadline) || (err == nil && info.NumPending == 0 && info.NumAckPending == 0) {
|
|
t.Fatalf("%d of %d followed events were never handed to the consumer (first at stream sequence %d), "+
|
|
"and it has nothing pending: the server moved past them", len(missing), len(sent), missing[0])
|
|
}
|
|
time.Sleep(200 * time.Millisecond)
|
|
}
|
|
}
|