diff --git a/internal-packages/run-engine/src/run-queue/tests/totalConcurrency.test.ts b/internal-packages/run-engine/src/run-queue/tests/totalConcurrency.test.ts index cc52b1aca..7ad89c8c5 100644 --- a/internal-packages/run-engine/src/run-queue/tests/totalConcurrency.test.ts +++ b/internal-packages/run-engine/src/run-queue/tests/totalConcurrency.test.ts @@ -286,7 +286,9 @@ describe("RunQueue total concurrency limit", () => { }); const r1Admitted = await waitFor(async () => { - const next = await queue.dequeueMessageFromWorkerQueue("consumer-1", "main"); + const next = await queue.dequeueMessageFromWorkerQueue("consumer-1", "main", { + blockingPop: false, + }); return next?.messageId === "r1"; }); expect(r1Admitted).toBe(true); @@ -350,7 +352,9 @@ describe("RunQueue total concurrency limit", () => { }); const r1Admitted = await waitFor(async () => { - const next = await queue.dequeueMessageFromWorkerQueue("consumer-1", "main"); + const next = await queue.dequeueMessageFromWorkerQueue("consumer-1", "main", { + blockingPop: false, + }); return next?.messageId === "r1"; }); expect(r1Admitted).toBe(true); diff --git a/packages/core/src/v3/types/queues.ts b/packages/core/src/v3/types/queues.ts index 01d28fc07..1fba786c8 100644 --- a/packages/core/src/v3/types/queues.ts +++ b/packages/core/src/v3/types/queues.ts @@ -54,6 +54,9 @@ export type QueueOptions = { * ``` * * Only enforced for runs triggered with a `concurrencyKey`, and requires server-side support. + * + * Omit for no total cap. Like `concurrencyLimit`, a value of `0` holds every keyed run in + * the queue rather than removing the cap. */ totalConcurrencyLimit?: number; };