diff --git a/.server-changes/2026-07-24-ck-fair-scheduling.md b/.server-changes/2026-07-24-ck-fair-scheduling.md new file mode 100644 index 00000000000..d9e04336819 --- /dev/null +++ b/.server-changes/2026-07-24-ck-fair-scheduling.md @@ -0,0 +1,6 @@ +--- +area: webapp +type: feature +--- + +The run queue can now serve concurrency keys fairly, so one key with a large backlog no longer starves runs waiting on other keys. It is opt-in and off by default. diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 24ed0833d8e..c118d484af9 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -978,6 +978,13 @@ const EnvironmentSchema = z RUN_ENGINE_TTL_CONSUMERS_DISABLED: BoolEnv.default(false), RUN_ENGINE_TTL_WORKER_BATCH_MAX_WAIT_MS: z.coerce.number().int().default(5_000), + // Fair (virtual-time) ordering across concurrency-key variants of a base queue. + // Off by default; when off the run queue behaves exactly as before. + RUN_ENGINE_CK_VTIME_SCHEDULING_ENABLED: BoolEnv.default(false), + RUN_ENGINE_CK_VTIME_QUANTUM: z.coerce.number().int().positive().default(1), + RUN_ENGINE_CK_VTIME_WINDOW_MULTIPLIER: z.coerce.number().int().positive().default(3), + RUN_ENGINE_CK_VTIME_STATE_TTL_SECONDS: z.coerce.number().int().positive().default(86400), + /** Optional maximum TTL for all runs (e.g. "14d"). If set, runs without an explicit TTL * will use this as their TTL, and runs with a TTL larger than this will be clamped. */ RUN_ENGINE_DEFAULT_MAX_TTL: z.string().optional(), diff --git a/apps/webapp/app/v3/runEngine.server.ts b/apps/webapp/app/v3/runEngine.server.ts index 85986933290..ba473be129c 100644 --- a/apps/webapp/app/v3/runEngine.server.ts +++ b/apps/webapp/app/v3/runEngine.server.ts @@ -108,6 +108,14 @@ function createRunEngine() { batchMaxSize: env.RUN_ENGINE_TTL_WORKER_BATCH_MAX_SIZE, batchMaxWaitMs: env.RUN_ENGINE_TTL_WORKER_BATCH_MAX_WAIT_MS, }, + ckVirtualTimeScheduling: env.RUN_ENGINE_CK_VTIME_SCHEDULING_ENABLED + ? { + enabled: true, + quantum: env.RUN_ENGINE_CK_VTIME_QUANTUM, + scanWindowMultiplier: env.RUN_ENGINE_CK_VTIME_WINDOW_MULTIPLIER, + stateTtlSeconds: env.RUN_ENGINE_CK_VTIME_STATE_TTL_SECONDS, + } + : undefined, }, runLock: { redis: { diff --git a/internal-packages/run-engine/src/engine/index.ts b/internal-packages/run-engine/src/engine/index.ts index b5131766559..2243a420df8 100644 --- a/internal-packages/run-engine/src/engine/index.ts +++ b/internal-packages/run-engine/src/engine/index.ts @@ -235,6 +235,7 @@ export class RunEngine { workerItemsSuffix: "ttl-worker:{queue:ttl-expiration:}items", visibilityTimeoutMs: options.queue?.ttlSystem?.visibilityTimeoutMs ?? 30_000, }, + ckVirtualTimeScheduling: options.queue?.ckVirtualTimeScheduling, }); this.worker = new Worker({ diff --git a/internal-packages/run-engine/src/engine/types.ts b/internal-packages/run-engine/src/engine/types.ts index f37ec7df50a..4360f889039 100644 --- a/internal-packages/run-engine/src/engine/types.ts +++ b/internal-packages/run-engine/src/engine/types.ts @@ -16,7 +16,7 @@ import { } from "@trigger.dev/redis-worker"; import type { ControlPlaneResolver } from "./controlPlaneResolver.js"; import type { FairQueueSelectionStrategyOptions } from "../run-queue/fairQueueSelectionStrategy.js"; -import type { RunQueueMetricsEmitter } from "../run-queue/index.js"; +import type { RunQueueMetricsEmitter, RunQueueOptions } from "../run-queue/index.js"; import type { MinimalAuthenticatedEnvironment } from "../shared/index.js"; import type { LockRetryConfig } from "./locking.js"; import type { workerCatalog } from "./workerCatalog.js"; @@ -126,6 +126,9 @@ export type RunEngineOptions = { /** Max time (ms) to wait for more items before flushing a batch (default: 5000) */ batchMaxWaitMs?: number; }; + /** Fair (virtual-time) ordering across concurrency-key variants of a base queue. + * Passed through to RunQueue; off by default (undefined = today's behaviour). */ + ckVirtualTimeScheduling?: RunQueueOptions["ckVirtualTimeScheduling"]; }; runLock: { redis: RedisOptions; diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index 58225cc5051..498d494bf50 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -214,6 +214,21 @@ export type RunQueueOptions = { /** Visibility timeout for TTL worker jobs (ms, default: 30000) */ visibilityTimeoutMs?: number; }; + /** + * Fair (virtual-time / SFQ) ordering across concurrency-key variants of a + * base queue. Off by default; when off, the exact pre-existing Lua commands + * run and no vtime keys are created. See the CK virtual-time scheduling + * design for the ordering model. + */ + ckVirtualTimeScheduling?: { + enabled: boolean; + /** Virtual-time advance per serve (dimensionless). Default 1. */ + quantum?: number; + /** Pass-1 candidate window = actualMaxCount * this. Default 3. */ + scanWindowMultiplier?: number; + /** EXPIRE applied to ckVtime/ckVtimeFloor on every write. Default 86400. */ + stateTtlSeconds?: number; + }; }; export interface ConcurrencySweeperCallback { @@ -298,10 +313,28 @@ export class RunQueue { private _observableWorkerQueues: Set = new Set(); private _meter: Meter; private _queueCooloffStates: Map = new Map(); + readonly #ckVtimeEnabled: boolean; + readonly #ckVtimeQuantum: number; + readonly #ckVtimeWindowMultiplier: number; + readonly #ckVtimeStateTtl: number; constructor(public readonly options: RunQueueOptions) { this.shardCount = options.shardCount ?? 2; this.counterTtlSeconds = options.counterTtlSeconds ?? 86400; + this.#ckVtimeEnabled = options.ckVirtualTimeScheduling?.enabled ?? false; + // Defense-in-depth: clamp so a directly-constructed RunQueue can't get bad + // values that would freeze tags (quantum <= 0) or force an O(N) scan / + // EX 0 error (multiplier / ttl <= 0). + const resolvedQuantum = options.ckVirtualTimeScheduling?.quantum ?? 1; + this.#ckVtimeQuantum = resolvedQuantum > 0 ? resolvedQuantum : 1; + this.#ckVtimeWindowMultiplier = Math.max( + 1, + Math.floor(options.ckVirtualTimeScheduling?.scanWindowMultiplier ?? 3) + ); + this.#ckVtimeStateTtl = Math.max( + 1, + Math.floor(options.ckVirtualTimeScheduling?.stateTtlSeconds ?? 86400) + ); this.retryOptions = options.retryOptions ?? defaultRetrySettings; this.redis = createRedisClient(options.redis, { onError: (error) => { @@ -2171,74 +2204,151 @@ export class RunQueue { const ckKeyPrefix = this.options.redis.keyPrefix ?? ""; if (ttlInfo) { - result = await this.redis.enqueueMessageWithTtlCkTracked( - // keys - masterQueueKey, - queueKey, - messageKey, - queueCurrentConcurrencyKey, - envCurrentConcurrencyKey, - queueCurrentDequeuedKey, - envCurrentDequeuedKey, - envQueueKey, - ttlInfo.ttlQueueKey, - ckIndexKey, - workerQueueKey, - queueConcurrencyLimitKey, - envConcurrencyLimitKey, - envConcurrencyLimitBurstFactorKey, - lengthCounterKey, - baseQueueKey, - // args - queueName, - messageId, - messageData, - messageScore, - ttlInfo.ttlMember, - String(ttlInfo.ttlExpiresAt), - ckWildcardName, - messageKeyValue, - defaultEnvConcurrencyLimit, - defaultEnvConcurrencyBurstFactor, - currentTime, - enableFastPathArg, - ckKeyPrefix, - String(this.counterTtlSeconds), - metricsGaugeArg - ); + result = this.#ckVtimeEnabled + ? await this.redis.enqueueMessageWithTtlCkVtimeTracked( + // keys + masterQueueKey, + queueKey, + messageKey, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + queueCurrentDequeuedKey, + envCurrentDequeuedKey, + envQueueKey, + ttlInfo.ttlQueueKey, + ckIndexKey, + workerQueueKey, + queueConcurrencyLimitKey, + envConcurrencyLimitKey, + envConcurrencyLimitBurstFactorKey, + lengthCounterKey, + baseQueueKey, + this.keys.ckVtimeKeyFromQueue(message.queue), + this.keys.ckVtimeFloorKeyFromQueue(message.queue), + // args + queueName, + messageId, + messageData, + messageScore, + ttlInfo.ttlMember, + String(ttlInfo.ttlExpiresAt), + ckWildcardName, + messageKeyValue, + defaultEnvConcurrencyLimit, + defaultEnvConcurrencyBurstFactor, + currentTime, + enableFastPathArg, + ckKeyPrefix, + String(this.counterTtlSeconds), + String(this.#ckVtimeStateTtl), + // Must stay last: the gauge fragment reads ARGV[#ARGV]. + metricsGaugeArg + ) + : await this.redis.enqueueMessageWithTtlCkTracked( + // keys + masterQueueKey, + queueKey, + messageKey, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + queueCurrentDequeuedKey, + envCurrentDequeuedKey, + envQueueKey, + ttlInfo.ttlQueueKey, + ckIndexKey, + workerQueueKey, + queueConcurrencyLimitKey, + envConcurrencyLimitKey, + envConcurrencyLimitBurstFactorKey, + lengthCounterKey, + baseQueueKey, + // args + queueName, + messageId, + messageData, + messageScore, + ttlInfo.ttlMember, + String(ttlInfo.ttlExpiresAt), + ckWildcardName, + messageKeyValue, + defaultEnvConcurrencyLimit, + defaultEnvConcurrencyBurstFactor, + currentTime, + enableFastPathArg, + ckKeyPrefix, + String(this.counterTtlSeconds), + metricsGaugeArg + ); } else { - result = await this.redis.enqueueMessageCkTracked( - // keys - masterQueueKey, - queueKey, - messageKey, - queueCurrentConcurrencyKey, - envCurrentConcurrencyKey, - queueCurrentDequeuedKey, - envCurrentDequeuedKey, - envQueueKey, - ckIndexKey, - workerQueueKey, - queueConcurrencyLimitKey, - envConcurrencyLimitKey, - envConcurrencyLimitBurstFactorKey, - lengthCounterKey, - baseQueueKey, - // args - queueName, - messageId, - messageData, - messageScore, - ckWildcardName, - messageKeyValue, - defaultEnvConcurrencyLimit, - defaultEnvConcurrencyBurstFactor, - currentTime, - enableFastPathArg, - ckKeyPrefix, - String(this.counterTtlSeconds), - metricsGaugeArg - ); + result = this.#ckVtimeEnabled + ? await this.redis.enqueueMessageCkVtimeTracked( + // keys + masterQueueKey, + queueKey, + messageKey, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + queueCurrentDequeuedKey, + envCurrentDequeuedKey, + envQueueKey, + ckIndexKey, + workerQueueKey, + queueConcurrencyLimitKey, + envConcurrencyLimitKey, + envConcurrencyLimitBurstFactorKey, + lengthCounterKey, + baseQueueKey, + this.keys.ckVtimeKeyFromQueue(message.queue), + this.keys.ckVtimeFloorKeyFromQueue(message.queue), + // args + queueName, + messageId, + messageData, + messageScore, + ckWildcardName, + messageKeyValue, + defaultEnvConcurrencyLimit, + defaultEnvConcurrencyBurstFactor, + currentTime, + enableFastPathArg, + ckKeyPrefix, + String(this.counterTtlSeconds), + String(this.#ckVtimeStateTtl), + // Must stay last: the gauge fragment reads ARGV[#ARGV]. + metricsGaugeArg + ) + : await this.redis.enqueueMessageCkTracked( + // keys + masterQueueKey, + queueKey, + messageKey, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + queueCurrentDequeuedKey, + envCurrentDequeuedKey, + envQueueKey, + ckIndexKey, + workerQueueKey, + queueConcurrencyLimitKey, + envConcurrencyLimitKey, + envConcurrencyLimitBurstFactorKey, + lengthCounterKey, + baseQueueKey, + // args + queueName, + messageId, + messageData, + messageScore, + ckWildcardName, + messageKeyValue, + defaultEnvConcurrencyLimit, + defaultEnvConcurrencyBurstFactor, + currentTime, + enableFastPathArg, + ckKeyPrefix, + String(this.counterTtlSeconds), + metricsGaugeArg + ); } } else if (ttlInfo) { // Use the TTL-aware enqueue that atomically adds to both queues @@ -2489,28 +2599,61 @@ export class RunQueue { const metricsGaugeArg = this.#queueMetricsGaugeArg(); - const reply = await this.redis.dequeueMessagesFromCkQueueTracked( - //keys - ckIndexKey, - queueConcurrencyLimitKey, - envConcurrencyLimitKey, - envConcurrencyLimitBurstFactorKey, - envCurrentConcurrencyKey, - messageKeyPrefix, - envQueueKey, - masterQueueKey, - ttlQueueKey, - lengthCounterKey, - runningCounterKey, - //args - ckWildcardQueue, - String(Date.now()), - String(this.options.defaultEnvConcurrency), - String(this.options.defaultEnvConcurrencyBurstFactor ?? 1), - this.options.redis.keyPrefix ?? "", - String(maxCount), - metricsGaugeArg - ); + if (this.#ckVtimeEnabled) { + span.setAttribute("ck_vtime_enabled", true); + } + + const reply = this.#ckVtimeEnabled + ? await this.redis.dequeueMessagesFromCkQueueVtimeTracked( + //keys + ckIndexKey, + queueConcurrencyLimitKey, + envConcurrencyLimitKey, + envConcurrencyLimitBurstFactorKey, + envCurrentConcurrencyKey, + messageKeyPrefix, + envQueueKey, + masterQueueKey, + ttlQueueKey, + lengthCounterKey, + runningCounterKey, + this.keys.ckVtimeKeyFromQueue(ckWildcardQueue), + this.keys.ckVtimeFloorKeyFromQueue(ckWildcardQueue), + //args + ckWildcardQueue, + String(Date.now()), + String(this.options.defaultEnvConcurrency), + String(this.options.defaultEnvConcurrencyBurstFactor ?? 1), + this.options.redis.keyPrefix ?? "", + String(maxCount), + String(this.#ckVtimeQuantum), + String(this.#ckVtimeWindowMultiplier), + String(this.#ckVtimeStateTtl), + // Must stay last: the gauge fragment reads ARGV[#ARGV]. + metricsGaugeArg + ) + : await this.redis.dequeueMessagesFromCkQueueTracked( + //keys + ckIndexKey, + queueConcurrencyLimitKey, + envConcurrencyLimitKey, + envConcurrencyLimitBurstFactorKey, + envCurrentConcurrencyKey, + messageKeyPrefix, + envQueueKey, + masterQueueKey, + ttlQueueKey, + lengthCounterKey, + runningCounterKey, + //args + ckWildcardQueue, + String(Date.now()), + String(this.options.defaultEnvConcurrency), + String(this.options.defaultEnvConcurrencyBurstFactor ?? 1), + this.options.redis.keyPrefix ?? "", + String(maxCount), + metricsGaugeArg + ); // Reply is [flatMessages|null, gauge|null]; the CK aggregate gauge rides here. const gauge = reply?.[1] ?? null; @@ -2874,28 +3017,56 @@ export class RunQueue { const lengthCounterKey = this.keys.queueLengthCounterKeyFromQueue(message.queue); const runningCounterKey = this.keys.queueRunningCounterKeyFromQueue(message.queue); - await this.redis.nackMessageCkTracked( - //keys - masterQueueKey, - messageKey, - messageQueue, - queueCurrentConcurrencyKey, - envCurrentConcurrencyKey, - queueCurrentDequeuedKey, - envCurrentDequeuedKey, - envQueueKey, - ckIndexKey, - lengthCounterKey, - runningCounterKey, - //args - messageId, - messageQueue, - JSON.stringify(message), - String(messageScore), - ckWildcardName, - this.options.redis.keyPrefix ?? "", - String(this.counterTtlSeconds) - ); + if (this.#ckVtimeEnabled) { + await this.redis.nackMessageCkVtimeTracked( + //keys + masterQueueKey, + messageKey, + messageQueue, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + queueCurrentDequeuedKey, + envCurrentDequeuedKey, + envQueueKey, + ckIndexKey, + lengthCounterKey, + runningCounterKey, + this.keys.ckVtimeKeyFromQueue(message.queue), + this.keys.ckVtimeFloorKeyFromQueue(message.queue), + //args + messageId, + messageQueue, + JSON.stringify(message), + String(messageScore), + ckWildcardName, + this.options.redis.keyPrefix ?? "", + String(this.counterTtlSeconds), + String(this.#ckVtimeStateTtl) + ); + } else { + await this.redis.nackMessageCkTracked( + //keys + masterQueueKey, + messageKey, + messageQueue, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + queueCurrentDequeuedKey, + envCurrentDequeuedKey, + envQueueKey, + ckIndexKey, + lengthCounterKey, + runningCounterKey, + //args + messageId, + messageQueue, + JSON.stringify(message), + String(messageScore), + ckWildcardName, + this.options.redis.keyPrefix ?? "", + String(this.counterTtlSeconds) + ); + } } else { await this.redis.nackMessage( //keys @@ -3967,72 +4138,345 @@ return __qmret(0) `, }); - // Expire TTL runs - atomically removes from TTL set, acknowledges from normal queue, and enqueues to TTL worker - this.redis.defineCommand("expireTtlRuns", { - numberOfKeys: 1, + // Vtime variant of enqueueMessageCkTracked (feature-flagged via + // ckVirtualTimeScheduling.enabled). Identical script body, plus slow-path + // registration of the variant into the :ckVtime ZSET at the floor (NX), so a + // brand-new key is present in the fair order from its first enqueue. The + // fast path (direct-to-worker-queue) neither registers nor advances. + this.redis.defineCommand("enqueueMessageCkVtimeTracked", { + numberOfKeys: 17, lua: ` -local ttlQueueKey = KEYS[1] -local keyPrefix = ARGV[1] -local currentTime = tonumber(ARGV[2]) -local batchSize = tonumber(ARGV[3]) -local shardCount = tonumber(ARGV[4]) -local workerQueueKey = ARGV[5] -local workerItemsKey = ARGV[6] -local visibilityTimeoutMs = tonumber(ARGV[7]) +local masterQueueKey = KEYS[1] +local queueKey = KEYS[2] +local messageKey = KEYS[3] +local queueCurrentConcurrencyKey = KEYS[4] +local envCurrentConcurrencyKey = KEYS[5] +local queueCurrentDequeuedKey = KEYS[6] +local envCurrentDequeuedKey = KEYS[7] +local envQueueKey = KEYS[8] +local ckIndexKey = KEYS[9] +-- Fast-path keys (KEYS 10-13) +local workerQueueKey = KEYS[10] +local queueConcurrencyLimitKey = KEYS[11] +local envConcurrencyLimitKey = KEYS[12] +local envConcurrencyLimitBurstFactorKey = KEYS[13] +-- Counter keys (KEYS 14-15) +local lengthCounterKey = KEYS[14] +local baseQueueKey = KEYS[15] +-- Virtual-time keys (KEYS 16-17) +local ckVtimeKey = KEYS[16] +local ckVtimeFloorKey = KEYS[17] --- Get expired runs from TTL sorted set (score <= currentTime) -local expiredMembers = redis.call('ZRANGEBYSCORE', ttlQueueKey, '-inf', currentTime, 'LIMIT', 0, batchSize) +local queueName = ARGV[1] +local messageId = ARGV[2] +local messageData = ARGV[3] +local messageScore = ARGV[4] +local ckWildcardName = ARGV[5] +-- Fast-path args (ARGV 6-10) +local messageKeyValue = ARGV[6] +local defaultEnvConcurrencyLimit = ARGV[7] +local defaultEnvConcurrencyBurstFactor = ARGV[8] +local currentTime = ARGV[9] +local enableFastPath = ARGV[10] +-- keyPrefix for prepending to variant names stored as values in ckIndex +local keyPrefix = ARGV[11] +-- TTL (seconds) applied to counter lazy-init SETs +local counterTtl = ARGV[12] +-- TTL (seconds) applied to ckVtime on registration +local stateTtl = ARGV[13] -if #expiredMembers == 0 then - return {} -end +${QUEUE_METRICS_GAUGE_PRELUDE} -local time = redis.call('TIME') -local nowMs = tonumber(time[1]) * 1000 + math.floor(tonumber(time[2]) / 1000) +-- Fast path: check if we can skip the queue and go directly to worker queue +if enableFastPath == '1' then + local available = redis.call('ZRANGEBYSCORE', queueKey, '-inf', currentTime, 'LIMIT', 0, 1) + if #available == 0 then + local envCurrent = tonumber(redis.call('SCARD', envCurrentConcurrencyKey) or '0') + local envLimit = tonumber(redis.call('GET', envConcurrencyLimitKey) or defaultEnvConcurrencyLimit) + local envBurstFactor = tonumber(redis.call('GET', envConcurrencyLimitBurstFactorKey) or defaultEnvConcurrencyBurstFactor) + local envLimitWithBurst = math.floor(envLimit * envBurstFactor) -local results = {} + if envCurrent < envLimitWithBurst then + local queueCurrent = tonumber(redis.call('SCARD', queueCurrentConcurrencyKey) or '0') + local queueLimit = math.min( + tonumber(redis.call('GET', queueConcurrencyLimitKey) or '1000000'), + envLimit + ) -for i, member in ipairs(expiredMembers) do - -- Parse member format: "queueKey|runId|orgId" - local pipePos1 = string.find(member, "|", 1, true) - if pipePos1 then - local pipePos2 = string.find(member, "|", pipePos1 + 1, true) - if pipePos2 then - local rawQueueKey = string.sub(member, 1, pipePos1 - 1) - local runId = string.sub(member, pipePos1 + 1, pipePos2 - 1) - local orgId = string.sub(member, pipePos2 + 1) + if queueCurrent < queueLimit then + redis.call('SET', messageKey, messageData) + redis.call('SADD', queueCurrentConcurrencyKey, messageId) + redis.call('SADD', envCurrentConcurrencyKey, messageId) + redis.call('RPUSH', workerQueueKey, messageKeyValue) + -- Fast-path skips the CK variant zset entirely; lengthCounter is unchanged. + -- runningCounter is bumped later by dequeueMessageFromKeyTracked when the + -- worker pulls the message from the worker queue. +${QUEUE_METRICS_CK_ENQUEUE_FASTPATH_GAUGE_LUA} + return __qmret(1) + end + end + end +end - -- Prefix the queue key so it matches the actual Redis keys - local queueKey = keyPrefix .. rawQueueKey +-- Slow path: normal enqueue +redis.call('SET', messageKey, messageData) - -- Remove from TTL set - redis.call('ZREM', ttlQueueKey, member) +-- Lazy-init lengthCounter from existing ckIndex variants (once per base queue per 24h). +-- The 24h TTL means the counter periodically re-anchors to truth, bounding any drift +-- that accumulated during rolling-deploy overlap windows. +-- Run BEFORE the ZADD so we capture pre-state; the subsequent INCR accounts for the new message. +-- The counter tracks ONLY CK-variant messages — the read path adds ZCARD(base) separately, +-- so the base zset is intentionally excluded here. +if redis.call('EXISTS', lengthCounterKey) == 0 then + local total = 0 + local variants = redis.call('ZRANGE', ckIndexKey, 0, -1) + for _, v in ipairs(variants) do + total = total + tonumber(redis.call('ZCARD', keyPrefix .. v) or '0') + end + redis.call('SET', lengthCounterKey, total, 'EX', counterTtl) +end - -- Construct keys for acknowledging the run from normal queue - -- Extract org from rawQueueKey: {org:orgId}:proj:... - local orgKeyStart = string.find(rawQueueKey, "{org:", 1, true) - local orgKeyEnd = string.find(rawQueueKey, "}", orgKeyStart, true) - local orgFromQueue = string.sub(rawQueueKey, orgKeyStart + 5, orgKeyEnd - 1) +-- INCR is gated on ZADD returning 1 (new entry). A duplicate enqueue (same messageId +-- already in the variant zset) returns 0 and must not bump the counter. +local added = redis.call('ZADD', queueKey, messageScore, messageId) +redis.call('ZADD', envQueueKey, messageScore, messageId) +if added == 1 then + redis.call('INCR', lengthCounterKey) +end - local messageKey = keyPrefix .. "{org:" .. orgFromQueue .. "}:message:" .. runId +-- Rebalance CK index +local earliest = redis.call('ZRANGE', queueKey, 0, 0, 'WITHSCORES') +if #earliest > 0 then + redis.call('ZADD', ckIndexKey, earliest[2], queueName) +end - -- Delete message key - redis.call('DEL', messageKey) +-- Register this variant in the virtual-time index at the floor. NX means an +-- already-advanced tag is never rewound. +local vfloor = redis.call('GET', ckVtimeFloorKey) or '0' +redis.call('ZADD', ckVtimeKey, 'NX', vfloor, queueName) +redis.call('EXPIRE', ckVtimeKey, stateTtl) +redis.call('EXPIRE', ckVtimeFloorKey, stateTtl) - -- Remove from queue sorted set - redis.call('ZREM', queueKey, runId) +-- Rebalance master queue with ck:* member +local earliestIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') +if #earliestIdx > 0 then + redis.call('ZADD', masterQueueKey, earliestIdx[2], ckWildcardName) +end - -- Remove from env queue (derive from rawQueueKey) - -- rawQueueKey format: {org:X}:proj:Y:env:Z:queue:Q[:ck:C] - local envMatch = string.match(rawQueueKey, ":env:([^:]+)") - if envMatch then - local envQueueKey = keyPrefix .. "{org:" .. orgFromQueue .. "}:env:" .. envMatch - redis.call('ZREM', envQueueKey, runId) - end +-- Remove old-format entry from master queue (transition cleanup) +redis.call('ZREM', masterQueueKey, queueName) - -- Remove from concurrency sets - local concurrencyKey = queueKey .. ":currentConcurrency" - local dequeuedKey = queueKey .. ":currentDequeued" +-- Update the concurrency keys +redis.call('SREM', queueCurrentConcurrencyKey, messageId) +redis.call('SREM', envCurrentConcurrencyKey, messageId) +redis.call('SREM', queueCurrentDequeuedKey, messageId) +redis.call('SREM', envCurrentDequeuedKey, messageId) + +${QUEUE_METRICS_CK_ENQUEUE_GAUGE_LUA} + +return __qmret(0) + `, + }); + + // Vtime variant of enqueueMessageWithTtlCkTracked. Same slow-path-only + // registration as enqueueMessageCkVtimeTracked above. + this.redis.defineCommand("enqueueMessageWithTtlCkVtimeTracked", { + numberOfKeys: 18, + lua: ` +local masterQueueKey = KEYS[1] +local queueKey = KEYS[2] +local messageKey = KEYS[3] +local queueCurrentConcurrencyKey = KEYS[4] +local envCurrentConcurrencyKey = KEYS[5] +local queueCurrentDequeuedKey = KEYS[6] +local envCurrentDequeuedKey = KEYS[7] +local envQueueKey = KEYS[8] +local ttlQueueKey = KEYS[9] +local ckIndexKey = KEYS[10] +-- Fast-path keys (KEYS 11-14) +local workerQueueKey = KEYS[11] +local queueConcurrencyLimitKey = KEYS[12] +local envConcurrencyLimitKey = KEYS[13] +local envConcurrencyLimitBurstFactorKey = KEYS[14] +-- Counter keys (KEYS 15-16) +local lengthCounterKey = KEYS[15] +local baseQueueKey = KEYS[16] +-- Virtual-time keys (KEYS 17-18) +local ckVtimeKey = KEYS[17] +local ckVtimeFloorKey = KEYS[18] + +local queueName = ARGV[1] +local messageId = ARGV[2] +local messageData = ARGV[3] +local messageScore = ARGV[4] +local ttlMember = ARGV[5] +local ttlScore = ARGV[6] +local ckWildcardName = ARGV[7] +-- Fast-path args (ARGV 8-12) +local messageKeyValue = ARGV[8] +local defaultEnvConcurrencyLimit = ARGV[9] +local defaultEnvConcurrencyBurstFactor = ARGV[10] +local currentTime = ARGV[11] +local enableFastPath = ARGV[12] +-- keyPrefix for prepending to variant names stored as values in ckIndex +local keyPrefix = ARGV[13] +-- TTL (seconds) applied to counter lazy-init SETs +local counterTtl = ARGV[14] +-- TTL (seconds) applied to ckVtime on registration +local stateTtl = ARGV[15] + +${QUEUE_METRICS_GAUGE_PRELUDE} + +-- Fast path: check if we can skip the queue and go directly to worker queue +if enableFastPath == '1' then + local available = redis.call('ZRANGEBYSCORE', queueKey, '-inf', currentTime, 'LIMIT', 0, 1) + if #available == 0 then + local envCurrent = tonumber(redis.call('SCARD', envCurrentConcurrencyKey) or '0') + local envLimit = tonumber(redis.call('GET', envConcurrencyLimitKey) or defaultEnvConcurrencyLimit) + local envBurstFactor = tonumber(redis.call('GET', envConcurrencyLimitBurstFactorKey) or defaultEnvConcurrencyBurstFactor) + local envLimitWithBurst = math.floor(envLimit * envBurstFactor) + + if envCurrent < envLimitWithBurst then + local queueCurrent = tonumber(redis.call('SCARD', queueCurrentConcurrencyKey) or '0') + local queueLimit = math.min( + tonumber(redis.call('GET', queueConcurrencyLimitKey) or '1000000'), + envLimit + ) + + if queueCurrent < queueLimit then + redis.call('SET', messageKey, messageData) + redis.call('SADD', queueCurrentConcurrencyKey, messageId) + redis.call('SADD', envCurrentConcurrencyKey, messageId) + redis.call('RPUSH', workerQueueKey, messageKeyValue) +${QUEUE_METRICS_CK_ENQUEUE_FASTPATH_GAUGE_LUA} + return __qmret(1) + end + end + end +end + +-- Slow path: normal enqueue +redis.call('SET', messageKey, messageData) + +-- Lazy-init lengthCounter from existing ckIndex variants (once per base queue per 24h). +-- See enqueueMessageCkTracked for the TTL rationale. +if redis.call('EXISTS', lengthCounterKey) == 0 then + local total = 0 + local variants = redis.call('ZRANGE', ckIndexKey, 0, -1) + for _, v in ipairs(variants) do + total = total + tonumber(redis.call('ZCARD', keyPrefix .. v) or '0') + end + redis.call('SET', lengthCounterKey, total, 'EX', counterTtl) +end + +-- INCR is gated on ZADD returning 1 (new entry). +local added = redis.call('ZADD', queueKey, messageScore, messageId) +redis.call('ZADD', envQueueKey, messageScore, messageId) +redis.call('ZADD', ttlQueueKey, ttlScore, ttlMember) +if added == 1 then + redis.call('INCR', lengthCounterKey) +end + +-- Rebalance CK index +local earliest = redis.call('ZRANGE', queueKey, 0, 0, 'WITHSCORES') +if #earliest > 0 then + redis.call('ZADD', ckIndexKey, earliest[2], queueName) +end + +-- Register this variant in the virtual-time index at the floor. NX means an +-- already-advanced tag is never rewound. +local vfloor = redis.call('GET', ckVtimeFloorKey) or '0' +redis.call('ZADD', ckVtimeKey, 'NX', vfloor, queueName) +redis.call('EXPIRE', ckVtimeKey, stateTtl) +redis.call('EXPIRE', ckVtimeFloorKey, stateTtl) + +-- Rebalance master queue with ck:* member +local earliestIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') +if #earliestIdx > 0 then + redis.call('ZADD', masterQueueKey, earliestIdx[2], ckWildcardName) +end + +-- Remove old-format entry from master queue (transition cleanup) +redis.call('ZREM', masterQueueKey, queueName) + +-- Update the concurrency keys +redis.call('SREM', queueCurrentConcurrencyKey, messageId) +redis.call('SREM', envCurrentConcurrencyKey, messageId) +redis.call('SREM', queueCurrentDequeuedKey, messageId) +redis.call('SREM', envCurrentDequeuedKey, messageId) + +${QUEUE_METRICS_CK_ENQUEUE_GAUGE_LUA} + +return __qmret(0) + `, + }); + + // Expire TTL runs - atomically removes from TTL set, acknowledges from normal queue, and enqueues to TTL worker + this.redis.defineCommand("expireTtlRuns", { + numberOfKeys: 1, + lua: ` +local ttlQueueKey = KEYS[1] +local keyPrefix = ARGV[1] +local currentTime = tonumber(ARGV[2]) +local batchSize = tonumber(ARGV[3]) +local shardCount = tonumber(ARGV[4]) +local workerQueueKey = ARGV[5] +local workerItemsKey = ARGV[6] +local visibilityTimeoutMs = tonumber(ARGV[7]) + +-- Get expired runs from TTL sorted set (score <= currentTime) +local expiredMembers = redis.call('ZRANGEBYSCORE', ttlQueueKey, '-inf', currentTime, 'LIMIT', 0, batchSize) + +if #expiredMembers == 0 then + return {} +end + +local time = redis.call('TIME') +local nowMs = tonumber(time[1]) * 1000 + math.floor(tonumber(time[2]) / 1000) + +local results = {} + +for i, member in ipairs(expiredMembers) do + -- Parse member format: "queueKey|runId|orgId" + local pipePos1 = string.find(member, "|", 1, true) + if pipePos1 then + local pipePos2 = string.find(member, "|", pipePos1 + 1, true) + if pipePos2 then + local rawQueueKey = string.sub(member, 1, pipePos1 - 1) + local runId = string.sub(member, pipePos1 + 1, pipePos2 - 1) + local orgId = string.sub(member, pipePos2 + 1) + + -- Prefix the queue key so it matches the actual Redis keys + local queueKey = keyPrefix .. rawQueueKey + + -- Remove from TTL set + redis.call('ZREM', ttlQueueKey, member) + + -- Construct keys for acknowledging the run from normal queue + -- Extract org from rawQueueKey: {org:orgId}:proj:... + local orgKeyStart = string.find(rawQueueKey, "{org:", 1, true) + local orgKeyEnd = string.find(rawQueueKey, "}", orgKeyStart, true) + local orgFromQueue = string.sub(rawQueueKey, orgKeyStart + 5, orgKeyEnd - 1) + + local messageKey = keyPrefix .. "{org:" .. orgFromQueue .. "}:message:" .. runId + + -- Delete message key + redis.call('DEL', messageKey) + + -- Remove from queue sorted set + redis.call('ZREM', queueKey, runId) + + -- Remove from env queue (derive from rawQueueKey) + -- rawQueueKey format: {org:X}:proj:Y:env:Z:queue:Q[:ck:C] + local envMatch = string.match(rawQueueKey, ":env:([^:]+)") + if envMatch then + local envQueueKey = keyPrefix .. "{org:" .. orgFromQueue .. "}:env:" .. envMatch + redis.call('ZREM', envQueueKey, runId) + end + + -- Remove from concurrency sets + local concurrencyKey = queueKey .. ":currentConcurrency" + local dequeuedKey = queueKey .. ":currentDequeued" redis.call('SREM', concurrencyKey, runId) redis.call('SREM', dequeuedKey, runId) @@ -4342,32 +4786,186 @@ local envConcurrencyLimitBurstFactor = tonumber(redis.call('GET', envConcurrency local envConcurrencyLimitWithBurstFactor = math.floor(envConcurrencyLimit * envConcurrencyLimitBurstFactor) if envCurrentConcurrency >= envConcurrencyLimitWithBurstFactor then - return nil + return nil +end + +-- Get base queue concurrency limit +local queueConcurrencyLimit = math.min(tonumber(redis.call('GET', queueConcurrencyLimitKey) or '1000000'), envConcurrencyLimit) + +-- Calculate env available capacity +local envAvailableCapacity = envConcurrencyLimitWithBurstFactor - envCurrentConcurrency +local actualMaxCount = math.min(maxCount, envAvailableCapacity) + +if actualMaxCount <= 0 then + return nil +end + +-- Get CK sub-queues from CK index with scores <= currentTime +local ckQueues = redis.call('ZRANGEBYSCORE', ckIndexKey, '-inf', tostring(currentTime), 'LIMIT', 0, actualMaxCount * 3) + +if #ckQueues == 0 then + -- Rebalance master queue in case CK index has future-scored entries + local anyIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') + if #anyIdx == 0 then + redis.call('ZREM', masterQueueKey, ckWildcardName) + else + redis.call('ZADD', masterQueueKey, anyIdx[2], ckWildcardName) + end + return nil +end + +local results = {} +local dequeuedCount = 0 + +for _, ckQueueName in ipairs(ckQueues) do + if dequeuedCount >= actualMaxCount then + break + end + + local fullQueueKey = keyPrefix .. ckQueueName + + -- Check CK-specific concurrency + local ckConcurrencyKey = fullQueueKey .. ':currentConcurrency' + local ckCurrentConcurrency = tonumber(redis.call('SCARD', ckConcurrencyKey) or '0') + + if ckCurrentConcurrency < queueConcurrencyLimit then + -- Try to dequeue one message from this CK sub-queue + local messages = redis.call('ZRANGEBYSCORE', fullQueueKey, '-inf', tostring(currentTime), 'WITHSCORES', 'LIMIT', 0, 1) + + if #messages >= 2 then + local messageId = messages[1] + local messageScore = messages[2] + + local messageKey = messageKeyPrefix .. messageId + local messagePayload = redis.call('GET', messageKey) + + if messagePayload then + local messageData = cjson.decode(messagePayload) + local ttlExpiresAt = messageData and messageData.ttlExpiresAt + + if ttlExpiresAt and ttlExpiresAt <= currentTime then + -- TTL expired - remove from queues + redis.call('ZREM', fullQueueKey, messageId) + redis.call('ZREM', envQueueKey, messageId) + else + -- Dequeue normally + redis.call('ZREM', fullQueueKey, messageId) + redis.call('ZREM', envQueueKey, messageId) + redis.call('SADD', ckConcurrencyKey, messageId) + redis.call('SADD', envCurrentConcurrencyKey, messageId) + + -- Remove from TTL set if applicable + if ttlQueueKey and ttlQueueKey ~= '' and ttlExpiresAt then + local ttlMember = ckQueueName .. '|' .. messageId .. '|' .. (messageData.orgId or '') + redis.call('ZREM', ttlQueueKey, ttlMember) + end + + table.insert(results, messageId) + table.insert(results, messageScore) + table.insert(results, messagePayload) + + dequeuedCount = dequeuedCount + 1 + end + else + -- Stale entry + redis.call('ZREM', fullQueueKey, messageId) + redis.call('ZREM', envQueueKey, messageId) + end + + -- Rebalance CK index for this sub-queue + local earliest = redis.call('ZRANGE', fullQueueKey, 0, 0, 'WITHSCORES') + if #earliest == 0 then + redis.call('ZREM', ckIndexKey, ckQueueName) + else + redis.call('ZADD', ckIndexKey, earliest[2], ckQueueName) + end + else + -- No messages available in score range, update CK index + local any = redis.call('ZRANGE', fullQueueKey, 0, 0, 'WITHSCORES') + if #any == 0 then + redis.call('ZREM', ckIndexKey, ckQueueName) + else + redis.call('ZADD', ckIndexKey, any[2], ckQueueName) + end + end + end +end + +-- Rebalance master queue +local earliestIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') +if #earliestIdx == 0 then + redis.call('ZREM', masterQueueKey, ckWildcardName) +else + redis.call('ZADD', masterQueueKey, earliestIdx[2], ckWildcardName) +end + +return results + `, + }); + + // Tracked variant: same as dequeueMessagesFromCkQueue plus DECR of the + // per-base-queue lengthCounter for every message removed from a CK variant + // (normal dequeue, TTL-expired, or stale-orphan path — all of which were + // counted at enqueue time). + this.redis.defineCommand("dequeueMessagesFromCkQueueTracked", { + numberOfKeys: 11, + lua: ` +local ckIndexKey = KEYS[1] +local queueConcurrencyLimitKey = KEYS[2] +local envConcurrencyLimitKey = KEYS[3] +local envConcurrencyLimitBurstFactorKey = KEYS[4] +local envCurrentConcurrencyKey = KEYS[5] +local messageKeyPrefix = KEYS[6] +local envQueueKey = KEYS[7] +local masterQueueKey = KEYS[8] +local ttlQueueKey = KEYS[9] +local lengthCounterKey = KEYS[10] +local runningCounterKey = KEYS[11] + +local ckWildcardName = ARGV[1] +local currentTime = tonumber(ARGV[2]) +local defaultEnvConcurrencyLimit = ARGV[3] +local defaultEnvConcurrencyBurstFactor = ARGV[4] +local keyPrefix = ARGV[5] +local maxCount = tonumber(ARGV[6] or '1') +${QUEUE_METRICS_GAUGE_PRELUDE} +${QUEUE_METRICS_CK_DEQUEUE_GAUGE_LUA} + +local function decrLengthCounter() + if tonumber(redis.call('GET', lengthCounterKey) or '0') > 0 then + redis.call('DECR', lengthCounterKey) + end +end + +-- Check env concurrency +local envCurrentConcurrency = tonumber(redis.call('SCARD', envCurrentConcurrencyKey) or '0') +local envConcurrencyLimit = tonumber(redis.call('GET', envConcurrencyLimitKey) or defaultEnvConcurrencyLimit) +local envConcurrencyLimitBurstFactor = tonumber(redis.call('GET', envConcurrencyLimitBurstFactorKey) or defaultEnvConcurrencyBurstFactor) +local envConcurrencyLimitWithBurstFactor = math.floor(envConcurrencyLimit * envConcurrencyLimitBurstFactor) + +if envCurrentConcurrency >= envConcurrencyLimitWithBurstFactor then + return __qmret(nil) end --- Get base queue concurrency limit local queueConcurrencyLimit = math.min(tonumber(redis.call('GET', queueConcurrencyLimitKey) or '1000000'), envConcurrencyLimit) --- Calculate env available capacity local envAvailableCapacity = envConcurrencyLimitWithBurstFactor - envCurrentConcurrency local actualMaxCount = math.min(maxCount, envAvailableCapacity) if actualMaxCount <= 0 then - return nil + return __qmret(nil) end --- Get CK sub-queues from CK index with scores <= currentTime local ckQueues = redis.call('ZRANGEBYSCORE', ckIndexKey, '-inf', tostring(currentTime), 'LIMIT', 0, actualMaxCount * 3) if #ckQueues == 0 then - -- Rebalance master queue in case CK index has future-scored entries local anyIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') if #anyIdx == 0 then redis.call('ZREM', masterQueueKey, ckWildcardName) else redis.call('ZADD', masterQueueKey, anyIdx[2], ckWildcardName) end - return nil + return __qmret(nil) end local results = {} @@ -4380,12 +4978,10 @@ for _, ckQueueName in ipairs(ckQueues) do local fullQueueKey = keyPrefix .. ckQueueName - -- Check CK-specific concurrency local ckConcurrencyKey = fullQueueKey .. ':currentConcurrency' local ckCurrentConcurrency = tonumber(redis.call('SCARD', ckConcurrencyKey) or '0') if ckCurrentConcurrency < queueConcurrencyLimit then - -- Try to dequeue one message from this CK sub-queue local messages = redis.call('ZRANGEBYSCORE', fullQueueKey, '-inf', tostring(currentTime), 'WITHSCORES', 'LIMIT', 0, 1) if #messages >= 2 then @@ -4400,17 +4996,16 @@ for _, ckQueueName in ipairs(ckQueues) do local ttlExpiresAt = messageData and messageData.ttlExpiresAt if ttlExpiresAt and ttlExpiresAt <= currentTime then - -- TTL expired - remove from queues redis.call('ZREM', fullQueueKey, messageId) redis.call('ZREM', envQueueKey, messageId) + decrLengthCounter() else - -- Dequeue normally redis.call('ZREM', fullQueueKey, messageId) redis.call('ZREM', envQueueKey, messageId) + decrLengthCounter() redis.call('SADD', ckConcurrencyKey, messageId) redis.call('SADD', envCurrentConcurrencyKey, messageId) - -- Remove from TTL set if applicable if ttlQueueKey and ttlQueueKey ~= '' and ttlExpiresAt then local ttlMember = ckQueueName .. '|' .. messageId .. '|' .. (messageData.orgId or '') redis.call('ZREM', ttlQueueKey, ttlMember) @@ -4423,12 +5018,11 @@ for _, ckQueueName in ipairs(ckQueues) do dequeuedCount = dequeuedCount + 1 end else - -- Stale entry redis.call('ZREM', fullQueueKey, messageId) redis.call('ZREM', envQueueKey, messageId) + decrLengthCounter() end - -- Rebalance CK index for this sub-queue local earliest = redis.call('ZRANGE', fullQueueKey, 0, 0, 'WITHSCORES') if #earliest == 0 then redis.call('ZREM', ckIndexKey, ckQueueName) @@ -4436,7 +5030,6 @@ for _, ckQueueName in ipairs(ckQueues) do redis.call('ZADD', ckIndexKey, earliest[2], ckQueueName) end else - -- No messages available in score range, update CK index local any = redis.call('ZRANGE', fullQueueKey, 0, 0, 'WITHSCORES') if #any == 0 then redis.call('ZREM', ckIndexKey, ckQueueName) @@ -4447,7 +5040,6 @@ for _, ckQueueName in ipairs(ckQueues) do end end --- Rebalance master queue local earliestIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') if #earliestIdx == 0 then redis.call('ZREM', masterQueueKey, ckWildcardName) @@ -4455,16 +5047,22 @@ else redis.call('ZADD', masterQueueKey, earliestIdx[2], ckWildcardName) end -return results +return __qmret(results) `, }); - // Tracked variant: same as dequeueMessagesFromCkQueue plus DECR of the - // per-base-queue lengthCounter for every message removed from a CK variant - // (normal dequeue, TTL-expired, or stale-orphan path — all of which were - // counted at enqueue time). - this.redis.defineCommand("dequeueMessagesFromCkQueueTracked", { - numberOfKeys: 11, + // Virtual-time (SFQ) variant of dequeueMessagesFromCkQueueTracked. + // Flag-selected via ckVirtualTimeScheduling. Orders concurrency-key variants + // by virtual-time tag (ckVtime ZSET) instead of head timestamp, layered under + // the existing per-variant concurrency gate. Two passes: pass 1 serves in fair + // (lowest-tag) order; pass 2 fills the batch + discovers unregistered variants + // in the existing age order (work conservation, mixed-deploy safety). Only the + // :ckVtime / :ckVtimeFloor keys hold virtual times; ckIndex and the master + // queue keep their timestamp score domain. The per-candidate serve body is a + // verbatim copy of dequeueMessagesFromCkQueueTracked's, with the marked NEW + // lines added (tag advance on serve, ZREM ckVtime on GC). + this.redis.defineCommand("dequeueMessagesFromCkQueueVtimeTracked", { + numberOfKeys: 13, lua: ` local ckIndexKey = KEYS[1] local queueConcurrencyLimitKey = KEYS[2] @@ -4477,6 +5075,8 @@ local masterQueueKey = KEYS[8] local ttlQueueKey = KEYS[9] local lengthCounterKey = KEYS[10] local runningCounterKey = KEYS[11] +local ckVtimeKey = KEYS[12] +local ckVtimeFloorKey = KEYS[13] local ckWildcardName = ARGV[1] local currentTime = tonumber(ARGV[2]) @@ -4484,6 +5084,9 @@ local defaultEnvConcurrencyLimit = ARGV[3] local defaultEnvConcurrencyBurstFactor = ARGV[4] local keyPrefix = ARGV[5] local maxCount = tonumber(ARGV[6] or '1') +local quantum = tonumber(ARGV[7] or '1') +local windowMultiplier = tonumber(ARGV[8] or '3') +local stateTtl = tonumber(ARGV[9] or '86400') ${QUEUE_METRICS_GAUGE_PRELUDE} ${QUEUE_METRICS_CK_DEQUEUE_GAUGE_LUA} @@ -4512,26 +5115,26 @@ if actualMaxCount <= 0 then return __qmret(nil) end -local ckQueues = redis.call('ZRANGEBYSCORE', ckIndexKey, '-inf', tostring(currentTime), 'LIMIT', 0, actualMaxCount * 3) +local window = actualMaxCount * windowMultiplier -if #ckQueues == 0 then - local anyIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') - if #anyIdx == 0 then - redis.call('ZREM', masterQueueKey, ckWildcardName) - else - redis.call('ZADD', masterQueueKey, anyIdx[2], ckWildcardName) +-- Monotonic floor, advanced to the minimum stored virtual-time tag +local floor = tonumber(redis.call('GET', ckVtimeFloorKey) or '0') +local minEntry = redis.call('ZRANGE', ckVtimeKey, 0, 0, 'WITHSCORES') +if #minEntry > 0 then + local minTag = tonumber(minEntry[2]) + if minTag > floor then + floor = minTag end - return __qmret(nil) end local results = {} local dequeuedCount = 0 +local attempted = {} -for _, ckQueueName in ipairs(ckQueues) do - if dequeuedCount >= actualMaxCount then - break - end - +-- Per-candidate serve. Body is dequeueMessagesFromCkQueueTracked's per-candidate +-- block, verbatim, with the marked NEW lines added. +local function tryServe(ckQueueName) + attempted[ckQueueName] = true local fullQueueKey = keyPrefix .. ckQueueName local ckConcurrencyKey = fullQueueKey .. ':currentConcurrency' @@ -4572,6 +5175,12 @@ for _, ckQueueName in ipairs(ckQueues) do table.insert(results, messagePayload) dequeuedCount = dequeuedCount + 1 + + -- NEW: advance this variant's virtual time (weight hook: fixed 1 today) + local weight = 1 + local tag = tonumber(redis.call('ZSCORE', ckVtimeKey, ckQueueName) or floor) + if tag < floor then tag = floor end + redis.call('ZADD', ckVtimeKey, tostring(tag + (quantum / weight)), ckQueueName) end else redis.call('ZREM', fullQueueKey, messageId) @@ -4582,6 +5191,7 @@ for _, ckQueueName in ipairs(ckQueues) do local earliest = redis.call('ZRANGE', fullQueueKey, 0, 0, 'WITHSCORES') if #earliest == 0 then redis.call('ZREM', ckIndexKey, ckQueueName) + redis.call('ZREM', ckVtimeKey, ckQueueName) -- NEW else redis.call('ZADD', ckIndexKey, earliest[2], ckQueueName) end @@ -4589,6 +5199,7 @@ for _, ckQueueName in ipairs(ckQueues) do local any = redis.call('ZRANGE', fullQueueKey, 0, 0, 'WITHSCORES') if #any == 0 then redis.call('ZREM', ckIndexKey, ckQueueName) + redis.call('ZREM', ckVtimeKey, ckQueueName) -- NEW else redis.call('ZADD', ckIndexKey, any[2], ckQueueName) end @@ -4596,6 +5207,34 @@ for _, ckQueueName in ipairs(ckQueues) do end end +-- Pass 1: fair order (lowest virtual start tag first) +local vtimeCandidates = redis.call('ZRANGE', ckVtimeKey, 0, window - 1) +for _, ckQueueName in ipairs(vtimeCandidates) do + if dequeuedCount >= actualMaxCount then break end + tryServe(ckQueueName) +end + +-- Pass 2: fill + discovery in age order (work conservation, mixed-deploy safety). +-- Never runs when pass 1 filled the batch. +if dequeuedCount < actualMaxCount then + -- Clamp to at least 3x so pass 2 never scans fewer index variants than the old command, preserving work conservation regardless of the configured multiplier + local pass2Window = math.max(window, actualMaxCount * 3) + local ckQueues = redis.call('ZRANGEBYSCORE', ckIndexKey, '-inf', tostring(currentTime), 'LIMIT', 0, pass2Window) + for _, ckQueueName in ipairs(ckQueues) do + if dequeuedCount >= actualMaxCount then break end + if not attempted[ckQueueName] then + tryServe(ckQueueName) + end + end +end + +-- NEW: persist floor and refresh TTLs +redis.call('SET', ckVtimeFloorKey, tostring(floor), 'EX', stateTtl) +if redis.call('EXISTS', ckVtimeKey) == 1 then + redis.call('EXPIRE', ckVtimeKey, stateTtl) +end + +-- Rebalance master queue (ckIndex keeps its timestamp domain) local earliestIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') if #earliestIdx == 0 then redis.call('ZREM', masterQueueKey, ckWildcardName) @@ -5198,6 +5837,110 @@ else redis.call('ZADD', masterQueueKey, earliestIdx[2], ckWildcardName) end +-- Remove old-format entry from master queue (transition cleanup) +redis.call('ZREM', masterQueueKey, messageQueueName) +`, + }); + + // Vtime variant of nackMessageCkTracked (feature-flagged via + // ckVirtualTimeScheduling.enabled). Identical script body, plus registration + // of the variant into the :ckVtime ZSET at the floor (NX), so a GC'd variant + // that a nack revives rejoins the fair order. + this.redis.defineCommand("nackMessageCkVtimeTracked", { + numberOfKeys: 13, + lua: ` +-- Keys: +local masterQueueKey = KEYS[1] +local messageKey = KEYS[2] +local messageQueueKey = KEYS[3] +local queueCurrentConcurrencyKey = KEYS[4] +local envCurrentConcurrencyKey = KEYS[5] +local queueCurrentDequeuedKey = KEYS[6] +local envCurrentDequeuedKey = KEYS[7] +local envQueueKey = KEYS[8] +local ckIndexKey = KEYS[9] +local lengthCounterKey = KEYS[10] +local runningCounterKey = KEYS[11] +-- Virtual-time keys (KEYS 12-13) +local ckVtimeKey = KEYS[12] +local ckVtimeFloorKey = KEYS[13] + +-- Args: +local messageId = ARGV[1] +local messageQueueName = ARGV[2] +local messageData = ARGV[3] +local messageScore = tonumber(ARGV[4]) +local ckWildcardName = ARGV[5] +-- keyPrefix for prepending to variant names stored as values in ckIndex (lazy-init only) +local keyPrefix = ARGV[6] +-- TTL (seconds) applied to counter lazy-init SETs +local counterTtl = ARGV[7] +-- TTL (seconds) applied to ckVtime on registration +local stateTtl = ARGV[8] + +local function decrFloored(key) + if tonumber(redis.call('GET', key) or '0') > 0 then + redis.call('DECR', key) + end +end + +-- Update the message data +redis.call('SET', messageKey, messageData) + +-- Update the concurrency keys. nack only DECRs runningCounter, never INCRs it, +-- so we skip the eager lazy-init here (unlike releaseConcurrencyTracked, which +-- mirrors the same DECR pattern with init). A post-TTL nack's floored DECR +-- no-ops; the next dequeueMessageFromKeyTracked reseeds from current state. +redis.call('SREM', queueCurrentConcurrencyKey, messageId) +redis.call('SREM', envCurrentConcurrencyKey, messageId) +local removedFromDequeued = redis.call('SREM', queueCurrentDequeuedKey, messageId) +redis.call('SREM', envCurrentDequeuedKey, messageId) +if removedFromDequeued == 1 then + decrFloored(runningCounterKey) +end + +-- Lazy-init lengthCounter if missing (e.g. expired via 24h TTL). nack re-queues a +-- message, which means lengthCounter must be present before we INCR. Without this, +-- a nack after counter expiry would create the counter at 1 and stay drifted until +-- next reset. +if redis.call('EXISTS', lengthCounterKey) == 0 then + local total = 0 + local variants = redis.call('ZRANGE', ckIndexKey, 0, -1) + for _, v in ipairs(variants) do + total = total + tonumber(redis.call('ZCARD', keyPrefix .. v) or '0') + end + redis.call('SET', lengthCounterKey, total, 'EX', counterTtl) +end + +-- Enqueue the message back into the CK-specific queue. INCR lengthCounter only if +-- it's a new entry (ZADD returns 1). +local added = redis.call('ZADD', messageQueueKey, messageScore, messageId) +redis.call('ZADD', envQueueKey, messageScore, messageId) +if added == 1 then + redis.call('INCR', lengthCounterKey) +end + +-- Rebalance CK index +local earliest = redis.call('ZRANGE', messageQueueKey, 0, 0, 'WITHSCORES') +if #earliest > 0 then + redis.call('ZADD', ckIndexKey, earliest[2], messageQueueName) +end + +-- Register this variant in the virtual-time index at the floor. NX means an +-- already-advanced tag is never rewound. +local vfloor = redis.call('GET', ckVtimeFloorKey) or '0' +redis.call('ZADD', ckVtimeKey, 'NX', vfloor, messageQueueName) +redis.call('EXPIRE', ckVtimeKey, stateTtl) +redis.call('EXPIRE', ckVtimeFloorKey, stateTtl) + +-- Rebalance master queue with ck:* member +local earliestIdx = redis.call('ZRANGE', ckIndexKey, 0, 0, 'WITHSCORES') +if #earliestIdx == 0 then + redis.call('ZREM', masterQueueKey, ckWildcardName) +else + redis.call('ZADD', masterQueueKey, earliestIdx[2], ckWildcardName) +end + -- Remove old-format entry from master queue (transition cleanup) redis.call('ZREM', masterQueueKey, messageQueueName) `, @@ -5917,6 +6660,79 @@ declare module "@internal/redis" { callback?: Callback<[number, number[] | null]> ): Result<[number, number[] | null], Context>; + enqueueMessageCkVtimeTracked( + masterQueueKey: string, + queue: string, + messageKey: string, + queueCurrentConcurrencyKey: string, + envCurrentConcurrencyKey: string, + queueCurrentDequeuedKey: string, + envCurrentDequeuedKey: string, + envQueueKey: string, + ckIndexKey: string, + workerQueueKey: string, + queueConcurrencyLimitKey: string, + envConcurrencyLimitKey: string, + envConcurrencyLimitBurstFactorKey: string, + lengthCounterKey: string, + baseQueueKey: string, + ckVtimeKey: string, + ckVtimeFloorKey: string, + queueName: string, + messageId: string, + messageData: string, + messageScore: string, + ckWildcardName: string, + messageKeyValue: string, + defaultEnvConcurrencyLimit: string, + defaultEnvConcurrencyBurstFactor: string, + currentTime: string, + enableFastPath: string, + keyPrefix: string, + counterTtl: string, + stateTtl: string, + metricsEnabled: string, + callback?: Callback<[number, number[] | null]> + ): Result<[number, number[] | null], Context>; + + enqueueMessageWithTtlCkVtimeTracked( + masterQueueKey: string, + queue: string, + messageKey: string, + queueCurrentConcurrencyKey: string, + envCurrentConcurrencyKey: string, + queueCurrentDequeuedKey: string, + envCurrentDequeuedKey: string, + envQueueKey: string, + ttlQueueKey: string, + ckIndexKey: string, + workerQueueKey: string, + queueConcurrencyLimitKey: string, + envConcurrencyLimitKey: string, + envConcurrencyLimitBurstFactorKey: string, + lengthCounterKey: string, + baseQueueKey: string, + ckVtimeKey: string, + ckVtimeFloorKey: string, + queueName: string, + messageId: string, + messageData: string, + messageScore: string, + ttlMember: string, + ttlScore: string, + ckWildcardName: string, + messageKeyValue: string, + defaultEnvConcurrencyLimit: string, + defaultEnvConcurrencyBurstFactor: string, + currentTime: string, + enableFastPath: string, + keyPrefix: string, + counterTtl: string, + stateTtl: string, + metricsEnabled: string, + callback?: Callback<[number, number[] | null]> + ): Result<[number, number[] | null], Context>; + dequeueMessagesFromCkQueueTracked( ckIndexKey: string, queueConcurrencyLimitKey: string, @@ -5939,6 +6755,33 @@ declare module "@internal/redis" { callback?: Callback<[string[] | null, number[] | null]> ): Result<[string[] | null, number[] | null], Context>; + dequeueMessagesFromCkQueueVtimeTracked( + ckIndexKey: string, + queueConcurrencyLimitKey: string, + envConcurrencyLimitKey: string, + envConcurrencyLimitBurstFactorKey: string, + envCurrentConcurrencyKey: string, + messageKeyPrefix: string, + envQueueKey: string, + masterQueueKey: string, + ttlQueueKey: string, + lengthCounterKey: string, + runningCounterKey: string, + ckVtimeKey: string, + ckVtimeFloorKey: string, + ckWildcardName: string, + currentTime: string, + defaultEnvConcurrencyLimit: string, + defaultEnvConcurrencyBurstFactor: string, + keyPrefix: string, + maxCount: string, + quantum: string, + windowMultiplier: string, + stateTtlSeconds: string, + metricsEnabled: string, + callback?: Callback<[string[] | null, number[] | null]> + ): Result<[string[] | null, number[] | null], Context>; + dequeueMessageFromKeyTracked( messageKey: string, keyPrefix: string, @@ -5989,6 +6832,31 @@ declare module "@internal/redis" { callback?: Callback ): Result; + nackMessageCkVtimeTracked( + masterQueueKey: string, + messageKey: string, + messageQueue: string, + queueCurrentConcurrencyKey: string, + envCurrentConcurrencyKey: string, + queueCurrentDequeuedKey: string, + envCurrentDequeuedKey: string, + envQueueKey: string, + ckIndexKey: string, + lengthCounterKey: string, + runningCounterKey: string, + ckVtimeKey: string, + ckVtimeFloorKey: string, + messageId: string, + messageQueueName: string, + messageData: string, + messageScore: string, + ckWildcardName: string, + keyPrefix: string, + counterTtl: string, + stateTtl: string, + callback?: Callback + ): Result; + moveToDeadLetterQueueCkTracked( masterQueueKey: string, messageKey: string, diff --git a/internal-packages/run-engine/src/run-queue/keyProducer.ts b/internal-packages/run-engine/src/run-queue/keyProducer.ts index 0609b3d719b..c00ee47aba4 100644 --- a/internal-packages/run-engine/src/run-queue/keyProducer.ts +++ b/internal-packages/run-engine/src/run-queue/keyProducer.ts @@ -22,6 +22,8 @@ const constants = { MASTER_QUEUE_PART: "masterQueue", WORKER_QUEUE_PART: "workerQueue", CK_INDEX_PART: "ckIndex", + CK_VTIME_PART: "ckVtime", + CK_VTIME_FLOOR_PART: "ckVtimeFloor", LENGTH_COUNTER_PART: "lengthCounter", RUNNING_COUNTER_PART: "runningCounter", } as const; @@ -315,6 +317,14 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer { return `${this.baseQueueKeyFromQueue(queue)}:${constants.CK_INDEX_PART}`; } + ckVtimeKeyFromQueue(queue: string): string { + return `${this.baseQueueKeyFromQueue(queue)}:${constants.CK_VTIME_PART}`; + } + + ckVtimeFloorKeyFromQueue(queue: string): string { + return `${this.baseQueueKeyFromQueue(queue)}:${constants.CK_VTIME_FLOOR_PART}`; + } + // indexOf instead of /:ck:.+$/ (queue names are user-controlled; polynomial regex). // Only strips when at least one character follows ":ck:", matching the old semantics. baseQueueKeyFromQueue(queue: string): string { diff --git a/internal-packages/run-engine/src/run-queue/tests/ckVtime.test.ts b/internal-packages/run-engine/src/run-queue/tests/ckVtime.test.ts new file mode 100644 index 00000000000..98b0ac73d17 --- /dev/null +++ b/internal-packages/run-engine/src/run-queue/tests/ckVtime.test.ts @@ -0,0 +1,1061 @@ +import { redisTest } from "@internal/testcontainers"; +import { trace } from "@internal/tracing"; +import { Logger } from "@trigger.dev/core/logger"; +import { Decimal } from "@trigger.dev/database"; +import { describe } from "node:test"; +import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js"; +import { RunQueue } from "../index.js"; +import { RunQueueFullKeyProducer } from "../keyProducer.js"; +import type { InputPayload } from "../types.js"; + +const testOptions = { + name: "rq", + tracer: trace.getTracer("rq"), + workers: 1, + defaultEnvConcurrency: 25, + logger: new Logger("RunQueue", "warn"), + 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: 10, + concurrencyLimitBurstFactor: new Decimal(2.0), + project: { id: "p1234" }, + organization: { id: "o1234" }, +}; + +type VtimeOverrides = { + enabled?: boolean; + quantum?: number; + scanWindowMultiplier?: number; + stateTtlSeconds?: number; +}; + +// vtime: overrides merged into an enabled ckVirtualTimeScheduling option, or +// null to omit the option entirely (flag off, the production default). +function createQueue(redisContainer: any, vtime: VtimeOverrides | null = {}) { + return new RunQueue({ + ...testOptions, + // These tests drive every op themselves (testDequeueFromMasterQueue + skipDequeueProcessing), + // so the autonomous master-queue consumers and background worker must not race them. + masterQueueConsumersDisabled: true, + workerOptions: { disabled: true }, + ...(vtime === null + ? {} + : { + ckVirtualTimeScheduling: { + enabled: true, + ...vtime, + }, + }), + queueSelectionStrategy: new FairQueueSelectionStrategy({ + redis: { + keyPrefix: "runqueue:test:", + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + keys: testOptions.keys, + }), + redis: { + keyPrefix: "runqueue:test:", + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + }); +} + +function makeMessage(overrides: Partial = {}): InputPayload { + return { + runId: "r1", + taskIdentifier: "task/my-task", + orgId: "o1234", + projectId: "p1234", + environmentId: "e1234", + environmentType: "DEVELOPMENT", + queue: "task/my-task", + timestamp: Date.now(), + attempt: 0, + ...overrides, + }; +} + +// The ckVtime/ckIndex member for a variant is the fully-qualified variant queue +// key (org:proj:env:queue:...:ck:), which is exactly what queueKey() produces. +function variantName(ck: string): string { + return testOptions.keys.queueKey(authenticatedEnvDev, "task/my-task", ck); +} + +const QUEUE = "task/my-task"; + +vi.setConfig({ testTimeout: 60_000 }); + +describe("CK virtual-time (SFQ) dequeue", () => { + redisTest("vtime order beats head-timestamp order", async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + + // 30 old messages on heavy, timestamps t0..t0+29 + for (let i = 0; i < 30; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `h${i}`, concurrencyKey: "heavy", timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + // 3 much newer messages on light + for (let i = 0; i < 3; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `l${i}`, concurrencyKey: "light", timestamp: t0 + 1000 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + const heavyVariant = variantName("heavy"); + const lightVariant = variantName("light"); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(heavyVariant); + + // Explicit seed is redundant now that enqueue registers variants at the + // floor itself; kept as a belt-and-braces fixture. + await queue.redis.zadd(ckVtimeKey, 0, heavyVariant, 0, lightVariant); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + const lightServedInCall: number[] = []; + const lightSeen = new Set(); + + for (let call = 0; call < 3; call++) { + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 10); + + // At most one message per variant per call. + const heavyCount = messages.filter((m) => m.message.concurrencyKey === "heavy").length; + const lightCount = messages.filter((m) => m.message.concurrencyKey === "light").length; + expect(heavyCount).toBeLessThanOrEqual(1); + expect(lightCount).toBeLessThanOrEqual(1); + + for (const m of messages) { + if (m.message.concurrencyKey === "light" && !lightSeen.has(m.messageId)) { + lightSeen.add(m.messageId); + lightServedInCall.push(call); + } + // ack to free concurrency between calls + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + + // All 3 light messages served within the first 3 calls (age order alone + // would have drained the 30 heavy messages first). + expect(lightSeen.size).toBe(3); + } finally { + await queue.quit(); + } + }); + + redisTest( + "vtime order wins over older head when maxCount forces a choice", + async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + + // Old messages on heavy: age order strictly favours heavy. + for (let i = 0; i < 5; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `h${i}`, concurrencyKey: "heavy", timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + // Newer messages on light (two, so light isn't GC'd after its serve). + for (let i = 0; i < 2; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `l${i}`, + concurrencyKey: "light", + timestamp: t0 + 1000, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + const heavyVariant = variantName("heavy"); + const lightVariant = variantName("light"); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(heavyVariant); + + // Seed vtime the OPPOSITE way to age order: heavy has the HIGH tag, + // light the LOW tag. + await queue.redis.zadd(ckVtimeKey, 10, heavyVariant, 0, lightVariant); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + // maxCount 1 with two ready variants: ordering alone decides who is + // served. Age order (the old command) would pick heavy's older head; + // vtime rank must pick light's lower tag. + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 1); + + expect(messages.length).toBe(1); + expect(messages[0]!.message.concurrencyKey).toBe("light"); + + // Light's tag advanced by the quantum (=1); heavy's is untouched. + const lightTag = Number(await queue.redis.zscore(ckVtimeKey, lightVariant)); + const heavyTag = Number(await queue.redis.zscore(ckVtimeKey, heavyVariant)); + expect(lightTag).toBe(1); + expect(heavyTag).toBe(10); + } finally { + await queue.quit(); + } + } + ); + + redisTest("tags advance per serve within one batched call", async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + const cks = ["a", "b", "c", "d", "e"]; + + // Two messages per variant so a single serve doesn't drain (and GC) it, + // letting us observe the advanced tag afterwards. + for (const ck of cks) { + for (let i = 0; i < 2; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `r-${ck}-${i}`, concurrencyKey: ck, timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + } + + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("a")); + const seedArgs: (string | number)[] = []; + for (const ck of cks) { + seedArgs.push(0, variantName(ck)); + } + await queue.redis.zadd(ckVtimeKey, ...(seedArgs as [number, string])); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 5); + + expect(messages.length).toBe(5); + + for (const ck of cks) { + const score = await queue.redis.zscore(ckVtimeKey, variantName(ck)); + expect(Number(score)).toBe(1); + } + } finally { + await queue.quit(); + } + }); + + redisTest("floor is monotonic and read-repairs", async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + const cks = ["a", "b"]; + + // enough messages per variant to keep serving for many calls + for (const ck of cks) { + for (let i = 0; i < 25; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `r-${ck}-${i}`, concurrencyKey: ck, timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + } + + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("a")); + const ckVtimeFloorKey = testOptions.keys.ckVtimeFloorKeyFromQueue(variantName("a")); + await queue.redis.zadd(ckVtimeKey, 0, variantName("a"), 0, variantName("b")); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + let prevFloor = 0; + for (let call = 0; call < 20; call++) { + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 2); + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + const floor = Number((await queue.redis.get(ckVtimeFloorKey)) ?? "0"); + expect(floor).toBeGreaterThanOrEqual(prevFloor); + prevFloor = floor; + } + + // Final settle: gate both variants (limit 1 + an occupied slot) so nothing + // is served (no advance), and the floor read-repairs up to the current min tag. + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, QUEUE, 1); + for (const ck of cks) { + await queue.redis.sadd( + testOptions.keys.queueCurrentConcurrencyKeyFromQueue(variantName(ck)), + "occupant" + ); + } + await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 2); + + const floorAfter = Number((await queue.redis.get(ckVtimeFloorKey)) ?? "0"); + expect(floorAfter).toBeGreaterThanOrEqual(prevFloor); + + const minEntry = await queue.redis.zrange(ckVtimeKey, 0, 0, "WITHSCORES"); + const minTag = Number(minEntry[1]); + expect(floorAfter).toBe(minTag); + expect(floorAfter).toBeGreaterThan(10); + } finally { + await queue.quit(); + } + }); + + redisTest( + "new key initialises at the floor, not zero and not behind the backlog", + async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + const cks = ["a", "b"]; + + for (const ck of cks) { + for (let i = 0; i < 25; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-${ck}-${i}`, + concurrencyKey: ck, + timestamp: t0 + i, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + } + + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("a")); + const ckVtimeFloorKey = testOptions.keys.ckVtimeFloorKeyFromQueue(variantName("a")); + // No direct ZADD seeding: the enqueues above register a and b at the + // initial floor (0) themselves. + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + // Drive tags up to ~20. + for (let call = 0; call < 20; call++) { + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 2); + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + + const floor = Number((await queue.redis.get(ckVtimeFloorKey)) ?? "0"); + expect(floor).toBeGreaterThan(10); + + // A fresh variant registers itself at the current floor via the + // enqueue-time registration (ZADD NX in the enqueue script). + const freshVariant = variantName("fresh"); + // two messages so the fresh variant isn't GC'd on its first serve + for (let i = 0; i < 2; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-fresh-${i}`, + concurrencyKey: "fresh", + timestamp: t0 + i, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 10); + const served = messages.some((m) => m.message.concurrencyKey === "fresh"); + expect(served).toBe(true); + + const freshTag = Number(await queue.redis.zscore(ckVtimeKey, freshVariant)); + // initialised at the floor and advanced by quantum (=1), not stuck at 1 + expect(freshTag).toBe(floor + 1); + } finally { + await queue.quit(); + } + } + ); + + // H1 regression: the floor key must not be allowed to expire while ckVtime + // survives. Before the fix, only the dequeue command refreshed the floor + // key's TTL, so a dequeue-quiescent + enqueue-active base queue let the floor + // key expire underneath a live ckVtime; a brand-new variant then read a + // missing floor as 0 and jumped ahead of the whole established backlog. The + // enqueue/nack registration paths now refresh the floor key TTL too. + redisTest( + "enqueue refreshes the floor key TTL and a new variant registers at the current floor", + async ({ redisContainer }) => { + const stateTtlSeconds = 3600; + const queue = createQueue(redisContainer, { stateTtlSeconds }); + try { + const t0 = Date.now() - 100_000; + const cks = ["a", "b"]; + + for (const ck of cks) { + for (let i = 0; i < 25; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-${ck}-${i}`, + concurrencyKey: ck, + timestamp: t0 + i, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + } + + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("a")); + const ckVtimeFloorKey = testOptions.keys.ckVtimeFloorKeyFromQueue(variantName("a")); + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + // Drive tags and the floor above 0 with a run of serves. + for (let call = 0; call < 20; call++) { + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 2); + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + + const floor = Number((await queue.redis.get(ckVtimeFloorKey)) ?? "0"); + expect(floor).toBeGreaterThan(10); + expect(await queue.redis.exists(ckVtimeKey)).toBe(1); + + // Simulate the floor key's TTL decaying toward expiry while dequeues are + // quiescent. Without the fix, only a dequeue would ever bump it back. + await queue.redis.pexpire(ckVtimeFloorKey, 2_000); + + // WITHOUT dequeuing, enqueue several more messages on an existing + // variant. The enqueue registration path must refresh the floor key TTL. + for (let i = 0; i < 5; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-a-more-${i}`, + concurrencyKey: "a", + timestamp: t0 + 500 + i, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + // The floor key TTL was pushed back up to (about) stateTtl, well above + // the 2s decay we forced. + const floorPttl = await queue.redis.pttl(ckVtimeFloorKey); + expect(floorPttl).toBeGreaterThan(2_000); + expect(floorPttl).toBeLessThanOrEqual(stateTtlSeconds * 1000); + + // Enqueue-only activity does not move the floor value itself. + const floorAfter = Number((await queue.redis.get(ckVtimeFloorKey)) ?? "0"); + expect(floorAfter).toBe(floor); + + // A brand-new variant enqueued now registers at the CURRENT floor, so it + // cannot leapfrog the established backlog back to 0. + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-fresh", concurrencyKey: "fresh", timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + const freshTag = Number(await queue.redis.zscore(ckVtimeKey, variantName("fresh"))); + expect(freshTag).toBe(floor); + } finally { + await queue.quit(); + } + } + ); + + redisTest("no service, no advance", async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + + // ck:a has messages but its concurrency slot will be occupied + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-a", concurrencyKey: "a", timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + for (const ck of ["b", "c"]) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `r-${ck}`, concurrencyKey: ck, timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("a")); + await queue.redis.zadd( + ckVtimeKey, + 0, + variantName("a"), + 0, + variantName("b"), + 0, + variantName("c") + ); + + // base queue concurrency limit of 1 + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, QUEUE, 1); + + // occupy ck:a's single slot (equivalent to a prior dequeue-without-ack) + await queue.redis.sadd( + testOptions.keys.queueCurrentConcurrencyKeyFromQueue(variantName("a")), + "occupant" + ); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 10); + + const servedCks = messages.map((m) => m.message.concurrencyKey); + expect(servedCks).not.toContain("a"); + expect(servedCks).toContain("b"); + expect(servedCks).toContain("c"); + + // ck:a's tag is unchanged (never served) + const aTag = Number(await queue.redis.zscore(ckVtimeKey, variantName("a"))); + expect(aTag).toBe(0); + } finally { + await queue.quit(); + } + }); + + redisTest( + "GC on empty variant removes it from ckIndex and ckVtime", + async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-a", concurrencyKey: "a", timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + // second variant so ckVtime/ckIndex don't fully disappear, keeping the test focused + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-b", concurrencyKey: "b", timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + + const aVariant = variantName("a"); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(aVariant); + const ckIndexKey = testOptions.keys.ckIndexKeyFromQueue(aVariant); + await queue.redis.zadd(ckVtimeKey, 0, aVariant, 0, variantName("b")); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 10); + expect(messages.some((m) => m.message.concurrencyKey === "a")).toBe(true); + + const inVtime = await queue.redis.zscore(ckVtimeKey, aVariant); + const inIndex = await queue.redis.zscore(ckIndexKey, aVariant); + expect(inVtime).toBeNull(); + expect(inIndex).toBeNull(); + } finally { + await queue.quit(); + } + } + ); + + redisTest("TTL is set and refreshed on ckVtime and ckVtimeFloor", async ({ redisContainer }) => { + const stateTtlSeconds = 3600; + const queue = createQueue(redisContainer, { stateTtlSeconds }); + try { + const t0 = Date.now() - 100_000; + + // two messages on one variant so ckVtime survives (not GC'd) after a serve + for (let i = 0; i < 2; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `r-a-${i}`, concurrencyKey: "a", timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + const aVariant = variantName("a"); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(aVariant); + const ckVtimeFloorKey = testOptions.keys.ckVtimeFloorKeyFromQueue(aVariant); + await queue.redis.zadd(ckVtimeKey, 0, aVariant); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 1); + + const vtimeTtl = await queue.redis.pttl(ckVtimeKey); + const floorTtl = await queue.redis.pttl(ckVtimeFloorKey); + + expect(vtimeTtl).toBeGreaterThan(0); + expect(vtimeTtl).toBeLessThanOrEqual(stateTtlSeconds * 1000); + expect(floorTtl).toBeGreaterThan(0); + expect(floorTtl).toBeLessThanOrEqual(stateTtlSeconds * 1000); + } finally { + await queue.quit(); + } + }); + + redisTest( + "pass 2 fill serves unregistered variants and registers them", + async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + + // Enqueue on a variant but do NOT register it in ckVtime (simulating an + // enqueue from old code that predates enqueue-time registration, e.g. + // during a rolling deploy). Two messages so the variant survives its + // first serve and we can observe it was registered. + for (let i = 0; i < 2; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `r-a-${i}`, concurrencyKey: "a", timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + const aVariant = variantName("a"); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(aVariant); + // ensure no ckVtime entry exists for it + await queue.redis.zrem(ckVtimeKey, aVariant); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 10); + + expect(messages.some((m) => m.message.concurrencyKey === "a")).toBe(true); + + const tag = await queue.redis.zscore(ckVtimeKey, aVariant); + expect(tag).not.toBeNull(); + } finally { + await queue.quit(); + } + } + ); + + redisTest("future-scheduled variants are skipped without advance", async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + + // a normal ready variant so the :ck:* wildcard is selected from the master queue + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-now", concurrencyKey: "now", timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + // a future-scheduled variant + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: "r-future", + concurrencyKey: "future", + timestamp: Date.now() + 60_000, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("now")); + const futureVariant = variantName("future"); + await queue.redis.zadd(ckVtimeKey, 0, variantName("now"), 5, futureVariant); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 10); + + expect(messages.some((m) => m.message.concurrencyKey === "future")).toBe(false); + + const futureTag = Number(await queue.redis.zscore(ckVtimeKey, futureVariant)); + expect(futureTag).toBe(5); + } finally { + await queue.quit(); + } + }); + + redisTest( + "enqueue registers the variant at the current floor with NX", + async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + + // two variants with enough messages that the drive loop never drains them + for (const ck of ["a", "b"]) { + for (let i = 0; i < 10; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-${ck}-${i}`, + concurrencyKey: ck, + timestamp: t0 + i, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + } + + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("a")); + const ckVtimeFloorKey = testOptions.keys.ckVtimeFloorKeyFromQueue(variantName("a")); + + // enqueue registered both variants at the initial floor (0), before any dequeue + expect(Number(await queue.redis.zscore(ckVtimeKey, variantName("a")))).toBe(0); + expect(Number(await queue.redis.zscore(ckVtimeKey, variantName("b")))).toBe(0); + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + // drive the floor up to ~5 via serves + for (let call = 0; call < 8; call++) { + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 2); + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + + const floor = Number((await queue.redis.get(ckVtimeFloorKey)) ?? "0"); + expect(floor).toBeGreaterThanOrEqual(5); + + // a fresh key enqueued now lands exactly at the floor (no dequeue in between) + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-fresh-0", concurrencyKey: "fresh", timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + expect(Number(await queue.redis.zscore(ckVtimeKey, variantName("fresh")))).toBe(floor); + + // NX: enqueueing on a key whose tag is already 9 never rewinds it + await queue.redis.zadd(ckVtimeKey, 9, variantName("nine")); + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-nine-0", concurrencyKey: "nine", timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + expect(Number(await queue.redis.zscore(ckVtimeKey, variantName("nine")))).toBe(9); + } finally { + await queue.quit(); + } + } + ); + + redisTest("fast path leaves vtime state untouched", async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const aVariant = variantName("a"); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(aVariant); + + // empty variant + free capacity: the fast path fires and skips the + // variant zset entirely + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-fast", concurrencyKey: "a" }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + enableFastPath: true, + }); + + // fast-path proof: nothing landed in the variant zset.. + expect(await queue.redis.zcard(aVariant)).toBe(0); + // ..and no vtime registration happened + expect(await queue.redis.zscore(ckVtimeKey, aVariant)).toBeNull(); + + // saturate capacity so the next enqueue takes the slow path + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, QUEUE, 1); + await queue.redis.sadd( + testOptions.keys.queueCurrentConcurrencyKeyFromQueue(aVariant), + "occupant" + ); + + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-slow", concurrencyKey: "a" }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + enableFastPath: true, + }); + + // slow path taken and the variant is registered + expect(await queue.redis.zcard(aVariant)).toBe(1); + expect(await queue.redis.zscore(ckVtimeKey, aVariant)).not.toBeNull(); + } finally { + await queue.quit(); + } + }); + + redisTest("nack re-registers a GC'd variant", async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const t0 = Date.now() - 100_000; + + // one message on ck:a, plenty on ck:b so serves keep flowing and the floor rises + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "r-a-0", concurrencyKey: "a", timestamp: t0 }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + for (let i = 0; i < 10; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `r-b-${i}`, concurrencyKey: "b", timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + const aVariant = variantName("a"); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(aVariant); + const ckVtimeFloorKey = testOptions.keys.ckVtimeFloorKeyFromQueue(aVariant); + const ckIndexKey = testOptions.keys.ckIndexKeyFromQueue(aVariant); + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + // first dequeue serves a's only message: a is drained and GC'd from both indexes + let aMessageId: string | undefined; + const first = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 2); + for (const m of first) { + if (m.message.concurrencyKey === "a") { + aMessageId = m.messageId; + } else { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + expect(aMessageId).toBeDefined(); + expect(await queue.redis.zscore(ckVtimeKey, aVariant)).toBeNull(); + expect(await queue.redis.zscore(ckIndexKey, aVariant)).toBeNull(); + + // drive the floor up via b serves + for (let call = 0; call < 6; call++) { + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 1); + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + const floor = Number((await queue.redis.get(ckVtimeFloorKey)) ?? "0"); + expect(floor).toBeGreaterThan(0); + + // nack with an immediate retry score so the revived message is servable now + await queue.nackMessage({ + orgId: authenticatedEnvDev.organization.id, + messageId: aMessageId!, + retryAt: Date.now(), + skipDequeueProcessing: true, + }); + + // the variant is back in ckIndex AND in ckVtime at the floor + expect(await queue.redis.zscore(ckIndexKey, aVariant)).not.toBeNull(); + expect(Number(await queue.redis.zscore(ckVtimeKey, aVariant))).toBe(floor); + + // and a subsequent dequeue serves it + const after = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 10); + expect(after.some((m) => m.messageId === aMessageId)).toBe(true); + } finally { + await queue.quit(); + } + }); + + redisTest("ckVtime membership tracks ckIndex membership", async ({ redisContainer }) => { + const queue = createQueue(redisContainer); + try { + const cks = ["a", "b", "c", "d", "e", "f", "g", "h"]; + const ckIndexKey = testOptions.keys.ckIndexKeyFromQueue(variantName("a")); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("a")); + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + // deterministic LCG so failures reproduce + let seed = 123456789; + const rand = () => { + seed = (seed * 1103515245 + 12345) % 2147483648; + return seed / 2147483648; + }; + const pick = (n: number) => Math.floor(rand() * n); + + const inFlight: string[] = []; + let nextRun = 0; + + for (let step = 0; step < 200; step++) { + const op = pick(4); + let opName = "noop"; + + if (op === 0) { + opName = "enqueue"; + const ck = cks[pick(cks.length)]!; + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r${nextRun++}`, + concurrencyKey: ck, + timestamp: Date.now() - 100_000, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } else if (op === 1) { + opName = "dequeue"; + const messages = await queue.testDequeueFromMasterQueue( + shard, + authenticatedEnvDev.id, + 1 + pick(4) + ); + for (const m of messages) { + inFlight.push(m.messageId); + } + } else if (op === 2 && inFlight.length > 0) { + opName = "ack"; + const [id] = inFlight.splice(pick(inFlight.length), 1); + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, id!, { + skipDequeueProcessing: true, + }); + } else if (op === 3 && inFlight.length > 0) { + opName = "nack"; + const [id] = inFlight.splice(pick(inFlight.length), 1); + await queue.nackMessage({ + orgId: authenticatedEnvDev.organization.id, + messageId: id!, + retryAt: Date.now(), + skipDequeueProcessing: true, + }); + } + + // closure invariant: every ckIndex member is a ckVtime member. The + // converse may transiently not hold (stale ckVtime entries GC on scan). + const members = await queue.redis.zrange(ckIndexKey, 0, -1); + for (const member of members) { + const tag = await queue.redis.zscore(ckVtimeKey, member); + expect( + tag, + `step ${step} (${opName}): ${member} in ckIndex but not ckVtime` + ).not.toBeNull(); + } + } + } finally { + await queue.quit(); + } + }); + + redisTest( + "flag off creates no vtime keys and matches head-timestamp order", + async ({ redisContainer }) => { + // ckVirtualTimeScheduling ABSENT: the off path calls the pre-existing + // command names (enqueueMessage*CkTracked, dequeueMessagesFromCkQueueTracked, + // nackMessageCkTracked) whose defineCommand script text this feature never + // edited, so the stronger same-script-SHA guarantee holds by construction. + // What a test CAN observe is asserted here: no vtime state is ever created, + // and the dequeue order is head-timestamp (age) order, matching the + // pre-existing ckIndex.test.ts expectation. + const queue = createQueue(redisContainer, null); + try { + const t0 = Date.now() - 100_000; + + // 3 variants with distinct head ages: old < mid < new, 3 messages each. + const heads: Record = { + old: t0, + mid: t0 + 10_000, + new: t0 + 20_000, + }; + for (const [ck, head] of Object.entries(heads)) { + for (let i = 0; i < 3; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-${ck}-${i}`, + concurrencyKey: ck, + timestamp: head + i, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + } + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + // The off-path command serves at most one message per variant per call, + // visiting variants in ckIndex (head-timestamp) order: oldest head first. + const first = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 5); + expect(first.map((m) => m.message.concurrencyKey)).toEqual(["old", "mid", "new"]); + + // nack old's head (immediate retry), ack the rest + const nackedId = first[0]!.messageId; + await queue.nackMessage({ + orgId: authenticatedEnvDev.organization.id, + messageId: nackedId, + retryAt: Date.now(), + skipDequeueProcessing: true, + }); + for (const m of first.slice(1)) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + + // Two more batched calls drain the original heads in age order each time. + for (let call = 0; call < 2; call++) { + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 5); + expect(messages.map((m) => m.message.concurrencyKey)).toEqual(["old", "mid", "new"]); + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + + // Only the nacked message remains; it is re-served. + const last = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 5); + expect(last.length).toBe(1); + expect(last[0]!.messageId).toBe(nackedId); + expect(last[0]!.message.concurrencyKey).toBe("old"); + + // After the whole mixed sequence (enqueues, batched dequeues, a nack, + // acks, one message still in flight so the keyspace is non-empty) no + // vtime state exists at all: no :ckVtime, no :ckVtimeFloor. The KEYS + // scan is safe here because redisTest runs flushall before each test, + // so the DB only holds this test's keys. + const allKeys = await queue.redis.keys("*"); + expect(allKeys.length).toBeGreaterThan(0); + expect(allKeys.filter((k) => k.includes("ckVtime"))).toEqual([]); + + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, nackedId, { + skipDequeueProcessing: true, + }); + } finally { + await queue.quit(); + } + } + ); +}); diff --git a/internal-packages/run-engine/src/run-queue/tests/ckVtimeConcurrency.test.ts b/internal-packages/run-engine/src/run-queue/tests/ckVtimeConcurrency.test.ts new file mode 100644 index 00000000000..9edbd772839 --- /dev/null +++ b/internal-packages/run-engine/src/run-queue/tests/ckVtimeConcurrency.test.ts @@ -0,0 +1,375 @@ +import { createRedisClient } from "@internal/redis"; +import { redisTest } from "@internal/testcontainers"; +import { trace } from "@internal/tracing"; +import { Logger } from "@trigger.dev/core/logger"; +import { Decimal } from "@trigger.dev/database"; +import { setTimeout as sleep } from "node:timers/promises"; +import { describe } from "node:test"; +import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js"; +import { RunQueue } from "../index.js"; +import { RunQueueFullKeyProducer } from "../keyProducer.js"; +import type { InputPayload } from "../types.js"; + +// Multi-consumer / multi-shard correctness for CK virtual-time scheduling, plus +// an op-count budget pinning the per-dequeue overhead of the vtime path. +// +// The correctness argument for concurrent consumers is that every ckVtime / +// ckIndex mutation happens inside a single Lua script and Redis serialises +// scripts. These tests check the scripts do not assume any cross-call state: +// two RunQueue instances hammering the same keyspace must still serve every +// message exactly once, never rewind a tag, and leave the vtime state clean. + +const testOptions = { + name: "rq", + tracer: trace.getTracer("rq"), + workers: 1, + defaultEnvConcurrency: 25, + logger: new Logger("RunQueue", "warn"), + 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: 10, + concurrencyLimitBurstFactor: new Decimal(2.0), + project: { id: "p1234" }, + organization: { id: "o1234" }, +}; + +function createQueue(redisContainer: any, keyPrefix: string, vtimeEnabled: boolean) { + return new RunQueue({ + ...testOptions, + // These tests drive every dequeue themselves (testDequeueFromMasterQueue + + // skipDequeueProcessing). The ONLY concurrency is the explicit consumer + // loops below, so the autonomous master-queue consumers and the background + // worker must stay off in every instance. + masterQueueConsumersDisabled: true, + workerOptions: { disabled: true }, + ckVirtualTimeScheduling: { + enabled: vtimeEnabled, + }, + queueSelectionStrategy: new FairQueueSelectionStrategy({ + redis: { + keyPrefix, + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + keys: testOptions.keys, + }), + redis: { + keyPrefix, + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + }); +} + +function makeMessage(overrides: Partial = {}): InputPayload { + return { + runId: "r1", + taskIdentifier: "task/my-task", + orgId: "o1234", + projectId: "p1234", + environmentId: "e1234", + environmentType: "DEVELOPMENT", + queue: "task/my-task", + timestamp: Date.now(), + attempt: 0, + ...overrides, + }; +} + +// The ckVtime/ckIndex member for a variant is the fully-qualified variant queue +// key (org:proj:env:queue:...:ck:), which is exactly what queueKey() produces. +function variantName(ck: string): string { + return testOptions.keys.queueKey(authenticatedEnvDev, "task/my-task", ck); +} + +vi.setConfig({ testTimeout: 120_000 }); + +describe("CK virtual-time concurrency and op-count budget", () => { + redisTest("two consumers, one base queue, no corruption", async ({ redisContainer }) => { + const keyPrefix = "rq15:"; + // one instance for enqueues, two more (same Redis, same key prefix) as consumers + const producer = createQueue(redisContainer, keyPrefix, true); + const consumerA = createQueue(redisContainer, keyPrefix, true); + const consumerB = createQueue(redisContainer, keyPrefix, true); + try { + const t0 = Date.now() - 100_000; + const cks = ["a", "b", "c", "d", "e", "f"]; + const perKey = 30; + + const enqueuedIds = new Set(); + for (const ck of cks) { + for (let i = 0; i < perKey; i++) { + const runId = `r-${ck}-${i}`; + enqueuedIds.add(runId); + await producer.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId, concurrencyKey: ck, timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + } + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(variantName("a")); + const ckVtimeFloorKey = testOptions.keys.ckVtimeFloorKeyFromQueue(variantName("a")); + + // shared across both consumer loops: messageId -> times served + const serveCounts = new Map(); + const floorSamples: number[] = []; + let floorRewind: { consumer: string; prev: number; next: number } | undefined; + + const runConsumer = async (name: string, queue: RunQueue) => { + let prevFloor = 0; + let iterations = 0; + while (serveCounts.size < enqueuedIds.size) { + iterations++; + if (iterations > 600) { + throw new Error( + `consumer ${name}: iteration cap hit with ${serveCounts.size}/${enqueuedIds.size} unique messages served` + ); + } + + const messages = await queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 5); + + // record serves immediately, so the exactly-once bookkeeping covers + // messages currently held by the other consumer too + for (const m of messages) { + serveCounts.set(m.messageId, (serveCounts.get(m.messageId) ?? 0) + 1); + } + + if (messages.length === 0) { + // nothing servable right now (the other consumer holds the slots); + // yield so its hold can elapse + await sleep(2); + } else { + // short hold before acking, so the two loops genuinely overlap + await sleep(3); + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + + // sample the floor between iterations: it must never decrease + const floor = Number((await queue.redis.get(ckVtimeFloorKey)) ?? "0"); + if (floor < prevFloor && !floorRewind) { + floorRewind = { consumer: name, prev: prevFloor, next: floor }; + } + prevFloor = floor; + floorSamples.push(floor); + } + }; + + await Promise.all([runConsumer("A", consumerA), runConsumer("B", consumerB)]); + + // exactly once: the union of served IDs equals the enqueued set, no duplicates + const duplicates = [...serveCounts.entries()].filter(([, count]) => count > 1); + expect(duplicates).toEqual([]); + expect(serveCounts.size).toBe(enqueuedIds.size); + expect(new Set(serveCounts.keys())).toEqual(enqueuedIds); + + // the floor never rewound in either consumer's sample sequence + expect(floorRewind).toBeUndefined(); + + // after drain: every variant was GC'd from ckVtime.. + expect(await consumerA.redis.zcard(ckVtimeKey)).toBe(0); + // ..and the floor sits at the max it ever reached + const finalFloor = Number((await consumerA.redis.get(ckVtimeFloorKey)) ?? "0"); + expect(finalFloor).toBe(Math.max(finalFloor, ...floorSamples)); + } finally { + await producer.quit(); + await consumerA.quit(); + await consumerB.quit(); + } + }); + + redisTest("concurrent enqueue during dequeue cannot rewind a tag", async ({ redisContainer }) => { + const queue = createQueue(redisContainer, "rq16:", true); + try { + const t0 = Date.now() - 100_000; + + // hot backlog large enough that it never drains (so it is never GC'd and + // re-registered, keeping the ZSCORE comparison meaningful), plus a + // competitor key so hot is not the only candidate + for (let i = 0; i < 12; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: `r-hot-${i}`, concurrencyKey: "hot", timestamp: t0 + i }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + for (let i = 0; i < 30; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-cold-${i}`, + concurrencyKey: "cold", + timestamp: t0 + i, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + + const hotVariant = variantName("hot"); + const ckVtimeKey = testOptions.keys.ckVtimeKeyFromQueue(hotVariant); + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + // enqueue registered hot at the initial floor + let prevTag = Number(await queue.redis.zscore(ckVtimeKey, hotVariant)); + expect(prevTag).toBe(0); + + let extra = 0; + for (let round = 0; round < 12; round++) { + // enqueues on the hot key racing a dequeue batch: the enqueue script's + // ZADD NX registration must never rewind the tag the dequeue script is + // advancing (advance-only writes) + const [, messages] = await Promise.all([ + (async () => { + for (let j = 0; j < 2; j++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-hot-extra-${extra++}`, + concurrencyKey: "hot", + timestamp: t0 + 1000 + round, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + })(), + queue.testDequeueFromMasterQueue(shard, authenticatedEnvDev.id, 3), + ]); + + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + + const tag = await queue.redis.zscore(ckVtimeKey, hotVariant); + // never drained, so never GC'd + expect(tag, `round ${round}: hot variant missing from ckVtime`).not.toBeNull(); + expect(Number(tag), `round ${round}: tag rewound`).toBeGreaterThanOrEqual(prevTag); + prevTag = Number(tag); + } + + // hot was actually served along the way (the invariant wasn't vacuous) + expect(prevTag).toBeGreaterThan(0); + } finally { + await queue.quit(); + } + }); + + redisTest("op-count budget: vtime dequeue overhead is bounded", async ({ redisContainer }) => { + const maxCount = 5; + const dequeueCalls = 50; + const cks = ["a", "b", "c", "d", "e", "f"]; + const perKey = 30; + + // second plain ioredis client (no key prefix) for CONFIG RESETSTAT / INFO. + // Redis command stats are server-wide, so each phase resets them after its + // enqueues and reads them right after its 50th dequeue call. + const statsClient = createRedisClient({ + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }); + + // Runs one phase: identical data under a fresh keyspace, then 50 identical + // dequeue calls (ack immediately, so env concurrency never gates a serve + // and both phases fully drain the same 180 messages inside the window). + const runPhase = async (keyPrefix: string, vtimeEnabled: boolean) => { + const queue = createQueue(redisContainer, keyPrefix, vtimeEnabled); + try { + const t0 = Date.now() - 100_000; + for (const ck of cks) { + for (let i = 0; i < perKey; i++) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r-${ck}-${i}`, + concurrencyKey: ck, + timestamp: t0 + i, + }), + workerQueue: authenticatedEnvDev.id, + skipDequeueProcessing: true, + }); + } + } + + const shard = testOptions.keys.masterQueueShardForEnvironment(authenticatedEnvDev.id, 2); + + await statsClient.call("CONFIG", "RESETSTAT"); + + let served = 0; + for (let call = 0; call < dequeueCalls; call++) { + const messages = await queue.testDequeueFromMasterQueue( + shard, + authenticatedEnvDev.id, + maxCount + ); + served += messages.length; + for (const m of messages) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, m.messageId, { + skipDequeueProcessing: true, + }); + } + } + + const info = await statsClient.info("commandstats"); + return { served, totalCalls: totalCommandCalls(info) }; + } finally { + await queue.quit(); + } + }; + + try { + const off = await runPhase("rq17off:", false); + const on = await runPhase("rq17on:", true); + + // both phases did identical work: the full 180 messages served and acked + expect(off.served).toBe(cks.length * perKey); + expect(on.served).toBe(cks.length * perKey); + + // Per dequeue call the vtime path adds at worst 7 fixed ops: GET floor, + // ZRANGE min, ZRANGE window, the pass-2 ZRANGEBYSCORE, SET floor, + // EXISTS ckVtime, EXPIRE ckVtime — plus per serve one ZSCORE and one ZADD. + const budget = dequeueCalls * (7 + 2 * maxCount); + expect( + on.totalCalls, + `on_total ${on.totalCalls} exceeds off_total ${off.totalCalls} + budget ${budget}` + ).toBeLessThanOrEqual(off.totalCalls + budget); + } finally { + await statsClient.quit(); + } + }); +}); + +// Sums calls= across every cmdstat_ line of INFO commandstats. Includes +// commands executed from inside Lua scripts, which is exactly what we want: +// the vtime overhead lives in the dequeue script body. +function totalCommandCalls(info: string): number { + let total = 0; + for (const line of info.split("\n")) { + const match = line.match(/^cmdstat_[^:]+:calls=(\d+)/); + if (match) { + total += Number(match[1]); + } + } + return total; +} diff --git a/internal-packages/run-engine/src/run-queue/tests/ckVtimeFairness.test.ts b/internal-packages/run-engine/src/run-queue/tests/ckVtimeFairness.test.ts new file mode 100644 index 00000000000..622b2cdd247 --- /dev/null +++ b/internal-packages/run-engine/src/run-queue/tests/ckVtimeFairness.test.ts @@ -0,0 +1,531 @@ +import { redisTest } from "@internal/testcontainers"; +import { trace } from "@internal/tracing"; +import { Logger } from "@trigger.dev/core/logger"; +import { Decimal } from "@trigger.dev/database"; +import { appendFileSync } from "node:fs"; +import { describe } from "node:test"; +import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js"; +import { RunQueue } from "../index.js"; +import { RunQueueFullKeyProducer } from "../keyProducer.js"; +import type { InputPayload } from "../types.js"; + +// Fairness scenarios driven through the REAL batched dequeue path (maxCount 10), +// closing the spike's maxCount=1 fidelity gap. Every assertion is a ratio between +// a flag-ON and a flag-OFF run of the same scenario (identical enqueue order and +// timestamps), so the tests are stable in CI. +// +// Scenario shapes are ported from the throwaway fairness spike (ckScenarios.ts / +// capsFairness.bench.test.ts): the message counts and head-age structure are +// copied as values, nothing is imported from the spike. +// +// Harness: a deterministic step loop. All messages are enqueued before step 0 +// with explicit past timestamps (a pre-existing backlog), so every message's +// logical arrival step is 0 and its wait is simply the step it was served at. +// Each step makes one dequeue call with maxCount 10, records the serves, then +// acks in-flight messages whose logical hold has elapsed (servedAt + hold <= +// step), which is how the env concurrency contends across steps. No wall-clock +// sleeps and no randomness anywhere. + +const testOptions = { + name: "rq", + tracer: trace.getTracer("rq"), + workers: 1, + defaultEnvConcurrency: 25, + logger: new Logger("RunQueue", "warn"), + 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: 10, + concurrencyLimitBurstFactor: new Decimal(2.0), + project: { id: "p1234" }, + organization: { id: "o1234" }, +}; + +function createQueue(redisContainer: any, keyPrefix: string, vtimeEnabled: boolean) { + return new RunQueue({ + ...testOptions, + // The step loop drives every op itself (testDequeueFromMasterQueue + + // skipDequeueProcessing), so the autonomous master-queue consumers and the + // background worker must not race it. + masterQueueConsumersDisabled: true, + workerOptions: { disabled: true }, + ckVirtualTimeScheduling: { + enabled: vtimeEnabled, + }, + queueSelectionStrategy: new FairQueueSelectionStrategy({ + redis: { + keyPrefix, + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + keys: testOptions.keys, + }), + redis: { + keyPrefix, + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + }); +} + +function makeMessage(overrides: Partial = {}): InputPayload { + return { + runId: "r1", + taskIdentifier: "task/my-task", + orgId: "o1234", + projectId: "p1234", + environmentId: "e1234", + environmentType: "DEVELOPMENT", + queue: "task/my-task", + timestamp: Date.now(), + attempt: 0, + ...overrides, + }; +} + +type ScenarioMessage = { runId: string; ck: string; timestamp: number }; + +type Scenario = { + name: string; + messages: ScenarioMessage[]; + // Effective env concurrency for the run (burst factor is pinned to 1.0). + // This is the contention knob: it caps how many serves fit in one dequeue + // call (actualMaxCount = min(maxCount, available env capacity)). + envConcurrencyLimit: number; + // Logical hold: a served message occupies its env slot until the end of + // step servedAt + holdSteps, when it is acked. + holdSteps: number; + // Safety cap so a work-conservation bug fails the count assertions instead + // of hanging the test. + maxSteps: number; +}; + +type ServeRecord = { step: number; ck: string; messageId: string }; + +type ScenarioResult = { + serves: ServeRecord[]; + // step at which the last message was served + drainStep: number; + // serves that happened in steps where >= 2 keys still had queued backlog + contentionServes: { total: number; byCk: Map }; +}; + +async function runScenario( + redisContainer: any, + scenario: Scenario, + vtimeEnabled: boolean +): Promise { + // Separate key prefix per run: the ON and OFF runs of a scenario share one + // Redis container but never share state. + const keyPrefix = `runqueue:test:${scenario.name}:${vtimeEnabled ? "on" : "off"}:`; + const queue = createQueue(redisContainer, keyPrefix, vtimeEnabled); + + try { + const env = { + ...authenticatedEnvDev, + maximumConcurrencyLimit: scenario.envConcurrencyLimit, + concurrencyLimitBurstFactor: new Decimal(1), + }; + await queue.updateEnvConcurrencyLimits(env); + + for (const msg of scenario.messages) { + await queue.enqueueMessage({ + env, + message: makeMessage({ + runId: msg.runId, + concurrencyKey: msg.ck, + timestamp: msg.timestamp, + }), + workerQueue: env.id, + skipDequeueProcessing: true, + }); + } + + const shard = testOptions.keys.masterQueueShardForEnvironment(env.id, 2); + const total = scenario.messages.length; + + const remaining = new Map(); + for (const m of scenario.messages) { + remaining.set(m.ck, (remaining.get(m.ck) ?? 0) + 1); + } + + const serves: ServeRecord[] = []; + const inFlight: { messageId: string; servedAtStep: number }[] = []; + const contentionServes = { total: 0, byCk: new Map() }; + let drainStep = -1; + + for (let step = 0; step < scenario.maxSteps && serves.length < total; step++) { + // evaluated before the dequeue: does this step have cross-key contention? + let keysWithBacklog = 0; + for (const count of remaining.values()) { + if (count > 0) keysWithBacklog++; + } + + const messages = await queue.testDequeueFromMasterQueue(shard, env.id, 10); + + for (const m of messages) { + const ck = m.message.concurrencyKey ?? ""; + serves.push({ step, ck, messageId: m.messageId }); + remaining.set(ck, (remaining.get(ck) ?? 0) - 1); + inFlight.push({ messageId: m.messageId, servedAtStep: step }); + if (keysWithBacklog >= 2) { + contentionServes.total++; + contentionServes.byCk.set(ck, (contentionServes.byCk.get(ck) ?? 0) + 1); + } + if (serves.length === total) { + drainStep = step; + } + } + + // release env/queue concurrency for serves whose hold has elapsed + for (let i = inFlight.length - 1; i >= 0; i--) { + const entry = inFlight[i]!; + if (entry.servedAtStep + scenario.holdSteps <= step) { + await queue.acknowledgeMessage(authenticatedEnvDev.organization.id, entry.messageId, { + skipDequeueProcessing: true, + }); + inFlight.splice(i, 1); + } + } + } + + return { serves, drainStep, contentionServes }; + } finally { + await queue.quit(); + } +} + +// Wait per message = serve step - arrival step, and arrival is step 0 for the +// whole pre-enqueued backlog, so the wait is just the serve step. +function meanWait(result: ScenarioResult, matches: (ck: string) => boolean): number { + const waits = result.serves.filter((s) => matches(s.ck)).map((s) => s.step); + expect(waits.length).toBeGreaterThan(0); + return waits.reduce((a, b) => a + b, 0) / waits.length; +} + +function firstServeStep(result: ScenarioResult, matches: (ck: string) => boolean): number { + const first = result.serves.find((s) => matches(s.ck)); + expect(first).toBeDefined(); + return first!.step; +} + +// No loss and no double-serve, in both runs. +function assertConservation(scenario: Scenario, on: ScenarioResult, off: ScenarioResult) { + expect(on.serves.length).toBe(scenario.messages.length); + expect(off.serves.length).toBe(scenario.messages.length); + expect(new Set(on.serves.map((s) => s.messageId)).size).toBe(scenario.messages.length); + expect(new Set(off.serves.map((s) => s.messageId)).size).toBe(scenario.messages.length); +} + +function debugLog(name: string, data: Record) { + if (process.env.CK_FAIRNESS_DEBUG) { + // the test reporter swallows console output, so append to a file instead + appendFileSync( + process.env.CK_FAIRNESS_DEBUG, + `[ckVtimeFairness] ${name} ${JSON.stringify(data)}\n` + ); + } +} + +vi.setConfig({ testTimeout: 120_000 }); + +describe("CK virtual-time fairness on the real batched dequeue path", () => { + // ckSkew (spike shape): one heavy key with a 120-message backlog on an old + // shared head, 4 light keys with 10 messages each on later heads. + // + // Contention regime: env limit 1, hold 3. The batched dequeue serves at most + // one message per variant per call, so at env limit 4 (the spike's driver + // setting) a single heavy key cannot crowd out 4 light keys at all: both + // flags serve every light key each round and the ON/OFF ratio sits near 1. + // The head-age starvation the spike measured appears on this path when env + // capacity serializes the calls (limit 1): flag OFF then always picks the + // globally oldest head, which is heavy for its whole backlog. + redisTest( + "ckSkew: light keys stop waiting behind the heavy backlog", + async ({ redisContainer }) => { + const t0 = Date.now() - 500_000; + const messages: ScenarioMessage[] = []; + for (let i = 0; i < 120; i++) { + messages.push({ runId: `heavy-${i}`, ck: "heavy", timestamp: t0 }); + } + for (let i = 0; i < 10; i++) { + for (let k = 0; k < 4; k++) { + messages.push({ + runId: `light${k}-${i}`, + ck: `light${k}`, + timestamp: t0 + 10_000 + i * 4 + k, + }); + } + } + const scenario: Scenario = { + name: "ckSkew", + messages, + envConcurrencyLimit: 1, + holdSteps: 3, + maxSteps: 1_000, + }; + + const on = await runScenario(redisContainer, scenario, true); + const off = await runScenario(redisContainer, scenario, false); + + assertConservation(scenario, on, off); + + const isLight = (ck: string) => ck.startsWith("light"); + const onWait = meanWait(on, isLight); + const offWait = meanWait(off, isLight); + debugLog("ckSkew", { onWait, offWait, ratio: onWait / offWait }); + + // Heavy's wait may rise under the fair order; that is expected and not + // asserted down. + expect(onWait).toBeLessThanOrEqual(0.3 * offWait); + } + ); + + // ckTrickle (spike shape): one bulk key with a 120-message backlog on an old + // shared head, two trickle keys with 15 messages each on later heads. Same + // serialized contention regime as ckSkew, same assertion. + redisTest( + "ckTrickle: trickle keys stop waiting behind the bulk backlog", + async ({ redisContainer }) => { + const t0 = Date.now() - 500_000; + const messages: ScenarioMessage[] = []; + for (let i = 0; i < 120; i++) { + messages.push({ runId: `bulk-${i}`, ck: "bulk", timestamp: t0 }); + } + for (let i = 0; i < 15; i++) { + for (let k = 0; k < 2; k++) { + messages.push({ + runId: `trickle${k}-${i}`, + ck: `trickle${k}`, + timestamp: t0 + 10_000 + i * 2 + k, + }); + } + } + const scenario: Scenario = { + name: "ckTrickle", + messages, + envConcurrencyLimit: 1, + holdSteps: 3, + maxSteps: 1_000, + }; + + const on = await runScenario(redisContainer, scenario, true); + const off = await runScenario(redisContainer, scenario, false); + + assertConservation(scenario, on, off); + + const isTrickle = (ck: string) => ck.startsWith("trickle"); + const onWait = meanWait(on, isTrickle); + const offWait = meanWait(off, isTrickle); + debugLog("ckTrickle", { onWait, offWait, ratio: onWait / offWait }); + + expect(onWait).toBeLessThanOrEqual(0.3 * offWait); + } + ); + + // ckSybil (spike shape, the case per-key caps cannot fix): 20 attacker keys + // with 8 messages each, all on older heads, and 1 light key with 10 newer + // messages. 21 variants against a batch of 10 exercises the batched path + // properly: flag OFF walks the age order and only reaches the light key when + // the attackers are nearly drained; flag ON serves the light key from the + // floor on its first fair round. + redisTest("ckSybil: many attacker keys cannot starve a light key", async ({ redisContainer }) => { + const t0 = Date.now() - 500_000; + const messages: ScenarioMessage[] = []; + for (let i = 0; i < 8; i++) { + for (let k = 0; k < 20; k++) { + const ck = `att${String(k).padStart(2, "0")}`; + messages.push({ runId: `${ck}-${i}`, ck, timestamp: t0 + i * 20 + k }); + } + } + for (let i = 0; i < 10; i++) { + messages.push({ runId: `light-${i}`, ck: "light", timestamp: t0 + 50_000 + i }); + } + const scenario: Scenario = { + name: "ckSybil", + messages, + envConcurrencyLimit: 25, + holdSteps: 3, + maxSteps: 300, + }; + + const on = await runScenario(redisContainer, scenario, true); + const off = await runScenario(redisContainer, scenario, false); + + assertConservation(scenario, on, off); + + const isLight = (ck: string) => ck === "light"; + + // Reachability at the floor: enqueue registered the light key at the + // floor, so it is served within the first 3 steps even though 20 attacker + // variants sit ahead of it in age order. + const onFirstServe = firstServeStep(on, isLight); + expect(onFirstServe).toBeLessThanOrEqual(2); + + const onWait = meanWait(on, isLight); + const offWait = meanWait(off, isLight); + + // Contention-window share (directional, per the spike's confounding + // caveat; the wait ratio is the headline): over the steps where >= 2 keys + // had queued backlog, light's served fraction is at least half its fair + // share of 1/21. + const lightContentionServes = on.contentionServes.byCk.get("light") ?? 0; + const lightShare = lightContentionServes / on.contentionServes.total; + + debugLog("ckSybil", { + onWait, + offWait, + ratio: onWait / offWait, + onFirstServe, + lightShare, + fairShare: 1 / 21, + }); + + expect(onWait).toBeLessThanOrEqual(0.7 * offWait); + expect(lightShare).toBeGreaterThanOrEqual(0.5 * (1 / 21)); + }); + + // ckBalanced (spike shape, no-harm check): 4 symmetric keys with 25 messages + // each. The fair order must not make the symmetric case worse. + redisTest( + "ckBalanced: fair order does not hurt the symmetric case", + async ({ redisContainer }) => { + const t0 = Date.now() - 500_000; + const cks = ["bal0", "bal1", "bal2", "bal3"]; + const messages: ScenarioMessage[] = []; + for (let i = 0; i < 25; i++) { + for (let k = 0; k < cks.length; k++) { + messages.push({ + runId: `${cks[k]}-${i}`, + ck: cks[k]!, + timestamp: t0 + i * 4 + k, + }); + } + } + const scenario: Scenario = { + name: "ckBalanced", + messages, + envConcurrencyLimit: 4, + holdSteps: 3, + maxSteps: 500, + }; + + const on = await runScenario(redisContainer, scenario, true); + const off = await runScenario(redisContainer, scenario, false); + + assertConservation(scenario, on, off); + + const maxPerKeyMeanWait = (result: ScenarioResult) => + Math.max(...cks.map((ck) => meanWait(result, (c) => c === ck))); + + const onMax = maxPerKeyMeanWait(on); + const offMax = maxPerKeyMeanWait(off); + debugLog("ckBalanced", { onMax, offMax, ratio: onMax / offMax }); + + // Observed ratio is 1.0 (the fair order is neutral on the symmetric case), + // so allow only modest headroom rather than the original 1.25. + expect(onMax).toBeLessThanOrEqual(1.1 * offMax); + } + ); + + // ckManyKeys (sharding coverage): cardinality ABOVE the pass-1 fair window. + // The batched dequeue uses maxCount 10, so window = actualMaxCount * 3 = 30. + // With ~60 attacker keys (all on the same old head) plus 1 light key on a + // newer head, 61 variants sit above the 30-wide pass-1 ZRANGE window, so no + // single fair pass can even see every key. The property to hold is that this + // does NOT permanently starve the light key: as attackers advance their tags + // out of the bottom of the window, the light key (still at the floor) rises + // into it and gets served, and every message drains exactly once. A bounded + // first-serve delay is fine; permanent starvation or a stuck drain is not. + redisTest( + "ckManyKeys: light key is not starved when cardinality exceeds the fair window", + async ({ redisContainer }) => { + const t0 = Date.now() - 500_000; + const messages: ScenarioMessage[] = []; + const attackerCount = 60; + for (let i = 0; i < 8; i++) { + for (let k = 0; k < attackerCount; k++) { + const ck = `att${String(k).padStart(2, "0")}`; + // All attackers share the same old head timestamp (tied heads). + messages.push({ runId: `${ck}-${i}`, ck, timestamp: t0 }); + } + } + for (let i = 0; i < 10; i++) { + messages.push({ runId: `light-${i}`, ck: "light", timestamp: t0 + 50_000 + i }); + } + const scenario: Scenario = { + name: "ckManyKeys", + messages, + envConcurrencyLimit: 25, + holdSteps: 3, + maxSteps: 1_000, + }; + + const on = await runScenario(redisContainer, scenario, true); + const off = await runScenario(redisContainer, scenario, false); + + // No loss and no double-serve in either run: the run terminates and every + // message (attackers + light) is served exactly once within maxSteps. + assertConservation(scenario, on, off); + + const isLight = (ck: string) => ck === "light"; + + // The light key IS eventually served (no permanent starvation) in both + // runs, and drains fully. + const onFirstServe = firstServeStep(on, isLight); + const offFirstServe = firstServeStep(off, isLight); + expect(on.drainStep).toBeGreaterThanOrEqual(0); + expect(off.drainStep).toBeGreaterThanOrEqual(0); + + debugLog("ckManyKeys", { + variants: attackerCount + 1, + onFirstServe, + offFirstServe, + onDrainStep: on.drainStep, + offDrainStep: off.drainStep, + }); + } + ); + + // ckHeavyIdle (spike shape, work conservation): a single key with 60 + // messages and nothing else contending. Any extra step to drain under the + // fair order is a work-conservation bug, so the step counts must be exactly + // equal. + redisTest( + "ckHeavyIdle: a lone key drains in exactly the same steps", + async ({ redisContainer }) => { + const t0 = Date.now() - 500_000; + const messages: ScenarioMessage[] = []; + for (let i = 0; i < 60; i++) { + messages.push({ runId: `solo-${i}`, ck: "solo", timestamp: t0 + i }); + } + const scenario: Scenario = { + name: "ckHeavyIdle", + messages, + envConcurrencyLimit: 25, + holdSteps: 3, + maxSteps: 300, + }; + + const on = await runScenario(redisContainer, scenario, true); + const off = await runScenario(redisContainer, scenario, false); + + assertConservation(scenario, on, off); + + debugLog("ckHeavyIdle", { onDrainStep: on.drainStep, offDrainStep: off.drainStep }); + + expect(on.drainStep).toBeGreaterThanOrEqual(0); + expect(on.drainStep).toBe(off.drainStep); + } + ); +}); diff --git a/internal-packages/run-engine/src/run-queue/tests/keyProducer.test.ts b/internal-packages/run-engine/src/run-queue/tests/keyProducer.test.ts index 3e31085d678..1f61ed1cfc7 100644 --- a/internal-packages/run-engine/src/run-queue/tests/keyProducer.test.ts +++ b/internal-packages/run-engine/src/run-queue/tests/keyProducer.test.ts @@ -432,4 +432,19 @@ describe("KeyProducer", () => { "{org:o1234}:proj:p1234:env:e1234:queue:task/foo:ck:*" ); }); + + it("produces ckVtime keys from a CK variant queue name", () => { + const keyProducer = new RunQueueFullKeyProducer(); + const q = "{org:o1}:proj:p1:env:e1:queue:task/my-task:ck:tenant-a"; + expect(keyProducer.ckVtimeKeyFromQueue(q)).toBe( + "{org:o1}:proj:p1:env:e1:queue:task/my-task:ckVtime" + ); + expect(keyProducer.ckVtimeFloorKeyFromQueue(q)).toBe( + "{org:o1}:proj:p1:env:e1:queue:task/my-task:ckVtimeFloor" + ); + // ck wildcard and base-queue inputs normalise the same way + expect(keyProducer.ckVtimeKeyFromQueue(q.replace(":ck:tenant-a", ":ck:*"))).toBe( + keyProducer.ckVtimeKeyFromQueue(q) + ); + }); }); diff --git a/internal-packages/run-engine/src/run-queue/types.ts b/internal-packages/run-engine/src/run-queue/types.ts index 8a7d3c93ec5..b82153bc99b 100644 --- a/internal-packages/run-engine/src/run-queue/types.ts +++ b/internal-packages/run-engine/src/run-queue/types.ts @@ -132,6 +132,8 @@ export interface RunQueueKeyProducer { // CK index methods ckIndexKeyFromQueue(queue: string): string; + ckVtimeKeyFromQueue(queue: string): string; + ckVtimeFloorKeyFromQueue(queue: string): string; baseQueueKeyFromQueue(queue: string): string; isCkWildcard(queue: string): boolean; toCkWildcard(queue: string): string;