From c294949f2ad132c2a96cc7fc3a385e9ca6b0bfd8 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 11:07:11 +0200 Subject: [PATCH] A worker of the wrong type on a history-keeping stream is re-made to deliver from now on, never from the start (hq issue 207) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Left for a hand, the hand re-made it with the server's default — everything the stream holds — and on 2026-10-03 that replayed every build ask since 1 October into the catalogue. Re-made with deliver-new instead: nothing acknowledged comes back; what was in flight is said and asked again. --- internal/broker/jetstream.go | 22 ++++++++++----- internal/broker/worker_type_change_test.go | 32 ++++++++++++++++++---- 2 files changed, 42 insertions(+), 12 deletions(-) diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 72cf05f..cf86519 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -235,14 +235,22 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { 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 + // **A stream that keeps its history is re-made from now on, never from the start.** + // Left for a hand, the hand re-makes it with the server's default — everything the + // stream holds — which on 2026-10-03 replayed every build ask since 1 October and + // re-registered nine modules from the past (novox/hq issue 207). What this consumer + // had not yet acknowledged is lost with it, and said: on a history stream that is + // the smaller cost, and the asks in flight are visible to whoever asked. + j.note("consumer %s on %s changes from %s to %s delivery on a stream that keeps its history: "+ + "re-made to deliver from now on, so nothing this one acknowledged comes back (novox/hq issue "+ + "207); %d ask(s) it had not acknowledged are not carried over and must be asked again", + c.Name, c.Stream, shape(havePush), shape(wantPush), have.NumPending+uint64(have.NumAckPending)) + want.DeliverPolicy = nats.DeliverNewPolicy + } else { + 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)) } - 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) } diff --git a/internal/broker/worker_type_change_test.go b/internal/broker/worker_type_change_test.go index dfe719c..641ebd5 100644 --- a/internal/broker/worker_type_change_test.go +++ b/internal/broker/worker_type_change_test.go @@ -108,9 +108,10 @@ func TestAWorkerOfTheWrongTypeIsRemadeOnAWorkQueueAndAPullThenBinds(t *testing.T } } -// 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) { +// On a stream that keeps its history, a worker of the wrong type is re-made to deliver from now on: +// re-making it from the start would replay what it acknowledged (novox/hq issue 156), and leaving it +// for a hand re-made it exactly that way on 2026-10-03 (issue 207). +func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsRemadeFromNowOn(t *testing.T) { url := os.Getenv("MESH_TEST_NATS") if url == "" { t.Skip("MESH_TEST_NATS unset") @@ -132,6 +133,12 @@ func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(t *testing.T) { }); err != nil { t.Fatal(err) } + // History the old consumer would have acknowledged long ago, and must not come back. + for i := 0; i < 3; i++ { + if _, err := js.js.Publish("mesh.t.event.old", []byte("old")); err != nil { + t.Fatal(err) + } + } if err := js.EnsureConsumer(Consumer{Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60}); err != nil { t.Fatal(err) } @@ -139,7 +146,22 @@ func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(t *testing.T) { 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") + if have.Config.DeliverSubject != "" { + t.Fatal("a history stream's consumer of the wrong type was left as it was") + } + if have.Config.DeliverPolicy != nats.DeliverNewPolicy || have.NumPending != 0 { + t.Fatalf("re-made consumer delivers %v with %d pending; it must deliver from now on with nothing of the past", have.Config.DeliverPolicy, have.NumPending) + } + // And what arrives from now on is delivered. + if _, err := js.js.Publish("mesh.t.event.new", []byte("new")); err != nil { + t.Fatal(err) + } + sub, err := js.js.PullSubscribe(filter, worker, nats.Bind(stream, worker)) + if err != nil { + t.Fatal(err) + } + got, err := sub.Fetch(1, nats.MaxWait(3*time.Second)) + if err != nil || len(got) != 1 || string(got[0].Data) != "new" { + t.Fatalf("the re-made consumer delivered %v, %v; want the one new message", got, err) } }