package broker import ( "fmt" "os" "strings" "testing" "time" "github.com/nats-io/nats.go" ) 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). // // docker run -d --rm --name t -p 14231:4222 nats:2.10-alpine -js // MESH_TEST_NATS=nats://127.0.0.1:14231 go test ./internal/broker/ -run TestUpgrading func TestUpgradingAConsumerWhoseDeliverySubjectMoved(t *testing.T) { url := os.Getenv("MESH_TEST_NATS") if url == "" { t.Skip("MESH_TEST_NATS unset") } 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 := os.Getenv("MESH_TEST_NATS") if url == "" { t.Skip("MESH_TEST_NATS unset") } 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)) } }