Skip to content

Commit 0aedf39

Browse files
authored
feat(run-engine): fair virtual-time scheduling for the concurrency-key dequeue (#4367)
Off by default. When many concurrency-key variants share one task queue, the dequeue serves the oldest waiting run first, so one key's large backlog is served to exhaustion while keys queued behind it wait for the whole pile to drain. This adds an opt-in fair order: each key gets a virtual clock, the dequeue serves the smallest clock and advances it, so keys take turns instead of one pile draining. With the flag off, the existing scripts run unchanged. ![The problem, and the fix](https://raw.githubusercontent.com/triggerdotdev/trigger.dev/84d4659b2565c7fe45eb31969d9e0858223a4d4e/internal-packages/run-engine/design/references/diagrams/fairness-problem-and-fix.png) ## How it works ![How the fix works](https://raw.githubusercontent.com/triggerdotdev/trigger.dev/84d4659b2565c7fe45eb31969d9e0858223a4d4e/internal-packages/run-engine/design/references/diagrams/fairness-how-it-works.png) - Three keys per base queue: `:ckVtime` holds a virtual-time tag per variant, `:ckVtimeFloor` is the current virtual time, and `:ckVtimeIdle` remembers the tags of variants that have drained. `ckIndex` keeps its head-timestamp domain, so time-eligibility, master-queue rebalancing and every other writer stay untouched, which is what makes it mixed-deploy safe. - A flag-selected two-pass dequeue. Pass 1 serves the lowest tags and charges each serve a quantum, bounded by a window of `maxCount * RUN_ENGINE_CK_VTIME_WINDOW_MULTIPLIER`. Pass 2 fills any leftover slots in today's age order and discovers variants that have no tag yet, so the command is a strict superset of today and stays work-conserving. - Only pass-1 serves move the floor, so a variant that cannot be served cannot drag the clock along behind it. - A draining variant parks its tag in the idle set and takes it back on its next enqueue or nack, so draining and returning isn't a way to reset your clock. All four routes out park it: the dequeue, ack, dead-letter and TTL expiry. - A brand-new variant joins one quantum behind the current leader rather than at the floor. Registering at the floor is what let a tenant sharding across fresh keys outrank everything already waiting, and it's now reserved for repair (pass-2 discovery and the gated batch), where everything in sight is established work that lost its tag. - A candidate that can't be served this call, because its head is scheduled in the future or it's at its per-key concurrency ceiling, says so instead of spending one of pass 1's window slots. - New behaviour lives only in new Lua command names. The existing enqueue/dequeue/nack scripts are byte-for-byte unchanged, which is checkable by hashing each script body, so flag-off is identical to today. ## Numbers Fairness, flag OFF against ON under the same load on the same box, wait measured in logical dequeue steps: | scenario | victim wait p99 | Jain fairness | | --- | --- | --- | | skewed backlog | 636 -> 196 (-69%) | 0.34 -> 1 | | trickle behind a backlog | 596 -> 176 (-70%) | 0.50 -> 1 | A fresh key minted per run against a 2000-deep backlog: the backlogged variant went from 1 slot in 600 to 120 in 600, which is what age order would have given it. Cost, measured saturated with the generator next to Redis, 1M invocations an arm: | path | flag-on cost | scaling | | --- | --- | --- | | enqueue | +1.8 usec | flat, 100 to 50k keys | | dequeue, serving | +3 to 4% call p95 | from the fairness runs above | | dequeue, every variant gated | 33.5 usec at 1k keys | 36.4 at 10k | | memory at rest | none | the idle set doesn't exist until something drains | | memory under churn | ~1.5MB per queue | capped by rank at 10k entries | Rule of thumb that fell out of the sweep, useful for costing anything else added to these scripts: about 1.38 usec fixed per EVALSHA plus 0.33 per `redis.call`. ## Testing 176 tests across the run-queue suite, green with the flag off and on. Fairness is proven on the real batched dequeue path rather than one message per call, plus multi-consumer exactly-once, a per-dequeue op-count budget, and behaviour tests for the floor, tag advance, GC, registration and the credit round trip through every drain path. Two mutation audits (33 mutations, one change to the production Lua at a time, rerun the suites, on the principle that a green run against a broken invariant is a hole). They found six things nothing was checking, all now covered: the pass-1 guard for a future-scheduled head, the idle park on ack, dead-letter and TTL expiry, nack's credit restore, the dead-letter caller, the TTL enqueue's own copy of the registration block, and any quantum other than 1. ## Rollout Off by default behind `RUN_ENGINE_CK_VTIME_SCHEDULING_ENABLED`, with `RUN_ENGINE_CK_VTIME_QUANTUM`, `RUN_ENGINE_CK_VTIME_WINDOW_MULTIPLIER` and `RUN_ENGINE_CK_VTIME_STATE_TTL_SECONDS` for the knobs. Enable on a staging cell, then production. Rollback is flipping the flag off, and leftover state expires within a day. During a rolling deploy, old instances serve in age order and are folded in by pass 2, so nothing is lost and no run is served twice. Every mutation is a single atomic Lua script, which is what makes those interleavings safe. That atomicity assumes the single-node Redis the run queue actually runs on: it has no cluster-mode setting (every other Redis in `env.server.ts` has one, `RUN_ENGINE_RUN_QUEUE_REDIS_*` doesn't), and the master queue key sits outside the base queue's hash slot exactly as it does in the command this one is modelled on. ## Known limitations **Registration order sets relative priority.** Within the fair pass a variant's tag comes from when it joined and how much it has been served, so message age doesn't break in. A key that arrives during a busy period sits a quantum behind the leader and keeps that place until the floor catches up. This is deliberate, since the alternative (strict fair share) makes minting fresh keys pay, and it's worth confirming against a real high-cardinality workload before the flag goes on anywhere. **Ties break lexically.** Variants on the same tag, at a cold start or when a collected variant re-registers, are served in queue-name order. It's a pre-existing effect of the old head-timestamp ordering and only affects who goes first, not long-run fairness. **The vtime scripts are copies.** Seven pairs, about 688 identical Lua lines, kept as copies so the flag-off path stays byte-identical. A change to one of the originals doesn't reach its copy, which already happened once on this branch (main added a TTL re-registration to `dequeueMessagesFromCkQueueTracked` and the copy went without it until 7f01aaf). A follow-up PR will pin each copy against its original so drift shows up as a red check.
1 parent 2fb4210 commit 0aedf39

22 files changed

Lines changed: 7666 additions & 1021 deletions
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: feature
4+
---
5+
6+
One concurrency key with a large backlog no longer holds up runs waiting on other keys on the same queue.

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

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1242,6 +1242,16 @@ const EnvironmentSchema = z
12421242
RUN_ENGINE_TTL_CONSUMERS_DISABLED: BoolEnv.default(false),
12431243
RUN_ENGINE_TTL_WORKER_BATCH_MAX_WAIT_MS: z.coerce.number().int().default(5_000),
12441244

1245+
// Fair (virtual-time) ordering across concurrency-key variants of a base queue.
1246+
// Off by default; when off the run queue behaves exactly as before.
1247+
RUN_ENGINE_CK_VTIME_SCHEDULING_ENABLED: BoolEnv.default(false),
1248+
// Fractional allowed: the vtime Lua serves weighted fair-queue tags, so a sub-1 quantum
1249+
// is a valid finer serve granularity. The weight hook is fixed at 1 today, so 1 stays the
1250+
// default, but the schema no longer blocks the capability the engine already has.
1251+
RUN_ENGINE_CK_VTIME_QUANTUM: z.coerce.number().finite().positive().default(1),
1252+
RUN_ENGINE_CK_VTIME_WINDOW_MULTIPLIER: z.coerce.number().int().positive().default(3),
1253+
RUN_ENGINE_CK_VTIME_STATE_TTL_SECONDS: z.coerce.number().int().positive().default(86400),
1254+
12451255
/** Optional maximum TTL for all runs (e.g. "14d"). If set, runs without an explicit TTL
12461256
* will use this as their TTL, and runs with a TTL larger than this will be clamped. */
12471257
RUN_ENGINE_DEFAULT_MAX_TTL: z.string().optional(),

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

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -116,6 +116,14 @@ function createRunEngine() {
116116
batchMaxSize: env.RUN_ENGINE_TTL_WORKER_BATCH_MAX_SIZE,
117117
batchMaxWaitMs: env.RUN_ENGINE_TTL_WORKER_BATCH_MAX_WAIT_MS,
118118
},
119+
ckVirtualTimeScheduling: env.RUN_ENGINE_CK_VTIME_SCHEDULING_ENABLED
120+
? {
121+
enabled: true,
122+
quantum: env.RUN_ENGINE_CK_VTIME_QUANTUM,
123+
scanWindowMultiplier: env.RUN_ENGINE_CK_VTIME_WINDOW_MULTIPLIER,
124+
stateTtlSeconds: env.RUN_ENGINE_CK_VTIME_STATE_TTL_SECONDS,
125+
}
126+
: undefined,
119127
},
120128
runLock: {
121129
redis: {

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -263,6 +263,7 @@ export class RunEngine {
263263
workerItemsSuffix: "ttl-worker:{queue:ttl-expiration:}items",
264264
visibilityTimeoutMs: options.queue?.ttlSystem?.visibilityTimeoutMs ?? 30_000,
265265
},
266+
ckVirtualTimeScheduling: options.queue?.ckVirtualTimeScheduling,
266267
});
267268

