From 8cae45ed245725d21ef143e03d20e89669eff4dc Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Fri, 28 Aug 2026 19:14:53 +0100 Subject: [PATCH] fix(run-engine,core): non-blocking test polls and document the zero total limit The waitFor polls dequeued with the default 10s blocking pop, which could blow past the helper deadline on a slow runner; poll non-blocking instead. Document that totalConcurrencyLimit: 0 holds every keyed run, matching concurrencyLimit's zero semantics. --- .../src/run-queue/tests/totalConcurrency.test.ts | 8 ++++++-- packages/core/src/v3/types/queues.ts | 3 +++ 2 files changed, 9 insertions(+), 2 deletions(-) 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; };