diff --git a/internal/broker/derived.go b/internal/broker/derived.go index f5f514a..b6bbb4d 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -40,7 +40,14 @@ type Consumer struct { AckWaitSeconds int // MaxDeliver before the message is dead-lettered; zero for the mesh's default. MaxDeliver int - Why string + // MaxAckPending is how many deliveries the server lets stand unacknowledged at once; zero for + // the server's default, which is many. **One, for a consumer handled one at a time** + // (novox/hq issue 175): a handler that builds for minutes keeps its own message alive with a + // heartbeat, but everything handed over behind it times out unacknowledged and comes back — + // 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 } // 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 6f81b5d..4fb76b6 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -179,6 +179,7 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { AckPolicy: nats.AckExplicitPolicy, AckWait: time.Duration(c.AckWaitSeconds) * time.Second, MaxDeliver: c.MaxDeliver, + MaxAckPending: c.MaxAckPending, DeliverGroup: c.Queue, DeliverSubject: "", Description: c.Why, diff --git a/internal/broker/streams.go b/internal/broker/streams.go index f6dcd2d..1be7620 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -244,8 +244,13 @@ func MeshConsumers() []Consumer { Push: true, AckWaitSeconds: 30, MaxDeliver: 5, - Why: "the two events the mesh's own controller reacts to; after max-deliver it " + - "dead-letters, because an announcement it cannot act on will not become actionable", + // One at a time (novox/hq issue 175): acting on a merge builds for minutes, and an + // 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, + 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", }, } } diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go index 9a49289..ae94143 100644 --- a/internal/broker/streams_test.go +++ b/internal/broker/streams_test.go @@ -260,3 +260,14 @@ func TestNoTwoConsumersDeliverOntoTheSameSubject(t *testing.T) { seen[subject] = c.Name + " on " + c.Stream } } + +// The controller's events consumer is handed one announcement at a time (novox/hq issue 175): a +// merge's handler builds for minutes, and what is queued behind it must wait on the server rather +// than time out on the client and be acted on twice. +func TestTheControllerTakesOneAnnouncementAtATime(t *testing.T) { + for _, c := range MeshConsumers() { + if c.Stream == "EVENTS" && c.Name == ControllerName && c.MaxAckPending != 1 { + t.Fatalf("the events consumer may have %d outstanding; one announcement at a time", c.MaxAckPending) + } + } +}