Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
65 commits
Select commit Hold shift + click to select a range
0375746
feat(run-engine): concurrency-key virtual-time key builders
1stvamp Jul 23, 2026
91fb435
feat(run-engine): ckVirtualTimeScheduling options flag
1stvamp Jul 23, 2026
3591bea
feat(run-engine): two-pass virtual-time CK dequeue command
1stvamp Jul 24, 2026
113904f
feat(run-engine): register CK variants in vtime index on enqueue
1stvamp Jul 24, 2026
1765b77
feat(run-engine): re-register CK variants in vtime index on nack
1stvamp Jul 24, 2026
f646766
test(run-engine): fairness scenarios on the real batched dequeue path
1stvamp Jul 24, 2026
9dac71b
test(run-engine): multi-consumer correctness + op-count budget
1stvamp Jul 24, 2026
24b6220
test(run-engine): default-off regression proof
1stvamp Jul 24, 2026
bddc82e
feat(run-engine,webapp): wire ckVirtualTimeScheduling env flag (code-…
1stvamp Jul 24, 2026
86977ff
chore(run-engine): ship note + comment/format cleanup
1stvamp Jul 24, 2026
2b1c389
fix(run-engine,webapp): address whole-branch adversarial review
1stvamp Jul 24, 2026
5cb0aa5
docs(run-engine): record CK vtime known limitations for GA decision
1stvamp Jul 24, 2026
e9da64c
docs(run-engine): plan + references for virtual-time CK fair scheduling
1stvamp Jul 23, 2026
dce0b72
docs(run-engine): add fairness explainer diagrams
1stvamp Jul 24, 2026
21211ae
fix(run-engine): move design docs out of the Mintlify docs/ tree; format
1stvamp Jul 24, 2026
e1444eb
docs(run-engine): tighten CK fairness server-changes note
1stvamp Jul 24, 2026
dea03c6
docs(run-engine): fix stale references in the CK fairness design docs
1stvamp Jul 26, 2026
c804278
docs(run-engine): add CK virtual-time A/B benchmark plan and harness
1stvamp Jul 27, 2026
da3c8fb
docs(run-engine): make the e2e bench harness work on self-hosted
1stvamp Jul 27, 2026
4994295
docs(run-engine): add CK virtual-time benchmark results
1stvamp Jul 27, 2026
3dd789f
chore(run-engine): gitignore benchmark output artifacts
1stvamp Jul 27, 2026
31606d3
docs(run-engine): benchmark Redis CPU and memory vs concurrency-key c…
1stvamp Jul 28, 2026
9015810
chore(run-engine): drop benchmark, e2e, and design docs from the branch
1stvamp Jul 28, 2026
bb07b0b
fix(run-engine): stop an unservable variant pinning the ck virtual-ti…
1stvamp Jul 31, 2026
819e021
fix(run-engine): advance the ck virtual-time floor from servable vari…
1stvamp Jul 31, 2026
b932442
fix(run-engine): keep unservable ck variants in the fair order
1stvamp Jul 31, 2026
d8c5201
fix(run-engine): make ck vtime pass 2 discover unregistered variants
1stvamp Aug 3, 2026
f538f8d
test(run-engine): bound the ckManyKeys first-serve claim
1stvamp Aug 3, 2026
5cd0c65
test(run-engine): pin the ck vtime window-freeze and stranded-entry b…
1stvamp Aug 6, 2026
c2fd2be
docs(run-engine): note the ck vtime retry-storm fallback in the relea…
1stvamp Aug 6, 2026
3298ee2
fix(run-engine): stop unready ck variants blocking the fair pass, and…
1stvamp Aug 15, 2026
6ce9d2a
fix(run-engine): guard the wildcard cleanup in the ck vtime scripts
1stvamp Aug 15, 2026
701d164
fix(run-engine): register a gated ck variant so it can rejoin the fai…
1stvamp Aug 16, 2026
d21cbb9
fix(run-engine): drop drained ck variants from the fair order on ack,…
1stvamp Aug 17, 2026
af407f5
fix(run-engine): remember a drained ck variant's virtual time
1stvamp Aug 18, 2026
50351eb
fix(run-engine): bound the ck idle set by rank, not just by floor
1stvamp Aug 18, 2026
dca1e9c
perf(run-engine): cut two Redis calls from the ck vtime registration …
1stvamp Aug 19, 2026
2b6bfee
perf(run-engine): stop re-registering ck variants already in the fair…
1stvamp Aug 19, 2026
f5c1d74
fix(run-engine): keep the ck vtime floor alive alongside the tags it …
1stvamp Aug 19, 2026
d018b28
docs(run-engine): cut the ck fair scheduling release note down to the…
1stvamp Aug 19, 2026
c11b5c8
fix(run-engine): register a brand-new ck variant behind the pack, not…
1stvamp Aug 20, 2026
bb4d67d
test(run-engine): pin the fresh-key fix and isolate the ck idle rank cap
1stvamp Aug 20, 2026
f3ad8d6
fix(run-engine): stop the ck arrival cap from being reachable
1stvamp Aug 20, 2026
64b2fc9
fix(run-engine): stop concurrency-gated ck variants eating the fair pass
1stvamp Aug 21, 2026
e8aa263
test(run-engine): cover the future-head window guard a mutation audit…
1stvamp Aug 21, 2026
18e6dce
test(run-engine): pin the ck vtime credit round trip and TTL registra…
1stvamp Aug 21, 2026
0a9998d
fix(run-engine): stop the gated ck batch rewinding a live virtual-tim…
1stvamp Aug 21, 2026
2b7a533
test(run-engine): pin the parked tag a gated variant registers with
1stvamp Aug 21, 2026
88127be
docs(run-engine): cut the ck vtime comments back to what the code can…
1stvamp Aug 21, 2026
3c23a3a
fix(run-engine): re-register the TTL entry when the vtime dequeue exp…
1stvamp Aug 21, 2026
e69f000
fix(run-engine): keep parked ck credit when pass 2 repairs a variant
1stvamp Sep 7, 2026
a5fa891
refactor(run-engine): share one copy of the ck ack, nack and dead-let…
1stvamp Sep 15, 2026
838896a
fix(run-engine): restore the snapshotRoute read in the vtime TTL sweep
1stvamp Sep 15, 2026
0584290
refactor(run-engine): share one copy of the remaining ck lua
1stvamp Sep 15, 2026
c5c6bc9
fix(run-engine): make a mistyped ck lua slot a compile error
1stvamp Sep 15, 2026
c8872e2
test(run-engine): pin the snapshotRoute the vtime TTL sweep hands the…
1stvamp Sep 15, 2026
43d0722
test(run-engine): load every run-queue script into redis
1stvamp Sep 15, 2026
565ad2c
merge: integrate origin/main into ck virtual-time scheduling
1stvamp Sep 18, 2026
c4b657d
test(run-engine): thread groupConcurrencyKey into the two direct-Lua …
1stvamp Sep 19, 2026
c4af9aa
docs(run-engine): correct the ckExpireTtlLua sharing comment
1stvamp Sep 19, 2026
b02f4f0
feat(run-engine): warn when ck vtime and queue-gates are both enabled
1stvamp Sep 19, 2026
a9a14be
feat(run-engine): enforce queue-gates in the vtime dequeue and TTL sweep
1stvamp Sep 19, 2026
9040133
feat(webapp): allow a fractional ck vtime quantum
1stvamp Sep 19, 2026
fab8175
fix(run-engine): reject a non-finite ck vtime quantum
1stvamp Sep 20, 2026
027415b
Merge remote-tracking branch 'origin/main' into feat/ck-virtual-time-…
1stvamp Sep 22, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
fix(run-engine): stop concurrency-gated ck variants eating the fair pass
A candidate parked at its per-key concurrency ceiling fell out of tryServe
returning nil, so pass 1 spent one of its window slots on a variant it could not
serve. A gated variant's tag also stops advancing, so it keeps sorting to the
front of ckVtime and is revisited first on every call. Enough of them and pass 1
serves nothing, ever, and the scheduler quietly runs on pass 2's age order
instead.

The original note on this said work conservation still held because pass 2 fills
the batch, and that was the reason it was left alone. It does not hold. Where the
gated variants are also the oldest, which is the ordinary case since a variant
that has been queued longest is likely to be both old and saturated, pass 2's own
window fills with the same variants and servable work behind them is reached by
neither pass. The test added here starts from that shape and serves nothing at all
before the fix, rather than serving in the wrong order.

So a gated candidate now reports 'notReady', exactly as a future-scheduled head
already did, and pass 1 reads past it without spending a slot. The read is bounded
by scanLimit, which is already the cap on how far pass 1 will look.

It is not free. On a fully gated call pass 1 now reads to scanLimit instead of
stopping at the window: measured 53 Redis operations before and 80 after, every
one of the 27 a SCARD, so roughly 10 usec. That is worth paying, because a fully
gated call serves nothing either way, while a partially gated one goes from
serving nothing to serving in fair order. A second test pins the op count against
scanLimit plus the pass-2 window so the read cannot start running away.

Reported by Devin on #4367.
  • Loading branch information
1stvamp committed Sep 18, 2026
commit 64b2fc9557bcee329b59686764510f9f5f2df2e1
10 changes: 10 additions & 0 deletions internal-packages/run-engine/src/run-queue/index.ts
Comment thread
1stvamp marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -5712,6 +5712,16 @@ local function tryServe(ckQueueName, mayRaiseFloor, knownRegistered)
if gatedPending == nil then gatedPending = {} end
table.insert(gatedPending, ckQueueName)
end
-- NEW: report the gate the same way a future head is reported, so pass 1 declines to
-- spend a window slot on a candidate it cannot serve. Previously this fell through
-- returning nil and cost a slot, and because a gated variant's tag stops advancing it
-- also keeps sorting to the front and being revisited first, so enough of them
-- permanently consumed the fair pass and the scheduler ran on pass 2's age order
-- instead. Worse than that in the shape where the gated variants also hold the oldest
-- heads: pass 2's own window fills with them too and servable work behind them is
-- reached by neither pass for as long as the gate holds. Skipping without spending is
-- bounded by scanLimit, which is already the cap on how far pass 1 will read.
return 'notReady'
end
end

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
import { redisTest } from "@internal/testcontainers";
import { trace } from "@internal/tracing";
import { Logger } from "@trigger.dev/core/logger";
import { Decimal } from "@trigger.dev/database";
import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js";
import { RunQueue } from "../index.js";
import { RunQueueFullKeyProducer } from "../keyProducer.js";
const keys = new RunQueueFullKeyProducer();
const baseEnv: any = {
id: "e1234",
type: "DEVELOPMENT",
maximumConcurrencyLimit: 100,
concurrencyLimitBurstFactor: new Decimal(1),
project: { id: "p1234" },
organization: { id: "o1234" },
};
const QUEUE = "task/my-task";
const mk = (o: any) => ({
runId: "r1",
taskIdentifier: QUEUE,
orgId: "o1234",
projectId: "p1234",
environmentId: "e1234",
environmentType: "DEVELOPMENT",
queue: QUEUE,
timestamp: Date.now(),
attempt: 0,
...o,
});
// Declining to spend a window slot on a gated candidate means pass 1 reads further, so a
// fully-gated call costs more than it used to: measured 53 ops before and 80 after, the
// whole difference being SCARDs. That is the price of not silently degrading to age order,
// and it is worth paying because a fully-gated call serves nothing either way. What must
// not happen is the read running away, so this pins it against scanLimit (window * 2)
// rather than against the measured number, which would only be a tripwire.
describe("op count: fully gated dequeue", () => {
redisTest(
"pass 1 reads further when everything is gated, but stays inside scanLimit",
async ({ redisContainer }) => {
const MAX = 10; // window = 30, scanLimit = 60
const GATED = 80; // more than scanLimit, so both bounds bind
const kp = "rq:opc:";
const q: any = new RunQueue({
name: "rq",
tracer: trace.getTracer("rq"),
workers: 1,
defaultEnvConcurrency: 100,
logger: new Logger("RunQueue", "error"),
retryOptions: {
maxAttempts: 5,
factor: 1.1,
minTimeoutInMs: 100,
maxTimeoutInMs: 1000,
randomize: true,
},
keys,
masterQueueConsumersDisabled: true,
workerOptions: { disabled: true },
ckVirtualTimeScheduling: { enabled: true, scanWindowMultiplier: 3 },
queueSelectionStrategy: new FairQueueSelectionStrategy({
redis: { keyPrefix: kp, host: redisContainer.getHost(), port: redisContainer.getPort() },
keys,
}),
redis: { keyPrefix: kp, host: redisContainer.getHost(), port: redisContainer.getPort() },
} as any);
const env = baseEnv;
await q.updateEnvConcurrencyLimits(env);
const shard = keys.masterQueueShardForEnvironment(env.id, 2);
const t0 = Date.now() - 5000000;
for (let i = 0; i < GATED; i++) {
await q.enqueueMessage({
env,
message: mk({ runId: "g" + i, concurrencyKey: "gated-" + i, timestamp: t0 + i }),
workerQueue: env.id,
skipDequeueProcessing: true,
});
const members = Array.from({ length: 105 }, (_, k) => "busy-" + i + "-" + k);
await q.redis.sadd(
keys.queueKey(env, QUEUE, "gated-" + i) + ":currentConcurrency",
...members
);
}
await q.redis.config("RESETSTAT");
const served = await q.testDequeueFromMasterQueue(shard, env.id, MAX);
const info = await q.redis.info("commandstats");
let total = 0;
const per: Record<string, number> = {};
for (const line of info.split("\n")) {
const m = line.match(/^cmdstat_([a-z|]+):calls=(\d+)/);
if (!m || ["info", "config"].includes(m[1])) continue;
per[m[1]] = parseInt(m[2], 10);
total += parseInt(m[2], 10);
}
// Nothing is servable, so the call is pure scan.
expect(served.length).toBe(0);
// window = MAX * 3 = 30, scanLimit = 60. Pass 1 may read up to scanLimit candidates and
// pass 2 up to its own window, one SCARD each, so that sum is the ceiling.
const scanLimit = MAX * 3 * 2;
const pass2Window = MAX * 3;
expect(per.scard ?? 0).toBeLessThanOrEqual(scanLimit + pass2Window);
// And it must read past the old window bound, or the fix is not in effect.
expect(per.scard ?? 0).toBeGreaterThan(MAX * 3);
await q.quit();
},
120000
);
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,153 @@
import { redisTest } from "@internal/testcontainers";
import { trace } from "@internal/tracing";
import { Logger } from "@trigger.dev/core/logger";
import { Decimal } from "@trigger.dev/database";
import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js";
import { RunQueue } from "../index.js";
import { RunQueueFullKeyProducer } from "../keyProducer.js";

// Devin's finding B on #4367: a concurrency-gated candidate returns from tryServe without
// the 'notReady' marker, so it spends one of pass 1's window slots even though it can
// never be served. A gated variant also stops advancing its tag, so it keeps sorting to
// the front and is revisited first on every call. Fill the window with them and pass 1
// serves nothing, every call, and the scheduler silently degrades to pass 2's age order.
//
// Work conservation survives that, which is why it was originally waved through. What does
// not survive is the feature's whole purpose: fair order. This pins the difference.

const testOptions = {
name: "rq",
tracer: trace.getTracer("rq"),
workers: 1,
defaultEnvConcurrency: 100,
logger: new Logger("RunQueue", "error"),
retryOptions: {
maxAttempts: 5,
factor: 1.1,
minTimeoutInMs: 100,
maxTimeoutInMs: 1000,
randomize: true,
},
keys: new RunQueueFullKeyProducer(),
};
const baseEnv: any = {
id: "e1234",
type: "DEVELOPMENT",
maximumConcurrencyLimit: 100,
concurrencyLimitBurstFactor: new Decimal(1),
project: { id: "p1234" },
organization: { id: "o1234" },
};
const QUEUE = "task/my-task";
const makeMessage = (o: any) => ({
runId: "r1",
taskIdentifier: QUEUE,
orgId: "o1234",
projectId: "p1234",
environmentId: "e1234",
environmentType: "DEVELOPMENT",
queue: QUEUE,
timestamp: Date.now(),
attempt: 0,
...o,
});
const variantName = (ck: string) => testOptions.keys.queueKey(baseEnv, QUEUE, ck);

function createQueue(rc: any, keyPrefix: string) {
return new RunQueue({
...testOptions,
masterQueueConsumersDisabled: true,
workerOptions: { disabled: true },
ckVirtualTimeScheduling: { enabled: true, scanWindowMultiplier: 3 },
queueSelectionStrategy: new FairQueueSelectionStrategy({
redis: { keyPrefix, host: rc.getHost(), port: rc.getPort() },
keys: testOptions.keys,
}),
redis: { keyPrefix, host: rc.getHost(), port: rc.getPort() },
} as any) as any;
}

describe("CK vtime: gated variants and the pass-1 window", () => {
redisTest(
"gated variants must not spend the fair pass's budget",
async ({ redisContainer }) => {
const MAX = 2; // window = MAX * 3 = 6
const GATED = 8; // more than the window, all sorting ahead on tag
const queue = createQueue(redisContainer, "runqueue:test:gatedwin:");

try {
const env = baseEnv;
await queue.updateEnvConcurrencyLimits(env);
const shard = testOptions.keys.masterQueueShardForEnvironment(env.id, 2);
const t0 = Date.now() - 5_000_000;

// Gated variants: queued work, but each parked at its per-key ceiling so it can
// never be served. Enqueued first so their heads are oldest too.
for (let i = 0; i < GATED; i++) {
await queue.enqueueMessage({
env,
message: makeMessage({
runId: `g-${i}`,
concurrencyKey: `gated-${i}`,
timestamp: t0 + i,
}),
workerQueue: env.id,
skipDequeueProcessing: true,
});
}

// "owed" is what fair order says to serve next: the lowest tag among servable
// variants. Its head is the NEWEST, so age order would put it last.
await queue.enqueueMessage({
env,
message: makeMessage({
runId: "owed-0",
concurrencyKey: "owed",
timestamp: t0 + 900_000,
}),
workerQueue: env.id,
skipDequeueProcessing: true,
});
// "old" has the OLDEST head of the servable pair but a higher tag, so age order
// serves it first and fair order serves it second.
await queue.enqueueMessage({
env,
message: makeMessage({ runId: "old-0", concurrencyKey: "old", timestamp: t0 + 100 }),
workerQueue: env.id,
skipDequeueProcessing: true,
});

// Park every gated variant at its ceiling.
const limit = 100;
for (let i = 0; i < GATED; i++) {
const members = Array.from({ length: limit + 5 }, (_, k) => `busy-${i}-${k}`);
await queue.redis.sadd(`${variantName(`gated-${i}`)}:currentConcurrency`, ...members);
}

// Tags: gated variants lowest so they lead pass 1, then owed, then old.
const ckv = testOptions.keys.ckVtimeKeyFromQueue(variantName("owed"));
for (let i = 0; i < GATED; i++) await queue.redis.zadd(ckv, 0, variantName(`gated-${i}`));
await queue.redis.zadd(ckv, 1, variantName("owed"));
await queue.redis.zadd(ckv, 5, variantName("old"));

const served: string[] = [];
for (let c = 0; c < 4 && served.length < 2; c++) {
for (const m of await queue.testDequeueFromMasterQueue(shard, env.id, MAX)) {
served.push(m.message.concurrencyKey as string);
await queue.acknowledgeMessage(env.organization.id, m.messageId, {
skipDequeueProcessing: true,
});
}
}

// Fair order is the point of the feature: lowest tag first. If gated variants have
// eaten the window, pass 1 served nothing and pass 2's age order ran instead,
// which puts "old" first.
expect(served[0]).toBe("owed");
} finally {
await queue.quit();
}
},
60_000
);
});
Original file line number Diff line number Diff line change
Expand Up @@ -230,17 +230,34 @@ describe("FairDequeuingStrategy", () => {
envId: "env-3",
});

const startDistribute1 = performance.now();
// Command counts rather than wall clock. What this test is really asserting is that
// the second call reuses the snapshot instead of rebuilding it, and the timing ratio
// it used to assert was a proxy for that: sub-millisecond durations compared as a
// ratio, which flakes the moment anything else is running on the box. Counting the
// commands the strategy issues measures the same thing and cannot be perturbed by
// load.
const counter = createRedisClient(redis);
const commandCount = async () => {
const info = await counter.info("commandstats");
let total = 0;
for (const line of info.split("\n")) {
const m = line.match(/^cmdstat_([a-z|]+):calls=(\d+)/);
if (!m || ["info", "config"].includes(m[1])) continue;
total += parseInt(m[2], 10);
}
return total;
};

await counter.config("RESETSTAT");
const before1 = await commandCount();

const envResult = await strategy.distributeFairQueuesFromParentQueue(
"parent-queue",
"consumer-1"
);
const result = flattenResults(envResult);

const distribute1Duration = performance.now() - startDistribute1;

console.log("First distribution took", distribute1Duration, "ms");
const distribute1Commands = (await commandCount()) - before1;

expect(result).toHaveLength(3);
// Should only get the two oldest queues
Expand All @@ -249,33 +266,32 @@ describe("FairDequeuingStrategy", () => {
const queue3 = keyProducer.queueKey("org-3", "proj-3", "env-3", "queue-3");
expect(result).toEqual([queue2, queue1, queue3]);

const startDistribute2 = performance.now();
const before2 = await commandCount();

const _result2 = await strategy.distributeFairQueuesFromParentQueue(
"parent-queue",
"consumer-1"
);

const distribute2Duration = performance.now() - startDistribute2;

console.log("Second distribution took", distribute2Duration, "ms");
const distribute2Commands = (await commandCount()) - before2;

// Make sure the second call is more than 2 times faster than the first
expect(distribute2Duration).toBeLessThan(distribute1Duration / 2);
// Reused snapshot: the second call does materially less Redis work than the first.
expect(distribute2Commands).toBeLessThan(distribute1Commands / 2);

const startDistribute3 = performance.now();
const before3 = await commandCount();

const _result3 = await strategy.distributeFairQueuesFromParentQueue(
"parent-queue",
"consumer-1"
);

const distribute3Duration = performance.now() - startDistribute3;
const distribute3Commands = (await commandCount()) - before3;

console.log("Third distribution took", distribute3Duration, "ms");
// The snapshot has aged out by now, so the third call rebuilds it and pays the full
// cost again rather than the cached one.
expect(distribute3Commands).toBeGreaterThan(distribute2Commands * 2);

// Make sure the third call is more than 4 times the second
expect(distribute3Duration).toBeGreaterThan(distribute2Duration * 2);
await counter.quit();
}
);

Expand Down