Skip to content

Commit b81e175

Browse files
authored
feat(run-engine,sdk,webapp): cap a queue's total concurrency across keyed and keyless runs (#4823)
## Summary Adds a `totalConcurrencyLimit` option to queues. On a queue used with `concurrencyKey`, `concurrencyLimit` applies to each key value independently, so ten active keys with a limit of 5 can run 50 at once and nothing bounds the queue as a whole short of the environment limit. `totalConcurrencyLimit` caps in-flight runs across all keys while each key still gets at most `concurrencyLimit`: ```ts export const perUserQueue = queue({ name: "per-user-queue", concurrencyLimit: 1, totalConcurrencyLimit: 10, }); ``` Enforcement is gated behind `RUN_ENGINE_TOTAL_CONCURRENCY_LIMITS_ENABLED` (default off) and applies to runs triggered with a `concurrencyKey`. With the gate off, admit paths are unchanged. ## Design The engine keeps a per-base-queue `groupConcurrency` set shared by every concurrency-key variant; its cardinality is the queue's total in-flight count. The concurrency-key dequeue script bounds each batch by `min(totalLimit, envLimit) - SCARD(group)` and adds admitted runs to the set. The enqueue fast path checks the same gate and falls back to a normal enqueue when the queue is at its total limit. Every release path (ack, nack, dead-letter, concurrency release, TTL expiry, sweeper clear) removes a run from the group set whenever it removes it from the per-key `currentConcurrency` set. Those removals run regardless of the gate, so the set stays correct if the gate is later turned off, and the existing reconciliation sweep self-heals the group set because it is always a subset of the per-key sets it acks against. The stored limit is the raw declared value; readers clamp to the environment concurrency limit. The option flows through the queue manifest into a new nullable `TaskQueue.totalConcurrencyLimit` column (additive migration, internal name) and syncs to the engine on deploy. The group set self-heals against release paths that miss the removal (an instance on an older build during a rollout, or a future release script such as the one [#4398](#4398) adds). Every terminal release deletes the run's message key, so when a queue sits at its total the dequeue gate prunes members whose message key no longer exists, throttled to one pass per interval per queue. Runs in flight before the gate is enabled are the opposite case: they are absent from the set, so a queue can transiently exceed its total by at most that count, converging as each one completes. Not included here, planned as follow-ups: a runtime override API for the total limit, dashboard and metrics surfacing, and skipping queues at their total during fair-queue selection. A follow-up in this stack renames the public option to `combinedConcurrencyLimit` before release; engine internals and storage keep these names.
1 parent 8067b1f commit b81e175

18 files changed

Lines changed: 791 additions & 32 deletions

File tree

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
1+
---
2+
"@trigger.dev/sdk": patch
3+
"@trigger.dev/core": patch
4+
---
5+
6+
Cap a queue's total concurrency across all of its `concurrencyKey` values with the new `totalConcurrencyLimit` queue option. On a keyed queue, `concurrencyLimit` applies to each key value independently, so ten active keys with a limit of 5 can run 50 at once. `totalConcurrencyLimit` bounds the whole queue while each key still gets at most `concurrencyLimit`.
7+
8+
```ts
9+
import { queue } from "@trigger.dev/sdk";
10+
11+
export const perUserQueue = queue({
12+
name: "per-user-queue",
13+
concurrencyLimit: 1,
14+
totalConcurrencyLimit: 10,
15+
});
16+
```
17+
18+
Enforcement happens server-side and only applies to runs triggered with a `concurrencyKey`. Servers that have not enabled total concurrency limits accept the option but do not enforce it yet.

‎apps/webapp/app/env.server.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1440,6 +1440,7 @@ const EnvironmentSchema = z
14401440
RUN_ENGINE_RUN_QUEUE_LOG_LEVEL: z
14411441
.enum(["log", "error", "warn", "info", "debug"])
14421442
.default("info"),
1443+
RUN_ENGINE_TOTAL_CONCURRENCY_LIMITS_ENABLED: z.string().default("1"),
14431444
RUN_ENGINE_TREAT_PRODUCTION_EXECUTION_STALLS_AS_OOM: z.string().default("0"),
14441445
RUN_ENGINE_READ_REPLICA_SNAPSHOTS_SINCE_ENABLED: z.string().default("0"),
14451446
RUN_ENGINE_SNAPSHOTS_SINCE_REPLICA_RETRY_MIN_MS: z.coerce.number().int().default(50),

