Repository navigation
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 the gated ck batch rewinding a live virtual-tim…
…e tag The gated-candidate block registers its batch with one variadic ZADD NX and then, if that added anything, walks every member of gatedPending applying its parked idle tag. The ZADD reports how many members it added but not which ones, so the correction lands on candidates it did not register. That reaches an already-registered variant whenever the pass-1 scan is truncated, since knownRegistered comes from that scan and it reads only scanLimit entries. A queue with more variants than that pushes registered ones into pass 2 as if they were new. Their idle entry also survives re-registration (the enqueue path reads the parked tag but never deletes it, and it is reaped only once the floor climbs past), so with floor < parked < live the XX write overwrites the live tag with the older one. The variant's clock winds back and it is served ahead of variants that are genuinely due, which is the opposite of what the feature is for and exactly what NX exists to prevent everywhere else. The current score is the discriminator: at the floor means the variant either just registered here or has no credit to lose, and above the floor means it has spent a turn and keeps its tag. Reading it first also skips the idle lookup for the advanced ones, so the branch gets cheaper rather than dearer, and the steady state is untouched because none of this runs unless something registered. Reported by Devin on #4367.
- Loading branch information
commit 0a9998dc5d15da71adaf1f5a14e00630fbe38494
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
136 changes: 136 additions & 0 deletions
136
internal-packages/run-engine/src/run-queue/tests/ckVtimeGatedRewind.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,136 @@ | ||
| 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"; | ||
| import type { InputPayload } from "../types.js"; | ||
|
|
||
| // Devin on #4367: the gated-candidate block corrects every member of gatedPending with its | ||
| // parked idle tag once the batched ZADD NX has added at least one member, rather than only | ||
| // the members it actually added. A candidate that is already registered with an advanced | ||
| // tag can therefore have that tag overwritten by an older parked one, which rewinds its | ||
| // virtual clock and hands it a turn it has already taken. | ||
| // | ||
| // Reaching it needs the pass-1 scan to be incomplete, because that scan is the only thing | ||
| // that decides knownRegistered. It reads scanLimit entries, so a queue holding more | ||
| // variants than that pushes already-registered ones into pass 2 as if they were new. | ||
|
|
||
| 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: 1_000, | ||
| randomize: true, | ||
| }, | ||
| keys: new RunQueueFullKeyProducer(), | ||
| }; | ||
|
|
||
| const authenticatedEnvDev = { | ||
| id: "e1234", | ||
| type: "DEVELOPMENT" as const, | ||
| maximumConcurrencyLimit: 100, | ||
| concurrencyLimitBurstFactor: new Decimal(1), | ||
| project: { id: "p1234" }, | ||
| organization: { id: "o1234" }, | ||
| }; | ||
|
|
||
| const QUEUE = "task/my-task"; | ||
|
|
||
| function createQueue(redisContainer: any): any { | ||
| const redis = { | ||
| keyPrefix: "runqueue:test:rewind:", | ||
| host: redisContainer.getHost(), | ||
| port: redisContainer.getPort(), | ||
| }; | ||
| return new RunQueue({ | ||
| ...testOptions, | ||
| masterQueueConsumersDisabled: true, | ||
| workerOptions: { disabled: true }, | ||
| // window = 1 * 1, so scanLimit is 2 and four variants overflow it. | ||
| ckVirtualTimeScheduling: { enabled: true, scanWindowMultiplier: 1 }, | ||
| queueSelectionStrategy: new FairQueueSelectionStrategy({ redis, keys: testOptions.keys }), | ||
| redis, | ||
| } as any) as any; | ||
| } | ||
|
|
||
| function makeMessage(overrides: Partial<InputPayload> = {}): InputPayload { | ||
| return { | ||
| runId: "r1", | ||
| taskIdentifier: QUEUE, | ||
| orgId: "o1234", | ||
| projectId: "p1234", | ||
| environmentId: "e1234", | ||
| environmentType: "DEVELOPMENT", | ||
| queue: QUEUE, | ||
| timestamp: Date.now(), | ||
| attempt: 0, | ||
| ...overrides, | ||
| }; | ||
| } | ||
|
|
||
| const variantName = (ck: string) => testOptions.keys.queueKey(authenticatedEnvDev, QUEUE, ck); | ||
|
|
||
| describe("CK vtime: the gated batch must not rewind a live tag", () => { | ||
| redisTest("a stale parked tag cannot undercut an advanced one", async ({ redisContainer }) => { | ||
| const queue = createQueue(redisContainer); | ||
| try { | ||
| const t0 = Date.now() - 100_000; | ||
|
|
||
| // Head age decides pass 2's order, so victim is visited first. | ||
| const cks = ["victim", "wnew", "aa", "bb"]; | ||
| for (let i = 0; i < cks.length; i++) { | ||
| await queue.enqueueMessage({ | ||
| env: authenticatedEnvDev, | ||
| message: makeMessage({ runId: `r-${cks[i]}`, concurrencyKey: cks[i], timestamp: t0 + i }), | ||
| workerQueue: authenticatedEnvDev.id, | ||
| skipDequeueProcessing: true, | ||
| }); | ||
| } | ||
|
|
||
| const victim = variantName("victim"); | ||
| const wnew = variantName("wnew"); | ||
| const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(victim); | ||
| const ckVtimeIdleKey = testOptions.keys.ckVtimeIdleKeyFromQueue(victim); | ||
|
|
||
| // aa and bb hold the two scan slots, so victim falls outside the scan and reaches | ||
| // pass 2 with knownRegistered false even though it is registered. | ||
| await queue.redis.zadd(ckVtimeKey, 0, variantName("aa"), 1, variantName("bb"), 10, victim); | ||
| // wnew is genuinely unregistered, so the batched ZADD NX adds one member and the | ||
| // correction loop runs at all. | ||
| await queue.redis.zrem(ckVtimeKey, wnew); | ||
| // Left behind by an earlier drain: the enqueue path reads this to restore credit but | ||
| // never deletes it, and it is only reaped once the floor climbs past it. | ||
| await queue.redis.zadd(ckVtimeIdleKey, 5, victim); | ||
|
|
||
| // Every variant parked at its per-key ceiling, so pass 1 serves nothing and pass 2 | ||
| // reaches the gated block. | ||
| await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, QUEUE, 1); | ||
| for (const ck of cks) { | ||
| await queue.redis.sadd( | ||
| testOptions.keys.queueCurrentConcurrencyKeyFromQueue(variantName(ck)), | ||
| "occupant" | ||
| ); | ||
| } | ||
|
|
||
| const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); | ||
| const served = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 1); | ||
| expect(served.length).toBe(0); | ||
|
|
||
| // wnew joined at the floor, which is the whole point of the block. | ||
| expect(await queue.redis.zscore(ckVtimeKey, wnew)).toBe("0"); | ||
| // victim spent its credit already and must keep its advanced tag. Rewinding it to | ||
| // the parked 5 would put it ahead of variants that are genuinely due. | ||
| expect(await queue.redis.zscore(ckVtimeKey, victim)).toBe("10"); | ||
| } finally { | ||
| await queue.quit(); | ||
| } | ||
| }); | ||
| }); |
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.