From 8d88a51bda8bdd4f0a5da9fc46ddd24ce79882a1 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 10:36:57 +0100 Subject: [PATCH 01/13] feat(run-engine): per-concurrency-key limit override storage and methods Sparse ckLimits HASH at the base queue whose fields are the exact ck-variant queue names, plus engine methods to set (atomic cardinality cap, default 1000 per queue), remove, and read the overrides. The admit-path gate wiring follows. --- .../run-engine/src/run-queue/index.ts | 95 +++++++++++++++++++ .../run-engine/src/run-queue/keyProducer.ts | 15 +++ .../run-engine/src/run-queue/types.ts | 2 + 3 files changed, 112 insertions(+) diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index 6ab8d90aaec..8d8781eca65 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -257,6 +257,13 @@ export interface RunQueueMetricsEmitter { emitGauge(shardKey: string, fields: Record): void; } +export class RunQueueConcurrencyKeyLimitExceededError extends Error { + constructor(message: string) { + super(message); + this.name = "RunQueueConcurrencyKeyLimitExceededError"; + } +} + export type RunQueueOptions = { name: string; tracer: Tracer; @@ -316,6 +323,8 @@ export type RunQueueOptions = { * the total cap covering releases from builds without the mirror. */ gatesEnabled?: boolean; + /** Cap on per-concurrency-key limit overrides stored per queue. Default 1000. */ + maxConcurrencyKeyOverridesPerQueue?: number; workerOptions?: { pollIntervalMs?: number; immediatePollIntervalMs?: number; @@ -429,6 +438,7 @@ export class RunQueue { private queueSelectionStrategy: RunQueueSelectionStrategy; private shardCount: number; private counterTtlSeconds: number; + private maxConcurrencyKeyOverridesPerQueue: number; private abortController: AbortController; private worker: Worker; private workerQueueResolver: WorkerQueueResolver; @@ -439,6 +449,7 @@ export class RunQueue { constructor(public readonly options: RunQueueOptions) { this.shardCount = options.shardCount ?? 2; this.counterTtlSeconds = options.counterTtlSeconds ?? 86400; + this.maxConcurrencyKeyOverridesPerQueue = options.maxConcurrencyKeyOverridesPerQueue ?? 1000; this.retryOptions = options.retryOptions ?? defaultRetrySettings; this.redis = createRedisClient(options.redis, { onError: (error) => { @@ -633,6 +644,62 @@ export class RunQueue { return this.redis.scard(this.keys.queueGroupConcurrencyKey(env, queue)); } + /** + * Sets a per-concurrency-key limit override for a queue. The stored value is the + * raw requested limit; admit paths clamp to the environment limit at read time. + * Throws RunQueueConcurrencyKeyLimitExceededError when a NEW key would push the + * queue past maxConcurrencyKeyOverridesPerQueue (updates to existing keys always + * succeed). + */ + public async updateQueueConcurrencyKeyLimit( + env: MinimalAuthenticatedEnvironment, + queue: string, + concurrencyKey: string, + limit: number + ) { + const result = await this.redis.setQueueConcurrencyKeyLimit( + this.keys.queueCkLimitsKey(env, queue), + this.keys.queueKey(env, queue, concurrencyKey), + String(limit), + String(this.maxConcurrencyKeyOverridesPerQueue) + ); + + if (result === 0) { + throw new RunQueueConcurrencyKeyLimitExceededError( + `Cannot add a concurrency key override to queue ${queue}: the queue already has ${this.maxConcurrencyKeyOverridesPerQueue} overrides` + ); + } + } + + public async removeQueueConcurrencyKeyLimit( + env: MinimalAuthenticatedEnvironment, + queue: string, + concurrencyKey: string + ) { + return this.redis.hdel( + this.keys.queueCkLimitsKey(env, queue), + this.keys.queueKey(env, queue, concurrencyKey) + ); + } + + /** Returns the raw per-concurrency-key limit overrides for a queue, keyed by concurrency key value. */ + public async getQueueConcurrencyKeyLimits( + env: MinimalAuthenticatedEnvironment, + queue: string + ): Promise> { + const raw = await this.redis.hgetall(this.keys.queueCkLimitsKey(env, queue)); + + const limits: Record = {}; + for (const [variantName, value] of Object.entries(raw)) { + const ckIndex = variantName.indexOf(":ck:"); + if (ckIndex === -1) { + continue; + } + limits[variantName.slice(ckIndex + 4)] = Number(value); + } + return limits; + } + public async updateEnvConcurrencyLimits(env: MinimalAuthenticatedEnvironment) { await this.#callUpdateEnvironmentConcurrencyLimits({ envConcurrencyLimitKey: this.keys.envConcurrencyLimitKey(env), @@ -5879,6 +5946,26 @@ __gatesRelease(keyPrefix, redis.call('GET', messageKey), messageId) `, }); + this.redis.defineCommand("setQueueConcurrencyKeyLimit", { + numberOfKeys: 1, + lua: ` +local ckLimitsKey = KEYS[1] + +local fieldName = ARGV[1] +local limit = ARGV[2] +local maxFields = tonumber(ARGV[3]) + +if redis.call('HEXISTS', ckLimitsKey, fieldName) == 0 then + if redis.call('HLEN', ckLimitsKey) >= maxFields then + return 0 + end +end + +redis.call('HSET', ckLimitsKey, fieldName, limit) +return 1 +`, + }); + this.redis.defineCommand("updateEnvironmentConcurrencyLimits", { numberOfKeys: 2, lua: ` @@ -6240,6 +6327,14 @@ declare module "@internal/redis" { callback?: Callback ): Result; + setQueueConcurrencyKeyLimit( + ckLimitsKey: string, + fieldName: string, + limit: string, + maxFields: string, + callback?: Callback + ): Result; + updateEnvironmentConcurrencyLimits( // keys envConcurrencyLimitKey: string, diff --git a/internal-packages/run-engine/src/run-queue/keyProducer.ts b/internal-packages/run-engine/src/run-queue/keyProducer.ts index 98028f5af7b..7b997043244 100644 --- a/internal-packages/run-engine/src/run-queue/keyProducer.ts +++ b/internal-packages/run-engine/src/run-queue/keyProducer.ts @@ -26,6 +26,7 @@ const constants = { RUNNING_COUNTER_PART: "runningCounter", GROUP_CONCURRENCY_PART: "groupConcurrency", TOTAL_CONCURRENCY_LIMIT_PART: "totalConcurrency", + CK_LIMITS_PART: "ckLimits", } as const; export class RunQueueFullKeyProducer implements RunQueueKeyProducer { @@ -366,6 +367,20 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer { return `${this.baseQueueKeyFromQueue(queue)}:${constants.TOTAL_CONCURRENCY_LIMIT_PART}`; } + /** + * HASH of per-concurrency-key limit overrides for a queue. Lives at the base + * queue; each field is the EXACT full ck-variant queue name (the ckIndex ZSET + * member), so reads need no parsing, and values are the raw requested limits + * (readers clamp to the environment limit). + */ + queueCkLimitsKey(env: RunQueueKeyProducerEnvironment, queue: string): string { + return `${this.queueKey(env, queue)}:${constants.CK_LIMITS_PART}`; + } + + queueCkLimitsKeyFromQueue(queue: string): string { + return `${this.baseQueueKeyFromQueue(queue)}:${constants.CK_LIMITS_PART}`; + } + isCkWildcard(queue: string): boolean { return queue.endsWith(":ck:*"); } diff --git a/internal-packages/run-engine/src/run-queue/types.ts b/internal-packages/run-engine/src/run-queue/types.ts index 2cbfe40c775..b21a3f66368 100644 --- a/internal-packages/run-engine/src/run-queue/types.ts +++ b/internal-packages/run-engine/src/run-queue/types.ts @@ -110,6 +110,8 @@ export interface RunQueueKeyProducer { queueGroupConcurrencyKeyFromQueue(queue: string): string; queueTotalConcurrencyLimitKey(env: RunQueueKeyProducerEnvironment, queue: string): string; queueTotalConcurrencyLimitKeyFromQueue(queue: string): string; + queueCkLimitsKey(env: RunQueueKeyProducerEnvironment, queue: string): string; + queueCkLimitsKeyFromQueue(queue: string): string; //env oncurrency envCurrentConcurrencyKey(env: EnvDescriptor): string; From ce66aaccf785eaf92e1aebf82e12fe5113e18392 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 10:46:04 +0100 Subject: [PATCH 02/13] feat(run-engine): enforce per-key limit overrides at admit time The ck dequeue admit and both enqueue fast paths read the queue's ckLimits HASH for the variant being admitted and use the env-clamped override in place of the queue's per-key limit, behind the totalConcurrencyEnabled flag. Covered by tests for lowered and raised keys, removal, the cardinality cap, and flag-off behavior. --- .../run-engine/src/run-queue/index.ts | 38 ++- .../tests/concurrencyKeyOverrides.test.ts | 270 ++++++++++++++++++ 2 files changed, 304 insertions(+), 4 deletions(-) create mode 100644 internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index 8d8781eca65..b774732f5d3 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -2433,6 +2433,7 @@ export class RunQueue { const totalConcurrencyLimitKey = this.keys.queueTotalConcurrencyLimitKeyFromQueue( message.queue ); + const ckLimitsKey = this.keys.queueCkLimitsKeyFromQueue(message.queue); const totalConcurrencyEnabledArg = this.options.totalConcurrencyEnabled ? "1" : "0"; if (ttlInfo) { @@ -2456,6 +2457,7 @@ export class RunQueue { baseQueueKey, groupConcurrencyKey, totalConcurrencyLimitKey, + ckLimitsKey, // args queueName, messageId, @@ -2495,6 +2497,7 @@ export class RunQueue { baseQueueKey, groupConcurrencyKey, totalConcurrencyLimitKey, + ckLimitsKey, // args queueName, messageId, @@ -2782,6 +2785,7 @@ export class RunQueue { runningCounterKey, this.keys.queueGroupConcurrencyKeyFromQueue(ckWildcardQueue), this.keys.queueTotalConcurrencyLimitKeyFromQueue(ckWildcardQueue), + this.keys.queueCkLimitsKeyFromQueue(ckWildcardQueue), //args ckWildcardQueue, String(Date.now()), @@ -4070,7 +4074,7 @@ return __qmret(0) // *Tracked variants of dequeueMessageFromKey and the ack/nack/dlq/release/clear // scripts. this.redis.defineCommand("enqueueMessageCkTracked", { - numberOfKeys: 17, + numberOfKeys: 18, lua: ` local masterQueueKey = KEYS[1] local queueKey = KEYS[2] @@ -4092,6 +4096,7 @@ local baseQueueKey = KEYS[15] -- Total-cap keys (KEYS 16-17) local groupConcurrencyKey = KEYS[16] local totalConcurrencyLimitKey = KEYS[17] +local ckLimitsKey = KEYS[18] local queueName = ARGV[1] local messageId = ARGV[2] @@ -4129,6 +4134,12 @@ if enableFastPath == '1' then tonumber(redis.call('GET', queueConcurrencyLimitKey) or '1000000'), envLimit ) + if totalConcurrencyEnabled then + local perKeyOverride = redis.call('HGET', ckLimitsKey, queueName) + if perKeyOverride then + queueLimit = math.min(tonumber(perKeyOverride), envLimit) + end + end if queueCurrent < queueLimit then -- Total-cap gate: a fast-path admit consumes a group slot, so it must @@ -4240,7 +4251,7 @@ return __qmret(0) }); this.redis.defineCommand("enqueueMessageWithTtlCkTracked", { - numberOfKeys: 18, + numberOfKeys: 19, lua: ` local masterQueueKey = KEYS[1] local queueKey = KEYS[2] @@ -4263,6 +4274,7 @@ local baseQueueKey = KEYS[16] -- Total-cap keys (KEYS 17-18) local groupConcurrencyKey = KEYS[17] local totalConcurrencyLimitKey = KEYS[18] +local ckLimitsKey = KEYS[19] local queueName = ARGV[1] local messageId = ARGV[2] @@ -4302,6 +4314,12 @@ if enableFastPath == '1' then tonumber(redis.call('GET', queueConcurrencyLimitKey) or '1000000'), envLimit ) + if totalConcurrencyEnabled then + local perKeyOverride = redis.call('HGET', ckLimitsKey, queueName) + if perKeyOverride then + queueLimit = math.min(tonumber(perKeyOverride), envLimit) + end + end if queueCurrent < queueLimit then -- Total-cap gate: see enqueueMessageCkTracked. @@ -4926,7 +4944,7 @@ return results // (normal dequeue, TTL-expired, or stale-orphan path — all of which were // counted at enqueue time). this.redis.defineCommand("dequeueMessagesFromCkQueueTracked", { - numberOfKeys: 13, + numberOfKeys: 14, lua: ` local ckIndexKey = KEYS[1] local queueConcurrencyLimitKey = KEYS[2] @@ -4941,6 +4959,7 @@ local lengthCounterKey = KEYS[10] local runningCounterKey = KEYS[11] local groupConcurrencyKey = KEYS[12] local totalConcurrencyLimitKey = KEYS[13] +local ckLimitsKey = KEYS[14] local ckWildcardName = ARGV[1] local currentTime = tonumber(ARGV[2]) @@ -5030,7 +5049,15 @@ for _, ckQueueName in ipairs(ckQueues) do local ckConcurrencyKey = fullQueueKey .. ':currentConcurrency' local ckCurrentConcurrency = tonumber(redis.call('SCARD', ckConcurrencyKey) or '0') - if ckCurrentConcurrency < queueConcurrencyLimit then + local perKeyLimit = queueConcurrencyLimit + if totalConcurrencyEnabled then + local perKeyOverride = redis.call('HGET', ckLimitsKey, ckQueueName) + if perKeyOverride then + perKeyLimit = math.min(tonumber(perKeyOverride), envConcurrencyLimit) + end + end + + if ckCurrentConcurrency < perKeyLimit then local messages = redis.call('ZRANGEBYSCORE', fullQueueKey, '-inf', tostring(currentTime), 'WITHSCORES', 'LIMIT', 0, 1) if #messages >= 2 then @@ -6521,6 +6548,7 @@ declare module "@internal/redis" { baseQueueKey: string, groupConcurrencyKey: string, totalConcurrencyLimitKey: string, + ckLimitsKey: string, queueName: string, messageId: string, messageData: string, @@ -6558,6 +6586,7 @@ declare module "@internal/redis" { baseQueueKey: string, groupConcurrencyKey: string, totalConcurrencyLimitKey: string, + ckLimitsKey: string, queueName: string, messageId: string, messageData: string, @@ -6592,6 +6621,7 @@ declare module "@internal/redis" { runningCounterKey: string, groupConcurrencyKey: string, totalConcurrencyLimitKey: string, + ckLimitsKey: string, ckWildcardName: string, currentTime: string, defaultEnvConcurrencyLimit: string, diff --git a/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts b/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts new file mode 100644 index 00000000000..7995369bc3c --- /dev/null +++ b/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts @@ -0,0 +1,270 @@ +import { redisTest } from "@internal/testcontainers"; +import { trace } from "@internal/tracing"; +import { setTimeout } from "node:timers/promises"; +import { describe } from "vitest"; +import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js"; +import { RunQueue, RunQueueConcurrencyKeyLimitExceededError } from "../index.js"; +import { RunQueueFullKeyProducer } from "../keyProducer.js"; +import type { InputPayload } from "../types.js"; +import { Decimal } from "@trigger.dev/database"; + +const testOptions = { + name: "rq", + tracer: trace.getTracer("rq"), + workers: 1, + defaultEnvConcurrency: 25, + 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, totalConcurrencyEnabled: boolean, maxOverrides?: number) { + return new RunQueue({ + ...testOptions, + totalConcurrencyEnabled, + maxConcurrencyKeyOverridesPerQueue: maxOverrides, + 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, + }; +} + +async function waitFor(condition: () => Promise, timeoutMs = 20_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (await condition()) { + return true; + } + await setTimeout(250); + } + return condition(); +} + +vi.setConfig({ testTimeout: 60_000 }); + +describe("RunQueue per-concurrency-key limit overrides", () => { + redisTest( + "a lowered key is capped while other keys keep the queue limit", + async ({ redisContainer }) => { + const queue = createQueue(redisContainer, true); + try { + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 2); + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 1); + + const now = Date.now(); + const messages = [ + ["ck-a", "a0"], + ["ck-a", "a1"], + ["ck-b", "b0"], + ["ck-b", "b1"], + ] as const; + for (const [i, [ck, id]] of messages.entries()) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: id, concurrencyKey: ck, timestamp: now - 1000 + i }), + workerQueue: "main", + }); + } + + const settled = await waitFor(async () => { + const a = await queue.currentConcurrencyOfQueue( + authenticatedEnvDev, + "task/my-task", + "ck-a" + ); + const b = await queue.currentConcurrencyOfQueue( + authenticatedEnvDev, + "task/my-task", + "ck-b" + ); + return a === 1 && b === 2; + }); + expect(settled).toBe(true); + + await setTimeout(2000); + expect( + await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "task/my-task", "ck-a") + ).toBe(1); + expect(await queue.lengthOfQueue(authenticatedEnvDev, "task/my-task")).toBe(1); + } finally { + await queue.quit(); + } + } + ); + + redisTest("a raised key admits past the queue limit", async ({ redisContainer }) => { + const queue = createQueue(redisContainer, true); + try { + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 1); + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 3); + + const now = Date.now(); + for (const i of [0, 1, 2]) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `a${i}`, + concurrencyKey: "ck-a", + timestamp: now - 1000 + i, + }), + workerQueue: "main", + }); + } + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "b0", concurrencyKey: "ck-b", timestamp: now - 500 }), + workerQueue: "main", + }); + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ runId: "b1", concurrencyKey: "ck-b", timestamp: now - 499 }), + workerQueue: "main", + }); + + const settled = await waitFor(async () => { + const a = await queue.currentConcurrencyOfQueue( + authenticatedEnvDev, + "task/my-task", + "ck-a" + ); + const b = await queue.currentConcurrencyOfQueue( + authenticatedEnvDev, + "task/my-task", + "ck-b" + ); + return a === 3 && b === 1; + }); + expect(settled).toBe(true); + + await setTimeout(2000); + expect(await queue.lengthOfQueue(authenticatedEnvDev, "task/my-task")).toBe(1); + } finally { + await queue.quit(); + } + }); + + redisTest("removing an override restores the queue limit", async ({ redisContainer }) => { + const queue = createQueue(redisContainer, true); + try { + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 2); + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 1); + expect(await queue.getQueueConcurrencyKeyLimits(authenticatedEnvDev, "task/my-task")).toEqual( + { "ck-a": 1 } + ); + + await queue.removeQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a"); + expect(await queue.getQueueConcurrencyKeyLimits(authenticatedEnvDev, "task/my-task")).toEqual( + {} + ); + + const now = Date.now(); + for (const i of [0, 1]) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `a${i}`, + concurrencyKey: "ck-a", + timestamp: now - 1000 + i, + }), + workerQueue: "main", + }); + } + + const settled = await waitFor( + async () => + (await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "task/my-task", "ck-a")) === 2 + ); + expect(settled).toBe(true); + } finally { + await queue.quit(); + } + }); + + redisTest("the per-queue override count is capped", async ({ redisContainer }) => { + const queue = createQueue(redisContainer, true, 2); + try { + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 1); + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-b", 1); + + await expect( + queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-c", 1) + ).rejects.toThrow(RunQueueConcurrencyKeyLimitExceededError); + + /** Updates to existing keys always succeed at the cap. */ + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 4); + expect(await queue.getQueueConcurrencyKeyLimits(authenticatedEnvDev, "task/my-task")).toEqual( + { "ck-a": 4, "ck-b": 1 } + ); + } finally { + await queue.quit(); + } + }); + + redisTest("overrides are ignored when disabled", async ({ redisContainer }) => { + const queue = createQueue(redisContainer, false); + try { + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 2); + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 1); + + const now = Date.now(); + for (const i of [0, 1]) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `a${i}`, + concurrencyKey: "ck-a", + timestamp: now - 1000 + i, + }), + workerQueue: "main", + }); + } + + const settled = await waitFor( + async () => + (await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "task/my-task", "ck-a")) === 2 + ); + expect(settled).toBe(true); + } finally { + await queue.quit(); + } + }); +}); From efdbc726a84a8327a29c798faf542b4d5df3a826 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 10:54:53 +0100 Subject: [PATCH 03/13] feat(database): total-override bookkeeping and per-key override table Three nullable TaskQueue columns record when, by whom, and from what declared base the total concurrency limit was overridden, and a new TaskQueueConcurrencyKeyOverride child table stores per-key limit overrides, unique per queue and key and cascading with the queue. --- .../migration.sql | 24 +++++++++++++++ .../database/prisma/schema.prisma | 30 ++++++++++++++++++- 2 files changed, 53 insertions(+), 1 deletion(-) create mode 100644 internal-packages/database/prisma/migrations/20260829150000_add_concurrency_overrides/migration.sql diff --git a/internal-packages/database/prisma/migrations/20260829150000_add_concurrency_overrides/migration.sql b/internal-packages/database/prisma/migrations/20260829150000_add_concurrency_overrides/migration.sql new file mode 100644 index 00000000000..d95778f276d --- /dev/null +++ b/internal-packages/database/prisma/migrations/20260829150000_add_concurrency_overrides/migration.sql @@ -0,0 +1,24 @@ +-- AlterTable +ALTER TABLE "TaskQueue" ADD COLUMN "totalConcurrencyLimitOverriddenAt" TIMESTAMP(3); +ALTER TABLE "TaskQueue" ADD COLUMN "totalConcurrencyLimitOverriddenBy" TEXT; +ALTER TABLE "TaskQueue" ADD COLUMN "totalConcurrencyLimitBase" INTEGER; + +-- CreateTable +CREATE TABLE "TaskQueueConcurrencyKeyOverride" ( + "id" TEXT NOT NULL, + "taskQueueId" TEXT NOT NULL, + "concurrencyKey" TEXT NOT NULL, + "concurrencyLimit" INTEGER NOT NULL, + "overriddenAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "overriddenBy" TEXT, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" TIMESTAMP(3) NOT NULL, + + CONSTRAINT "TaskQueueConcurrencyKeyOverride_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE UNIQUE INDEX "TaskQueueConcurrencyKeyOverride_taskQueueId_concurrencyKey_key" ON "TaskQueueConcurrencyKeyOverride"("taskQueueId", "concurrencyKey"); + +-- AddForeignKey +ALTER TABLE "TaskQueueConcurrencyKeyOverride" ADD CONSTRAINT "TaskQueueConcurrencyKeyOverride_taskQueueId_fkey" FOREIGN KEY ("taskQueueId") REFERENCES "TaskQueue"("id") ON DELETE CASCADE ON UPDATE CASCADE; diff --git a/internal-packages/database/prisma/schema.prisma b/internal-packages/database/prisma/schema.prisma index c88ba5b5887..65a84ff1645 100644 --- a/internal-packages/database/prisma/schema.prisma +++ b/internal-packages/database/prisma/schema.prisma @@ -1983,7 +1983,13 @@ model TaskQueue { concurrencyLimitOverridePercent Decimal? @db.Decimal(5, 2) /// Caps total concurrent runs across ALL concurrencyKey values of this queue /// (concurrencyLimit applies per key value). Null = no total cap. - totalConcurrencyLimit Int? + totalConcurrencyLimit Int? + /// When the total concurrency limit was overridden + totalConcurrencyLimitOverriddenAt DateTime? + /// Who overrode the total concurrency limit (null when overridden via the API) + totalConcurrencyLimitOverriddenBy String? + /// If totalConcurrencyLimit is overridden, the declared value it reverts to on reset + totalConcurrencyLimitBase Int? rateLimit Json? paused Boolean @default(false) @@ -1995,9 +2001,31 @@ model TaskQueue { tasks BackgroundWorkerTask[] workers BackgroundWorker[] + concurrencyKeyOverrides TaskQueueConcurrencyKeyOverride[] + @@unique([runtimeEnvironmentId, name]) } +/// A per-concurrency-key limit override for a queue: the named key value gets this +/// limit instead of the queue's concurrencyLimit. Deleting the row resets the key. +model TaskQueueConcurrencyKeyOverride { + id String @id @default(cuid()) + + taskQueue TaskQueue @relation(fields: [taskQueueId], references: [id], onDelete: Cascade, onUpdate: Cascade) + taskQueueId String + + concurrencyKey String + concurrencyLimit Int + + overriddenAt DateTime @default(now()) + overriddenBy String? + + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + @@unique([taskQueueId, concurrencyKey]) +} + enum TaskQueueType { VIRTUAL NAMED From 9bcaee382579bf497f21dd46df977f945f468d21 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 11:02:54 +0100 Subject: [PATCH 04/13] feat(sdk,core,webapp): runtime overrides for total and per-key limits queues.overrideConcurrencyLimit accepts a named concurrencyKey to adjust one key's limit independently, and new overrideTotalConcurrencyLimit and resetTotalConcurrencyLimit calls adjust the cap across all keys. Four API routes back them; the concurrency system validates against the environment limit, captures the declared base on first override, persists per-key overrides in the child table alongside the engine hash, and deploys keep an overridden total instead of clobbering it from the manifest. --- .changeset/queue-concurrency-overrides.md | 18 ++ ...es.$queueParam.concurrency.key.override.ts | 103 ++++++++ ...ueues.$queueParam.concurrency.key.reset.ts | 95 +++++++ ....$queueParam.concurrency.total.override.ts | 96 +++++++ ...ues.$queueParam.concurrency.total.reset.ts | 97 +++++++ .../v3/services/concurrencySystem.server.ts | 244 +++++++++++++++++- .../services/createBackgroundWorker.server.ts | 4 +- internal-packages/run-engine/src/index.ts | 1 + packages/core/src/v3/apiClient/index.ts | 97 +++++++ packages/trigger-sdk/src/v3/queues.ts | 98 ++++++- 10 files changed, 847 insertions(+), 6 deletions(-) create mode 100644 .changeset/queue-concurrency-overrides.md create mode 100644 apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts create mode 100644 apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts create mode 100644 apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts create mode 100644 apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts diff --git a/.changeset/queue-concurrency-overrides.md b/.changeset/queue-concurrency-overrides.md new file mode 100644 index 00000000000..65579b0d81a --- /dev/null +++ b/.changeset/queue-concurrency-overrides.md @@ -0,0 +1,18 @@ +--- +"@trigger.dev/sdk": patch +"@trigger.dev/core": patch +--- + +Adjust queue concurrency at runtime, per key and in total. `queues.overrideConcurrencyLimit` accepts a `concurrencyKey` to raise or lower one key's limit without touching the rest of the queue, and the new `queues.overrideTotalConcurrencyLimit` and `queues.resetTotalConcurrencyLimit` adjust the cap across all keys. + +```ts +import { queues } from "@trigger.dev/sdk"; + +await queues.overrideConcurrencyLimit("my-queue", 20, { concurrencyKey: "tenant-123" }); +await queues.resetConcurrencyLimit("my-queue", { concurrencyKey: "tenant-123" }); + +await queues.overrideTotalConcurrencyLimit("my-queue", 100); +await queues.resetTotalConcurrencyLimit("my-queue"); +``` + +Overrides survive deploys and reset back to the declared configuration. Enforcement happens server-side on servers with total concurrency limits enabled. diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts new file mode 100644 index 00000000000..9a72155d4fc --- /dev/null +++ b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts @@ -0,0 +1,103 @@ +import { json } from "@remix-run/server-runtime"; +import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3"; +import { z } from "zod"; +import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server"; +import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server"; +import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server"; + +const BodySchema = z.object({ + type: RetrieveQueueType.default("id"), + concurrencyKey: z.string().min(1).max(128), + concurrencyLimit: z.number().int().min(0).max(100000), +}); + +const route = createActionApiRoute( + { + body: BodySchema, + params: z.object({ + queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")), + }), + authorization: { + action: "write", + resource: () => ({ type: "queues" }), + }, + }, + async ({ params, body, authentication }) => { + const input: RetrieveQueueParam = + body.type === "id" + ? params.queueParam + : { + type: body.type, + name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"), + }; + + return concurrencySystem.queues + .overrideConcurrencyKeyLimit( + authentication.environment, + input, + body.concurrencyKey, + body.concurrencyLimit + ) + .match( + (queue) => { + return json( + toQueueItem({ + friendlyId: queue.friendlyId, + name: queue.name, + type: queue.type, + running: queue.running, + queued: queue.queued, + concurrencyLimit: queue.concurrencyLimit, + concurrencyLimitBase: queue.concurrencyLimitBase, + concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt, + concurrencyLimitOverriddenBy: null, + paused: queue.paused, + }), + { status: 200 } + ); + }, + (error) => { + switch (error.type) { + case "queue_not_found": { + return json({ error: "Queue not found" }, { status: 404 }); + } + case "invalid_override": + case "concurrency_limit_exceeds_maximum": + case "too_many_key_overrides": { + return json({ error: error.message }, { status: 400 }); + } + case "queue_update_failed": { + return json( + { error: "Failed to update queue concurrency key limit" }, + { status: 500 } + ); + } + case "sync_queue_concurrency_to_engine_failed": { + return json({ error: "Failed to sync the concurrency key limit" }, { status: 500 }); + } + case "get_queue_stats_failed": { + return json({ error: "Failed to read queue stats" }, { status: 500 }); + } + case "other": { + return json( + { error: "Failed to update queue concurrency key limit" }, + { + status: 500, + } + ); + } + default: { + return json( + { error: "Failed to update queue concurrency key limit" }, + { + status: 500, + } + ); + } + } + } + ); + } +); + +export const { action } = route; diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts new file mode 100644 index 00000000000..3024ac84b72 --- /dev/null +++ b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts @@ -0,0 +1,95 @@ +import { json } from "@remix-run/server-runtime"; +import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3"; +import { z } from "zod"; +import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server"; +import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server"; +import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server"; + +const BodySchema = z.object({ + type: RetrieveQueueType.default("id"), + concurrencyKey: z.string().min(1).max(128), +}); + +const route = createActionApiRoute( + { + body: BodySchema, + params: z.object({ + queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")), + }), + authorization: { + action: "write", + resource: () => ({ type: "queues" }), + }, + }, + async ({ params, body, authentication }) => { + const input: RetrieveQueueParam = + body.type === "id" + ? params.queueParam + : { + type: body.type, + name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"), + }; + + return concurrencySystem.queues + .resetConcurrencyKeyLimit(authentication.environment, input, body.concurrencyKey) + .match( + (queue) => { + return json( + toQueueItem({ + friendlyId: queue.friendlyId, + name: queue.name, + type: queue.type, + running: queue.running, + queued: queue.queued, + concurrencyLimit: queue.concurrencyLimit, + concurrencyLimitBase: queue.concurrencyLimitBase, + concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt, + concurrencyLimitOverriddenBy: null, + paused: queue.paused, + }), + { status: 200 } + ); + }, + (error) => { + switch (error.type) { + case "queue_not_found": { + return json({ error: "Queue not found" }, { status: 404 }); + } + case "queue_not_overridden": { + return json( + { error: "This concurrency key does not have an override" }, + { status: 400 } + ); + } + case "queue_update_failed": { + return json({ error: "Failed to reset the concurrency key limit" }, { status: 500 }); + } + case "sync_queue_concurrency_to_engine_failed": { + return json({ error: "Failed to sync the concurrency key limit" }, { status: 500 }); + } + case "get_queue_stats_failed": { + return json({ error: "Failed to read queue stats" }, { status: 500 }); + } + case "other": { + return json( + { error: "Failed to reset the concurrency key limit" }, + { + status: 500, + } + ); + } + default: { + return json( + { error: "Failed to reset the concurrency key limit" }, + { + status: 500, + } + ); + } + } + } + ); + } +); + +export const { action } = route; diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts new file mode 100644 index 00000000000..b0c34ae1379 --- /dev/null +++ b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts @@ -0,0 +1,96 @@ +import { json } from "@remix-run/server-runtime"; +import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3"; +import { z } from "zod"; +import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server"; +import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server"; +import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server"; + +const BodySchema = z.object({ + type: RetrieveQueueType.default("id"), + concurrencyLimit: z.number().int().min(0).max(100000), +}); + +const route = createActionApiRoute( + { + body: BodySchema, + params: z.object({ + queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")), + }), + authorization: { + action: "write", + resource: () => ({ type: "queues" }), + }, + }, + async ({ params, body, authentication }) => { + const input: RetrieveQueueParam = + body.type === "id" + ? params.queueParam + : { + type: body.type, + name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"), + }; + + return concurrencySystem.queues + .overrideTotalConcurrencyLimit(authentication.environment, input, body.concurrencyLimit) + .match( + (queue) => { + return json( + toQueueItem({ + friendlyId: queue.friendlyId, + name: queue.name, + type: queue.type, + running: queue.running, + queued: queue.queued, + concurrencyLimit: queue.concurrencyLimit, + concurrencyLimitBase: queue.concurrencyLimitBase, + concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt, + concurrencyLimitOverriddenBy: null, + paused: queue.paused, + }), + { status: 200 } + ); + }, + (error) => { + switch (error.type) { + case "queue_not_found": { + return json({ error: "Queue not found" }, { status: 404 }); + } + case "invalid_override": + case "concurrency_limit_exceeds_maximum": { + return json({ error: error.message }, { status: 400 }); + } + case "queue_update_failed": { + return json( + { error: "Failed to update queue total concurrency limit" }, + { status: 500 } + ); + } + case "sync_queue_concurrency_to_engine_failed": { + return json({ error: "Failed to sync the total concurrency limit" }, { status: 500 }); + } + case "get_queue_stats_failed": { + return json({ error: "Failed to read queue stats" }, { status: 500 }); + } + case "other": { + return json( + { error: "Failed to update queue total concurrency limit" }, + { + status: 500, + } + ); + } + default: { + return json( + { error: "Failed to update queue total concurrency limit" }, + { + status: 500, + } + ); + } + } + } + ); + } +); + +export const { action } = route; diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts new file mode 100644 index 00000000000..2eacbb3ed6a --- /dev/null +++ b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts @@ -0,0 +1,97 @@ +import { json } from "@remix-run/server-runtime"; +import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3"; +import { z } from "zod"; +import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server"; +import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server"; +import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server"; + +const BodySchema = z.object({ + type: RetrieveQueueType.default("id"), +}); + +const route = createActionApiRoute( + { + body: BodySchema, + params: z.object({ + queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")), + }), + authorization: { + action: "write", + resource: () => ({ type: "queues" }), + }, + }, + async ({ params, body, authentication }) => { + const input: RetrieveQueueParam = + body.type === "id" + ? params.queueParam + : { + type: body.type, + name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"), + }; + + return concurrencySystem.queues + .resetTotalConcurrencyLimit(authentication.environment, input) + .match( + (queue) => { + return json( + toQueueItem({ + friendlyId: queue.friendlyId, + name: queue.name, + type: queue.type, + running: queue.running, + queued: queue.queued, + concurrencyLimit: queue.concurrencyLimit, + concurrencyLimitBase: queue.concurrencyLimitBase, + concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt, + concurrencyLimitOverriddenBy: null, + paused: queue.paused, + }), + { status: 200 } + ); + }, + (error) => { + switch (error.type) { + case "queue_not_found": { + return json({ error: "Queue not found" }, { status: 404 }); + } + case "queue_not_overridden": { + return json( + { error: "The queue total concurrency limit is not overridden" }, + { status: 400 } + ); + } + case "queue_update_failed": { + return json( + { error: "Failed to reset the queue total concurrency limit" }, + { status: 500 } + ); + } + case "sync_queue_concurrency_to_engine_failed": { + return json({ error: "Failed to sync the total concurrency limit" }, { status: 500 }); + } + case "get_queue_stats_failed": { + return json({ error: "Failed to read queue stats" }, { status: 500 }); + } + case "other": { + return json( + { error: "Failed to reset the queue total concurrency limit" }, + { + status: 500, + } + ); + } + default: { + return json( + { error: "Failed to reset the queue total concurrency limit" }, + { + status: 500, + } + ); + } + } + } + ); + } +); + +export const { action } = route; diff --git a/apps/webapp/app/v3/services/concurrencySystem.server.ts b/apps/webapp/app/v3/services/concurrencySystem.server.ts index 51c51674234..3945c087053 100644 --- a/apps/webapp/app/v3/services/concurrencySystem.server.ts +++ b/apps/webapp/app/v3/services/concurrencySystem.server.ts @@ -3,7 +3,13 @@ import { errAsync, fromPromise, okAsync } from "neverthrow"; import type { PrismaClientOrTransaction } from "~/db.server"; import type { AuthenticatedEnvironment } from "~/services/apiAuth.server"; import { logger } from "~/services/logger.server"; -import { removeQueueConcurrencyLimits, updateQueueConcurrencyLimits } from "../runQueue.server"; +import { + removeQueueConcurrencyLimits, + removeQueueTotalConcurrencyLimits, + updateQueueConcurrencyLimits, + updateQueueTotalConcurrencyLimits, +} from "../runQueue.server"; +import { RunQueueConcurrencyKeyLimitExceededError } from "@internal/run-engine"; import { engine } from "../runEngine.server"; export type ConcurrencySystemOptions = { @@ -77,6 +83,60 @@ export class ConcurrencySystem { .andThen((queue) => syncQueueConcurrencyToEngine(environment, queue)) .andThen((queue) => getQueueStats(environment, queue)); }, + overrideTotalConcurrencyLimit: ( + environment: AuthenticatedEnvironment, + queue: QueueInput, + totalConcurrencyLimit: number, + overriddenBy?: User + ) => { + return findQueueFromInput(this.db, environment, queue) + .andThen((queue) => + overrideQueueTotalConcurrencyLimit( + this.db, + environment, + queue, + totalConcurrencyLimit, + overriddenBy + ) + ) + .andThen((queue) => syncQueueTotalConcurrencyToEngine(environment, queue)) + .andThen((queue) => getQueueStats(environment, queue)); + }, + resetTotalConcurrencyLimit: (environment: AuthenticatedEnvironment, queue: QueueInput) => { + return findQueueFromInput(this.db, environment, queue) + .andThen((queue) => resetQueueTotalConcurrencyLimit(this.db, queue)) + .andThen((queue) => syncQueueTotalConcurrencyToEngine(environment, queue)) + .andThen((queue) => getQueueStats(environment, queue)); + }, + overrideConcurrencyKeyLimit: ( + environment: AuthenticatedEnvironment, + queue: QueueInput, + concurrencyKey: string, + concurrencyLimit: number, + overriddenBy?: User + ) => { + return findQueueFromInput(this.db, environment, queue) + .andThen((queue) => + overrideQueueConcurrencyKeyLimit( + this.db, + environment, + queue, + concurrencyKey, + concurrencyLimit, + overriddenBy + ) + ) + .andThen((queue) => getQueueStats(environment, queue)); + }, + resetConcurrencyKeyLimit: ( + environment: AuthenticatedEnvironment, + queue: QueueInput, + concurrencyKey: string + ) => { + return findQueueFromInput(this.db, environment, queue) + .andThen((queue) => resetQueueConcurrencyKeyLimit(this.db, environment, queue, concurrencyKey)) + .andThen((queue) => getQueueStats(environment, queue)); + }, /** * Recalculates the materialized limit of every percent-based override in the environment * against its CURRENT maximumConcurrencyLimit and syncs changed queues to the run engine. @@ -316,6 +376,188 @@ function syncQueueConcurrencyToEngine(environment: AuthenticatedEnvironment, que } } +function overrideQueueTotalConcurrencyLimit( + db: PrismaClientOrTransaction, + environment: AuthenticatedEnvironment, + queue: TaskQueue, + totalConcurrencyLimit: number, + overriddenBy?: User +) { + const maximum = environment.maximumConcurrencyLimit; + + if (!Number.isFinite(totalConcurrencyLimit) || totalConcurrencyLimit < 0) { + return errAsync({ + type: "invalid_override" as const, + message: "Total concurrency limit must be a non-negative number", + }); + } + + if (totalConcurrencyLimit > maximum) { + return errAsync({ + type: "concurrency_limit_exceeds_maximum" as const, + message: `Total concurrency limit (${totalConcurrencyLimit}) cannot exceed the environment limit (${maximum})`, + }); + } + + const totalConcurrencyLimitBase = queue.totalConcurrencyLimitOverriddenAt + ? queue.totalConcurrencyLimitBase + : queue.totalConcurrencyLimit; + + return fromPromise( + db.taskQueue.update({ + where: { id: queue.id }, + data: { + totalConcurrencyLimit, + totalConcurrencyLimitBase: totalConcurrencyLimitBase ?? null, + totalConcurrencyLimitOverriddenAt: new Date(), + totalConcurrencyLimitOverriddenBy: overriddenBy?.id ?? null, + }, + }), + (error) => ({ + type: "queue_update_failed" as const, + cause: error, + }) + ); +} + +function resetQueueTotalConcurrencyLimit(db: PrismaClientOrTransaction, queue: TaskQueue) { + if (queue.totalConcurrencyLimitOverriddenAt === null) { + return errAsync({ type: "queue_not_overridden" as const }); + } + + return fromPromise( + db.taskQueue.update({ + where: { id: queue.id }, + data: { + totalConcurrencyLimit: queue.totalConcurrencyLimitBase, + totalConcurrencyLimitBase: null, + totalConcurrencyLimitOverriddenAt: null, + totalConcurrencyLimitOverriddenBy: null, + }, + }), + (error) => ({ + type: "queue_update_failed" as const, + cause: error, + }) + ); +} + +/** + * The total limit key is separate from the per-queue limit key that pause zeroes, + * so it syncs regardless of the paused state. + */ +function syncQueueTotalConcurrencyToEngine( + environment: AuthenticatedEnvironment, + queue: TaskQueue +) { + if (typeof queue.totalConcurrencyLimit === "number") { + return fromPromise( + updateQueueTotalConcurrencyLimits(environment, queue.name, queue.totalConcurrencyLimit), + (error) => ({ + type: "sync_queue_concurrency_to_engine_failed" as const, + cause: error, + }) + ).andThen(() => okAsync(queue)); + } + + return fromPromise(removeQueueTotalConcurrencyLimits(environment, queue.name), (error) => ({ + type: "sync_queue_concurrency_to_engine_failed" as const, + cause: error, + })).andThen(() => okAsync(queue)); +} + +function overrideQueueConcurrencyKeyLimit( + db: PrismaClientOrTransaction, + environment: AuthenticatedEnvironment, + queue: TaskQueue, + concurrencyKey: string, + concurrencyLimit: number, + overriddenBy?: User +) { + const maximum = environment.maximumConcurrencyLimit; + + if (!Number.isFinite(concurrencyLimit) || concurrencyLimit < 0) { + return errAsync({ + type: "invalid_override" as const, + message: "Concurrency limit must be a non-negative number", + }); + } + + if (concurrencyLimit > maximum) { + return errAsync({ + type: "concurrency_limit_exceeds_maximum" as const, + message: `Concurrency limit (${concurrencyLimit}) cannot exceed the environment limit (${maximum})`, + }); + } + + return fromPromise( + engine.runQueue.updateQueueConcurrencyKeyLimit( + environment, + queue.name, + concurrencyKey, + concurrencyLimit + ), + (error) => { + if (error instanceof RunQueueConcurrencyKeyLimitExceededError) { + return { type: "too_many_key_overrides" as const, message: error.message }; + } + return { type: "sync_queue_concurrency_to_engine_failed" as const, cause: error }; + } + ) + .andThen(() => + fromPromise( + db.taskQueueConcurrencyKeyOverride.upsert({ + where: { taskQueueId_concurrencyKey: { taskQueueId: queue.id, concurrencyKey } }, + create: { + taskQueueId: queue.id, + concurrencyKey, + concurrencyLimit, + overriddenBy: overriddenBy?.id ?? null, + }, + update: { + concurrencyLimit, + overriddenAt: new Date(), + overriddenBy: overriddenBy?.id ?? null, + }, + }), + (error) => ({ + type: "queue_update_failed" as const, + cause: error, + }) + ) + ) + .andThen(() => okAsync(queue)); +} + +function resetQueueConcurrencyKeyLimit( + db: PrismaClientOrTransaction, + environment: AuthenticatedEnvironment, + queue: TaskQueue, + concurrencyKey: string +) { + return fromPromise( + db.taskQueueConcurrencyKeyOverride.deleteMany({ + where: { taskQueueId: queue.id, concurrencyKey }, + }), + (error) => ({ + type: "queue_update_failed" as const, + cause: error, + }) + ).andThen((deleted) => { + if (deleted.count === 0) { + return errAsync({ type: "queue_not_overridden" as const }); + } + + return fromPromise( + engine.runQueue.removeQueueConcurrencyKeyLimit(environment, queue.name, concurrencyKey), + (error) => ({ + type: "sync_queue_concurrency_to_engine_failed" as const, + cause: error, + }) + ).andThen(() => okAsync(queue)); + }); +} + function getQueueStats(environment: AuthenticatedEnvironment, queue: TaskQueue) { return fromPromise( Promise.all([ diff --git a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts index 386f8f038d6..147248bd376 100644 --- a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts +++ b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts @@ -660,6 +660,7 @@ async function upsertWorkerQueueRecord( }); } else { const hasOverride = taskQueue.concurrencyLimitOverriddenAt !== null; + const hasTotalOverride = taskQueue.totalConcurrencyLimitOverriddenAt !== null; taskQueue = await prisma.taskQueue.update({ where: { @@ -672,7 +673,8 @@ async function upsertWorkerQueueRecord( // If overridden, keep current limit and update base; otherwise update limit normally concurrencyLimit: hasOverride ? undefined : concurrencyLimit, concurrencyLimitBase: hasOverride ? concurrencyLimit : undefined, - totalConcurrencyLimit, + totalConcurrencyLimit: hasTotalOverride ? undefined : totalConcurrencyLimit, + totalConcurrencyLimitBase: hasTotalOverride ? totalConcurrencyLimit : undefined, }, }); } diff --git a/internal-packages/run-engine/src/index.ts b/internal-packages/run-engine/src/index.ts index 2c98edf6866..c9b4535b280 100644 --- a/internal-packages/run-engine/src/index.ts +++ b/internal-packages/run-engine/src/index.ts @@ -1,4 +1,5 @@ export { RunEngine } from "./engine/index.js"; +export { RunQueueConcurrencyKeyLimitExceededError } from "./run-queue/index.js"; export { RunDuplicateIdempotencyKeyError, RunOneTimeUseTokenError, diff --git a/packages/core/src/v3/apiClient/index.ts b/packages/core/src/v3/apiClient/index.ts index 585aeaed6c2..e13bb3f76d3 100644 --- a/packages/core/src/v3/apiClient/index.ts +++ b/packages/core/src/v3/apiClient/index.ts @@ -1704,6 +1704,103 @@ export class ApiClient { ); } + overrideQueueTotalConcurrencyLimit( + queue: RetrieveQueueParam, + concurrencyLimit: number, + requestOptions?: ZodFetchOptions + ) { + const type = typeof queue === "string" ? "id" : queue.type; + const value = typeof queue === "string" ? queue : queue.name; + + const encodedValue = encodeURIComponent(value.replace(/\//g, "%2F")); + + return zodfetch( + QueueItem, + `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/total/override`, + { + method: "POST", + headers: this.#getHeaders(false), + body: JSON.stringify({ + type, + concurrencyLimit, + }), + }, + mergeRequestOptions(this.defaultRequestOptions, requestOptions) + ); + } + + resetQueueTotalConcurrencyLimit(queue: RetrieveQueueParam, requestOptions?: ZodFetchOptions) { + const type = typeof queue === "string" ? "id" : queue.type; + const value = typeof queue === "string" ? queue : queue.name; + + const encodedValue = encodeURIComponent(value.replace(/\//g, "%2F")); + + return zodfetch( + QueueItem, + `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/total/reset`, + { + method: "POST", + headers: this.#getHeaders(false), + body: JSON.stringify({ + type, + }), + }, + mergeRequestOptions(this.defaultRequestOptions, requestOptions) + ); + } + + overrideQueueConcurrencyKeyLimit( + queue: RetrieveQueueParam, + concurrencyKey: string, + concurrencyLimit: number, + requestOptions?: ZodFetchOptions + ) { + const type = typeof queue === "string" ? "id" : queue.type; + const value = typeof queue === "string" ? queue : queue.name; + + const encodedValue = encodeURIComponent(value.replace(/\//g, "%2F")); + + return zodfetch( + QueueItem, + `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/key/override`, + { + method: "POST", + headers: this.#getHeaders(false), + body: JSON.stringify({ + type, + concurrencyKey, + concurrencyLimit, + }), + }, + mergeRequestOptions(this.defaultRequestOptions, requestOptions) + ); + } + + resetQueueConcurrencyKeyLimit( + queue: RetrieveQueueParam, + concurrencyKey: string, + requestOptions?: ZodFetchOptions + ) { + const type = typeof queue === "string" ? "id" : queue.type; + const value = typeof queue === "string" ? queue : queue.name; + + const encodedValue = encodeURIComponent(value.replace(/\//g, "%2F")); + + return zodfetch( + QueueItem, + `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/key/reset`, + { + method: "POST", + headers: this.#getHeaders(false), + body: JSON.stringify({ + type, + concurrencyKey, + }), + }, + mergeRequestOptions(this.defaultRequestOptions, requestOptions) + ); + } + subscribeToRun( runId: string, options?: { diff --git a/packages/trigger-sdk/src/v3/queues.ts b/packages/trigger-sdk/src/v3/queues.ts index 7e76c5f940b..e4f3a9b6641 100644 --- a/packages/trigger-sdk/src/v3/queues.ts +++ b/packages/trigger-sdk/src/v3/queues.ts @@ -144,9 +144,10 @@ export function pause( export function overrideConcurrencyLimit( queue: RetrieveQueueParam, concurrencyLimit: number, - requestOptions?: ApiRequestOptions + options?: ApiRequestOptions & { concurrencyKey?: string } ): ApiPromise { const apiClient = apiClientManager.clientOrThrow(); + const { concurrencyKey, ...requestOptions } = options ?? {}; const $requestOptions = mergeRequestOptions( { @@ -154,7 +155,7 @@ export function overrideConcurrencyLimit( name: "queues.overrideConcurrencyLimit()", icon: "queue", attributes: { - ...flattenAttributes({ queue }), + ...flattenAttributes({ queue, concurrencyKey }), ...accessoryAttributes({ items: [ { @@ -169,9 +170,93 @@ export function overrideConcurrencyLimit( requestOptions ); + if (concurrencyKey) { + return apiClient.overrideQueueConcurrencyKeyLimit( + queue, + concurrencyKey, + concurrencyLimit, + $requestOptions + ); + } + return apiClient.overrideQueueConcurrencyLimit(queue, concurrencyLimit, $requestOptions); } +/** + * Overrides the total concurrency limit of a queue: the cap on concurrent runs across + * all of its `concurrencyKey` values. + * + * @param queue - The ID of the queue, or the type and name + * @param concurrencyLimit - The total concurrency limit to apply + * @returns The updated queue state + */ +export function overrideTotalConcurrencyLimit( + queue: RetrieveQueueParam, + concurrencyLimit: number, + requestOptions?: ApiRequestOptions +): ApiPromise { + const apiClient = apiClientManager.clientOrThrow(); + + const $requestOptions = mergeRequestOptions( + { + tracer, + name: "queues.overrideTotalConcurrencyLimit()", + icon: "queue", + attributes: { + ...flattenAttributes({ queue }), + ...accessoryAttributes({ + items: [ + { + text: typeof queue === "string" ? queue : queue.name, + variant: "normal", + }, + ], + style: "codepath", + }), + }, + }, + requestOptions + ); + + return apiClient.overrideQueueTotalConcurrencyLimit(queue, concurrencyLimit, $requestOptions); +} + +/** + * Resets the total concurrency limit of a queue back to its declared value. + * + * @param queue - The ID of the queue, or the type and name + * @returns The updated queue state + */ +export function resetTotalConcurrencyLimit( + queue: RetrieveQueueParam, + requestOptions?: ApiRequestOptions +): ApiPromise { + const apiClient = apiClientManager.clientOrThrow(); + + const $requestOptions = mergeRequestOptions( + { + tracer, + name: "queues.resetTotalConcurrencyLimit()", + icon: "queue", + attributes: { + ...flattenAttributes({ queue }), + ...accessoryAttributes({ + items: [ + { + text: typeof queue === "string" ? queue : queue.name, + variant: "normal", + }, + ], + style: "codepath", + }), + }, + }, + requestOptions + ); + + return apiClient.resetQueueTotalConcurrencyLimit(queue, $requestOptions); +} + /** * Resets the concurrency limit of a queue to the base value. * @@ -180,9 +265,10 @@ export function overrideConcurrencyLimit( */ export function resetConcurrencyLimit( queue: RetrieveQueueParam, - requestOptions?: ApiRequestOptions + options?: ApiRequestOptions & { concurrencyKey?: string } ): ApiPromise { const apiClient = apiClientManager.clientOrThrow(); + const { concurrencyKey, ...requestOptions } = options ?? {}; const $requestOptions = mergeRequestOptions( { @@ -190,7 +276,7 @@ export function resetConcurrencyLimit( name: "queues.resetConcurrencyLimit()", icon: "queue", attributes: { - ...flattenAttributes({ queue }), + ...flattenAttributes({ queue, concurrencyKey }), ...accessoryAttributes({ items: [ { @@ -205,6 +291,10 @@ export function resetConcurrencyLimit( requestOptions ); + if (concurrencyKey) { + return apiClient.resetQueueConcurrencyKeyLimit(queue, concurrencyKey, $requestOptions); + } + return apiClient.resetQueueConcurrencyLimit(queue, $requestOptions); } From ded4316200b69bf17e1e7bb2712842e805ca807f Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 11:15:06 +0100 Subject: [PATCH 05/13] fix(run-engine,webapp,sdk): converge override failures and unpin blocked keys A variant blocked at its per-key limit or by a gate now backs off in the candidate index instead of pinning the bounded window, so zero-limit keys can never starve runnable keys behind them; acks and nacks rebalance the score back the moment capacity frees. Override writes persist before enforcing and resets enforce before clearing, so a failure on either side leaves a state a retry converges from, with the cap-rejected row compensated away. Deploys restore per-key overrides from their durable rows into the engine, and an empty concurrencyKey in the SDK no longer falls through to the queue-wide endpoint. --- .../v3/services/concurrencySystem.server.ts | 119 ++++++++++++++---- .../services/createBackgroundWorker.server.ts | 25 ++++ .../run-engine/src/run-queue/index.ts | 28 ++++- .../tests/concurrencyKeyOverrides.test.ts | 63 +++++++++- packages/trigger-sdk/src/v3/queues.ts | 4 +- 5 files changed, 209 insertions(+), 30 deletions(-) diff --git a/apps/webapp/app/v3/services/concurrencySystem.server.ts b/apps/webapp/app/v3/services/concurrencySystem.server.ts index 3945c087053..28fb32b868d 100644 --- a/apps/webapp/app/v3/services/concurrencySystem.server.ts +++ b/apps/webapp/app/v3/services/concurrencySystem.server.ts @@ -104,6 +104,7 @@ export class ConcurrencySystem { }, resetTotalConcurrencyLimit: (environment: AuthenticatedEnvironment, queue: QueueInput) => { return findQueueFromInput(this.db, environment, queue) + .andThen((queue) => syncQueueTotalConcurrencyResetToEngine(environment, queue)) .andThen((queue) => resetQueueTotalConcurrencyLimit(this.db, queue)) .andThen((queue) => syncQueueTotalConcurrencyToEngine(environment, queue)) .andThen((queue) => getQueueStats(environment, queue)); @@ -134,7 +135,9 @@ export class ConcurrencySystem { concurrencyKey: string ) => { return findQueueFromInput(this.db, environment, queue) - .andThen((queue) => resetQueueConcurrencyKeyLimit(this.db, environment, queue, concurrencyKey)) + .andThen((queue) => + resetQueueConcurrencyKeyLimit(this.db, environment, queue, concurrencyKey) + ) .andThen((queue) => getQueueStats(environment, queue)); }, /** @@ -420,6 +423,35 @@ function overrideQueueTotalConcurrencyLimit( ); } +/** + * Enforce first, then persist: syncs the engine to the declared base BEFORE clearing + * the override marker, so an engine failure leaves the marker set and a retry + * converges instead of being rejected while the overridden limit stays enforced. + */ +function syncQueueTotalConcurrencyResetToEngine( + environment: AuthenticatedEnvironment, + queue: TaskQueue +) { + if (queue.totalConcurrencyLimitOverriddenAt === null) { + return errAsync({ type: "queue_not_overridden" as const }); + } + + if (typeof queue.totalConcurrencyLimitBase === "number") { + return fromPromise( + updateQueueTotalConcurrencyLimits(environment, queue.name, queue.totalConcurrencyLimitBase), + (error) => ({ + type: "sync_queue_concurrency_to_engine_failed" as const, + cause: error, + }) + ).andThen(() => okAsync(queue)); + } + + return fromPromise(removeQueueTotalConcurrencyLimits(environment, queue.name), (error) => ({ + type: "sync_queue_concurrency_to_engine_failed" as const, + cause: error, + })).andThen(() => okAsync(queue)); +} + function resetQueueTotalConcurrencyLimit(db: PrismaClientOrTransaction, queue: TaskQueue) { if (queue.totalConcurrencyLimitOverriddenAt === null) { return errAsync({ type: "queue_not_overridden" as const }); @@ -466,6 +498,13 @@ function syncQueueTotalConcurrencyToEngine( })).andThen(() => okAsync(queue)); } +/** + * Persist first, then enforce: a database failure leaves nothing enforced and the + * request errors cleanly, while an engine failure after persistence leaves a durable + * record and a retry converges (the upsert is idempotent). When the engine rejects a + * NEW key for exceeding the per-queue override cap, the just-created row is removed + * again so the record never claims an override the engine refused. + */ function overrideQueueConcurrencyKeyLimit( db: PrismaClientOrTransaction, environment: AuthenticatedEnvironment, @@ -491,20 +530,13 @@ function overrideQueueConcurrencyKeyLimit( } return fromPromise( - engine.runQueue.updateQueueConcurrencyKeyLimit( - environment, - queue.name, - concurrencyKey, - concurrencyLimit - ), - (error) => { - if (error instanceof RunQueueConcurrencyKeyLimitExceededError) { - return { type: "too_many_key_overrides" as const, message: error.message }; - } - return { type: "sync_queue_concurrency_to_engine_failed" as const, cause: error }; - } + db.taskQueueConcurrencyKeyOverride.findFirst({ + where: { taskQueueId: queue.id, concurrencyKey }, + select: { id: true }, + }), + (error) => ({ type: "other" as const, cause: error }) ) - .andThen(() => + .andThen((existing) => fromPromise( db.taskQueueConcurrencyKeyOverride.upsert({ where: { taskQueueId_concurrencyKey: { taskQueueId: queue.id, concurrencyKey } }, @@ -524,11 +556,42 @@ function overrideQueueConcurrencyKeyLimit( type: "queue_update_failed" as const, cause: error, }) - ) + ).map(() => existing) + ) + .andThen((existing) => + fromPromise( + engine.runQueue.updateQueueConcurrencyKeyLimit( + environment, + queue.name, + concurrencyKey, + concurrencyLimit + ), + (error) => { + if (error instanceof RunQueueConcurrencyKeyLimitExceededError) { + return { type: "too_many_key_overrides" as const, message: error.message }; + } + return { type: "sync_queue_concurrency_to_engine_failed" as const, cause: error }; + } + ).orElse((error) => { + if (!existing && error.type === "too_many_key_overrides") { + return fromPromise( + db.taskQueueConcurrencyKeyOverride.deleteMany({ + where: { taskQueueId: queue.id, concurrencyKey }, + }), + () => error + ).andThen(() => errAsync(error)); + } + return errAsync(error); + }) ) .andThen(() => okAsync(queue)); } +/** + * Enforce first, then persist: removing the engine limit is idempotent, so an engine + * failure leaves the override row in place and a retry converges instead of being + * rejected as not overridden while the old limit is still enforced. + */ function resetQueueConcurrencyKeyLimit( db: PrismaClientOrTransaction, environment: AuthenticatedEnvironment, @@ -536,15 +599,13 @@ function resetQueueConcurrencyKeyLimit( concurrencyKey: string ) { return fromPromise( - db.taskQueueConcurrencyKeyOverride.deleteMany({ + db.taskQueueConcurrencyKeyOverride.findFirst({ where: { taskQueueId: queue.id, concurrencyKey }, + select: { id: true }, }), - (error) => ({ - type: "queue_update_failed" as const, - cause: error, - }) - ).andThen((deleted) => { - if (deleted.count === 0) { + (error) => ({ type: "other" as const, cause: error }) + ).andThen((existing) => { + if (!existing) { return errAsync({ type: "queue_not_overridden" as const }); } @@ -554,7 +615,19 @@ function resetQueueConcurrencyKeyLimit( type: "sync_queue_concurrency_to_engine_failed" as const, cause: error, }) - ).andThen(() => okAsync(queue)); + ) + .andThen(() => + fromPromise( + db.taskQueueConcurrencyKeyOverride.deleteMany({ + where: { taskQueueId: queue.id, concurrencyKey }, + }), + (error) => ({ + type: "queue_update_failed" as const, + cause: error, + }) + ) + ) + .andThen(() => okAsync(queue)); }); } diff --git a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts index 147248bd376..6799c925af8 100644 --- a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts +++ b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts @@ -582,6 +582,31 @@ async function createWorkerQueue( await removeQueueTotalConcurrencyLimits(environment, taskQueue.name); } + /** + * Restore per-key limit overrides into the engine so a fresh or flushed Redis + * converges back to the durable records on the next deploy. Row-level failures + * are logged rather than failing the deploy; the next deploy retries them. + */ + const keyOverrides = await prisma.taskQueueConcurrencyKeyOverride.findMany({ + where: { taskQueueId: taskQueue.id }, + }); + for (const keyOverride of keyOverrides) { + try { + await engine.runQueue.updateQueueConcurrencyKeyLimit( + environment, + taskQueue.name, + keyOverride.concurrencyKey, + keyOverride.concurrencyLimit + ); + } catch (error) { + logger.error("createWorkerQueue: failed to restore concurrency key override", { + taskQueueId: taskQueue.id, + concurrencyKey: keyOverride.concurrencyKey, + error, + }); + } + } + if (!taskQueue.paused) { if (typeof newConcurrencyLimit === "number") { logger.debug("createWorkerQueue: updating concurrency limit", { diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index b774732f5d3..43614d13c8c 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -312,6 +312,10 @@ export type RunQueueOptions = { * that dead-lettered or suspended through a mirror-less path. Enabling only after * every instance runs this build avoids the noise but is no longer load-bearing * for correctness. + * + * Per-concurrency-key limit overrides are part of the same concurrency-limits + * feature and are deliberately enforced behind this flag too: writes are always + * accepted and durable, and enforcement of both arrives together. */ totalConcurrencyEnabled?: boolean; /** @@ -5045,6 +5049,7 @@ for _, ckQueueName in ipairs(ckQueues) do end local fullQueueKey = keyPrefix .. ckQueueName + local blockedByGates = false local ckConcurrencyKey = fullQueueKey .. ':currentConcurrency' local ckCurrentConcurrency = tonumber(redis.call('SCARD', ckConcurrencyKey) or '0') @@ -5057,6 +5062,14 @@ for _, ckQueueName in ipairs(ckQueues) do end end + if ckCurrentConcurrency >= perKeyLimit then + -- Back a blocked variant off so it cannot pin the bounded candidate window + -- and starve later keys (acute with a zero per-key override, which never + -- self-clears). Acks and nacks rebalance the score back to the oldest + -- message, so the key is eligible again the moment capacity frees. + redis.call('ZADD', ckIndexKey, currentTime + 1000, ckQueueName) + end + if ckCurrentConcurrency < perKeyLimit then local messages = redis.call('ZRANGEBYSCORE', fullQueueKey, '-inf', tostring(currentTime), 'WITHSCORES', 'LIMIT', 0, 1) @@ -5084,6 +5097,9 @@ for _, ckQueueName in ipairs(ckQueues) do if gatesEnabled then gatesAllow = __gatesHaveCapacity(keyPrefix, messageData, messageId, envConcurrencyLimit, messageKeyPrefix) end + if not gatesAllow then + blockedByGates = true + end local alreadyInGroup = false local totalAllows = true @@ -5126,11 +5142,15 @@ for _, ckQueueName in ipairs(ckQueues) do decrLengthCounter() end - local earliest = redis.call('ZRANGE', fullQueueKey, 0, 0, 'WITHSCORES') - if #earliest == 0 then - redis.call('ZREM', ckIndexKey, ckQueueName) + if blockedByGates then + redis.call('ZADD', ckIndexKey, currentTime + 1000, ckQueueName) else - redis.call('ZADD', ckIndexKey, earliest[2], ckQueueName) + 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 end else local any = redis.call('ZRANGE', fullQueueKey, 0, 0, 'WITHSCORES') diff --git a/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts b/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts index 7995369bc3c..22a52a46625 100644 --- a/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts +++ b/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts @@ -32,11 +32,17 @@ const authenticatedEnvDev = { organization: { id: "o1234" }, }; -function createQueue(redisContainer: any, totalConcurrencyEnabled: boolean, maxOverrides?: number) { +function createQueue( + redisContainer: any, + totalConcurrencyEnabled: boolean, + maxOverrides?: number, + dequeueCount?: number +) { return new RunQueue({ ...testOptions, totalConcurrencyEnabled, maxConcurrencyKeyOverridesPerQueue: maxOverrides, + masterQueueConsumerDequeueCount: dequeueCount, queueSelectionStrategy: new FairQueueSelectionStrategy({ redis: { keyPrefix: "runqueue:test:", @@ -239,6 +245,61 @@ describe("RunQueue per-concurrency-key limit overrides", () => { } }); + redisTest( + "blocked keys cannot pin the candidate window and starve later keys", + async ({ redisContainer }) => { + /** dequeueCount 2 makes the candidate window 6 variants wide. */ + const queue = createQueue(redisContainer, true, undefined, 2); + try { + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 5); + + /** + * Ten zero-limit keys with OLDER messages fill the window many times over; + * without the blocked-key backoff the runnable key behind them would never + * be examined. + */ + const now = Date.now(); + for (let i = 0; i < 10; i++) { + const ck = `blocked-${i}`; + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", ck, 0); + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `b${i}`, + concurrencyKey: ck, + timestamp: now - 10_000 + i, + }), + workerQueue: "main", + }); + } + + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: "good-0", + concurrencyKey: "ck-good", + timestamp: now - 500, + }), + workerQueue: "main", + }); + + const goodAdmitted = await waitFor( + async () => + (await queue.currentConcurrencyOfQueue( + authenticatedEnvDev, + "task/my-task", + "ck-good" + )) === 1, + 30_000 + ); + expect(goodAdmitted).toBe(true); + expect(await queue.lengthOfQueue(authenticatedEnvDev, "task/my-task")).toBe(10); + } finally { + await queue.quit(); + } + } + ); + redisTest("overrides are ignored when disabled", async ({ redisContainer }) => { const queue = createQueue(redisContainer, false); try { diff --git a/packages/trigger-sdk/src/v3/queues.ts b/packages/trigger-sdk/src/v3/queues.ts index e4f3a9b6641..7f16d96e58b 100644 --- a/packages/trigger-sdk/src/v3/queues.ts +++ b/packages/trigger-sdk/src/v3/queues.ts @@ -170,7 +170,7 @@ export function overrideConcurrencyLimit( requestOptions ); - if (concurrencyKey) { + if (concurrencyKey !== undefined) { return apiClient.overrideQueueConcurrencyKeyLimit( queue, concurrencyKey, @@ -291,7 +291,7 @@ export function resetConcurrencyLimit( requestOptions ); - if (concurrencyKey) { + if (concurrencyKey !== undefined) { return apiClient.resetQueueConcurrencyKeyLimit(queue, concurrencyKey, $requestOptions); } From 4c261ead941ba3f5312e2a51c79d8a7242823879 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 11:54:58 +0100 Subject: [PATCH 06/13] fix(run-engine,webapp): gate admission honors per-key overrides; harden races Gate capacity now reads the gate queue's ckLimits hash, so an override on a key applies whether runs meet it as their own queue or as a gate. A reset deletes only the exact row generation it read, so a concurrent override's newer record survives, and a cap rejection deletes its row unconditionally since the cap can only reject keys absent from the engine hash. --- .../v3/services/concurrencySystem.server.ts | 21 +++++++--- .../run-engine/src/run-queue/index.ts | 6 +++ .../tests/concurrencyKeyOverrides.test.ts | 40 ++++++++++++++++++- 3 files changed, 60 insertions(+), 7 deletions(-) diff --git a/apps/webapp/app/v3/services/concurrencySystem.server.ts b/apps/webapp/app/v3/services/concurrencySystem.server.ts index 28fb32b868d..94417abf7df 100644 --- a/apps/webapp/app/v3/services/concurrencySystem.server.ts +++ b/apps/webapp/app/v3/services/concurrencySystem.server.ts @@ -501,9 +501,10 @@ function syncQueueTotalConcurrencyToEngine( /** * Persist first, then enforce: a database failure leaves nothing enforced and the * request errors cleanly, while an engine failure after persistence leaves a durable - * record and a retry converges (the upsert is idempotent). When the engine rejects a - * NEW key for exceeding the per-queue override cap, the just-created row is removed - * again so the record never claims an override the engine refused. + * record and a retry converges (the upsert is idempotent). A cap rejection deletes + * the row unconditionally: the cap only rejects keys absent from the engine hash + * (updates to present keys always succeed), so a cap-rejected row is never backing + * an enforced limit and must not survive as an authoritative override. */ function overrideQueueConcurrencyKeyLimit( db: PrismaClientOrTransaction, @@ -573,7 +574,7 @@ function overrideQueueConcurrencyKeyLimit( return { type: "sync_queue_concurrency_to_engine_failed" as const, cause: error }; } ).orElse((error) => { - if (!existing && error.type === "too_many_key_overrides") { + if (error.type === "too_many_key_overrides") { return fromPromise( db.taskQueueConcurrencyKeyOverride.deleteMany({ where: { taskQueueId: queue.id, concurrencyKey }, @@ -601,7 +602,7 @@ function resetQueueConcurrencyKeyLimit( return fromPromise( db.taskQueueConcurrencyKeyOverride.findFirst({ where: { taskQueueId: queue.id, concurrencyKey }, - select: { id: true }, + select: { id: true, overriddenAt: true }, }), (error) => ({ type: "other" as const, cause: error }) ).andThen((existing) => { @@ -618,8 +619,16 @@ function resetQueueConcurrencyKeyLimit( ) .andThen(() => fromPromise( + /** + * Deletes only the exact row generation this reset read, so a concurrent + * override that re-wrote the row after the reset began keeps its record + * (its next write, or the deploy-time restore, re-syncs the engine). + */ db.taskQueueConcurrencyKeyOverride.deleteMany({ - where: { taskQueueId: queue.id, concurrencyKey }, + where: { + id: existing.id, + overriddenAt: existing.overriddenAt, + }, }), (error) => ({ type: "queue_update_failed" as const, diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index 43614d13c8c..6b11034645c 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -120,6 +120,12 @@ local function __gatesHaveCapacity(gatesKeyPrefix, msg, messageId, envLimit, msg local base, variant, gateKey = __gateKeys(gatesKeyPrefix, msg, gate) local occupancy = tonumber(redis.call('SCARD', variant .. ':currentConcurrency') or '0') local perKeyLimit = math.min(tonumber(redis.call('GET', base .. ':concurrency') or '1000000'), envLimit) + if gateKey and gateKey ~= '' then + local gateOverride = redis.call('HGET', base .. ':ckLimits', string.sub(variant, #gatesKeyPrefix + 1)) + if gateOverride then + perKeyLimit = math.min(tonumber(gateOverride), envLimit) + end + end if occupancy >= perKeyLimit and redis.call('SISMEMBER', variant .. ':currentConcurrency', messageId) == 0 then __gateReconcile(variant .. ':currentConcurrency', msgKeyPrefix, gatesKeyPrefix) return false diff --git a/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts b/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts index 22a52a46625..8c10a8ac2bf 100644 --- a/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts +++ b/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts @@ -36,11 +36,13 @@ function createQueue( redisContainer: any, totalConcurrencyEnabled: boolean, maxOverrides?: number, - dequeueCount?: number + dequeueCount?: number, + gatesEnabled?: boolean ) { return new RunQueue({ ...testOptions, totalConcurrencyEnabled, + gatesEnabled, maxConcurrencyKeyOverridesPerQueue: maxOverrides, masterQueueConsumerDequeueCount: dequeueCount, queueSelectionStrategy: new FairQueueSelectionStrategy({ @@ -300,6 +302,42 @@ describe("RunQueue per-concurrency-key limit overrides", () => { } ); + redisTest("gate admission honors the gate queue per-key override", async ({ redisContainer }) => { + const queue = createQueue(redisContainer, true, undefined, undefined, true); + try { + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 5); + await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "tenant", 1); + await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "tenant", "acme", 2); + + const now = Date.now(); + for (const [i, ck] of ["ck-a", "ck-b", "ck-c"].entries()) { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: makeMessage({ + runId: `r${i}`, + concurrencyKey: ck, + timestamp: now - 1000 + i, + gates: [{ queue: "tenant", concurrencyKey: "acme" }], + }), + workerQueue: "main", + }); + } + + /** The declared gate limit is 1; the override raises acme to 2. */ + const twoAdmitted = await waitFor( + async () => + (await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "tenant", "acme")) === 2 + ); + expect(twoAdmitted).toBe(true); + + await setTimeout(2000); + expect(await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "tenant", "acme")).toBe(2); + expect(await queue.lengthOfQueue(authenticatedEnvDev, "task/my-task")).toBe(1); + } finally { + await queue.quit(); + } + }); + redisTest("overrides are ignored when disabled", async ({ redisContainer }) => { const queue = createQueue(redisContainer, false); try { From 2f78e447e731c23af3bfd21345dd703df01d07aa Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 12:03:40 +0100 Subject: [PATCH 07/13] fix(run-engine,webapp): flag-consistent gate overrides; generation-safe cap cleanup Gate admission reads a per-key override only when concurrency limit enforcement is enabled, matching the primary admit paths (the flag now threads through the unkeyed enqueue and dequeue scripts too). The cap-rejection cleanup deletes only the exact row generation the rejected request wrote, so a concurrent request that succeeded after capacity freed keeps its durable record. --- .../v3/services/concurrencySystem.server.ts | 13 +++++++--- .../run-engine/src/run-queue/index.ts | 25 +++++++++++++------ 2 files changed, 26 insertions(+), 12 deletions(-) diff --git a/apps/webapp/app/v3/services/concurrencySystem.server.ts b/apps/webapp/app/v3/services/concurrencySystem.server.ts index 94417abf7df..571ed5daf64 100644 --- a/apps/webapp/app/v3/services/concurrencySystem.server.ts +++ b/apps/webapp/app/v3/services/concurrencySystem.server.ts @@ -537,7 +537,7 @@ function overrideQueueConcurrencyKeyLimit( }), (error) => ({ type: "other" as const, cause: error }) ) - .andThen((existing) => + .andThen(() => fromPromise( db.taskQueueConcurrencyKeyOverride.upsert({ where: { taskQueueId_concurrencyKey: { taskQueueId: queue.id, concurrencyKey } }, @@ -557,9 +557,9 @@ function overrideQueueConcurrencyKeyLimit( type: "queue_update_failed" as const, cause: error, }) - ).map(() => existing) + ) ) - .andThen((existing) => + .andThen((written) => fromPromise( engine.runQueue.updateQueueConcurrencyKeyLimit( environment, @@ -576,8 +576,13 @@ function overrideQueueConcurrencyKeyLimit( ).orElse((error) => { if (error.type === "too_many_key_overrides") { return fromPromise( + /** + * Deletes only the exact row generation this request wrote, so a + * concurrent request that succeeded after capacity freed keeps its + * durable record. + */ db.taskQueueConcurrencyKeyOverride.deleteMany({ - where: { taskQueueId: queue.id, concurrencyKey }, + where: { id: written.id, overriddenAt: written.overriddenAt }, }), () => error ).andThen(() => errAsync(error)); diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index 6b11034645c..7917f3a48be 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -114,13 +114,13 @@ local function __gateReconcile(setKey, msgKeyPrefix, reconcileKeyPrefix) end end -local function __gatesHaveCapacity(gatesKeyPrefix, msg, messageId, envLimit, msgKeyPrefix) +local function __gatesHaveCapacity(gatesKeyPrefix, msg, messageId, envLimit, msgKeyPrefix, ckOverridesEnabled) if not msg.gates then return true end for _, gate in ipairs(msg.gates) do local base, variant, gateKey = __gateKeys(gatesKeyPrefix, msg, gate) local occupancy = tonumber(redis.call('SCARD', variant .. ':currentConcurrency') or '0') local perKeyLimit = math.min(tonumber(redis.call('GET', base .. ':concurrency') or '1000000'), envLimit) - if gateKey and gateKey ~= '' then + if ckOverridesEnabled and gateKey and gateKey ~= '' then local gateOverride = redis.call('HGET', base .. ':ckLimits', string.sub(variant, #gatesKeyPrefix + 1)) if gateOverride then perKeyLimit = math.min(tonumber(gateOverride), envLimit) @@ -2557,6 +2557,7 @@ export class RunQueue { enableFastPathArg, this.options.redis.keyPrefix ?? "", this.options.gatesEnabled ? "1" : "0", + this.options.totalConcurrencyEnabled ? "1" : "0", metricsGaugeArg ); } else { @@ -2586,6 +2587,7 @@ export class RunQueue { enableFastPathArg, this.options.redis.keyPrefix ?? "", this.options.gatesEnabled ? "1" : "0", + this.options.totalConcurrencyEnabled ? "1" : "0", metricsGaugeArg ); } @@ -2667,6 +2669,7 @@ export class RunQueue { this.options.redis.keyPrefix ?? "", String(maxCount), this.options.gatesEnabled ? "1" : "0", + this.options.totalConcurrencyEnabled ? "1" : "0", metricsGaugeArg ); @@ -3664,6 +3667,7 @@ local currentTime = ARGV[8] local enableFastPath = ARGV[9] local keyPrefix = ARGV[10] local gatesEnabled = ARGV[11] == '1' +local totalConcurrencyEnabled = ARGV[12] == '1' ${QUEUE_METRICS_GAUGE_PRELUDE} ${QUEUE_GATES_LUA_HELPERS} @@ -3691,7 +3695,7 @@ if enableFastPath == '1' then local okDecode, decoded = pcall(cjson.decode, messageData) if okDecode and type(decoded) == 'table' and decoded.gates then gateMsg = decoded - gatesAllowFastPath = __gatesHaveCapacity(keyPrefix, decoded, messageId, envLimit, nil) + gatesAllowFastPath = __gatesHaveCapacity(keyPrefix, decoded, messageId, envLimit, nil, totalConcurrencyEnabled) end end @@ -3778,6 +3782,7 @@ local currentTime = ARGV[10] local enableFastPath = ARGV[11] local keyPrefix = ARGV[12] local gatesEnabled = ARGV[13] == '1' +local totalConcurrencyEnabled = ARGV[14] == '1' ${QUEUE_METRICS_GAUGE_PRELUDE} ${QUEUE_GATES_LUA_HELPERS} @@ -3805,7 +3810,7 @@ if enableFastPath == '1' then local okDecode, decoded = pcall(cjson.decode, messageData) if okDecode and type(decoded) == 'table' and decoded.gates then gateMsg = decoded - gatesAllowFastPath = __gatesHaveCapacity(keyPrefix, decoded, messageId, envLimit, nil) + gatesAllowFastPath = __gatesHaveCapacity(keyPrefix, decoded, messageId, envLimit, nil, totalConcurrencyEnabled) end end @@ -4172,7 +4177,7 @@ if enableFastPath == '1' then local okDecode, decoded = pcall(cjson.decode, messageData) if okDecode and type(decoded) == 'table' and decoded.gates then gateMsg = decoded - gatesAllowFastPath = __gatesHaveCapacity(keyPrefix, decoded, messageId, envLimit, nil) + gatesAllowFastPath = __gatesHaveCapacity(keyPrefix, decoded, messageId, envLimit, nil, totalConcurrencyEnabled) end end @@ -4350,7 +4355,7 @@ if enableFastPath == '1' then local okDecode, decoded = pcall(cjson.decode, messageData) if okDecode and type(decoded) == 'table' and decoded.gates then gateMsg = decoded - gatesAllowFastPath = __gatesHaveCapacity(keyPrefix, decoded, messageId, envLimit, nil) + gatesAllowFastPath = __gatesHaveCapacity(keyPrefix, decoded, messageId, envLimit, nil, totalConcurrencyEnabled) end end @@ -4680,6 +4685,7 @@ local defaultEnvConcurrencyBurstFactor = ARGV[4] local keyPrefix = ARGV[5] local maxCount = tonumber(ARGV[6] or '1') local gatesEnabled = ARGV[7] == '1' +local totalConcurrencyEnabled = ARGV[8] == '1' ${QUEUE_METRICS_GAUGE_PRELUDE} ${QUEUE_GATES_LUA_HELPERS} ${QUEUE_METRICS_GAUGE_LUA} @@ -4752,7 +4758,7 @@ for i = 1, #messages, 2 do else local gatesAllow = true if gatesEnabled then - gatesAllow = __gatesHaveCapacity(keyPrefix, messageData, messageId, envConcurrencyLimit, messageKeyPrefix) + gatesAllow = __gatesHaveCapacity(keyPrefix, messageData, messageId, envConcurrencyLimit, messageKeyPrefix, totalConcurrencyEnabled) end if gatesAllow then @@ -5101,7 +5107,7 @@ for _, ckQueueName in ipairs(ckQueues) do else local gatesAllow = true if gatesEnabled then - gatesAllow = __gatesHaveCapacity(keyPrefix, messageData, messageId, envConcurrencyLimit, messageKeyPrefix) + gatesAllow = __gatesHaveCapacity(keyPrefix, messageData, messageId, envConcurrencyLimit, messageKeyPrefix, totalConcurrencyEnabled) end if not gatesAllow then blockedByGates = true @@ -6209,6 +6215,7 @@ declare module "@internal/redis" { enableFastPath: string, keyPrefix: string, gatesEnabled: string, + totalConcurrencyEnabled: string, metricsEnabled: string, callback?: Callback<[number, number[] | null]> ): Result<[number, number[] | null], Context>; @@ -6242,6 +6249,7 @@ declare module "@internal/redis" { enableFastPath: string, keyPrefix: string, gatesEnabled: string, + totalConcurrencyEnabled: string, metricsEnabled: string, callback?: Callback<[number, number[] | null]> ): Result<[number, number[] | null], Context>; @@ -6280,6 +6288,7 @@ declare module "@internal/redis" { keyPrefix: string, maxCount: string, gatesEnabled: string, + totalConcurrencyEnabled: string, metricsEnabled: string, callback?: Callback<[string[] | null, number[] | null]> ): Result<[string[] | null, number[] | null], Context>; From 3987ca058f76bd5ac0781bcfc3f81995a2357874 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 16:52:12 +0100 Subject: [PATCH 08/13] fix(webapp): export queue concurrency route handlers by property access The client build's export analyzer cannot statically resolve destructured route exports, so it treats the module's exports as depending on server-only code and the build fails. Also export the builder's loader so non-POST methods get a 405, matching the other concurrency routes. --- .../api.v1.queues.$queueParam.concurrency.key.override.ts | 4 +++- .../routes/api.v1.queues.$queueParam.concurrency.key.reset.ts | 4 +++- .../api.v1.queues.$queueParam.concurrency.total.override.ts | 4 +++- .../api.v1.queues.$queueParam.concurrency.total.reset.ts | 4 +++- 4 files changed, 12 insertions(+), 4 deletions(-) diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts index 9a72155d4fc..5a37b4526ec 100644 --- a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts +++ b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts @@ -100,4 +100,6 @@ const route = createActionApiRoute( } ); -export const { action } = route; +export const action = route.action; +/** The builder's loader answers non-POST methods with a 405. */ +export const loader = route.loader; diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts index 3024ac84b72..51d14642e2c 100644 --- a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts +++ b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts @@ -92,4 +92,6 @@ const route = createActionApiRoute( } ); -export const { action } = route; +export const action = route.action; +/** The builder's loader answers non-POST methods with a 405. */ +export const loader = route.loader; diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts index b0c34ae1379..c643b77965a 100644 --- a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts +++ b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts @@ -93,4 +93,6 @@ const route = createActionApiRoute( } ); -export const { action } = route; +export const action = route.action; +/** The builder's loader answers non-POST methods with a 405. */ +export const loader = route.loader; diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts index 2eacbb3ed6a..b2841f1efe6 100644 --- a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts +++ b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts @@ -94,4 +94,6 @@ const route = createActionApiRoute( } ); -export const { action } = route; +export const action = route.action; +/** The builder's loader answers non-POST methods with a 405. */ +export const loader = route.loader; From c39e28625789db3285bca534eb0a318a8c82d1c4 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 29 Aug 2026 19:10:38 +0100 Subject: [PATCH 09/13] refactor(sdk,core,webapp): combined concurrency override API names The override and reset endpoints, client methods and SDK functions now say combinedConcurrencyLimit, matching the queue option. --- .changeset/queue-concurrency-overrides.md | 8 ++++---- ...ueueParam.concurrency.combined.override.ts} | 0 ....$queueParam.concurrency.combined.reset.ts} | 0 packages/core/src/v3/apiClient/index.ts | 8 ++++---- packages/trigger-sdk/src/v3/queues.ts | 18 +++++++++--------- 5 files changed, 17 insertions(+), 17 deletions(-) rename apps/webapp/app/routes/{api.v1.queues.$queueParam.concurrency.total.override.ts => api.v1.queues.$queueParam.concurrency.combined.override.ts} (100%) rename apps/webapp/app/routes/{api.v1.queues.$queueParam.concurrency.total.reset.ts => api.v1.queues.$queueParam.concurrency.combined.reset.ts} (100%) diff --git a/.changeset/queue-concurrency-overrides.md b/.changeset/queue-concurrency-overrides.md index 65579b0d81a..6bbe8f909ee 100644 --- a/.changeset/queue-concurrency-overrides.md +++ b/.changeset/queue-concurrency-overrides.md @@ -3,7 +3,7 @@ "@trigger.dev/core": patch --- -Adjust queue concurrency at runtime, per key and in total. `queues.overrideConcurrencyLimit` accepts a `concurrencyKey` to raise or lower one key's limit without touching the rest of the queue, and the new `queues.overrideTotalConcurrencyLimit` and `queues.resetTotalConcurrencyLimit` adjust the cap across all keys. +Adjust queue concurrency at runtime, per key and combined. `queues.overrideConcurrencyLimit` accepts a `concurrencyKey` to raise or lower one key's limit without touching the rest of the queue, and the new `queues.overrideCombinedConcurrencyLimit` and `queues.resetCombinedConcurrencyLimit` adjust the cap across all keys. ```ts import { queues } from "@trigger.dev/sdk"; @@ -11,8 +11,8 @@ import { queues } from "@trigger.dev/sdk"; await queues.overrideConcurrencyLimit("my-queue", 20, { concurrencyKey: "tenant-123" }); await queues.resetConcurrencyLimit("my-queue", { concurrencyKey: "tenant-123" }); -await queues.overrideTotalConcurrencyLimit("my-queue", 100); -await queues.resetTotalConcurrencyLimit("my-queue"); +await queues.overrideCombinedConcurrencyLimit("my-queue", 100); +await queues.resetCombinedConcurrencyLimit("my-queue"); ``` -Overrides survive deploys and reset back to the declared configuration. Enforcement happens server-side on servers with total concurrency limits enabled. +Overrides survive deploys and reset back to the declared configuration. Enforcement happens server-side on servers with combined concurrency limits enabled. diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.combined.override.ts similarity index 100% rename from apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.override.ts rename to apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.combined.override.ts diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.combined.reset.ts similarity index 100% rename from apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.total.reset.ts rename to apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.combined.reset.ts diff --git a/packages/core/src/v3/apiClient/index.ts b/packages/core/src/v3/apiClient/index.ts index e13bb3f76d3..13eea0ef2c1 100644 --- a/packages/core/src/v3/apiClient/index.ts +++ b/packages/core/src/v3/apiClient/index.ts @@ -1704,7 +1704,7 @@ export class ApiClient { ); } - overrideQueueTotalConcurrencyLimit( + overrideQueueCombinedConcurrencyLimit( queue: RetrieveQueueParam, concurrencyLimit: number, requestOptions?: ZodFetchOptions @@ -1716,7 +1716,7 @@ export class ApiClient { return zodfetch( QueueItem, - `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/total/override`, + `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/combined/override`, { method: "POST", headers: this.#getHeaders(false), @@ -1729,7 +1729,7 @@ export class ApiClient { ); } - resetQueueTotalConcurrencyLimit(queue: RetrieveQueueParam, requestOptions?: ZodFetchOptions) { + resetQueueCombinedConcurrencyLimit(queue: RetrieveQueueParam, requestOptions?: ZodFetchOptions) { const type = typeof queue === "string" ? "id" : queue.type; const value = typeof queue === "string" ? queue : queue.name; @@ -1737,7 +1737,7 @@ export class ApiClient { return zodfetch( QueueItem, - `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/total/reset`, + `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/combined/reset`, { method: "POST", headers: this.#getHeaders(false), diff --git a/packages/trigger-sdk/src/v3/queues.ts b/packages/trigger-sdk/src/v3/queues.ts index 7f16d96e58b..f2ec3dd6339 100644 --- a/packages/trigger-sdk/src/v3/queues.ts +++ b/packages/trigger-sdk/src/v3/queues.ts @@ -183,14 +183,14 @@ export function overrideConcurrencyLimit( } /** - * Overrides the total concurrency limit of a queue: the cap on concurrent runs across + * Overrides the combined concurrency limit of a queue: the cap on concurrent runs across * all of its `concurrencyKey` values. * * @param queue - The ID of the queue, or the type and name - * @param concurrencyLimit - The total concurrency limit to apply + * @param concurrencyLimit - The combined concurrency limit to apply * @returns The updated queue state */ -export function overrideTotalConcurrencyLimit( +export function overrideCombinedConcurrencyLimit( queue: RetrieveQueueParam, concurrencyLimit: number, requestOptions?: ApiRequestOptions @@ -200,7 +200,7 @@ export function overrideTotalConcurrencyLimit( const $requestOptions = mergeRequestOptions( { tracer, - name: "queues.overrideTotalConcurrencyLimit()", + name: "queues.overrideCombinedConcurrencyLimit()", icon: "queue", attributes: { ...flattenAttributes({ queue }), @@ -218,16 +218,16 @@ export function overrideTotalConcurrencyLimit( requestOptions ); - return apiClient.overrideQueueTotalConcurrencyLimit(queue, concurrencyLimit, $requestOptions); + return apiClient.overrideQueueCombinedConcurrencyLimit(queue, concurrencyLimit, $requestOptions); } /** - * Resets the total concurrency limit of a queue back to its declared value. + * Resets the combined concurrency limit of a queue back to its declared value. * * @param queue - The ID of the queue, or the type and name * @returns The updated queue state */ -export function resetTotalConcurrencyLimit( +export function resetCombinedConcurrencyLimit( queue: RetrieveQueueParam, requestOptions?: ApiRequestOptions ): ApiPromise { @@ -236,7 +236,7 @@ export function resetTotalConcurrencyLimit( const $requestOptions = mergeRequestOptions( { tracer, - name: "queues.resetTotalConcurrencyLimit()", + name: "queues.resetCombinedConcurrencyLimit()", icon: "queue", attributes: { ...flattenAttributes({ queue }), @@ -254,7 +254,7 @@ export function resetTotalConcurrencyLimit( requestOptions ); - return apiClient.resetQueueTotalConcurrencyLimit(queue, $requestOptions); + return apiClient.resetQueueCombinedConcurrencyLimit(queue, $requestOptions); } /** From 87e8cd1d0f6f224de1bbd2b9417457bc5d1c7e04 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Mon, 31 Aug 2026 12:57:07 +0100 Subject: [PATCH 10/13] refactor(run-engine,webapp,sdk,database): remove per-key concurrency limit overrides Per-key runtime overrides are deferred: composing a per-key hash with pause semantics and keeping it convergent with durable rows across deploys and resets needs its own design. Combined-limit overrides stay. Declared per-key behavior is unchanged (concurrencyLimit applies per key as before). --- .changeset/queue-concurrency-overrides.md | 7 +- .../v3/services/concurrencySystem.server.ts | 179 --------- .../services/createBackgroundWorker.server.ts | 25 -- .../migration.sql | 20 - .../database/prisma/schema.prisma | 20 - internal-packages/run-engine/src/index.ts | 1 - .../run-engine/src/run-queue/keyProducer.ts | 15 - .../tests/concurrencyKeyOverrides.test.ts | 369 ------------------ .../run-engine/src/run-queue/types.ts | 2 - packages/core/src/v3/apiClient/index.ts | 52 --- packages/trigger-sdk/src/v3/queues.ts | 23 +- 11 files changed, 6 insertions(+), 707 deletions(-) delete mode 100644 internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts diff --git a/.changeset/queue-concurrency-overrides.md b/.changeset/queue-concurrency-overrides.md index 6bbe8f909ee..22f79c5e72c 100644 --- a/.changeset/queue-concurrency-overrides.md +++ b/.changeset/queue-concurrency-overrides.md @@ -3,16 +3,13 @@ "@trigger.dev/core": patch --- -Adjust queue concurrency at runtime, per key and combined. `queues.overrideConcurrencyLimit` accepts a `concurrencyKey` to raise or lower one key's limit without touching the rest of the queue, and the new `queues.overrideCombinedConcurrencyLimit` and `queues.resetCombinedConcurrencyLimit` adjust the cap across all keys. +Adjust a queue's combined concurrency limit at runtime. `queues.overrideCombinedConcurrencyLimit` raises or lowers the cap on concurrent runs across all of a queue's `concurrencyKey` values, and `queues.resetCombinedConcurrencyLimit` reverts to the declared configuration. ```ts import { queues } from "@trigger.dev/sdk"; -await queues.overrideConcurrencyLimit("my-queue", 20, { concurrencyKey: "tenant-123" }); -await queues.resetConcurrencyLimit("my-queue", { concurrencyKey: "tenant-123" }); - await queues.overrideCombinedConcurrencyLimit("my-queue", 100); await queues.resetCombinedConcurrencyLimit("my-queue"); ``` -Overrides survive deploys and reset back to the declared configuration. Enforcement happens server-side on servers with combined concurrency limits enabled. +Overrides survive deploys. Enforcement happens server-side on servers with combined concurrency limits enabled. diff --git a/apps/webapp/app/v3/services/concurrencySystem.server.ts b/apps/webapp/app/v3/services/concurrencySystem.server.ts index 571ed5daf64..e1b285136fa 100644 --- a/apps/webapp/app/v3/services/concurrencySystem.server.ts +++ b/apps/webapp/app/v3/services/concurrencySystem.server.ts @@ -9,7 +9,6 @@ import { updateQueueConcurrencyLimits, updateQueueTotalConcurrencyLimits, } from "../runQueue.server"; -import { RunQueueConcurrencyKeyLimitExceededError } from "@internal/run-engine"; import { engine } from "../runEngine.server"; export type ConcurrencySystemOptions = { @@ -109,37 +108,6 @@ export class ConcurrencySystem { .andThen((queue) => syncQueueTotalConcurrencyToEngine(environment, queue)) .andThen((queue) => getQueueStats(environment, queue)); }, - overrideConcurrencyKeyLimit: ( - environment: AuthenticatedEnvironment, - queue: QueueInput, - concurrencyKey: string, - concurrencyLimit: number, - overriddenBy?: User - ) => { - return findQueueFromInput(this.db, environment, queue) - .andThen((queue) => - overrideQueueConcurrencyKeyLimit( - this.db, - environment, - queue, - concurrencyKey, - concurrencyLimit, - overriddenBy - ) - ) - .andThen((queue) => getQueueStats(environment, queue)); - }, - resetConcurrencyKeyLimit: ( - environment: AuthenticatedEnvironment, - queue: QueueInput, - concurrencyKey: string - ) => { - return findQueueFromInput(this.db, environment, queue) - .andThen((queue) => - resetQueueConcurrencyKeyLimit(this.db, environment, queue, concurrencyKey) - ) - .andThen((queue) => getQueueStats(environment, queue)); - }, /** * Recalculates the materialized limit of every percent-based override in the environment * against its CURRENT maximumConcurrencyLimit and syncs changed queues to the run engine. @@ -498,153 +466,6 @@ function syncQueueTotalConcurrencyToEngine( })).andThen(() => okAsync(queue)); } -/** - * Persist first, then enforce: a database failure leaves nothing enforced and the - * request errors cleanly, while an engine failure after persistence leaves a durable - * record and a retry converges (the upsert is idempotent). A cap rejection deletes - * the row unconditionally: the cap only rejects keys absent from the engine hash - * (updates to present keys always succeed), so a cap-rejected row is never backing - * an enforced limit and must not survive as an authoritative override. - */ -function overrideQueueConcurrencyKeyLimit( - db: PrismaClientOrTransaction, - environment: AuthenticatedEnvironment, - queue: TaskQueue, - concurrencyKey: string, - concurrencyLimit: number, - overriddenBy?: User -) { - const maximum = environment.maximumConcurrencyLimit; - - if (!Number.isFinite(concurrencyLimit) || concurrencyLimit < 0) { - return errAsync({ - type: "invalid_override" as const, - message: "Concurrency limit must be a non-negative number", - }); - } - - if (concurrencyLimit > maximum) { - return errAsync({ - type: "concurrency_limit_exceeds_maximum" as const, - message: `Concurrency limit (${concurrencyLimit}) cannot exceed the environment limit (${maximum})`, - }); - } - - return fromPromise( - db.taskQueueConcurrencyKeyOverride.findFirst({ - where: { taskQueueId: queue.id, concurrencyKey }, - select: { id: true }, - }), - (error) => ({ type: "other" as const, cause: error }) - ) - .andThen(() => - fromPromise( - db.taskQueueConcurrencyKeyOverride.upsert({ - where: { taskQueueId_concurrencyKey: { taskQueueId: queue.id, concurrencyKey } }, - create: { - taskQueueId: queue.id, - concurrencyKey, - concurrencyLimit, - overriddenBy: overriddenBy?.id ?? null, - }, - update: { - concurrencyLimit, - overriddenAt: new Date(), - overriddenBy: overriddenBy?.id ?? null, - }, - }), - (error) => ({ - type: "queue_update_failed" as const, - cause: error, - }) - ) - ) - .andThen((written) => - fromPromise( - engine.runQueue.updateQueueConcurrencyKeyLimit( - environment, - queue.name, - concurrencyKey, - concurrencyLimit - ), - (error) => { - if (error instanceof RunQueueConcurrencyKeyLimitExceededError) { - return { type: "too_many_key_overrides" as const, message: error.message }; - } - return { type: "sync_queue_concurrency_to_engine_failed" as const, cause: error }; - } - ).orElse((error) => { - if (error.type === "too_many_key_overrides") { - return fromPromise( - /** - * Deletes only the exact row generation this request wrote, so a - * concurrent request that succeeded after capacity freed keeps its - * durable record. - */ - db.taskQueueConcurrencyKeyOverride.deleteMany({ - where: { id: written.id, overriddenAt: written.overriddenAt }, - }), - () => error - ).andThen(() => errAsync(error)); - } - return errAsync(error); - }) - ) - .andThen(() => okAsync(queue)); -} - -/** - * Enforce first, then persist: removing the engine limit is idempotent, so an engine - * failure leaves the override row in place and a retry converges instead of being - * rejected as not overridden while the old limit is still enforced. - */ -function resetQueueConcurrencyKeyLimit( - db: PrismaClientOrTransaction, - environment: AuthenticatedEnvironment, - queue: TaskQueue, - concurrencyKey: string -) { - return fromPromise( - db.taskQueueConcurrencyKeyOverride.findFirst({ - where: { taskQueueId: queue.id, concurrencyKey }, - select: { id: true, overriddenAt: true }, - }), - (error) => ({ type: "other" as const, cause: error }) - ).andThen((existing) => { - if (!existing) { - return errAsync({ type: "queue_not_overridden" as const }); - } - - return fromPromise( - engine.runQueue.removeQueueConcurrencyKeyLimit(environment, queue.name, concurrencyKey), - (error) => ({ - type: "sync_queue_concurrency_to_engine_failed" as const, - cause: error, - }) - ) - .andThen(() => - fromPromise( - /** - * Deletes only the exact row generation this reset read, so a concurrent - * override that re-wrote the row after the reset began keeps its record - * (its next write, or the deploy-time restore, re-syncs the engine). - */ - db.taskQueueConcurrencyKeyOverride.deleteMany({ - where: { - id: existing.id, - overriddenAt: existing.overriddenAt, - }, - }), - (error) => ({ - type: "queue_update_failed" as const, - cause: error, - }) - ) - ) - .andThen(() => okAsync(queue)); - }); -} - function getQueueStats(environment: AuthenticatedEnvironment, queue: TaskQueue) { return fromPromise( Promise.all([ diff --git a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts index 6799c925af8..147248bd376 100644 --- a/apps/webapp/app/v3/services/createBackgroundWorker.server.ts +++ b/apps/webapp/app/v3/services/createBackgroundWorker.server.ts @@ -582,31 +582,6 @@ async function createWorkerQueue( await removeQueueTotalConcurrencyLimits(environment, taskQueue.name); } - /** - * Restore per-key limit overrides into the engine so a fresh or flushed Redis - * converges back to the durable records on the next deploy. Row-level failures - * are logged rather than failing the deploy; the next deploy retries them. - */ - const keyOverrides = await prisma.taskQueueConcurrencyKeyOverride.findMany({ - where: { taskQueueId: taskQueue.id }, - }); - for (const keyOverride of keyOverrides) { - try { - await engine.runQueue.updateQueueConcurrencyKeyLimit( - environment, - taskQueue.name, - keyOverride.concurrencyKey, - keyOverride.concurrencyLimit - ); - } catch (error) { - logger.error("createWorkerQueue: failed to restore concurrency key override", { - taskQueueId: taskQueue.id, - concurrencyKey: keyOverride.concurrencyKey, - error, - }); - } - } - if (!taskQueue.paused) { if (typeof newConcurrencyLimit === "number") { logger.debug("createWorkerQueue: updating concurrency limit", { diff --git a/internal-packages/database/prisma/migrations/20260829150000_add_concurrency_overrides/migration.sql b/internal-packages/database/prisma/migrations/20260829150000_add_concurrency_overrides/migration.sql index d95778f276d..4b8b6ec6077 100644 --- a/internal-packages/database/prisma/migrations/20260829150000_add_concurrency_overrides/migration.sql +++ b/internal-packages/database/prisma/migrations/20260829150000_add_concurrency_overrides/migration.sql @@ -2,23 +2,3 @@ ALTER TABLE "TaskQueue" ADD COLUMN "totalConcurrencyLimitOverriddenAt" TIMESTAMP(3); ALTER TABLE "TaskQueue" ADD COLUMN "totalConcurrencyLimitOverriddenBy" TEXT; ALTER TABLE "TaskQueue" ADD COLUMN "totalConcurrencyLimitBase" INTEGER; - --- CreateTable -CREATE TABLE "TaskQueueConcurrencyKeyOverride" ( - "id" TEXT NOT NULL, - "taskQueueId" TEXT NOT NULL, - "concurrencyKey" TEXT NOT NULL, - "concurrencyLimit" INTEGER NOT NULL, - "overriddenAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, - "overriddenBy" TEXT, - "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, - "updatedAt" TIMESTAMP(3) NOT NULL, - - CONSTRAINT "TaskQueueConcurrencyKeyOverride_pkey" PRIMARY KEY ("id") -); - --- CreateIndex -CREATE UNIQUE INDEX "TaskQueueConcurrencyKeyOverride_taskQueueId_concurrencyKey_key" ON "TaskQueueConcurrencyKeyOverride"("taskQueueId", "concurrencyKey"); - --- AddForeignKey -ALTER TABLE "TaskQueueConcurrencyKeyOverride" ADD CONSTRAINT "TaskQueueConcurrencyKeyOverride_taskQueueId_fkey" FOREIGN KEY ("taskQueueId") REFERENCES "TaskQueue"("id") ON DELETE CASCADE ON UPDATE CASCADE; diff --git a/internal-packages/database/prisma/schema.prisma b/internal-packages/database/prisma/schema.prisma index 65a84ff1645..406caec57e5 100644 --- a/internal-packages/database/prisma/schema.prisma +++ b/internal-packages/database/prisma/schema.prisma @@ -2001,30 +2001,10 @@ model TaskQueue { tasks BackgroundWorkerTask[] workers BackgroundWorker[] - concurrencyKeyOverrides TaskQueueConcurrencyKeyOverride[] @@unique([runtimeEnvironmentId, name]) } -/// A per-concurrency-key limit override for a queue: the named key value gets this -/// limit instead of the queue's concurrencyLimit. Deleting the row resets the key. -model TaskQueueConcurrencyKeyOverride { - id String @id @default(cuid()) - - taskQueue TaskQueue @relation(fields: [taskQueueId], references: [id], onDelete: Cascade, onUpdate: Cascade) - taskQueueId String - - concurrencyKey String - concurrencyLimit Int - - overriddenAt DateTime @default(now()) - overriddenBy String? - - createdAt DateTime @default(now()) - updatedAt DateTime @updatedAt - - @@unique([taskQueueId, concurrencyKey]) -} enum TaskQueueType { VIRTUAL diff --git a/internal-packages/run-engine/src/index.ts b/internal-packages/run-engine/src/index.ts index c9b4535b280..2c98edf6866 100644 --- a/internal-packages/run-engine/src/index.ts +++ b/internal-packages/run-engine/src/index.ts @@ -1,5 +1,4 @@ export { RunEngine } from "./engine/index.js"; -export { RunQueueConcurrencyKeyLimitExceededError } from "./run-queue/index.js"; export { RunDuplicateIdempotencyKeyError, RunOneTimeUseTokenError, diff --git a/internal-packages/run-engine/src/run-queue/keyProducer.ts b/internal-packages/run-engine/src/run-queue/keyProducer.ts index 7b997043244..98028f5af7b 100644 --- a/internal-packages/run-engine/src/run-queue/keyProducer.ts +++ b/internal-packages/run-engine/src/run-queue/keyProducer.ts @@ -26,7 +26,6 @@ const constants = { RUNNING_COUNTER_PART: "runningCounter", GROUP_CONCURRENCY_PART: "groupConcurrency", TOTAL_CONCURRENCY_LIMIT_PART: "totalConcurrency", - CK_LIMITS_PART: "ckLimits", } as const; export class RunQueueFullKeyProducer implements RunQueueKeyProducer { @@ -367,20 +366,6 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer { return `${this.baseQueueKeyFromQueue(queue)}:${constants.TOTAL_CONCURRENCY_LIMIT_PART}`; } - /** - * HASH of per-concurrency-key limit overrides for a queue. Lives at the base - * queue; each field is the EXACT full ck-variant queue name (the ckIndex ZSET - * member), so reads need no parsing, and values are the raw requested limits - * (readers clamp to the environment limit). - */ - queueCkLimitsKey(env: RunQueueKeyProducerEnvironment, queue: string): string { - return `${this.queueKey(env, queue)}:${constants.CK_LIMITS_PART}`; - } - - queueCkLimitsKeyFromQueue(queue: string): string { - return `${this.baseQueueKeyFromQueue(queue)}:${constants.CK_LIMITS_PART}`; - } - isCkWildcard(queue: string): boolean { return queue.endsWith(":ck:*"); } diff --git a/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts b/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts deleted file mode 100644 index 8c10a8ac2bf..00000000000 --- a/internal-packages/run-engine/src/run-queue/tests/concurrencyKeyOverrides.test.ts +++ /dev/null @@ -1,369 +0,0 @@ -import { redisTest } from "@internal/testcontainers"; -import { trace } from "@internal/tracing"; -import { setTimeout } from "node:timers/promises"; -import { describe } from "vitest"; -import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js"; -import { RunQueue, RunQueueConcurrencyKeyLimitExceededError } from "../index.js"; -import { RunQueueFullKeyProducer } from "../keyProducer.js"; -import type { InputPayload } from "../types.js"; -import { Decimal } from "@trigger.dev/database"; - -const testOptions = { - name: "rq", - tracer: trace.getTracer("rq"), - workers: 1, - defaultEnvConcurrency: 25, - 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, - totalConcurrencyEnabled: boolean, - maxOverrides?: number, - dequeueCount?: number, - gatesEnabled?: boolean -) { - return new RunQueue({ - ...testOptions, - totalConcurrencyEnabled, - gatesEnabled, - maxConcurrencyKeyOverridesPerQueue: maxOverrides, - masterQueueConsumerDequeueCount: dequeueCount, - 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, - }; -} - -async function waitFor(condition: () => Promise, timeoutMs = 20_000): Promise { - const deadline = Date.now() + timeoutMs; - while (Date.now() < deadline) { - if (await condition()) { - return true; - } - await setTimeout(250); - } - return condition(); -} - -vi.setConfig({ testTimeout: 60_000 }); - -describe("RunQueue per-concurrency-key limit overrides", () => { - redisTest( - "a lowered key is capped while other keys keep the queue limit", - async ({ redisContainer }) => { - const queue = createQueue(redisContainer, true); - try { - await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 2); - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 1); - - const now = Date.now(); - const messages = [ - ["ck-a", "a0"], - ["ck-a", "a1"], - ["ck-b", "b0"], - ["ck-b", "b1"], - ] as const; - for (const [i, [ck, id]] of messages.entries()) { - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ runId: id, concurrencyKey: ck, timestamp: now - 1000 + i }), - workerQueue: "main", - }); - } - - const settled = await waitFor(async () => { - const a = await queue.currentConcurrencyOfQueue( - authenticatedEnvDev, - "task/my-task", - "ck-a" - ); - const b = await queue.currentConcurrencyOfQueue( - authenticatedEnvDev, - "task/my-task", - "ck-b" - ); - return a === 1 && b === 2; - }); - expect(settled).toBe(true); - - await setTimeout(2000); - expect( - await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "task/my-task", "ck-a") - ).toBe(1); - expect(await queue.lengthOfQueue(authenticatedEnvDev, "task/my-task")).toBe(1); - } finally { - await queue.quit(); - } - } - ); - - redisTest("a raised key admits past the queue limit", async ({ redisContainer }) => { - const queue = createQueue(redisContainer, true); - try { - await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 1); - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 3); - - const now = Date.now(); - for (const i of [0, 1, 2]) { - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ - runId: `a${i}`, - concurrencyKey: "ck-a", - timestamp: now - 1000 + i, - }), - workerQueue: "main", - }); - } - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ runId: "b0", concurrencyKey: "ck-b", timestamp: now - 500 }), - workerQueue: "main", - }); - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ runId: "b1", concurrencyKey: "ck-b", timestamp: now - 499 }), - workerQueue: "main", - }); - - const settled = await waitFor(async () => { - const a = await queue.currentConcurrencyOfQueue( - authenticatedEnvDev, - "task/my-task", - "ck-a" - ); - const b = await queue.currentConcurrencyOfQueue( - authenticatedEnvDev, - "task/my-task", - "ck-b" - ); - return a === 3 && b === 1; - }); - expect(settled).toBe(true); - - await setTimeout(2000); - expect(await queue.lengthOfQueue(authenticatedEnvDev, "task/my-task")).toBe(1); - } finally { - await queue.quit(); - } - }); - - redisTest("removing an override restores the queue limit", async ({ redisContainer }) => { - const queue = createQueue(redisContainer, true); - try { - await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 2); - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 1); - expect(await queue.getQueueConcurrencyKeyLimits(authenticatedEnvDev, "task/my-task")).toEqual( - { "ck-a": 1 } - ); - - await queue.removeQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a"); - expect(await queue.getQueueConcurrencyKeyLimits(authenticatedEnvDev, "task/my-task")).toEqual( - {} - ); - - const now = Date.now(); - for (const i of [0, 1]) { - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ - runId: `a${i}`, - concurrencyKey: "ck-a", - timestamp: now - 1000 + i, - }), - workerQueue: "main", - }); - } - - const settled = await waitFor( - async () => - (await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "task/my-task", "ck-a")) === 2 - ); - expect(settled).toBe(true); - } finally { - await queue.quit(); - } - }); - - redisTest("the per-queue override count is capped", async ({ redisContainer }) => { - const queue = createQueue(redisContainer, true, 2); - try { - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 1); - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-b", 1); - - await expect( - queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-c", 1) - ).rejects.toThrow(RunQueueConcurrencyKeyLimitExceededError); - - /** Updates to existing keys always succeed at the cap. */ - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 4); - expect(await queue.getQueueConcurrencyKeyLimits(authenticatedEnvDev, "task/my-task")).toEqual( - { "ck-a": 4, "ck-b": 1 } - ); - } finally { - await queue.quit(); - } - }); - - redisTest( - "blocked keys cannot pin the candidate window and starve later keys", - async ({ redisContainer }) => { - /** dequeueCount 2 makes the candidate window 6 variants wide. */ - const queue = createQueue(redisContainer, true, undefined, 2); - try { - await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 5); - - /** - * Ten zero-limit keys with OLDER messages fill the window many times over; - * without the blocked-key backoff the runnable key behind them would never - * be examined. - */ - const now = Date.now(); - for (let i = 0; i < 10; i++) { - const ck = `blocked-${i}`; - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", ck, 0); - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ - runId: `b${i}`, - concurrencyKey: ck, - timestamp: now - 10_000 + i, - }), - workerQueue: "main", - }); - } - - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ - runId: "good-0", - concurrencyKey: "ck-good", - timestamp: now - 500, - }), - workerQueue: "main", - }); - - const goodAdmitted = await waitFor( - async () => - (await queue.currentConcurrencyOfQueue( - authenticatedEnvDev, - "task/my-task", - "ck-good" - )) === 1, - 30_000 - ); - expect(goodAdmitted).toBe(true); - expect(await queue.lengthOfQueue(authenticatedEnvDev, "task/my-task")).toBe(10); - } finally { - await queue.quit(); - } - } - ); - - redisTest("gate admission honors the gate queue per-key override", async ({ redisContainer }) => { - const queue = createQueue(redisContainer, true, undefined, undefined, true); - try { - await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 5); - await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "tenant", 1); - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "tenant", "acme", 2); - - const now = Date.now(); - for (const [i, ck] of ["ck-a", "ck-b", "ck-c"].entries()) { - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ - runId: `r${i}`, - concurrencyKey: ck, - timestamp: now - 1000 + i, - gates: [{ queue: "tenant", concurrencyKey: "acme" }], - }), - workerQueue: "main", - }); - } - - /** The declared gate limit is 1; the override raises acme to 2. */ - const twoAdmitted = await waitFor( - async () => - (await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "tenant", "acme")) === 2 - ); - expect(twoAdmitted).toBe(true); - - await setTimeout(2000); - expect(await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "tenant", "acme")).toBe(2); - expect(await queue.lengthOfQueue(authenticatedEnvDev, "task/my-task")).toBe(1); - } finally { - await queue.quit(); - } - }); - - redisTest("overrides are ignored when disabled", async ({ redisContainer }) => { - const queue = createQueue(redisContainer, false); - try { - await queue.updateQueueConcurrencyLimits(authenticatedEnvDev, "task/my-task", 2); - await queue.updateQueueConcurrencyKeyLimit(authenticatedEnvDev, "task/my-task", "ck-a", 1); - - const now = Date.now(); - for (const i of [0, 1]) { - await queue.enqueueMessage({ - env: authenticatedEnvDev, - message: makeMessage({ - runId: `a${i}`, - concurrencyKey: "ck-a", - timestamp: now - 1000 + i, - }), - workerQueue: "main", - }); - } - - const settled = await waitFor( - async () => - (await queue.currentConcurrencyOfQueue(authenticatedEnvDev, "task/my-task", "ck-a")) === 2 - ); - expect(settled).toBe(true); - } finally { - await queue.quit(); - } - }); -}); diff --git a/internal-packages/run-engine/src/run-queue/types.ts b/internal-packages/run-engine/src/run-queue/types.ts index b21a3f66368..2cbfe40c775 100644 --- a/internal-packages/run-engine/src/run-queue/types.ts +++ b/internal-packages/run-engine/src/run-queue/types.ts @@ -110,8 +110,6 @@ export interface RunQueueKeyProducer { queueGroupConcurrencyKeyFromQueue(queue: string): string; queueTotalConcurrencyLimitKey(env: RunQueueKeyProducerEnvironment, queue: string): string; queueTotalConcurrencyLimitKeyFromQueue(queue: string): string; - queueCkLimitsKey(env: RunQueueKeyProducerEnvironment, queue: string): string; - queueCkLimitsKeyFromQueue(queue: string): string; //env oncurrency envCurrentConcurrencyKey(env: EnvDescriptor): string; diff --git a/packages/core/src/v3/apiClient/index.ts b/packages/core/src/v3/apiClient/index.ts index 13eea0ef2c1..4acb1c40089 100644 --- a/packages/core/src/v3/apiClient/index.ts +++ b/packages/core/src/v3/apiClient/index.ts @@ -1749,58 +1749,6 @@ export class ApiClient { ); } - overrideQueueConcurrencyKeyLimit( - queue: RetrieveQueueParam, - concurrencyKey: string, - concurrencyLimit: number, - requestOptions?: ZodFetchOptions - ) { - const type = typeof queue === "string" ? "id" : queue.type; - const value = typeof queue === "string" ? queue : queue.name; - - const encodedValue = encodeURIComponent(value.replace(/\//g, "%2F")); - - return zodfetch( - QueueItem, - `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/key/override`, - { - method: "POST", - headers: this.#getHeaders(false), - body: JSON.stringify({ - type, - concurrencyKey, - concurrencyLimit, - }), - }, - mergeRequestOptions(this.defaultRequestOptions, requestOptions) - ); - } - - resetQueueConcurrencyKeyLimit( - queue: RetrieveQueueParam, - concurrencyKey: string, - requestOptions?: ZodFetchOptions - ) { - const type = typeof queue === "string" ? "id" : queue.type; - const value = typeof queue === "string" ? queue : queue.name; - - const encodedValue = encodeURIComponent(value.replace(/\//g, "%2F")); - - return zodfetch( - QueueItem, - `${this.baseUrl}/api/v1/queues/${encodedValue}/concurrency/key/reset`, - { - method: "POST", - headers: this.#getHeaders(false), - body: JSON.stringify({ - type, - concurrencyKey, - }), - }, - mergeRequestOptions(this.defaultRequestOptions, requestOptions) - ); - } - subscribeToRun( runId: string, options?: { diff --git a/packages/trigger-sdk/src/v3/queues.ts b/packages/trigger-sdk/src/v3/queues.ts index f2ec3dd6339..1292cde1a3c 100644 --- a/packages/trigger-sdk/src/v3/queues.ts +++ b/packages/trigger-sdk/src/v3/queues.ts @@ -144,10 +144,9 @@ export function pause( export function overrideConcurrencyLimit( queue: RetrieveQueueParam, concurrencyLimit: number, - options?: ApiRequestOptions & { concurrencyKey?: string } + requestOptions?: ApiRequestOptions ): ApiPromise { const apiClient = apiClientManager.clientOrThrow(); - const { concurrencyKey, ...requestOptions } = options ?? {}; const $requestOptions = mergeRequestOptions( { @@ -155,7 +154,7 @@ export function overrideConcurrencyLimit( name: "queues.overrideConcurrencyLimit()", icon: "queue", attributes: { - ...flattenAttributes({ queue, concurrencyKey }), + ...flattenAttributes({ queue }), ...accessoryAttributes({ items: [ { @@ -170,15 +169,6 @@ export function overrideConcurrencyLimit( requestOptions ); - if (concurrencyKey !== undefined) { - return apiClient.overrideQueueConcurrencyKeyLimit( - queue, - concurrencyKey, - concurrencyLimit, - $requestOptions - ); - } - return apiClient.overrideQueueConcurrencyLimit(queue, concurrencyLimit, $requestOptions); } @@ -265,10 +255,9 @@ export function resetCombinedConcurrencyLimit( */ export function resetConcurrencyLimit( queue: RetrieveQueueParam, - options?: ApiRequestOptions & { concurrencyKey?: string } + requestOptions?: ApiRequestOptions ): ApiPromise { const apiClient = apiClientManager.clientOrThrow(); - const { concurrencyKey, ...requestOptions } = options ?? {}; const $requestOptions = mergeRequestOptions( { @@ -276,7 +265,7 @@ export function resetConcurrencyLimit( name: "queues.resetConcurrencyLimit()", icon: "queue", attributes: { - ...flattenAttributes({ queue, concurrencyKey }), + ...flattenAttributes({ queue }), ...accessoryAttributes({ items: [ { @@ -291,10 +280,6 @@ export function resetConcurrencyLimit( requestOptions ); - if (concurrencyKey !== undefined) { - return apiClient.resetQueueConcurrencyKeyLimit(queue, concurrencyKey, $requestOptions); - } - return apiClient.resetQueueConcurrencyLimit(queue, $requestOptions); } From fac9403b124933384884642bcdce3c5d5bf812c8 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Mon, 31 Aug 2026 12:58:37 +0100 Subject: [PATCH 11/13] refactor(run-engine,webapp): drop per-key override reads and endpoints Removes the per-key override endpoints, the override-aware admit and gauge reads, and the per-key limit column, following the removal of runtime per-key overrides from this stack. --- ...es.$queueParam.concurrency.key.override.ts | 105 ------------------ ...ueues.$queueParam.concurrency.key.reset.ts | 97 ---------------- 2 files changed, 202 deletions(-) delete mode 100644 apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts delete mode 100644 apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts deleted file mode 100644 index 5a37b4526ec..00000000000 --- a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.override.ts +++ /dev/null @@ -1,105 +0,0 @@ -import { json } from "@remix-run/server-runtime"; -import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3"; -import { z } from "zod"; -import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server"; -import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server"; -import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server"; - -const BodySchema = z.object({ - type: RetrieveQueueType.default("id"), - concurrencyKey: z.string().min(1).max(128), - concurrencyLimit: z.number().int().min(0).max(100000), -}); - -const route = createActionApiRoute( - { - body: BodySchema, - params: z.object({ - queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")), - }), - authorization: { - action: "write", - resource: () => ({ type: "queues" }), - }, - }, - async ({ params, body, authentication }) => { - const input: RetrieveQueueParam = - body.type === "id" - ? params.queueParam - : { - type: body.type, - name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"), - }; - - return concurrencySystem.queues - .overrideConcurrencyKeyLimit( - authentication.environment, - input, - body.concurrencyKey, - body.concurrencyLimit - ) - .match( - (queue) => { - return json( - toQueueItem({ - friendlyId: queue.friendlyId, - name: queue.name, - type: queue.type, - running: queue.running, - queued: queue.queued, - concurrencyLimit: queue.concurrencyLimit, - concurrencyLimitBase: queue.concurrencyLimitBase, - concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt, - concurrencyLimitOverriddenBy: null, - paused: queue.paused, - }), - { status: 200 } - ); - }, - (error) => { - switch (error.type) { - case "queue_not_found": { - return json({ error: "Queue not found" }, { status: 404 }); - } - case "invalid_override": - case "concurrency_limit_exceeds_maximum": - case "too_many_key_overrides": { - return json({ error: error.message }, { status: 400 }); - } - case "queue_update_failed": { - return json( - { error: "Failed to update queue concurrency key limit" }, - { status: 500 } - ); - } - case "sync_queue_concurrency_to_engine_failed": { - return json({ error: "Failed to sync the concurrency key limit" }, { status: 500 }); - } - case "get_queue_stats_failed": { - return json({ error: "Failed to read queue stats" }, { status: 500 }); - } - case "other": { - return json( - { error: "Failed to update queue concurrency key limit" }, - { - status: 500, - } - ); - } - default: { - return json( - { error: "Failed to update queue concurrency key limit" }, - { - status: 500, - } - ); - } - } - } - ); - } -); - -export const action = route.action; -/** The builder's loader answers non-POST methods with a 405. */ -export const loader = route.loader; diff --git a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts b/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts deleted file mode 100644 index 51d14642e2c..00000000000 --- a/apps/webapp/app/routes/api.v1.queues.$queueParam.concurrency.key.reset.ts +++ /dev/null @@ -1,97 +0,0 @@ -import { json } from "@remix-run/server-runtime"; -import { type RetrieveQueueParam, RetrieveQueueType } from "@trigger.dev/core/v3"; -import { z } from "zod"; -import { toQueueItem } from "~/presenters/v3/QueueRetrievePresenter.server"; -import { createActionApiRoute } from "~/services/routeBuilders/apiBuilder.server"; -import { concurrencySystem } from "~/v3/services/concurrencySystemInstance.server"; - -const BodySchema = z.object({ - type: RetrieveQueueType.default("id"), - concurrencyKey: z.string().min(1).max(128), -}); - -const route = createActionApiRoute( - { - body: BodySchema, - params: z.object({ - queueParam: z.string().transform((val) => val.replace(/%2F/g, "/")), - }), - authorization: { - action: "write", - resource: () => ({ type: "queues" }), - }, - }, - async ({ params, body, authentication }) => { - const input: RetrieveQueueParam = - body.type === "id" - ? params.queueParam - : { - type: body.type, - name: decodeURIComponent(params.queueParam).replace(/%2F/g, "/"), - }; - - return concurrencySystem.queues - .resetConcurrencyKeyLimit(authentication.environment, input, body.concurrencyKey) - .match( - (queue) => { - return json( - toQueueItem({ - friendlyId: queue.friendlyId, - name: queue.name, - type: queue.type, - running: queue.running, - queued: queue.queued, - concurrencyLimit: queue.concurrencyLimit, - concurrencyLimitBase: queue.concurrencyLimitBase, - concurrencyLimitOverriddenAt: queue.concurrencyLimitOverriddenAt, - concurrencyLimitOverriddenBy: null, - paused: queue.paused, - }), - { status: 200 } - ); - }, - (error) => { - switch (error.type) { - case "queue_not_found": { - return json({ error: "Queue not found" }, { status: 404 }); - } - case "queue_not_overridden": { - return json( - { error: "This concurrency key does not have an override" }, - { status: 400 } - ); - } - case "queue_update_failed": { - return json({ error: "Failed to reset the concurrency key limit" }, { status: 500 }); - } - case "sync_queue_concurrency_to_engine_failed": { - return json({ error: "Failed to sync the concurrency key limit" }, { status: 500 }); - } - case "get_queue_stats_failed": { - return json({ error: "Failed to read queue stats" }, { status: 500 }); - } - case "other": { - return json( - { error: "Failed to reset the concurrency key limit" }, - { - status: 500, - } - ); - } - default: { - return json( - { error: "Failed to reset the concurrency key limit" }, - { - status: 500, - } - ); - } - } - } - ); - } -); - -export const action = route.action; -/** The builder's loader answers non-POST methods with a 405. */ -export const loader = route.loader; From 205b6f264eacb48e69f4a0be76844da3cb900123 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Mon, 31 Aug 2026 13:10:59 +0100 Subject: [PATCH 12/13] fix(run-engine): keep the ck-limits key builders while the Lua reads remain The per-key admit reads are removed in the PR above alongside the rest of the engine-internal plumbing; the key builders they reference stay until then so every level of the stack compiles. --- internal-packages/run-engine/src/run-queue/keyProducer.ts | 8 ++++++++ internal-packages/run-engine/src/run-queue/types.ts | 3 +++ 2 files changed, 11 insertions(+) diff --git a/internal-packages/run-engine/src/run-queue/keyProducer.ts b/internal-packages/run-engine/src/run-queue/keyProducer.ts index 98028f5af7b..120e04f8c38 100644 --- a/internal-packages/run-engine/src/run-queue/keyProducer.ts +++ b/internal-packages/run-engine/src/run-queue/keyProducer.ts @@ -366,6 +366,14 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer { return `${this.baseQueueKeyFromQueue(queue)}:${constants.TOTAL_CONCURRENCY_LIMIT_PART}`; } + queueCkLimitsKey(env: RunQueueKeyProducerEnvironment, queue: string): string { + return `${this.queueKey(env, queue)}:ckLimits`; + } + + queueCkLimitsKeyFromQueue(queue: string): string { + return `${this.baseQueueKeyFromQueue(queue)}:ckLimits`; + } + isCkWildcard(queue: string): boolean { return queue.endsWith(":ck:*"); } diff --git a/internal-packages/run-engine/src/run-queue/types.ts b/internal-packages/run-engine/src/run-queue/types.ts index 2cbfe40c775..2961b642314 100644 --- a/internal-packages/run-engine/src/run-queue/types.ts +++ b/internal-packages/run-engine/src/run-queue/types.ts @@ -111,6 +111,9 @@ export interface RunQueueKeyProducer { queueTotalConcurrencyLimitKey(env: RunQueueKeyProducerEnvironment, queue: string): string; queueTotalConcurrencyLimitKeyFromQueue(queue: string): string; + queueCkLimitsKey(env: RunQueueKeyProducerEnvironment, queue: string): string; + queueCkLimitsKeyFromQueue(queue: string): string; + //env oncurrency envCurrentConcurrencyKey(env: EnvDescriptor): string; envCurrentConcurrencyKey(env: RunQueueKeyProducerEnvironment): string; From 43c724ed76859d45e10d302573a058b79607d5e2 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Mon, 31 Aug 2026 14:02:52 +0100 Subject: [PATCH 13/13] fix(webapp): combined override error messages use the public name --- apps/webapp/app/v3/services/concurrencySystem.server.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/apps/webapp/app/v3/services/concurrencySystem.server.ts b/apps/webapp/app/v3/services/concurrencySystem.server.ts index e1b285136fa..09cc5f9ce11 100644 --- a/apps/webapp/app/v3/services/concurrencySystem.server.ts +++ b/apps/webapp/app/v3/services/concurrencySystem.server.ts @@ -359,14 +359,14 @@ function overrideQueueTotalConcurrencyLimit( if (!Number.isFinite(totalConcurrencyLimit) || totalConcurrencyLimit < 0) { return errAsync({ type: "invalid_override" as const, - message: "Total concurrency limit must be a non-negative number", + message: "Combined concurrency limit must be a non-negative number", }); } if (totalConcurrencyLimit > maximum) { return errAsync({ type: "concurrency_limit_exceeds_maximum" as const, - message: `Total concurrency limit (${totalConcurrencyLimit}) cannot exceed the environment limit (${maximum})`, + message: `Combined concurrency limit (${totalConcurrencyLimit}) cannot exceed the environment limit (${maximum})`, }); }