‎apps/webapp/app/v3/runEngine.server.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,7 @@ function createRunEngine() {
6262
queue: {
6363
defaultEnvConcurrency: env.DEFAULT_ENV_EXECUTION_CONCURRENCY_LIMIT,
6464
defaultEnvConcurrencyBurstFactor: env.DEFAULT_ENV_EXECUTION_CONCURRENCY_BURST_FACTOR,
65+
totalConcurrencyEnabled: env.RUN_ENGINE_TOTAL_CONCURRENCY_LIMITS_ENABLED === "1",
6566
logLevel: env.RUN_ENGINE_RUN_QUEUE_LOG_LEVEL,
6667
redis: {
6768
keyPrefix: "engine:",

‎apps/webapp/app/v3/runQueue.server.ts‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,23 @@ export async function updateQueueConcurrencyLimits(
4040
await engine.runQueue.updateQueueConcurrencyLimits(environment, queueName, concurrency);
4141
}
4242

43+
/** Updates the RunQueue total concurrency limit for a queue (the cap across all concurrency-key values) */
44+
export async function updateQueueTotalConcurrencyLimits(
45+
environment: AuthenticatedEnvironment,
46+
queueName: string,
47+
totalConcurrency: number
48+
) {
49+
await engine.runQueue.updateQueueTotalConcurrencyLimits(environment, queueName, totalConcurrency);
50+
}
51+
52+
/** Removes the RunQueue total concurrency limit for a queue */
53+
export async function removeQueueTotalConcurrencyLimits(
54+
environment: AuthenticatedEnvironment,
55+
queueName: string
56+
) {
57+
await engine.runQueue.removeQueueTotalConcurrencyLimits(environment, queueName);
58+
}
59+
4360
/** Removes the RunQueue limits for a queue */
4461
export async function removeQueueConcurrencyLimits(
4562
environment: AuthenticatedEnvironment,

‎apps/webapp/app/v3/services/createBackgroundWorker.server.ts‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,8 +35,10 @@ import { generateFriendlyId } from "../friendlyIdentifiers";
3535
import { engine } from "../runEngine.server";
3636
import {
3737
removeQueueConcurrencyLimits,
38+
removeQueueTotalConcurrencyLimits,
3839
updateEnvConcurrencyLimits,
3940
updateQueueConcurrencyLimits,
41+
updateQueueTotalConcurrencyLimits,
4042
} from "../runQueue.server";
4143
import { resolveScheduleWindow } from "@internal/schedule-engine";
4244
import { scheduleEngine } from "../scheduleEngine.server";
@@ -430,6 +432,7 @@ async function createWorkerTask(
430432
{
431433
name: task.queue?.name ?? `task/${task.id}`,
432434
concurrencyLimit: task.queue?.concurrencyLimit,
435+
totalConcurrencyLimit: task.queue?.totalConcurrencyLimit,
433436
},
434437
task.id,
435438
task.queue?.name ? "NAMED" : "VIRTUAL",
@@ -577,6 +580,7 @@ async function createWorkerQueue(
577580
const taskQueue = await upsertWorkerQueueRecord(
578581
queueName,
579582
baseConcurrencyLimit ?? null,
583+
queue.totalConcurrencyLimit ?? null,
580584
orderableName,
581585
queueType,
582586
worker,
@@ -585,6 +589,21 @@ async function createWorkerQueue(
585589

586590
const newConcurrencyLimit = taskQueue.concurrencyLimit;
587591

592+
/**
593+
* The total limit key is separate from the per-queue limit key that pause zeroes,
594+
* so it is safe to sync it regardless of the paused state. The engine clamps it
595+
* to the environment limit at read time, so the raw declared value is stored.
596+
*/
597+
if (typeof taskQueue.totalConcurrencyLimit === "number") {
598+
await updateQueueTotalConcurrencyLimits(
599+
environment,
600+
taskQueue.name,
601+
taskQueue.totalConcurrencyLimit
602+
);
603+
} else {
604+
await removeQueueTotalConcurrencyLimits(environment, taskQueue.name);
605+
}
606+
588607
if (!taskQueue.paused) {
589608
if (typeof newConcurrencyLimit === "number") {
590609
logger.debug("createWorkerQueue: updating concurrency limit", {
@@ -623,6 +642,7 @@ async function createWorkerQueue(
623642
async function upsertWorkerQueueRecord(
624643
queueName: string,
625644
concurrencyLimit: number | null,
645+
totalConcurrencyLimit: number | null,
626646
orderableName: string,
627647
queueType: TaskQueueType,
628648
worker: BackgroundWorker,
@@ -649,6 +669,7 @@ async function upsertWorkerQueueRecord(
649669
name: queueName,
650670
orderableName,
651671
concurrencyLimit,
672+
totalConcurrencyLimit,
652673
runtimeEnvironmentId: worker.runtimeEnvironmentId,
653674
projectId: worker.projectId,
654675
type: queueType,
@@ -673,6 +694,7 @@ async function upsertWorkerQueueRecord(
673694
// If overridden, keep current limit and update base; otherwise update limit normally
674695
concurrencyLimit: hasOverride ? undefined : concurrencyLimit,
675696
concurrencyLimitBase: hasOverride ? concurrencyLimit : undefined,
697+
totalConcurrencyLimit,
676698
},
677699
});
678700
}
@@ -684,6 +706,7 @@ async function upsertWorkerQueueRecord(
684706
return await upsertWorkerQueueRecord(
685707
queueName,
686708
concurrencyLimit,
709+
totalConcurrencyLimit,
687710
orderableName,
688711
queueType,
689712
worker,
Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,2 @@
1+
-- AlterTable
2+
ALTER TABLE "TaskQueue" ADD COLUMN "totalConcurrencyLimit" INTEGER;

‎internal-packages/database/prisma/schema.prisma‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1984,6 +1984,9 @@ model TaskQueue {
19841984
/// percentage (the source of truth). The absolute concurrencyLimit is materialized from it.
19851985
/// Decimal(5,2) allows fractional percentages like 12.50% (0.01–100.00).
19861986
concurrencyLimitOverridePercent Decimal? @db.Decimal(5, 2)
1987+
/// Caps total concurrent runs across ALL concurrencyKey values of this queue
1988+
/// (concurrencyLimit applies per key value). Null = no total cap.
1989+
totalConcurrencyLimit Int?
19871990
rateLimit Json?
19881991
19891992
paused Boolean @default(false)

‎internal-packages/run-engine/src/engine/index.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -225,6 +225,7 @@ export class RunEngine {
225225
queueSelectionStrategy: new FairQueueSelectionStrategy(queueSelectionStrategyOptions),
226226
defaultEnvConcurrency: options.queue?.defaultEnvConcurrency ?? 10,
227227
defaultEnvConcurrencyBurstFactor: options.queue?.defaultEnvConcurrencyBurstFactor,
228+
totalConcurrencyEnabled: options.queue?.totalConcurrencyEnabled,
228229
logger: new Logger("RunQueue", options.queue?.logLevel ?? "info"),
229230
redis: { ...options.queue.redis, keyPrefix: `${options.queue.redis.keyPrefix}runqueue:` },
230231
retryOptions: options.queue?.retryOptions,

‎internal-packages/run-engine/src/engine/types.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,8 @@ export type RunEngineOptions = {
9191
defaultEnvConcurrency?: number;
9292
defaultEnvConcurrencyBurstFactor?: number;
9393
logLevel?: LogLevel;
94+
/** Enforce per-queue total concurrency limits across concurrency-key variants. See RunQueueOptions.totalConcurrencyEnabled. */
95+
totalConcurrencyEnabled?: boolean;
9496
/** Optional queue-metrics emitter; enables gauge + counter emission from the RunQueue. */
9597
queueMetrics?: RunQueueMetricsEmitter;
9698
queueSelectionStrategyOptions?: Pick<

0 commit comments

Comments
 (0)