The controller takes one announcement at a time #173

Merged
jschoubben merged 1 commits from fix/one-announcement-at-a-time into main 2026-09-30 19:37:57 +00:00
4 changed files with 27 additions and 3 deletions
+7
View File
@@ -40,6 +40,13 @@ type Consumer struct {
AckWaitSeconds int AckWaitSeconds int
// MaxDeliver before the message is dead-lettered; zero for the mesh's default. // MaxDeliver before the message is dead-lettered; zero for the mesh's default.
MaxDeliver int MaxDeliver int
// 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 Why string
} }
+1
View File
@@ -179,6 +179,7 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
AckPolicy: nats.AckExplicitPolicy, AckPolicy: nats.AckExplicitPolicy,
AckWait: time.Duration(c.AckWaitSeconds) * time.Second, AckWait: time.Duration(c.AckWaitSeconds) * time.Second,
MaxDeliver: c.MaxDeliver, MaxDeliver: c.MaxDeliver,
MaxAckPending: c.MaxAckPending,
DeliverGroup: c.Queue, DeliverGroup: c.Queue,
DeliverSubject: "", DeliverSubject: "",
Description: c.Why, Description: c.Why,
+7 -2
View File
@@ -244,8 +244,13 @@ func MeshConsumers() []Consumer {
Push: true, Push: true,
AckWaitSeconds: 30, AckWaitSeconds: 30,
MaxDeliver: 5, MaxDeliver: 5,
Why: "the two events the mesh's own controller reacts to; after max-deliver it " + // One at a time (novox/hq issue 175): acting on a merge builds for minutes, and an
"dead-letters, because an announcement it cannot act on will not become actionable", // 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",
}, },
} }
} }
+11
View File
@@ -260,3 +260,14 @@ func TestNoTwoConsumersDeliverOntoTheSameSubject(t *testing.T) {
seen[subject] = c.Name + " on " + c.Stream 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)
}
}
}