-
-
Notifications
You must be signed in to change notification settings - Fork 1.5k
feat(run-engine): fair virtual-time scheduling for the concurrency-key dequeue #4367
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
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 91fb435
feat(run-engine): ckVirtualTimeScheduling options flag
1stvamp 3591bea
feat(run-engine): two-pass virtual-time CK dequeue command
1stvamp 113904f
feat(run-engine): register CK variants in vtime index on enqueue
1stvamp 1765b77
feat(run-engine): re-register CK variants in vtime index on nack
1stvamp f646766
test(run-engine): fairness scenarios on the real batched dequeue path
1stvamp 9dac71b
test(run-engine): multi-consumer correctness + op-count budget
1stvamp 24b6220
test(run-engine): default-off regression proof
1stvamp bddc82e
feat(run-engine,webapp): wire ckVirtualTimeScheduling env flag (code-…
1stvamp 86977ff
chore(run-engine): ship note + comment/format cleanup
1stvamp 2b1c389
fix(run-engine,webapp): address whole-branch adversarial review
1stvamp 5cb0aa5
docs(run-engine): record CK vtime known limitations for GA decision
1stvamp e9da64c
docs(run-engine): plan + references for virtual-time CK fair scheduling
1stvamp dce0b72
docs(run-engine): add fairness explainer diagrams
1stvamp 21211ae
fix(run-engine): move design docs out of the Mintlify docs/ tree; format
1stvamp e1444eb
docs(run-engine): tighten CK fairness server-changes note
1stvamp dea03c6
docs(run-engine): fix stale references in the CK fairness design docs
1stvamp c804278
docs(run-engine): add CK virtual-time A/B benchmark plan and harness
1stvamp da3c8fb
docs(run-engine): make the e2e bench harness work on self-hosted
1stvamp 4994295
docs(run-engine): add CK virtual-time benchmark results
1stvamp 3dd789f
chore(run-engine): gitignore benchmark output artifacts
1stvamp 31606d3
docs(run-engine): benchmark Redis CPU and memory vs concurrency-key c…
1stvamp 9015810
chore(run-engine): drop benchmark, e2e, and design docs from the branch
1stvamp bb07b0b
fix(run-engine): stop an unservable variant pinning the ck virtual-ti…
1stvamp 819e021
fix(run-engine): advance the ck virtual-time floor from servable vari…
1stvamp b932442
fix(run-engine): keep unservable ck variants in the fair order
1stvamp d8c5201
fix(run-engine): make ck vtime pass 2 discover unregistered variants
1stvamp f538f8d
test(run-engine): bound the ckManyKeys first-serve claim
1stvamp 5cd0c65
test(run-engine): pin the ck vtime window-freeze and stranded-entry b…
1stvamp c2fd2be
docs(run-engine): note the ck vtime retry-storm fallback in the relea…
1stvamp 3298ee2
fix(run-engine): stop unready ck variants blocking the fair pass, and…
1stvamp 6ce9d2a
fix(run-engine): guard the wildcard cleanup in the ck vtime scripts
1stvamp 701d164
fix(run-engine): register a gated ck variant so it can rejoin the fai…
1stvamp d21cbb9
fix(run-engine): drop drained ck variants from the fair order on ack,…
1stvamp af407f5
fix(run-engine): remember a drained ck variant's virtual time
1stvamp 50351eb
fix(run-engine): bound the ck idle set by rank, not just by floor
1stvamp dca1e9c
perf(run-engine): cut two Redis calls from the ck vtime registration …
1stvamp 2b6bfee
perf(run-engine): stop re-registering ck variants already in the fair…
1stvamp f5c1d74
fix(run-engine): keep the ck vtime floor alive alongside the tags it …
1stvamp d018b28
docs(run-engine): cut the ck fair scheduling release note down to the…
1stvamp c11b5c8
fix(run-engine): register a brand-new ck variant behind the pack, not…
1stvamp bb4d67d
test(run-engine): pin the fresh-key fix and isolate the ck idle rank cap
1stvamp f3ad8d6
fix(run-engine): stop the ck arrival cap from being reachable
1stvamp 64b2fc9
fix(run-engine): stop concurrency-gated ck variants eating the fair pass
1stvamp e8aa263
test(run-engine): cover the future-head window guard a mutation audit…
1stvamp 18e6dce
test(run-engine): pin the ck vtime credit round trip and TTL registra…
1stvamp 0a9998d
fix(run-engine): stop the gated ck batch rewinding a live virtual-tim…
1stvamp 2b7a533
test(run-engine): pin the parked tag a gated variant registers with
1stvamp 88127be
docs(run-engine): cut the ck vtime comments back to what the code can…
1stvamp 3c23a3a
fix(run-engine): re-register the TTL entry when the vtime dequeue exp…
1stvamp e69f000
fix(run-engine): keep parked ck credit when pass 2 repairs a variant
1stvamp a5fa891
refactor(run-engine): share one copy of the ck ack, nack and dead-let…
1stvamp 838896a
fix(run-engine): restore the snapshotRoute read in the vtime TTL sweep
1stvamp 0584290
refactor(run-engine): share one copy of the remaining ck lua
1stvamp c5c6bc9
fix(run-engine): make a mistyped ck lua slot a compile error
1stvamp c8872e2
test(run-engine): pin the snapshotRoute the vtime TTL sweep hands the…
1stvamp 43d0722
test(run-engine): load every run-queue script into redis
1stvamp 565ad2c
merge: integrate origin/main into ck virtual-time scheduling
1stvamp c4b657d
test(run-engine): thread groupConcurrencyKey into the two direct-Lua …
1stvamp c4af9aa
docs(run-engine): correct the ckExpireTtlLua sharing comment
1stvamp b02f4f0
feat(run-engine): warn when ck vtime and queue-gates are both enabled
1stvamp a9a14be
feat(run-engine): enforce queue-gates in the vtime dequeue and TTL sweep
1stvamp 9040133
feat(webapp): allow a fractional ck vtime quantum
1stvamp fab8175
fix(run-engine): reject a non-finite ck vtime quantum
1stvamp 027415b
Merge remote-tracking branch 'origin/main' into feat/ck-virtual-time-…
1stvamp File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
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
commit 64b2fc9557bcee329b59686764510f9f5f2df2e1
Some comments aren't visible on the classic Files Changed page.
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
107 changes: 107 additions & 0 deletions
107
internal-packages/run-engine/src/run-queue/tests/ckVtimeGatedOpBudget.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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 | ||
| ); | ||
| }); |
153 changes: 153 additions & 0 deletions
153
internal-packages/run-engine/src/run-queue/tests/ckVtimeGatedWindow.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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 | ||
| ); | ||
| }); |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.