268269
this.worker = new Worker({

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

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ import {
1616
} from "@trigger.dev/redis-worker";
1717
import type { ControlPlaneResolver } from "./controlPlaneResolver.js";
1818
import type { FairQueueSelectionStrategyOptions } from "../run-queue/fairQueueSelectionStrategy.js";
19-
import type { RunQueueMetricsEmitter } from "../run-queue/index.js";
19+
import type { RunQueueMetricsEmitter, RunQueueOptions } from "../run-queue/index.js";
2020
import type { MinimalAuthenticatedEnvironment } from "../shared/index.js";
2121
import type { LockRetryConfig } from "./locking.js";
2222
import type { workerCatalog } from "./workerCatalog.js";
@@ -138,6 +138,9 @@ export type RunEngineOptions = {
138138
/** Max time (ms) to wait for more items before flushing a batch (default: 5000) */
139139
batchMaxWaitMs?: number;
140140
};
141+
/** Fair (virtual-time) ordering across concurrency-key variants of a base queue.
142+
* Passed through to RunQueue; off by default (undefined = today's behaviour). */
143+
ckVirtualTimeScheduling?: RunQueueOptions["ckVirtualTimeScheduling"];
141144
};
142145
runLock: {
143146
redis: RedisOptions;

0 commit comments

Comments
 (0)