Files
mesh-lab/replays/bus_test.go
T
jochen 70f62ee17b Replay the core incidents, and prove each fails before its fix and passes on it (hq to-be 45 §9)
Phase 5 is done when the replays of 236, 262, 263 and 266 fail on the commit before their fix and
pass after. The register names every replay with its issue and fix; the bus replay (266) runs a
consumer filtered like the controller's against a bus of a given release, the resolver replay (262)
renders the catalogue's machine list and asks every machine's name by getaddrinfo under musl and
glibc, and the prover runs each at both commits: all five (236, 262, 263, 266, 273) proved.
2026-10-06 21:01:39 +02:00

142 lines
4.6 KiB
Go

package replays
import (
"fmt"
"math/rand"
"os"
"sync"
"testing"
"time"
"github.com/nats-io/nats.go"
)
// **R266 — 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 several filters. On nats 2.10.29 such a consumer was moved past a message now and then
// without handing it over — nothing pending, nothing redelivered — and on 2026-10-06 that message was a
// merge nobody acted on. This is that consumer under traffic shaped like the mesh's, against the bus
// MESH_TEST_NATS names: in a merge check, a throwaway of the release the mesh runs; for the prover, the
// release the catalogue pinned at the commit it proves. Fails on 2.10.29, passes on 2.11.17.
//
// Everything it makes is under a prefix of its own, so it can run on a bus other things use.
func TestReplay266AConsumerWithSeveralFiltersIsHandedEveryMessage(t *testing.T) {
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS names no bus to replay against")
}
p := fmt.Sprintf("r266x%d", time.Now().UnixNano())
admin, err := nats.Connect(url)
if err != nil {
t.Fatal(err)
}
defer admin.Close()
js, err := admin.JetStream()
if err != nil {
t.Fatal(err)
}
stream := "R266_" + p
if _, err := js.AddStream(&nats.StreamConfig{Name: stream, Subjects: []string{p + ".mod.*.event.>", p + ".seat.*.event.>"},
MaxAge: time.Hour, MaxMsgsPerSubject: 10000, Storage: nats.FileStorage}); err != nil {
t.Fatal(err)
}
defer func() { _ = js.DeleteStream(stream) }()
// The controller's consumer, as the mesh declared it on the day: seven filters, one at a time.
if _, err := js.AddConsumer(stream, &nats.ConsumerConfig{Durable: "controller", DeliverSubject: "_DELIVER." + p,
FilterSubjects: []string{
p + ".mod.mesh-catalog.event.upgraded",
p + ".mod.mesh-catalog.event.catching-up",
p + ".seat.node-build-agent.event.built",
p + ".mod.gitea.event.pull.merged",
p + ".seat.mesh-build-machine.event.built",
p + ".mod.*.event.provisioner.failing",
p + ".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(url)
if err != nil {
t.Fatal(err)
}
defer controller.Close()
cjs, _ := controller.JetStream()
handed := make(chan *nats.Msg, 64)
sub, err := cjs.ChanSubscribe("", handed, nats.Bind(stream, "controller"))
if err != nil {
t.Fatal(err)
}
defer func() { _ = sub.Unsubscribe() }()
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{p + ".mod.gitea.event.pull.merged", p + ".seat.node-build-agent.event.built",
p + ".mod.postgres.event.provisioner.recovered", p + ".mod.mesh-catalog.event.catching-up"}
var mu sync.Mutex
var sent []uint64
stop := time.Now().Add(6 * 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(url)
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 {
// A build's log, on a subject of its own: what the stream mostly holds.
_, _ = pjs.Publish(fmt.Sprintf("%s.seat.node-build-agent.event.log.build-%d-%d", p, w, i/50), make([]byte, 200))
}
if rand.Intn(100) == 0 {
time.Sleep(time.Duration(rand.Intn(300)) * time.Millisecond)
}
}
}(w)
}
wg.Wait()
if len(sent) == 0 {
t.Fatal("no followed event was published")
}
deadline := time.Now().Add(30 * 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(stream, "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 — issue 266, on %s", len(missing), len(sent),
missing[0], admin.ConnectedServerVersion())
}
time.Sleep(200 * time.Millisecond)
}
}