diff --git a/internal/assignmentqueue/jetstream.go b/internal/assignmentqueue/jetstream.go index 93525a8..c3d8f14 100644 --- a/internal/assignmentqueue/jetstream.go +++ b/internal/assignmentqueue/jetstream.go @@ -167,7 +167,7 @@ func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectB AckPolicy: jetstream.AckExplicitPolicy, AckWait: 5 * time.Minute, MaxAckPending: capacity, - MaxDeliver: 20, + MaxDeliver: 1000, }) if err != nil { return nil, fmt.Errorf("open %s assignment consumer: %w", backend, err) diff --git a/internal/assignmentqueue/jetstream_test.go b/internal/assignmentqueue/jetstream_test.go index 6ff7254..4bc39bf 100644 --- a/internal/assignmentqueue/jetstream_test.go +++ b/internal/assignmentqueue/jetstream_test.go @@ -205,7 +205,7 @@ func TestOpenConsumerUsesIndependentDurablePerBackend(t *testing.T) { if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.BackendPod, 4); err != nil { t.Fatal(err) } - if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 { + if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 || manager.config.MaxDeliver != 1000 { t.Fatalf("config = %#v", manager.config) } }