Use work queue policy, deliver all on restart

This commit is contained in:
Neil Alexander 2022-01-24 11:59:28 +00:00
parent 03a989d5c9
commit 9ddb8749c1
No known key found for this signature in database
GPG key ID: A02A2019A2BB0944
2 changed files with 4 additions and 4 deletions

View file

@ -105,9 +105,9 @@ func (r *Inputer) Start() error {
nats.MaxDeliver(0),
// Use a durable named consumer.
r.Durable,
// Only process one message at a time, rather than have NATS flood us with
// more messages when we're still busy working on the last one.
nats.MaxAckPending(1),
// If we've missed things in the stream, e.g. we restarted, then replay
// all of the queued messages that were waiting for us.
nats.DeliverAll(),
)
return err
}

View file

@ -24,7 +24,7 @@ var (
var streams = []*nats.StreamConfig{
{
Name: InputRoomEvent,
Retention: nats.InterestPolicy,
Retention: nats.WorkQueuePolicy,
Storage: nats.FileStorage,
},
{