From 0275c2eeacbb0a1fca3d091b3d12dc6eab1a0c7e Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 6 Oct 2026 01:29:20 +0200 Subject: [PATCH] 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. --- modules/nats/Dockerfile | 9 ++ modules/nats/cmd/nats-tools/server_test.go | 169 +++++++++++++++++++++ modules/nats/go.mod | 19 ++- modules/nats/go.sum | 25 +++ modules/nats/module.json | 2 +- 5 files changed, 222 insertions(+), 2 deletions(-) create mode 100644 modules/nats/cmd/nats-tools/server_test.go diff --git a/modules/nats/Dockerfile b/modules/nats/Dockerfile index 79ad25a..418fa11 100644 --- a/modules/nats/Dockerfile +++ b/modules/nats/Dockerfile @@ -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 diff --git a/modules/nats/cmd/nats-tools/server_test.go b/modules/nats/cmd/nats-tools/server_test.go new file mode 100644 index 0000000..773469e --- /dev/null +++ b/modules/nats/cmd/nats-tools/server_test.go @@ -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) + } +} diff --git a/modules/nats/go.mod b/modules/nats/go.mod index a5442db..ac9888b 100644 --- a/modules/nats/go.mod +++ b/modules/nats/go.mod @@ -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 +) diff --git a/modules/nats/go.sum b/modules/nats/go.sum index b474419..b28aa2f 100644 --- a/modules/nats/go.sum +++ b/modules/nats/go.sum @@ -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= diff --git a/modules/nats/module.json b/modules/nats/module.json index 5e48301..8f8798b 100644 --- a/modules/nats/module.json +++ b/modules/nats/module.json @@ -88,7 +88,7 @@ "on": [ { "arg": "NATS_BASE", - "image": "nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927" + "image": "nats@sha256:e4bf19f15fd3218814a4e3c9e0064e1334bd8aa20d5984b9f1a0afd084f8cc00" } ], "artifacts": [