A consumer a machine is bound to keeps the subject that works

novox/hq 04-ISSUES/156. Issue 146 put the stream into a push consumer's
delivery subject. The server will not move that subject while a
subscriber is bound, and answers `consumer name already in use` — a
message about the name, for a conflict about the subject. A node is bound
to its declaration consumer the whole time it is up: that IS a node
listening. So every node consumer in a running mesh became one the
assertion could not bring to match, and the control plane crash-looped on
the assertion it makes before it serves. A fresh mesh showed nothing,
because nothing was bound.

Kept rather than deleted and re-made. Re-making moves the subject, and a
holder may not be allowed to subscribe to the new one yet: the wider
grant travels in the bus's user list, which this same control plane
composes and a machine applies minutes later. On the live mesh the nodes
are granted `_DELIVER.<node>` and not `_DELIVER.<node>.>`, so re-making
would have silenced every machine — worse than the collision it fixes,
and harder to undo.

Kept rather than fatal, which is what 146's change intended and did not
do. The bare subject still delivers, and collides only where one holder
has two consumers of one name. That is the controller's own pair, and the
controller is not bound to them while it asserts, so those do move.

Also: an existing consumer's deliver policy is carried across rather than
reasserted, because the server refuses to change it and where a consumer
starts is its history.

Two tests against a real server: a consumer with a subscriber bound keeps
its subject, is reported, and still delivers; one with nothing bound
moves, so 146's fix still applies where it matters.
This commit is contained in:
2026-09-29 23:55:56 +02:00
parent 6c5dfd0c25
commit e6ddc59cde
3 changed files with 230 additions and 1 deletions
+6
View File
@@ -717,6 +717,12 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
broker.BareAddress(address), err)
}
defer js.Close()
// What the raise decided not to fail over. Said, for the reason everything else here is said:
// a consumer kept as it was is a difference between what the mesh asked for and what the bus
// holds, and one nobody would find by reading either (novox/hq 04-ISSUES/156).
js.Note = func(format string, args ...any) {
fmt.Printf(" "+format+"\n", args...)
}
// **Its own user, before anything else.** The controller's account is created by the installer at
// a bootstrap password, before there is a controller to mint one — so nothing recorded a hash for
+178
View File
@@ -0,0 +1,178 @@
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))
}
}
+46 -1
View File
@@ -26,6 +26,16 @@ import (
type JetStream struct {
conn *nats.Conn
js nats.JetStreamContext
// Note is how this says something it decided not to fail over. Nil is silent, which is only
// right for a caller that has no way to report; the controller sets it.
Note func(string, ...any)
}
// note reports without requiring a caller to have set one.
func (j *JetStream) note(format string, args ...any) {
if j.Note != nil {
j.Note(format, args...)
}
}
// Dial connects and returns the controller's JetStream handle.
@@ -200,9 +210,44 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
want.DeliverSubject = DeliverSubjectFor(c)
}
switch _, err := j.js.ConsumerInfo(c.Stream, c.Name); {
switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); {
case err == nil:
// Where an existing consumer starts is its history, not something an assertion may move:
// the server refuses a changed deliver policy outright. Carried across, so asserting twice
// is the no-op a restart depends on.
want.DeliverPolicy = have.Config.DeliverPolicy
want.OptStartSeq = have.Config.OptStartSeq
want.OptStartTime = have.Config.OptStartTime
if _, err := j.js.UpdateConsumer(c.Stream, want); err != nil {
// **A consumer that works is not replaced to make its name tidier**
// (novox/hq 04-ISSUES/156).
//
// The server will not move a push consumer's delivery subject while a subscriber is
// bound to it, and answers `consumer name already in use` — a message about the name,
// for a conflict about the subject. A node is bound to its declaration consumer the
// whole time it is up; that IS a node listening. So when 04-ISSUES/146 put the stream
// into the subject, every node consumer in a running mesh became one this could not
// bring to match, and the control plane crash-looped on the assertion it makes before
// it serves. A fresh mesh showed nothing: nothing was bound.
//
// Kept rather than deleted and re-made. Re-making moves the subject, and a holder may
// not be allowed to subscribe to the new one yet — the wider grant travels in the bus's
// user list, which this same control plane composes and a machine applies minutes
// later. Re-making here would have silenced every machine in the mesh, which is worse
// than the collision it was fixing and harder to undo.
//
// Kept rather than fatal, which is what 146's change intended and did not do: the bare
// subject it replaces still delivers, and it collides only where one holder has two
// consumers of one name. That is the controller's own pair, and the controller is not
// bound to them while it asserts, so those do move. A node has one consumer and nothing
// to collide with.
if have.Config.DeliverSubject != want.DeliverSubject {
j.note("consumer %s on %s still delivers to %q and not %q: %v. It keeps working; "+
"the subject moves on an assertion made while nothing is bound to it",
c.Name, c.Stream, have.Config.DeliverSubject, want.DeliverSubject, err)
return nil
}
return fmt.Errorf("bringing consumer %s on %s to match: %w", c.Name, c.Stream, err)
}
return nil