fix(run-queue): prevent concurrency keys from bloating master queue shards (#3219)
Queues with concurrency keys now appear as a single entry in the master queue instead of one entry per key. This prevents high-CK-count tenants from consuming the entire `parentQueueLimit` window and starving other tenants on the same shard. A new per-queue **CK index** (sorted set) tracks active concurrency key sub-queues. The master queue gets one `:ck:*` wildcard entry per base queue. Dequeuing from that entry round-robins across sub-queues, maintaining per-CK concurrency tracking and fairness. All existing operations (enqueue, dequeue, ack, nack, DLQ, TTL expiry) are CK-index-aware and keep the index consistent. Old-format entries drain naturally during rollout — no migration step needed, single deploy.
This commit is contained in:
@@ -0,0 +1,6 @@
|
||||
---
|
||||
area: webapp
|
||||
type: fix
|
||||
---
|
||||
|
||||
Concurrency-keyed queues now use a single master queue entry per base queue instead of one entry per key. Prevents high-CK-count tenants from consuming the entire parentQueueLimit window and starving other tenants on the same shard.
|
||||
File diff suppressed because it is too large
Load Diff
@@ -17,6 +17,7 @@ const constants = {
|
||||
DEAD_LETTER_QUEUE_PART: "deadLetter",
|
||||
MASTER_QUEUE_PART: "masterQueue",
|
||||
WORKER_QUEUE_PART: "workerQueue",
|
||||
CK_INDEX_PART: "ckIndex",
|
||||
} as const;
|
||||
|
||||
export class RunQueueFullKeyProducer implements RunQueueKeyProducer {
|
||||
@@ -305,6 +306,23 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer {
|
||||
return ["ttl", "shard", shard.toString()].join(":");
|
||||
}
|
||||
|
||||
ckIndexKeyFromQueue(queue: string): string {
|
||||
const baseQueue = queue.replace(/:ck:.+$/, "");
|
||||
return `${baseQueue}:${constants.CK_INDEX_PART}`;
|
||||
}
|
||||
|
||||
baseQueueKeyFromQueue(queue: string): string {
|
||||
return queue.replace(/:ck:.+$/, "");
|
||||
}
|
||||
|
||||
isCkWildcard(queue: string): boolean {
|
||||
return queue.endsWith(":ck:*");
|
||||
}
|
||||
|
||||
toCkWildcard(queue: string): string {
|
||||
return queue.replace(/:ck:.+$/, ":ck:*");
|
||||
}
|
||||
|
||||
descriptorFromQueue(queue: string): QueueDescriptor {
|
||||
const parts = queue.split(":");
|
||||
return {
|
||||
|
||||
@@ -0,0 +1,521 @@
|
||||
import { assertNonNullable, redisTest } from "@internal/testcontainers";
|
||||
import { trace } from "@internal/tracing";
|
||||
import { Logger } from "@trigger.dev/core/logger";
|
||||
import { describe } from "node:test";
|
||||
import { setTimeout } from "node:timers/promises";
|
||||
import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js";
|
||||
import { RunQueue } from "../index.js";
|
||||
import { RunQueueFullKeyProducer } from "../keyProducer.js";
|
||||
import { InputPayload } from "../types.js";
|
||||
import { Decimal } from "@trigger.dev/database";
|
||||
|
||||
const testOptions = {
|
||||
name: "rq",
|
||||
tracer: trace.getTracer("rq"),
|
||||
workers: 1,
|
||||
defaultEnvConcurrency: 25,
|
||||
logger: new Logger("RunQueue", "warn"),
|
||||
retryOptions: {
|
||||
maxAttempts: 5,
|
||||
factor: 1.1,
|
||||
minTimeoutInMs: 100,
|
||||
maxTimeoutInMs: 1_000,
|
||||
randomize: true,
|
||||
},
|
||||
keys: new RunQueueFullKeyProducer(),
|
||||
};
|
||||
|
||||
const authenticatedEnvDev = {
|
||||
id: "e1234",
|
||||
type: "DEVELOPMENT" as const,
|
||||
maximumConcurrencyLimit: 10,
|
||||
concurrencyLimitBurstFactor: new Decimal(2.0),
|
||||
project: { id: "p1234" },
|
||||
organization: { id: "o1234" },
|
||||
};
|
||||
|
||||
function createQueue(redisContainer: any) {
|
||||
return new RunQueue({
|
||||
...testOptions,
|
||||
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> = {}): 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,
|
||||
};
|
||||
}
|
||||
|
||||
vi.setConfig({ testTimeout: 60_000 });
|
||||
|
||||
describe("CK Index", () => {
|
||||
redisTest(
|
||||
"enqueue with CK creates CK index entry and :ck:* master queue entry",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = createQueue(redisContainer);
|
||||
try {
|
||||
const msg = makeMessage({ runId: "r1", concurrencyKey: "ck-a" });
|
||||
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: msg,
|
||||
workerQueue: authenticatedEnvDev.id,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
|
||||
// Check that the CK-specific sorted set has the message
|
||||
const queueLength = await queue.lengthOfQueue(
|
||||
authenticatedEnvDev,
|
||||
msg.queue,
|
||||
msg.concurrencyKey
|
||||
);
|
||||
expect(queueLength).toBe(1);
|
||||
|
||||
// Check master queue: should have :ck:* entry, not :ck:ck-a
|
||||
const masterQueueKey = testOptions.keys.masterQueueKeyForShard(
|
||||
testOptions.keys.masterQueueShardForEnvironment(msg.environmentId, 2)
|
||||
);
|
||||
const masterMembers = await queue.redis.zrange(
|
||||
masterQueueKey,
|
||||
0,
|
||||
-1,
|
||||
"WITHSCORES"
|
||||
);
|
||||
// Should have exactly one member ending with :ck:*
|
||||
const ckWildcardMembers = masterMembers.filter(
|
||||
(m, i) => i % 2 === 0 && m.endsWith(":ck:*")
|
||||
);
|
||||
expect(ckWildcardMembers.length).toBe(1);
|
||||
|
||||
// Should NOT have :ck:ck-a member
|
||||
const oldFormatMembers = masterMembers.filter(
|
||||
(m, i) => i % 2 === 0 && m.endsWith(":ck:ck-a")
|
||||
);
|
||||
expect(oldFormatMembers.length).toBe(0);
|
||||
|
||||
// Check CK index has the CK queue
|
||||
const ckIndexKey = testOptions.keys.ckIndexKeyFromQueue(
|
||||
testOptions.keys.queueKey(authenticatedEnvDev, msg.queue, msg.concurrencyKey)
|
||||
);
|
||||
const ckIndexMembers = await queue.redis.zrange(ckIndexKey, 0, -1);
|
||||
expect(ckIndexMembers.length).toBe(1);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"multiple CKs result in single master queue entry",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = createQueue(redisContainer);
|
||||
try {
|
||||
const now = Date.now();
|
||||
const msg1 = makeMessage({
|
||||
runId: "r1",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now,
|
||||
});
|
||||
const msg2 = makeMessage({
|
||||
runId: "r2",
|
||||
concurrencyKey: "ck-b",
|
||||
timestamp: now + 100,
|
||||
});
|
||||
const msg3 = makeMessage({
|
||||
runId: "r3",
|
||||
concurrencyKey: "ck-c",
|
||||
timestamp: now + 200,
|
||||
});
|
||||
|
||||
for (const msg of [msg1, msg2, msg3]) {
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: msg,
|
||||
workerQueue: authenticatedEnvDev.id,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
}
|
||||
|
||||
// Master queue should have exactly ONE entry (the :ck:* wildcard)
|
||||
const masterQueueKey = testOptions.keys.masterQueueKeyForShard(
|
||||
testOptions.keys.masterQueueShardForEnvironment(msg1.environmentId, 2)
|
||||
);
|
||||
const masterMembers = await queue.redis.zrange(
|
||||
masterQueueKey,
|
||||
0,
|
||||
-1
|
||||
);
|
||||
// Filter to only members for our queue
|
||||
const ourMembers = masterMembers.filter((m) =>
|
||||
m.includes("queue:task/my-task")
|
||||
);
|
||||
expect(ourMembers.length).toBe(1);
|
||||
expect(ourMembers[0]).toContain(":ck:*");
|
||||
|
||||
// CK index should have 3 entries
|
||||
const ckIndexKey = testOptions.keys.ckIndexKeyFromQueue(
|
||||
testOptions.keys.queueKey(authenticatedEnvDev, msg1.queue, msg1.concurrencyKey)
|
||||
);
|
||||
const ckIndexMembers = await queue.redis.zrange(ckIndexKey, 0, -1);
|
||||
expect(ckIndexMembers.length).toBe(3);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"dequeue from CK queue distributes across sub-queues",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = createQueue(redisContainer);
|
||||
try {
|
||||
const now = Date.now() - 1000; // In the past so they're ready
|
||||
const msg1 = makeMessage({
|
||||
runId: "r1",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now,
|
||||
});
|
||||
const msg2 = makeMessage({
|
||||
runId: "r2",
|
||||
concurrencyKey: "ck-b",
|
||||
timestamp: now + 1,
|
||||
});
|
||||
const msg3 = makeMessage({
|
||||
runId: "r3",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now + 2,
|
||||
});
|
||||
|
||||
for (const msg of [msg1, msg2, msg3]) {
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: msg,
|
||||
workerQueue: authenticatedEnvDev.id,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
}
|
||||
|
||||
// Dequeue via the master queue consumer
|
||||
const shard = testOptions.keys.masterQueueShardForEnvironment(
|
||||
msg1.environmentId,
|
||||
2
|
||||
);
|
||||
const messages = await queue.testDequeueFromMasterQueue(shard, msg1.environmentId, 10);
|
||||
|
||||
// Should dequeue messages from both CK sub-queues
|
||||
expect(messages).toBeDefined();
|
||||
// We should get at least 2 messages (one from each CK)
|
||||
// The exact order depends on CK index scoring
|
||||
expect(messages!.length).toBeGreaterThanOrEqual(2);
|
||||
|
||||
const dequeuedRunIds = messages!.map((m: any) => m.messageId);
|
||||
// r1 (ck-a, oldest) and r2 (ck-b) should be dequeued
|
||||
expect(dequeuedRunIds).toContain("r1");
|
||||
expect(dequeuedRunIds).toContain("r2");
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"empty CK sub-queue is removed from CK index",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = createQueue(redisContainer);
|
||||
try {
|
||||
const now = Date.now() - 1000;
|
||||
const msg1 = makeMessage({
|
||||
runId: "r1",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now,
|
||||
});
|
||||
const msg2 = makeMessage({
|
||||
runId: "r2",
|
||||
concurrencyKey: "ck-b",
|
||||
timestamp: now + 1,
|
||||
});
|
||||
|
||||
for (const msg of [msg1, msg2]) {
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: msg,
|
||||
workerQueue: authenticatedEnvDev.id,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
}
|
||||
|
||||
// CK index should have 2 entries initially
|
||||
const ckIndexKey = testOptions.keys.ckIndexKeyFromQueue(
|
||||
testOptions.keys.queueKey(authenticatedEnvDev, msg1.queue, msg1.concurrencyKey)
|
||||
);
|
||||
let ckIndexMembers = await queue.redis.zrange(ckIndexKey, 0, -1);
|
||||
expect(ckIndexMembers.length).toBe(2);
|
||||
|
||||
// Dequeue both messages
|
||||
const shard = testOptions.keys.masterQueueShardForEnvironment(
|
||||
msg1.environmentId,
|
||||
2
|
||||
);
|
||||
await queue.testDequeueFromMasterQueue(shard, msg1.environmentId, 10);
|
||||
|
||||
// CK index should be empty (both sub-queues drained)
|
||||
ckIndexMembers = await queue.redis.zrange(ckIndexKey, 0, -1);
|
||||
expect(ckIndexMembers.length).toBe(0);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"empty CK index removes :ck:* from master queue",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = createQueue(redisContainer);
|
||||
try {
|
||||
const now = Date.now() - 1000;
|
||||
const msg = makeMessage({
|
||||
runId: "r1",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now,
|
||||
});
|
||||
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: msg,
|
||||
workerQueue: authenticatedEnvDev.id,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
|
||||
const masterQueueKey = testOptions.keys.masterQueueKeyForShard(
|
||||
testOptions.keys.masterQueueShardForEnvironment(msg.environmentId, 2)
|
||||
);
|
||||
|
||||
// Master queue should have :ck:* entry
|
||||
let masterMembers = await queue.redis.zrange(masterQueueKey, 0, -1);
|
||||
expect(masterMembers.length).toBe(1);
|
||||
|
||||
// Dequeue the message
|
||||
const shard = testOptions.keys.masterQueueShardForEnvironment(
|
||||
msg.environmentId,
|
||||
2
|
||||
);
|
||||
await queue.testDequeueFromMasterQueue(shard, msg.environmentId, 10);
|
||||
|
||||
// Master queue should be empty
|
||||
masterMembers = await queue.redis.zrange(masterQueueKey, 0, -1);
|
||||
expect(masterMembers.length).toBe(0);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"mixed CK and non-CK queues in same shard",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = createQueue(redisContainer);
|
||||
try {
|
||||
const now = Date.now() - 1000;
|
||||
|
||||
// Non-CK message
|
||||
const msgNoCk = makeMessage({
|
||||
runId: "r-no-ck",
|
||||
timestamp: now,
|
||||
});
|
||||
|
||||
// CK messages
|
||||
const msgCk1 = makeMessage({
|
||||
runId: "r-ck-1",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now + 1,
|
||||
});
|
||||
const msgCk2 = makeMessage({
|
||||
runId: "r-ck-2",
|
||||
concurrencyKey: "ck-b",
|
||||
timestamp: now + 2,
|
||||
});
|
||||
|
||||
for (const msg of [msgNoCk, msgCk1, msgCk2]) {
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: msg,
|
||||
workerQueue: authenticatedEnvDev.id,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
}
|
||||
|
||||
// Master queue should have 2 entries: one non-CK queue and one :ck:*
|
||||
const masterQueueKey = testOptions.keys.masterQueueKeyForShard(
|
||||
testOptions.keys.masterQueueShardForEnvironment(msgNoCk.environmentId, 2)
|
||||
);
|
||||
const masterMembers = await queue.redis.zrange(masterQueueKey, 0, -1);
|
||||
expect(masterMembers.length).toBe(2);
|
||||
|
||||
// One should be the non-CK queue, one should be :ck:*
|
||||
const ckWildcard = masterMembers.filter((m) => m.endsWith(":ck:*"));
|
||||
const nonCk = masterMembers.filter(
|
||||
(m) => !m.includes(":ck:")
|
||||
);
|
||||
expect(ckWildcard.length).toBe(1);
|
||||
expect(nonCk.length).toBe(1);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"acknowledge CK message rebalances CK index and master queue",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = createQueue(redisContainer);
|
||||
try {
|
||||
const now = Date.now() - 1000;
|
||||
const msg1 = makeMessage({
|
||||
runId: "r1",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now,
|
||||
});
|
||||
const msg2 = makeMessage({
|
||||
runId: "r2",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now + 100,
|
||||
});
|
||||
|
||||
for (const msg of [msg1, msg2]) {
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: msg,
|
||||
workerQueue: authenticatedEnvDev.id,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
}
|
||||
|
||||
// Dequeue one message
|
||||
const shard = testOptions.keys.masterQueueShardForEnvironment(
|
||||
msg1.environmentId,
|
||||
2
|
||||
);
|
||||
const messages = await queue.testDequeueFromMasterQueue(shard, msg1.environmentId, 1);
|
||||
expect(messages!.length).toBe(1);
|
||||
expect(messages![0].messageId).toBe("r1");
|
||||
|
||||
// Acknowledge the dequeued message
|
||||
await queue.acknowledgeMessage(msg1.orgId, "r1", {
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
|
||||
// CK index should still have the ck-a entry (r2 is still queued)
|
||||
const ckIndexKey = testOptions.keys.ckIndexKeyFromQueue(
|
||||
testOptions.keys.queueKey(authenticatedEnvDev, msg1.queue, msg1.concurrencyKey)
|
||||
);
|
||||
const ckIndexMembers = await queue.redis.zrange(ckIndexKey, 0, -1);
|
||||
expect(ckIndexMembers.length).toBe(1);
|
||||
|
||||
// Master queue should still have the :ck:* entry
|
||||
const masterQueueKey = testOptions.keys.masterQueueKeyForShard(shard);
|
||||
const masterMembers = await queue.redis.zrange(masterQueueKey, 0, -1);
|
||||
expect(masterMembers.length).toBe(1);
|
||||
expect(masterMembers[0]).toContain(":ck:*");
|
||||
|
||||
// Dequeue and ack the last message
|
||||
const messages2 = await queue.testDequeueFromMasterQueue(shard, msg1.environmentId, 1);
|
||||
expect(messages2!.length).toBe(1);
|
||||
await queue.acknowledgeMessage(msg2.orgId, "r2", {
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
|
||||
// CK index should be empty
|
||||
const ckIndexMembers2 = await queue.redis.zrange(ckIndexKey, 0, -1);
|
||||
expect(ckIndexMembers2.length).toBe(0);
|
||||
|
||||
// Master queue should be empty
|
||||
const masterMembers2 = await queue.redis.zrange(masterQueueKey, 0, -1);
|
||||
expect(masterMembers2.length).toBe(0);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
redisTest(
|
||||
"nack CK message rebalances CK index",
|
||||
async ({ redisContainer }) => {
|
||||
const queue = createQueue(redisContainer);
|
||||
try {
|
||||
const now = Date.now() - 1000;
|
||||
const msg = makeMessage({
|
||||
runId: "r1",
|
||||
concurrencyKey: "ck-a",
|
||||
timestamp: now,
|
||||
});
|
||||
|
||||
await queue.enqueueMessage({
|
||||
env: authenticatedEnvDev,
|
||||
message: msg,
|
||||
workerQueue: authenticatedEnvDev.id,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
|
||||
// Dequeue the message
|
||||
const shard = testOptions.keys.masterQueueShardForEnvironment(
|
||||
msg.environmentId,
|
||||
2
|
||||
);
|
||||
const messages = await queue.testDequeueFromMasterQueue(shard, msg.environmentId, 1);
|
||||
expect(messages!.length).toBe(1);
|
||||
|
||||
// Nack the message (re-enqueue)
|
||||
await queue.nackMessage({
|
||||
orgId: msg.orgId,
|
||||
messageId: "r1",
|
||||
retryAt: Date.now() + 5000,
|
||||
incrementAttemptCount: false,
|
||||
skipDequeueProcessing: true,
|
||||
});
|
||||
|
||||
// CK index should have the ck-a entry (message re-enqueued)
|
||||
const ckIndexKey = testOptions.keys.ckIndexKeyFromQueue(
|
||||
testOptions.keys.queueKey(authenticatedEnvDev, msg.queue, msg.concurrencyKey)
|
||||
);
|
||||
const ckIndexMembers = await queue.redis.zrange(ckIndexKey, 0, -1);
|
||||
expect(ckIndexMembers.length).toBe(1);
|
||||
|
||||
// Master queue should have the :ck:* entry
|
||||
const masterQueueKey = testOptions.keys.masterQueueKeyForShard(shard);
|
||||
const masterMembers = await queue.redis.zrange(masterQueueKey, 0, -1);
|
||||
expect(masterMembers.length).toBe(1);
|
||||
expect(masterMembers[0]).toContain(":ck:*");
|
||||
|
||||
// No old-format entries
|
||||
const oldFormatMembers = masterMembers.filter(
|
||||
(m) => m.includes(":ck:") && !m.endsWith(":ck:*")
|
||||
);
|
||||
expect(oldFormatMembers.length).toBe(0);
|
||||
} finally {
|
||||
await queue.quit();
|
||||
}
|
||||
}
|
||||
);
|
||||
});
|
||||
@@ -359,4 +359,73 @@ describe("KeyProducer", () => {
|
||||
concurrencyKey: "c1234",
|
||||
});
|
||||
});
|
||||
|
||||
it("ckIndexKeyFromQueue", () => {
|
||||
const keyProducer = new RunQueueFullKeyProducer();
|
||||
const queueKey = keyProducer.queueKey(
|
||||
{
|
||||
id: "e1234",
|
||||
type: "PRODUCTION",
|
||||
project: { id: "p1234" },
|
||||
organization: { id: "o1234" },
|
||||
},
|
||||
"task/task-name",
|
||||
"c1234"
|
||||
);
|
||||
const key = keyProducer.ckIndexKeyFromQueue(queueKey);
|
||||
expect(key).toBe("{org:o1234}:proj:p1234:env:e1234:queue:task/task-name:ckIndex");
|
||||
});
|
||||
|
||||
it("ckIndexKeyFromQueue (from wildcard)", () => {
|
||||
const keyProducer = new RunQueueFullKeyProducer();
|
||||
const key = keyProducer.ckIndexKeyFromQueue(
|
||||
"{org:o1234}:proj:p1234:env:e1234:queue:task/task-name:ck:*"
|
||||
);
|
||||
expect(key).toBe("{org:o1234}:proj:p1234:env:e1234:queue:task/task-name:ckIndex");
|
||||
});
|
||||
|
||||
it("baseQueueKeyFromQueue (with CK)", () => {
|
||||
const keyProducer = new RunQueueFullKeyProducer();
|
||||
const queueKey = keyProducer.queueKey(
|
||||
{
|
||||
id: "e1234",
|
||||
type: "PRODUCTION",
|
||||
project: { id: "p1234" },
|
||||
organization: { id: "o1234" },
|
||||
},
|
||||
"task/task-name",
|
||||
"c1234"
|
||||
);
|
||||
const key = keyProducer.baseQueueKeyFromQueue(queueKey);
|
||||
expect(key).toBe("{org:o1234}:proj:p1234:env:e1234:queue:task/task-name");
|
||||
});
|
||||
|
||||
it("baseQueueKeyFromQueue (no CK)", () => {
|
||||
const keyProducer = new RunQueueFullKeyProducer();
|
||||
const queueKey = keyProducer.queueKey(
|
||||
{
|
||||
id: "e1234",
|
||||
type: "PRODUCTION",
|
||||
project: { id: "p1234" },
|
||||
organization: { id: "o1234" },
|
||||
},
|
||||
"task/task-name"
|
||||
);
|
||||
const key = keyProducer.baseQueueKeyFromQueue(queueKey);
|
||||
expect(key).toBe("{org:o1234}:proj:p1234:env:e1234:queue:task/task-name");
|
||||
});
|
||||
|
||||
it("isCkWildcard", () => {
|
||||
const keyProducer = new RunQueueFullKeyProducer();
|
||||
expect(keyProducer.isCkWildcard("{org:o1234}:proj:p1234:env:e1234:queue:task/foo:ck:*")).toBe(true);
|
||||
expect(keyProducer.isCkWildcard("{org:o1234}:proj:p1234:env:e1234:queue:task/foo:ck:bar")).toBe(false);
|
||||
expect(keyProducer.isCkWildcard("{org:o1234}:proj:p1234:env:e1234:queue:task/foo")).toBe(false);
|
||||
});
|
||||
|
||||
it("toCkWildcard", () => {
|
||||
const keyProducer = new RunQueueFullKeyProducer();
|
||||
expect(keyProducer.toCkWildcard("{org:o1234}:proj:p1234:env:e1234:queue:task/foo:ck:bar")).toBe(
|
||||
"{org:o1234}:proj:p1234:env:e1234:queue:task/foo:ck:*"
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -125,6 +125,12 @@ export interface RunQueueKeyProducer {
|
||||
|
||||
// TTL system methods
|
||||
ttlQueueKeyForShard(shard: number): string;
|
||||
|
||||
// CK index methods
|
||||
ckIndexKeyFromQueue(queue: string): string;
|
||||
baseQueueKeyFromQueue(queue: string): string;
|
||||
isCkWildcard(queue: string): boolean;
|
||||
toCkWildcard(queue: string): string;
|
||||
}
|
||||
|
||||
export type EnvQueues = {
|
||||
|
||||
@@ -0,0 +1,201 @@
|
||||
import { batch, logger, queue, task } from "@trigger.dev/sdk";
|
||||
import { setTimeout } from "node:timers/promises";
|
||||
|
||||
// Queue with concurrency limit for CK tests
|
||||
const ckQueue = queue({
|
||||
name: "ck-test-queue",
|
||||
concurrencyLimit: 2,
|
||||
});
|
||||
|
||||
// Worker task: simulates work with a concurrency key
|
||||
export const ckWorkerTask = task({
|
||||
id: "ck-worker-task",
|
||||
queue: ckQueue,
|
||||
retry: { maxAttempts: 1 },
|
||||
run: async (payload: { id: string; waitMs: number }) => {
|
||||
const startedAt = Date.now();
|
||||
logger.info(`CK worker ${payload.id} started`);
|
||||
await setTimeout(payload.waitMs);
|
||||
const completedAt = Date.now();
|
||||
logger.info(`CK worker ${payload.id} completed`);
|
||||
return { id: payload.id, startedAt, completedAt };
|
||||
},
|
||||
});
|
||||
|
||||
// Test 1: Multiple CKs should each get their own concurrency slot
|
||||
export const ckBasicTest = task({
|
||||
id: "ck-basic-test",
|
||||
retry: { maxAttempts: 1 },
|
||||
maxDuration: 120,
|
||||
run: async () => {
|
||||
logger.info("Testing basic CK behavior: multiple CKs run concurrently");
|
||||
|
||||
// Trigger 3 runs with different CKs - all should be able to run
|
||||
// because each CK gets its own concurrency tracking
|
||||
const results = await batch.triggerAndWait([
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "user-1", waitMs: 3000 },
|
||||
options: { concurrencyKey: "user-1" },
|
||||
},
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "user-2", waitMs: 3000 },
|
||||
options: { concurrencyKey: "user-2" },
|
||||
},
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "user-3", waitMs: 3000 },
|
||||
options: { concurrencyKey: "user-3" },
|
||||
},
|
||||
]);
|
||||
|
||||
if (!results.runs.every((r) => r.ok)) {
|
||||
throw new Error("Not all CK runs completed successfully");
|
||||
}
|
||||
|
||||
const executions = results.runs
|
||||
.map((r) => r.output)
|
||||
.sort((a, b) => a.startedAt - b.startedAt);
|
||||
|
||||
logger.info("CK basic test executions", { executions });
|
||||
|
||||
return { executions };
|
||||
},
|
||||
});
|
||||
|
||||
// Test 2: Same CK should respect concurrency limit
|
||||
export const ckSameConcurrencyTest = task({
|
||||
id: "ck-same-concurrency-test",
|
||||
retry: { maxAttempts: 1 },
|
||||
maxDuration: 120,
|
||||
run: async () => {
|
||||
logger.info("Testing same CK concurrency: runs with same CK respect queue limit");
|
||||
|
||||
// Trigger 4 runs all with the same CK
|
||||
// Queue limit is 2, so at most 2 should run concurrently
|
||||
const results = await batch.triggerAndWait([
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "same-1", waitMs: 4000 },
|
||||
options: { concurrencyKey: "shared-key" },
|
||||
},
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "same-2", waitMs: 4000 },
|
||||
options: { concurrencyKey: "shared-key" },
|
||||
},
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "same-3", waitMs: 4000 },
|
||||
options: { concurrencyKey: "shared-key" },
|
||||
},
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "same-4", waitMs: 4000 },
|
||||
options: { concurrencyKey: "shared-key" },
|
||||
},
|
||||
]);
|
||||
|
||||
if (!results.runs.every((r) => r.ok)) {
|
||||
throw new Error("Not all same-CK runs completed successfully");
|
||||
}
|
||||
|
||||
const executions = results.runs
|
||||
.map((r) => r.output)
|
||||
.sort((a, b) => a.startedAt - b.startedAt);
|
||||
|
||||
// Check max concurrent: with same CK and limit 2, should be <= 2
|
||||
let maxConcurrent = 0;
|
||||
for (const current of executions) {
|
||||
const concurrent = executions.filter(
|
||||
(e) =>
|
||||
e.startedAt <= current.startedAt &&
|
||||
e.completedAt > current.startedAt
|
||||
).length;
|
||||
maxConcurrent = Math.max(maxConcurrent, concurrent);
|
||||
}
|
||||
|
||||
logger.info("Same CK concurrency result", { maxConcurrent, executions });
|
||||
|
||||
if (maxConcurrent > 2) {
|
||||
throw new Error(`Expected max 2 concurrent with same CK, got ${maxConcurrent}`);
|
||||
}
|
||||
|
||||
return { executions, maxConcurrent };
|
||||
},
|
||||
});
|
||||
|
||||
// Test 3: Many CKs - the scenario that motivated the CK index
|
||||
export const ckManyKeysTest = task({
|
||||
id: "ck-many-keys-test",
|
||||
retry: { maxAttempts: 1 },
|
||||
maxDuration: 180,
|
||||
run: async () => {
|
||||
logger.info("Testing many CKs: all should complete without starving");
|
||||
|
||||
// Trigger 20 runs each with a different CK
|
||||
const items = Array.from({ length: 20 }, (_, i) => ({
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: `prospect-${i}`, waitMs: 2000 },
|
||||
options: { concurrencyKey: `prospect-${i}` },
|
||||
}));
|
||||
|
||||
const results = await batch.triggerAndWait(items);
|
||||
|
||||
const succeeded = results.runs.filter((r) => r.ok).length;
|
||||
const failed = results.runs.filter((r) => !r.ok).length;
|
||||
|
||||
logger.info("Many CKs test result", { succeeded, failed, total: results.runs.length });
|
||||
|
||||
if (failed > 0) {
|
||||
throw new Error(`${failed} of ${results.runs.length} runs failed`);
|
||||
}
|
||||
|
||||
return { succeeded, total: results.runs.length };
|
||||
},
|
||||
});
|
||||
|
||||
// Test 4: Mixed CK and non-CK triggers on same queue
|
||||
export const ckMixedTest = task({
|
||||
id: "ck-mixed-test",
|
||||
retry: { maxAttempts: 1 },
|
||||
maxDuration: 120,
|
||||
run: async () => {
|
||||
logger.info("Testing mixed CK and non-CK on same queue");
|
||||
|
||||
const results = await batch.triggerAndWait([
|
||||
// Non-CK runs
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "no-ck-1", waitMs: 2000 },
|
||||
},
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "no-ck-2", waitMs: 2000 },
|
||||
},
|
||||
// CK runs
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "with-ck-1", waitMs: 2000 },
|
||||
options: { concurrencyKey: "tenant-a" },
|
||||
},
|
||||
{
|
||||
id: ckWorkerTask.id,
|
||||
payload: { id: "with-ck-2", waitMs: 2000 },
|
||||
options: { concurrencyKey: "tenant-b" },
|
||||
},
|
||||
]);
|
||||
|
||||
const succeeded = results.runs.filter((r) => r.ok).length;
|
||||
const failed = results.runs.filter((r) => !r.ok).length;
|
||||
|
||||
logger.info("Mixed test result", { succeeded, failed });
|
||||
|
||||
if (failed > 0) {
|
||||
throw new Error(`${failed} runs failed in mixed CK/non-CK test`);
|
||||
}
|
||||
|
||||
return { succeeded, total: results.runs.length };
|
||||
},
|
||||
});
|
||||
Reference in New Issue
Block a user