Make the controller's event consumer from now, and give a stuck one a reset

A consumer made with the server's default replays everything a stream that
keeps history holds: the controller's EVENTS consumer, re-made that way,
replayed a week of merges and builds one at a time and held every new one
behind them. FromNow makes it start at the end; broker consumer-reset
re-makes a stuck one from now, refusing a work queue. hq issue 244.
This commit is contained in:
jochen
2026-10-05 17:15:25 +02:00
parent 869fb6d6bf
commit 89ec48b9a9
5 changed files with 180 additions and 2 deletions
+33 -1
View File
@@ -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 <directory> | broker accounts --into <file>")
return errors.New("broker show | broker certificate [--check] --into <directory> | broker accounts --into <file> | " +
"broker consumer-reset <stream> <consumer>")
}
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 <stream> <consumer>, 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
+91
View File
@@ -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")
}
}
+7 -1
View File
@@ -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
+48
View File
@@ -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:
+1
View File
@@ -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",