diff --git a/cmd/mesh-controller/nodes.go b/cmd/mesh-controller/nodes.go index 04a38e4..ec08a4b 100644 --- a/cmd/mesh-controller/nodes.go +++ b/cmd/mesh-controller/nodes.go @@ -386,8 +386,12 @@ func brokerCommand(ctx context.Context, args []string) error { if len(args) > 0 && args[0] == "accounts" { return busAccounts(ctx, args[1:]) } + if len(args) > 0 && args[0] == "consumer-reset" { + return consumerReset(args[1:]) + } if len(args) == 0 || args[0] != "show" { - return errors.New("broker show | broker certificate [--check] --into | broker accounts --into ") + return errors.New("broker show | broker certificate [--check] --into | broker accounts --into | " + + "broker consumer-reset ") } known, err := broker.FromEnvironment() if errors.Is(err, broker.ErrNotConfigured) { @@ -407,6 +411,34 @@ func brokerCommand(ctx context.Context, args []string) error { return nil } +// consumerReset re-makes one consumer on a stream that keeps history to start from now (novox/hq issue +// 244): the way out of a consumer replaying a week of announcements, said rather than done by hand. A +// person's act — what was pending is dropped — so it is a command, and nothing calls it on its own. +func consumerReset(args []string) error { + if len(args) != 2 { + return errors.New("broker consumer-reset , e.g. broker consumer-reset EVENTS controller") + } + address, err := broker.BusAddress() + if err != nil { + return err + } + js, err := broker.Dial(address) + if err != nil { + return fmt.Errorf("cannot reach the bus: %w", err) + } + defer js.Close() + before, after, err := js.ResetConsumer(args[0], args[1]) + if err != nil { + return err + } + fmt.Printf("consumer %s on %s re-made to deliver from now\n", args[1], args[0]) + fmt.Printf(" before: delivers %s, delivered to %d, acknowledged to %d, %d pending, %d unacknowledged\n", + before.DeliverPolicy, before.Delivered, before.AckFloor, before.Pending, before.AckPending) + fmt.Printf(" after: delivers %s, %d pending; what was pending is dropped. A holder bound to it may need its "+ + "process restarted to bind again\n", after.DeliverPolicy, after.Pending) + return nil +} + // heardFrom says when a node was last heard from, in a form somebody can act on. // // "never" and "an hour ago" are different answers and are kept different. A node that has never diff --git a/internal/broker/consumer_from_now_test.go b/internal/broker/consumer_from_now_test.go new file mode 100644 index 0000000..94508f5 --- /dev/null +++ b/internal/broker/consumer_from_now_test.go @@ -0,0 +1,91 @@ +package broker + +import ( + "fmt" + "os" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// A consumer that reacts to announcements, made on a stream that keeps a week of them, starts from now: +// made from the start it replays every merge and build of that week (novox/hq issue 244). And a reset +// re-makes a stuck one from now, its configuration otherwise kept, and refuses a work queue. +// +// 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 TestAConsumerMadeFromNow +func TestAConsumerMadeFromNowNeverReplaysTheStreamsHistory(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 = "HISTORY_TEST" + _ = js.js.DeleteStream(stream) + if _, err := js.js.AddStream(&nats.StreamConfig{Name: stream, Subjects: []string{"history.>"}, Storage: nats.MemoryStorage}); err != nil { + t.Fatal(err) + } + defer func() { _ = js.js.DeleteStream(stream) }() + for i := 0; i < 5; i++ { + if _, err := js.js.Publish("history.merged", []byte(fmt.Sprint(i))); err != nil { + t.Fatal(err) + } + } + + fresh := Consumer{Name: "fresh", Stream: stream, Filters: []string{"history.merged"}, Push: true, + AckWaitSeconds: 30, MaxDeliver: 5, MaxAckPending: 1, FromNow: true} + if err := js.EnsureConsumer(fresh); err != nil { + t.Fatal(err) + } + info, _ := js.js.ConsumerInfo(stream, "fresh") + if info.NumPending != 0 || info.Config.DeliverPolicy != nats.DeliverNewPolicy { + t.Fatalf("a consumer made from now holds %d of the stream's past (policy %v)", info.NumPending, info.Config.DeliverPolicy) + } + _, _ = js.js.Publish("history.merged", []byte("new")) + time.Sleep(100 * time.Millisecond) + if info, _ = js.js.ConsumerInfo(stream, "fresh"); info.NumPending != 1 { + t.Fatalf("a new announcement is not pending: %d", info.NumPending) + } + // Asserted again, it keeps where it is. + if err := js.EnsureConsumer(fresh); err != nil { + t.Fatal(err) + } + + // One made from the start, as the server's default makes it, and stuck behind its history. + old := fresh + old.Name, old.FromNow = "old", false + if err := js.EnsureConsumer(old); err != nil { + t.Fatal(err) + } + if info, _ = js.js.ConsumerInfo(stream, "old"); info.NumPending != 6 { + t.Fatalf("the default should replay all six: %d", info.NumPending) + } + before, after, err := js.ResetConsumer(stream, "old") + if err != nil { + t.Fatal(err) + } + info, _ = js.js.ConsumerInfo(stream, "old") + if before.Pending != 6 || after.Pending != 0 || info.Config.MaxAckPending != 1 || info.Config.MaxDeliver != 5 || + info.Config.AckWait != 30*time.Second || info.Config.DeliverSubject == "" { + t.Fatalf("before %+v after %+v config %+v", before, after, info.Config) + } + + // Never a work queue: what is pending there is work. + const queue = "QUEUE_TEST" + _ = js.js.DeleteStream(queue) + if _, err := js.js.AddStream(&nats.StreamConfig{Name: queue, Subjects: []string{"queue.>"}, Retention: nats.WorkQueuePolicy, Storage: nats.MemoryStorage}); err != nil { + t.Fatal(err) + } + defer func() { _ = js.js.DeleteStream(queue) }() + if _, err := js.js.AddConsumer(queue, &nats.ConsumerConfig{Durable: "w", AckPolicy: nats.AckExplicitPolicy}); err != nil { + t.Fatal(err) + } + if _, _, err := js.ResetConsumer(queue, "w"); err == nil { + t.Fatal("a work queue's consumer was reset") + } +} diff --git a/internal/broker/derived.go b/internal/broker/derived.go index e83d662..7a96be6 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -47,7 +47,13 @@ type Consumer struct { // and a merge that came back rebuilt what it had just built, five times over on 2026-09-30. // With one outstanding, the server holds the rest, and the heartbeat is keeping the message. MaxAckPending int - Why string + // FromNow makes a consumer that does not exist yet start at the stream's end rather than its + // beginning (novox/hq issue 244). For a consumer that reacts to announcements — a merge, a build's + // outcome — on a stream that keeps a week of them: the server's default, everything the stream + // holds, replays every merge and every build of that week as if it had just happened. A consumer + // that exists keeps where it is, whatever this says; only its making is decided here. + FromNow bool + Why string } // seatStreamName is the stream holding a seat's inbound work. Named after the seat rather than diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 71a50e2..4c5680e 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -308,6 +308,9 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { } return nil case errors.Is(err, nats.ErrConsumerNotFound): + if c.FromNow { + want.DeliverPolicy = nats.DeliverNewPolicy + } if _, err := j.js.AddConsumer(c.Stream, want); err != nil { return fmt.Errorf("creating consumer %s on %s: %w", c.Name, c.Stream, err) } @@ -317,6 +320,51 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { } } +// ConsumerState is what a reset says about a consumer, before and after. +type ConsumerState struct { + DeliverPolicy string + Delivered uint64 + AckFloor uint64 + Pending uint64 + AckPending int +} + +func stateOf(info *nats.ConsumerInfo) ConsumerState { + policy, _ := info.Config.DeliverPolicy.MarshalJSON() + return ConsumerState{DeliverPolicy: strings.Trim(string(policy), `"`), Delivered: info.Delivered.Stream, + AckFloor: info.AckFloor.Stream, Pending: info.NumPending, AckPending: info.NumAckPending} +} + +// ResetConsumer re-makes a consumer to start from now, its configuration otherwise unchanged (novox/hq +// issue 244): what it had not yet delivered or acknowledged is dropped, which is the point — on a stream +// that keeps history, a consumer replaying a week of announcements does nothing anyone wants. Refused on +// a work queue, where what is pending is work nobody else will do. +func (j *JetStream) ResetConsumer(stream, name string) (before, after ConsumerState, err error) { + info, err := j.js.StreamInfo(stream) + if err != nil { + return before, after, fmt.Errorf("asking about stream %s: %w", stream, err) + } + if info.Config.Retention == nats.WorkQueuePolicy { + return before, after, fmt.Errorf("%s is a work queue: what its consumer has pending is work, and a reset would drop it", stream) + } + have, err := j.js.ConsumerInfo(stream, name) + if err != nil { + return before, after, fmt.Errorf("asking about consumer %s on %s: %w", name, stream, err) + } + before = stateOf(have) + want := have.Config + want.DeliverPolicy = nats.DeliverNewPolicy + want.OptStartSeq, want.OptStartTime = 0, nil + if err := j.js.DeleteConsumer(stream, name); err != nil { + return before, after, fmt.Errorf("removing consumer %s on %s: %w", name, stream, err) + } + made, err := j.js.AddConsumer(stream, &want) + if err != nil { + return before, after, fmt.Errorf("re-making consumer %s on %s — it is gone until the controller asserts it at its next start: %w", name, stream, err) + } + return before, stateOf(made), nil +} + func retentionOf(r Retention) nats.RetentionPolicy { switch r { case RetentionWorkQueue: diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 526cd65..87a32fa 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -270,6 +270,7 @@ func MeshConsumers() []Consumer { // announcement handed over behind it must wait on the server, not time out on the // client and come back to be acted on again. MaxAckPending: 1, + FromNow: true, Why: "the two events the mesh's own controller reacts to, one at a time; after " + "max-deliver it dead-letters, because an announcement it cannot act on will not " + "become actionable",