Skip to content

Commit 4df6cda

Browse files
committed
fix(redis-worker): stop fair queue leaking concurrency slots
A slot was only released when the message's in-flight record could still be read, and the release ran after that record had already been deleted. Any failure in between left the slot held with nothing able to reclaim it. Once a tenant leaked its whole limit, every queue it owned stalled for good.
1 parent 337dda1 commit 4df6cda

3 files changed

Lines changed: 84 additions & 10 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@trigger.dev/redis-worker": patch
3+
---
4+
5+
Fair queue consumers no longer leak concurrency slots. A slot is now always released when a message completes or is put back on the queue, even when its in-flight record has already gone. Leaked slots were never reclaimed, so enough of them would permanently stall every queue belonging to that tenant.

packages/redis-worker/src/fair-queue/index.ts

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1253,14 +1253,14 @@ export class FairQueue<TPayloadSchema extends z.ZodTypeAny = z.ZodUnknown> {
12531253
})
12541254
: { id: queueId, tenantId: this.keys.extractTenantId(queueId), metadata: {} };
12551255

1256-
// Complete in visibility manager
1257-
await this.visibilityManager.complete(messageId, queueId);
1258-
12591256
// Release concurrency
1260-
if (this.concurrencyManager && storedMessage) {
1257+
if (this.concurrencyManager) {
12611258
await this.concurrencyManager.release(descriptor, messageId);
12621259
}
12631260

1261+
// Complete in visibility manager
1262+
await this.visibilityManager.complete(messageId, queueId);
1263+
12641264
// Update both old and new indexes, clean up caches if queue is empty
12651265
const { queueEmpty } = await this.#updateAllIndexesAfterDequeue(queueId, descriptor.tenantId);
12661266
if (queueEmpty) {
@@ -1308,6 +1308,11 @@ export class FairQueue<TPayloadSchema extends z.ZodTypeAny = z.ZodUnknown> {
13081308
})
13091309
: { id: queueId, tenantId: this.keys.extractTenantId(queueId), metadata: {} };
13101310

1311+
// Release concurrency
1312+
if (this.concurrencyManager) {
1313+
await this.concurrencyManager.release(descriptor, messageId);
1314+
}
1315+
13111316
// Release back to queue (visibility manager updates dispatch indexes atomically)
13121317
// Dispatch shard is tenant-based, not queue-based
13131318
const dispatchShardId = this.tenantDispatch.getShardForTenant(descriptor.tenantId);
@@ -1324,11 +1329,6 @@ export class FairQueue<TPayloadSchema extends z.ZodTypeAny = z.ZodUnknown> {
13241329
Date.now() // Put at back of queue
13251330
);
13261331

1327-
// Release concurrency
1328-
if (this.concurrencyManager && storedMessage) {
1329-
await this.concurrencyManager.release(descriptor, messageId);
1330-
}
1331-
13321332
this.logger.debug("Message released", {
13331333
messageId,
13341334
queueId,

packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts

Lines changed: 70 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ import {
1010
WorkerQueueManager,
1111
} from "../index.js";
1212
import type { FairQueueKeyProducer, FairQueueOptions } from "../types.js";
13-
import type { RedisOptions } from "@internal/redis";
13+
import { createRedisClient, type RedisOptions } from "@internal/redis";
1414

1515
// Define a common payload schema for tests
1616
const TestPayloadSchema = z.object({ value: z.string() });
@@ -1370,4 +1370,73 @@ describe("FairQueue", () => {
13701370
}
13711371
);
13721372
});
1373+
1374+
describe("concurrency slot release", () => {
1375+
redisTest(
1376+
"should release the concurrency slot when the in-flight record is already gone",
1377+
{ timeout: 15000 },
1378+
async ({ redisOptions }) => {
1379+
const processed: string[] = [];
1380+
keys = new DefaultFairQueueKeyProducer({ prefix: "test" });
1381+
1382+
const scheduler = new DRRScheduler({
1383+
redis: redisOptions,
1384+
keys,
1385+
quantum: 10,
1386+
maxDeficit: 100,
1387+
});
1388+
1389+
const queue = new TestFairQueueHelper(redisOptions, keys, {
1390+
scheduler,
1391+
payloadSchema: TestPayloadSchema,
1392+
shardCount: 1,
1393+
consumerCount: 1,
1394+
consumerIntervalMs: 20,
1395+
visibilityTimeoutMs: 60000,
1396+
concurrencyGroups: [
1397+
{
1398+
name: "tenant",
1399+
extractGroupId: (q) => q.tenantId,
1400+
getLimit: async () => 1,
1401+
defaultLimit: 1,
1402+
},
1403+
],
1404+
startConsumers: false,
1405+
});
1406+
1407+
const redis = createRedisClient(redisOptions);
1408+
1409+
queue.onMessage(async (ctx) => {
1410+
if (ctx.message.payload.value === "msg-0") {
1411+
await redis.hdel(keys.inflightDataKey(0), ctx.message.id);
1412+
}
1413+
processed.push(ctx.message.payload.value);
1414+
await ctx.complete();
1415+
});
1416+
1417+
for (let i = 0; i < 2; i++) {
1418+
await queue.enqueue({
1419+
queueId: "tenant:t1:queue:q1",
1420+
tenantId: "t1",
1421+
payload: { value: `msg-${i}` },
1422+
});
1423+
}
1424+
1425+
queue.start();
1426+
1427+
await vi.waitFor(
1428+
() => {
1429+
expect(processed).toHaveLength(2);
1430+
},
1431+
{ timeout: 10000 }
1432+
);
1433+
1434+
const held = await redis.scard(keys.concurrencyKey("tenant", "t1"));
1435+
expect(held).toBe(0);
1436+
1437+
await redis.quit();
1438+
await queue.close();
1439+
}
1440+
);
1441+
});
13731442
});

0 commit comments

Comments
 (0)