nats: run 2.11.17, which hands a consumer with several filters every message (hq issue 266)
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.
This commit is contained in:
@@ -7,6 +7,15 @@
|
||||
# that names a manifest rather than the mistake. This is the index — `docker pull` reports the same
|
||||
# one, and `RepoDigests` confirms it.
|
||||
#
|
||||
# **The release the digest is, said here because a digest does not say it:**
|
||||
#
|
||||
# upstream: nats 2.11.17-alpine
|
||||
#
|
||||
# Kept equal to the server version cmd/nats-tools tests against (its go.mod), and a test there fails
|
||||
# when they differ: that test is what says the server delivers every message to a consumer with
|
||||
# several filters. 2.10.29 did not — it moved such a consumer past a message now and then without
|
||||
# handing it over, and the controller never heard of a merge (novox/hq issue 266).
|
||||
#
|
||||
# Unlike every other module's Dockerfile, this builds no TypeScript and uses no mesh base image:
|
||||
# the module's code is the server, which upstream already built. There is no BUILD_BASE here on
|
||||
# purpose — nothing is compiled. The upstream image is declared in the manifest under build.on and
|
||||
|
||||
@@ -0,0 +1,169 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
+18
-1
@@ -2,4 +2,21 @@ module nats-tools
|
||||
|
||||
go 1.25.0
|
||||
|
||||
require git.novox.be/novox/mesh-sdk/go v0.1.7
|
||||
require (
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.7
|
||||
github.com/nats-io/nats-server/v2 v2.11.17
|
||||
github.com/nats-io/nats.go v1.51.0
|
||||
)
|
||||
|
||||
require (
|
||||
github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op // indirect
|
||||
github.com/google/go-tpm v0.9.8 // indirect
|
||||
github.com/klauspost/compress v1.18.5 // indirect
|
||||
github.com/minio/highwayhash v1.0.4 // indirect
|
||||
github.com/nats-io/jwt/v2 v2.8.1 // indirect
|
||||
github.com/nats-io/nkeys v0.4.15 // indirect
|
||||
github.com/nats-io/nuid v1.0.1 // indirect
|
||||
golang.org/x/crypto v0.50.0 // indirect
|
||||
golang.org/x/sys v0.43.0 // indirect
|
||||
golang.org/x/time v0.15.0 // indirect
|
||||
)
|
||||
|
||||
@@ -1,2 +1,27 @@
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.7 h1:C0sTQmtTiyYH7bnqZb7PusXnqA37gKuT7Nqjn9gG47w=
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.7/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY=
|
||||
github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op h1:Z/MZK75wC/NSrkgqeNIa7jexam9uWzhLmFTSCPI/kn0=
|
||||
github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op/go.mod h1:FQyySiasQQM8735Ddel3MRojmy4dA1IqCeyJ5jmPMbI=
|
||||
github.com/google/go-tpm v0.9.8 h1:slArAR9Ft+1ybZu0lBwpSmpwhRXaa85hWtMinMyRAWo=
|
||||
github.com/google/go-tpm v0.9.8/go.mod h1:h9jEsEECg7gtLis0upRBQU+GhYVH6jMjrFxI8u6bVUY=
|
||||
github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE=
|
||||
github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
|
||||
github.com/minio/highwayhash v1.0.4 h1:asJizugGgchQod2ja9NJlGOWq4s7KsAWr5XUc9Clgl4=
|
||||
github.com/minio/highwayhash v1.0.4/go.mod h1:GGYsuwP/fPD6Y9hMiXuapVvlIUEhFhMTh0rxU3ik1LQ=
|
||||
github.com/nats-io/jwt/v2 v2.8.1 h1:V0xpGuD/N8Mi+fQNDynXohVvp7ZztevW5io8CUWlPmU=
|
||||
github.com/nats-io/jwt/v2 v2.8.1/go.mod h1:nWnOEEiVMiKHQpnAy4eXlizVEtSfzacZ1Q43LIRavZg=
|
||||
github.com/nats-io/nats-server/v2 v2.11.17 h1:GKEghcFK6A+aFx11Yf1LjgLC3txAwvyhnYzhBIQZA8I=
|
||||
github.com/nats-io/nats-server/v2 v2.11.17/go.mod h1:B1sFVz4StNosQ903ak4N1G01Fl/9f8e06mXpFIE2K24=
|
||||
github.com/nats-io/nats.go v1.51.0 h1:ByW84XTz6W03GSSsygsZcA+xgKK8vPGaa/FCAAEHnAI=
|
||||
github.com/nats-io/nats.go v1.51.0/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno=
|
||||
github.com/nats-io/nkeys v0.4.15 h1:JACV5jRVO9V856KOapQ7x+EY8Jo3qw1vJt/9Jpwzkk4=
|
||||
github.com/nats-io/nkeys v0.4.15/go.mod h1:CpMchTXC9fxA5zrMo4KpySxNjiDVvr8ANOSZdiNfUrs=
|
||||
github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
|
||||
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
|
||||
golang.org/x/crypto v0.50.0 h1:zO47/JPrL6vsNkINmLoo/PH1gcxpls50DNogFvB5ZGI=
|
||||
golang.org/x/crypto v0.50.0/go.mod h1:3muZ7vA7PBCE6xgPX7nkzzjiUq87kRItoJQM1Yo8S+Q=
|
||||
golang.org/x/sys v0.21.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI=
|
||||
golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U=
|
||||
golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno=
|
||||
|
||||
@@ -88,7 +88,7 @@
|
||||
"on": [
|
||||
{
|
||||
"arg": "NATS_BASE",
|
||||
"image": "nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927"
|
||||
"image": "nats@sha256:e4bf19f15fd3218814a4e3c9e0064e1334bd8aa20d5984b9f1a0afd084f8cc00"
|
||||
}
|
||||
],
|
||||
"artifacts": [
|
||||
|
||||
Reference in New Issue
Block a user