The controller owns a worker's shape, type included: one of the wrong type is re-made on a work queue (hq issue 206)
A holder built for a pull worker cannot bind a push one — `cannot pull subscribe to push based consumer` — and on 2026-10-03 the build machine rolled before the controller that would have redefined its worker, restarted on that for an hour, and nothing could build the controller that would have ended it. The server cannot change a consumer's type in place, so the assertion re-makes one of the wrong type: on a work queue nothing is lost, because what was acknowledged is gone from the stream and what was not is delivered again from the start. On a stream that keeps its history it is said and left, since a re-made consumer replays what this one acknowledged (issue 156), and that is a person's call. Proven against a real bus: a push worker with one ask acknowledged and two pending is re-made as pull, a pull subscription binds, and takes exactly the two.
This commit is contained in:
@@ -214,6 +214,43 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
|
|||||||
|
|
||||||
switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); {
|
switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); {
|
||||||
case err == nil:
|
case err == nil:
|
||||||
|
// **The controller owns the worker's shape, type included** (novox/hq issue 206). A holder
|
||||||
|
// built for a pull worker cannot bind a push one — `cannot pull subscribe to push based
|
||||||
|
// consumer` — and on 2026-10-03 the build machine rolled before the controller that would
|
||||||
|
// have redefined its worker, restarted on that for an hour, and nothing could build the
|
||||||
|
// controller that would have ended it. The server cannot change a consumer's type in place,
|
||||||
|
// so one of the wrong type is re-made: on a work queue nothing is lost, because what was
|
||||||
|
// acknowledged is gone from the stream and what was not is delivered again from the start.
|
||||||
|
// On any other stream a re-made consumer would replay what this one acknowledged (issue
|
||||||
|
// 156), so there it is said and left, and the person re-makes it knowing the cost.
|
||||||
|
if havePush, wantPush := have.Config.DeliverSubject != "", want.DeliverSubject != ""; havePush != wantPush {
|
||||||
|
shape := func(push bool) string {
|
||||||
|
if push {
|
||||||
|
return "push"
|
||||||
|
}
|
||||||
|
return "pull"
|
||||||
|
}
|
||||||
|
info, err := j.js.StreamInfo(c.Stream)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("asking about stream %s to re-make consumer %s: %w", c.Stream, c.Name, err)
|
||||||
|
}
|
||||||
|
if info.Config.Retention != nats.WorkQueuePolicy {
|
||||||
|
j.note("consumer %s on %s is %s and should be %s; not re-made, because %s keeps its history "+
|
||||||
|
"and a re-made consumer replays what this one acknowledged (novox/hq issue 156). Re-make it by hand",
|
||||||
|
c.Name, c.Stream, shape(havePush), shape(wantPush), c.Stream)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
j.note("consumer %s on %s changes from %s to %s delivery: re-made where it left off, nothing "+
|
||||||
|
"acknowledged comes back and nothing pending is lost (novox/hq issue 206); a holder bound to "+
|
||||||
|
"the old shape binds again", c.Name, c.Stream, shape(havePush), shape(wantPush))
|
||||||
|
if err := j.js.DeleteConsumer(c.Stream, c.Name); err != nil {
|
||||||
|
return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err)
|
||||||
|
}
|
||||||
|
if _, err := j.js.AddConsumer(c.Stream, want); err != nil {
|
||||||
|
return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
// Where an existing consumer starts is its history, not something an assertion may move:
|
// 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
|
// the server refuses a changed deliver policy outright. Carried across, so asserting twice
|
||||||
// is the no-op a restart depends on.
|
// is the no-op a restart depends on.
|
||||||
|
|||||||
@@ -0,0 +1,145 @@
|
|||||||
|
package broker
|
||||||
|
|
||||||
|
import (
|
||||||
|
"os"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/nats-io/nats.go"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A seat's worker that changed from push to pull delivery strands a holder built for the new shape
|
||||||
|
// (novox/hq issue 206): the server refuses a pull subscription on a push consumer, and the controller
|
||||||
|
// that would redefine it was the build that nobody could take. The controller owns the worker's
|
||||||
|
// shape, type included: on a work queue it re-makes one of the wrong type, losing nothing, and a
|
||||||
|
// pull subscription then binds and takes what was pending.
|
||||||
|
//
|
||||||
|
// 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 TestAWorker
|
||||||
|
func TestAWorkerOfTheWrongTypeIsRemadeOnAWorkQueueAndAPullThenBinds(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, worker, filter = "SEAT_T_SHELF", "SEAT_T_SHELF_worker", "mesh.seat.t-shelf.accept.>"
|
||||||
|
_ = js.js.DeleteStream(stream)
|
||||||
|
if _, err := js.js.AddStream(&nats.StreamConfig{
|
||||||
|
Name: stream, Subjects: []string{filter}, Retention: nats.WorkQueuePolicy, Storage: nats.MemoryStorage,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer func() { _ = js.js.DeleteStream(stream) }()
|
||||||
|
|
||||||
|
// The worker as the previous controller defined it: push, in a queue group.
|
||||||
|
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
|
||||||
|
Durable: worker, AckPolicy: nats.AckExplicitPolicy, AckWait: 60 * time.Second, MaxDeliver: 5,
|
||||||
|
FilterSubject: filter, DeliverSubject: "_DELIVER." + worker, DeliverGroup: "holders",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
for _, body := range []string{"one", "two", "three"} {
|
||||||
|
if _, err := js.js.Publish("mesh.seat.t-shelf.accept.build", []byte(body)); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// The old holder took and acknowledged the first ask, then went away.
|
||||||
|
old, err := js.js.QueueSubscribeSync(filter, "holders", nats.Bind(stream, worker))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
m, err := old.NextMsg(twoSeconds)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if string(m.Data) != "one" {
|
||||||
|
t.Fatalf("the first ask is %q", m.Data)
|
||||||
|
}
|
||||||
|
if err := m.AckSync(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := old.Unsubscribe(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The new controller asserts the worker as the mesh derives it now: pull.
|
||||||
|
if err := js.EnsureConsumer(Consumer{
|
||||||
|
Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60, MaxDeliver: 5,
|
||||||
|
Why: "the test's worker",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
have, err := js.js.ConsumerInfo(stream, worker)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if have.Config.DeliverSubject != "" || have.Config.DeliverGroup != "" {
|
||||||
|
t.Fatalf("the worker is still push: %+v", have.Config)
|
||||||
|
}
|
||||||
|
|
||||||
|
// A holder built for the new shape binds, and takes exactly what the old one left.
|
||||||
|
sub, err := js.js.PullSubscribe(filter, worker, nats.Bind(stream, worker), nats.ManualAck())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("a pull subscription does not bind the re-made worker: %v", err)
|
||||||
|
}
|
||||||
|
got, err := sub.Fetch(3, nats.MaxWait(twoSeconds))
|
||||||
|
if err != nil && len(got) == 0 {
|
||||||
|
t.Fatalf("nothing pending was delivered: %v", err)
|
||||||
|
}
|
||||||
|
var bodies []string
|
||||||
|
for _, g := range got {
|
||||||
|
bodies = append(bodies, string(g.Data))
|
||||||
|
_ = g.Ack()
|
||||||
|
}
|
||||||
|
if len(bodies) != 2 || bodies[0] != "two" || bodies[1] != "three" {
|
||||||
|
t.Fatalf("the pending asks after the acknowledged one, in order: %v", bodies)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Asserted again, the pull worker is the no-op a restart depends on.
|
||||||
|
if err := js.EnsureConsumer(Consumer{
|
||||||
|
Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60, MaxDeliver: 5,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// On a stream that keeps its history, a worker of the wrong type is said and left: re-making it would
|
||||||
|
// replay what it acknowledged (novox/hq issue 156), and that is a person's call.
|
||||||
|
func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(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, worker, filter = "EVENTS_T", "EVENTS_T_reader", "mesh.t.event.>"
|
||||||
|
_ = js.js.DeleteStream(stream)
|
||||||
|
if _, err := js.js.AddStream(&nats.StreamConfig{Name: stream, Subjects: []string{filter}, Storage: nats.MemoryStorage}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer func() { _ = js.js.DeleteStream(stream) }()
|
||||||
|
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
|
||||||
|
Durable: worker, AckPolicy: nats.AckExplicitPolicy, AckWait: 60 * time.Second,
|
||||||
|
FilterSubject: filter, DeliverSubject: "_DELIVER." + worker,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := js.EnsureConsumer(Consumer{Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
have, err := js.js.ConsumerInfo(stream, worker)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if have.Config.DeliverSubject == "" {
|
||||||
|
t.Fatal("a history stream's consumer was re-made, which replays what it acknowledged")
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user