diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 4870bfa..c2e4194 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -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 diff --git a/internal/broker/consumer_upgrade_test.go b/internal/broker/consumer_upgrade_test.go new file mode 100644 index 0000000..6ade5c0 --- /dev/null +++ b/internal/broker/consumer_upgrade_test.go @@ -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)) + } +} diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 5bbee69..6f81b5d 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -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