The live tests reached one shared bus and assert, read and remove the mesh's own objects by their fixed names, so packages run in parallel deleted what each other read and the suite passed only one package at a time; a red suite read as noise. internal/testbus starts a server per test, linked in at the nats-server release go.mod pins, and a test holds that pin to the catalogue's bus image and to the facts snapshot's bus when there is one, so the tests never run a bus the mesh does not. The waiter test read a timing (the most connections held at one look) and now reads the state it means (the fewest held across the wait). make check runs the packages in parallel under the race detector, with a timeout.
173 lines
5.9 KiB
Go
173 lines
5.9 KiB
Go
package broker
|
|
|
|
import (
|
|
"fmt"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/testbus"
|
|
)
|
|
|
|
const twoSeconds = 2 * time.Second
|
|
|
|
// A running mesh already holds consumers made before the delivery subject carried the stream
|
|
// (novox/hq 04-ISSUES/146). The server will not change a push consumer's delivery subject in place,
|
|
// so bringing one to match must replace it — and must not replay what it already acknowledged
|
|
// (novox/hq 04-ISSUES/156).
|
|
//
|
|
// go test ./internal/broker/ -run TestUpgrading (each test on a bus of its own: internal/testbus)
|
|
func TestUpgradingAConsumerWhoseDeliverySubjectMoved(t *testing.T) {
|
|
url := testbus.URL(t)
|
|
js, err := Dial(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer js.Close()
|
|
|
|
// The stream exactly as the mesh's own is — one declaration per node, always the newest.
|
|
// Reproduced rather than approximated: the first version of this test used a plain stream and
|
|
// a plain consumer, and the server accepted the update it refuses in a running mesh, so the
|
|
// test passed against the very code that was crash-looping on the control node.
|
|
const stream, name = "NODES", "novox"
|
|
subject := "mesh.node." + name + ".declare"
|
|
_ = js.js.DeleteStream(stream)
|
|
if _, err := js.js.AddStream(&nats.StreamConfig{
|
|
Name: stream, Subjects: []string{"mesh.node.*.declare"},
|
|
MaxMsgsPerSubject: 1, Storage: nats.MemoryStorage,
|
|
}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = js.js.DeleteStream(stream) }()
|
|
|
|
for i := 0; i < 6; i++ {
|
|
if _, err := js.js.Publish(subject, []byte(fmt.Sprint(i))); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
// The consumer as a running mesh holds it: made before the subject carried the stream, and
|
|
// otherwise exactly what NodeConsumer asks for.
|
|
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
|
|
Durable: name, AckPolicy: nats.AckExplicitPolicy,
|
|
AckWait: 300 * time.Second, MaxDeliver: -1,
|
|
FilterSubject: subject,
|
|
DeliverSubject: "_DELIVER." + name,
|
|
}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
// It acknowledged the first four. Those must not come back.
|
|
sub, err := js.js.SubscribeSync(subject, nats.Bind(stream, name))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
for i := 0; i < 1; i++ {
|
|
m, err := sub.NextMsg(twoSeconds)
|
|
if err != nil {
|
|
t.Fatalf("message %d never arrived: %v", i, err)
|
|
}
|
|
if err := m.AckSync(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
// **The subscription stays up.** In a running mesh the machine is attached to this consumer
|
|
// the whole time — that is what a node listening for its declaration IS. The first version of
|
|
// this test unsubscribed first, and the server then accepted an update it refuses while a
|
|
// subscriber is bound, so the test passed against the code that was crash-looping.
|
|
defer func() { _ = sub.Unsubscribe() }()
|
|
|
|
// Now the upgrade: the consumer the controller asserts on every start, with the subject that
|
|
// carries the stream.
|
|
want := NodeConsumer(name)
|
|
var notes []string
|
|
js.Note = func(f string, a ...any) { notes = append(notes, fmt.Sprintf(f, a...)) }
|
|
|
|
if err := js.EnsureConsumer(want); err != nil {
|
|
t.Fatalf("a consumer the mesh already held could not be brought to match, which is the "+
|
|
"control plane failing to start: %v", err)
|
|
}
|
|
|
|
info, err := js.js.ConsumerInfo(stream, name)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
// It KEEPS the subject it had. Moving it would need the holder's grant to have widened first,
|
|
// and that grant travels in the bus's user list, which a machine applies minutes later.
|
|
if got := info.Config.DeliverSubject; got != "_DELIVER."+name {
|
|
t.Fatalf("the consumer a machine is bound to was moved to %q; a machine not yet allowed "+
|
|
"to subscribe there is a machine that hears nothing", got)
|
|
}
|
|
if len(notes) != 1 {
|
|
t.Fatalf("keeping it was not reported, so it would be invisible: %v", notes)
|
|
}
|
|
if !strings.Contains(notes[0], "keeps working") {
|
|
t.Fatalf("the note does not say the consumer still works: %q", notes[0])
|
|
}
|
|
|
|
// And the machine bound to it is still being delivered to — the point of keeping it.
|
|
if _, err := js.js.Publish(subject, []byte("after the assertion")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
m, err := sub.NextMsg(twoSeconds)
|
|
if err != nil {
|
|
t.Fatalf("the machine stopped hearing its declarations after the assertion: %v", err)
|
|
}
|
|
if string(m.Data) != "after the assertion" {
|
|
t.Fatalf("delivered %q", m.Data)
|
|
}
|
|
|
|
// Asserting again is a no-op, or the controller crash-loops on its own restart.
|
|
if err := js.EnsureConsumer(want); err != nil {
|
|
t.Fatalf("the second assertion failed: %v", err)
|
|
}
|
|
}
|
|
|
|
// And where nothing is bound, the subject DOES move — that is 04-ISSUES/146's fix, which this must
|
|
// not undo. The controller's own two consumers are in exactly this position: it asserts them before
|
|
// it subscribes.
|
|
func TestAConsumerNothingIsBoundToDoesMove(t *testing.T) {
|
|
url := testbus.URL(t)
|
|
js, err := Dial(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer js.Close()
|
|
|
|
const stream, name = "NODES", "shanks"
|
|
subject := "mesh.node." + name + ".declare"
|
|
_ = js.js.DeleteStream(stream)
|
|
if _, err := js.js.AddStream(&nats.StreamConfig{
|
|
Name: stream, Subjects: []string{"mesh.node.*.declare"},
|
|
MaxMsgsPerSubject: 1, Storage: nats.MemoryStorage,
|
|
}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = js.js.DeleteStream(stream) }()
|
|
|
|
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
|
|
Durable: name, AckPolicy: nats.AckExplicitPolicy,
|
|
AckWait: 300 * time.Second, MaxDeliver: -1,
|
|
FilterSubject: subject,
|
|
DeliverSubject: "_DELIVER." + name,
|
|
}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
want := NodeConsumer(name)
|
|
if err := js.EnsureConsumer(want); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
info, err := js.js.ConsumerInfo(stream, name)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if got := info.Config.DeliverSubject; got != DeliverSubjectFor(want) {
|
|
t.Fatalf("delivery subject is %q, wanted %q -- issue 146's fix no longer applies to a "+
|
|
"consumer nothing is holding", got, DeliverSubjectFor(want))
|
|
}
|
|
}
|