The controller takes one announcement at a time #173
@@ -40,7 +40,14 @@ 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
|
||||||
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
|
// seatStreamName is the stream holding a seat's inbound work. Named after the seat rather than
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -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",
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user