From 906e6a2e18cc2dcfa841537a0d5df6e678357fcf Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Mon, 21 Sep 2026 08:42:01 +0000 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=9D=E6=8C=81=E5=90=8E=E7=AB=AF=20a?= =?UTF-8?q?ssignment=20=E5=8F=AF=E9=87=8D=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/assignmentqueue/jetstream.go | 2 +- internal/assignmentqueue/jetstream_test.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) 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) } }