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