diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 8b4a26e..72cf05f 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -214,6 +214,43 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); { 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: // the server refuses a changed deliver policy outright. Carried across, so asserting twice // is the no-op a restart depends on. diff --git a/internal/broker/worker_type_change_test.go b/internal/broker/worker_type_change_test.go new file mode 100644 index 0000000..dfe719c --- /dev/null +++ b/internal/broker/worker_type_change_test.go @@ -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") + } +}