From 4df6cda09dbe42868bd2f38bfb21b6993319a917 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 8 Aug 2026 21:34:07 +0100 Subject: [PATCH 1/2] 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. --- .../fair-queue-concurrency-slot-leak.md | 5 ++ packages/redis-worker/src/fair-queue/index.ts | 18 ++--- .../src/fair-queue/tests/fairQueue.test.ts | 71 ++++++++++++++++++- 3 files changed, 84 insertions(+), 10 deletions(-) create mode 100644 .changeset/fair-queue-concurrency-slot-leak.md diff --git a/.changeset/fair-queue-concurrency-slot-leak.md b/.changeset/fair-queue-concurrency-slot-leak.md new file mode 100644 index 00000000000..4ccced1c9c9 --- /dev/null +++ b/.changeset/fair-queue-concurrency-slot-leak.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/redis-worker": patch +--- + +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. diff --git a/packages/redis-worker/src/fair-queue/index.ts b/packages/redis-worker/src/fair-queue/index.ts index e24de876d85..fa3b0b3da5a 100644 --- a/packages/redis-worker/src/fair-queue/index.ts +++ b/packages/redis-worker/src/fair-queue/index.ts @@ -1253,14 +1253,14 @@ export class FairQueue { }) : { id: queueId, tenantId: this.keys.extractTenantId(queueId), metadata: {} }; - // Complete in visibility manager - await this.visibilityManager.complete(messageId, queueId); - // Release concurrency - if (this.concurrencyManager && storedMessage) { + if (this.concurrencyManager) { await this.concurrencyManager.release(descriptor, messageId); } + // Complete in visibility manager + await this.visibilityManager.complete(messageId, queueId); + // Update both old and new indexes, clean up caches if queue is empty const { queueEmpty } = await this.#updateAllIndexesAfterDequeue(queueId, descriptor.tenantId); if (queueEmpty) { @@ -1308,6 +1308,11 @@ export class FairQueue { }) : { id: queueId, tenantId: this.keys.extractTenantId(queueId), metadata: {} }; + // Release concurrency + if (this.concurrencyManager) { + await this.concurrencyManager.release(descriptor, messageId); + } + // Release back to queue (visibility manager updates dispatch indexes atomically) // Dispatch shard is tenant-based, not queue-based const dispatchShardId = this.tenantDispatch.getShardForTenant(descriptor.tenantId); @@ -1324,11 +1329,6 @@ export class FairQueue { Date.now() // Put at back of queue ); - // Release concurrency - if (this.concurrencyManager && storedMessage) { - await this.concurrencyManager.release(descriptor, messageId); - } - this.logger.debug("Message released", { messageId, queueId, diff --git a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts index 5ae1f390f9f..ab71198a2e7 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -10,7 +10,7 @@ import { WorkerQueueManager, } from "../index.js"; import type { FairQueueKeyProducer, FairQueueOptions } from "../types.js"; -import type { RedisOptions } from "@internal/redis"; +import { createRedisClient, type RedisOptions } from "@internal/redis"; // Define a common payload schema for tests const TestPayloadSchema = z.object({ value: z.string() }); @@ -1370,4 +1370,73 @@ describe("FairQueue", () => { } ); }); + + describe("concurrency slot release", () => { + redisTest( + "should release the concurrency slot when the in-flight record is already gone", + { timeout: 15000 }, + async ({ redisOptions }) => { + const processed: string[] = []; + keys = new DefaultFairQueueKeyProducer({ prefix: "test" }); + + const scheduler = new DRRScheduler({ + redis: redisOptions, + keys, + quantum: 10, + maxDeficit: 100, + }); + + const queue = new TestFairQueueHelper(redisOptions, keys, { + scheduler, + payloadSchema: TestPayloadSchema, + shardCount: 1, + consumerCount: 1, + consumerIntervalMs: 20, + visibilityTimeoutMs: 60000, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 1, + defaultLimit: 1, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + + queue.onMessage(async (ctx) => { + if (ctx.message.payload.value === "msg-0") { + await redis.hdel(keys.inflightDataKey(0), ctx.message.id); + } + processed.push(ctx.message.payload.value); + await ctx.complete(); + }); + + for (let i = 0; i < 2; i++) { + await queue.enqueue({ + queueId: "tenant:t1:queue:q1", + tenantId: "t1", + payload: { value: `msg-${i}` }, + }); + } + + queue.start(); + + await vi.waitFor( + () => { + expect(processed).toHaveLength(2); + }, + { timeout: 10000 } + ); + + const held = await redis.scard(keys.concurrencyKey("tenant", "t1")); + expect(held).toBe(0); + + await redis.quit(); + await queue.close(); + } + ); + }); }); From 35568a32e7ad4762ba9c44eab12d370476c15eae Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 8 Aug 2026 21:51:23 +0100 Subject: [PATCH 2/2] fix(redis-worker): keep metadata-derived concurrency groups on the release path Concurrency groups can key on queue metadata, not just the tenant. The fallback descriptor dropped metadata and skipped the descriptor cache, so those groups released against the wrong Redis set and leaked exactly as before. Prefer the cached descriptor in both branches, and move the retry path's release ahead of the re-queue so a redelivery cannot have its fresh reservation deleted by the previous holder. --- packages/redis-worker/src/fair-queue/index.ts | 34 +++-- .../src/fair-queue/tests/fairQueue.test.ts | 122 ++++++++++++++---- 2 files changed, 115 insertions(+), 41 deletions(-) diff --git a/packages/redis-worker/src/fair-queue/index.ts b/packages/redis-worker/src/fair-queue/index.ts index fa3b0b3da5a..d8a2cd00523 100644 --- a/packages/redis-worker/src/fair-queue/index.ts +++ b/packages/redis-worker/src/fair-queue/index.ts @@ -1245,13 +1245,11 @@ export class FairQueue { } } - const descriptor: QueueDescriptor = storedMessage - ? (this.queueDescriptorCache.get(queueId) ?? { - id: queueId, - tenantId: storedMessage.tenantId, - metadata: storedMessage.metadata ?? {}, - }) - : { id: queueId, tenantId: this.keys.extractTenantId(queueId), metadata: {} }; + const descriptor: QueueDescriptor = this.queueDescriptorCache.get(queueId) ?? { + id: queueId, + tenantId: storedMessage?.tenantId ?? this.keys.extractTenantId(queueId), + metadata: storedMessage?.metadata ?? {}, + }; // Release concurrency if (this.concurrencyManager) { @@ -1300,13 +1298,11 @@ export class FairQueue { } } - const descriptor: QueueDescriptor = storedMessage - ? (this.queueDescriptorCache.get(queueId) ?? { - id: queueId, - tenantId: storedMessage.tenantId, - metadata: storedMessage.metadata ?? {}, - }) - : { id: queueId, tenantId: this.keys.extractTenantId(queueId), metadata: {} }; + const descriptor: QueueDescriptor = this.queueDescriptorCache.get(queueId) ?? { + id: queueId, + tenantId: storedMessage?.tenantId ?? this.keys.extractTenantId(queueId), + metadata: storedMessage?.metadata ?? {}, + }; // Release concurrency if (this.concurrencyManager) { @@ -1411,6 +1407,11 @@ export class FairQueue { attempt: storedMessage.attempt + 1, }; + // Release concurrency + if (this.concurrencyManager) { + await this.concurrencyManager.release(descriptor, storedMessage.id); + } + // Release with delay, passing the updated message data so the Lua script // atomically writes the incremented attempt count when re-queuing. const tenantQueueIndexKey = this.keys.tenantQueueIndexKey(descriptor.tenantId); @@ -1427,11 +1428,6 @@ export class FairQueue { JSON.stringify(updatedMessage) ); - // Release concurrency - if (this.concurrencyManager) { - await this.concurrencyManager.release(descriptor, storedMessage.id); - } - this.telemetry.recordRetry(); this.logger.debug("Message scheduled for retry", { diff --git a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts index ab71198a2e7..f0ca5ef758f 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -1406,36 +1406,114 @@ describe("FairQueue", () => { const redis = createRedisClient(redisOptions); - queue.onMessage(async (ctx) => { - if (ctx.message.payload.value === "msg-0") { - await redis.hdel(keys.inflightDataKey(0), ctx.message.id); + try { + queue.onMessage(async (ctx) => { + if (ctx.message.payload.value === "msg-0") { + await redis.hdel(keys.inflightDataKey(0), ctx.message.id); + } + processed.push(ctx.message.payload.value); + await ctx.complete(); + }); + + for (let i = 0; i < 2; i++) { + await queue.enqueue({ + queueId: "tenant:t1:queue:q1", + tenantId: "t1", + payload: { value: `msg-${i}` }, + }); } - processed.push(ctx.message.payload.value); - await ctx.complete(); + + queue.start(); + + await vi.waitFor( + () => { + expect(processed).toHaveLength(2); + }, + { timeout: 10000 } + ); + + const held = await redis.scard(keys.concurrencyKey("tenant", "t1")); + expect(held).toBe(0); + } finally { + await redis.quit(); + await queue.close(); + } + } + ); + + redisTest( + "should release metadata-derived concurrency groups when the in-flight record is gone", + { timeout: 15000 }, + async ({ redisOptions }) => { + const processed: string[] = []; + keys = new DefaultFairQueueKeyProducer({ prefix: "test" }); + + const scheduler = new DRRScheduler({ + redis: redisOptions, + keys, + quantum: 10, + maxDeficit: 100, }); - for (let i = 0; i < 2; i++) { - await queue.enqueue({ - queueId: "tenant:t1:queue:q1", - tenantId: "t1", - payload: { value: `msg-${i}` }, + const queue = new TestFairQueueHelper(redisOptions, keys, { + scheduler, + payloadSchema: TestPayloadSchema, + shardCount: 1, + consumerCount: 1, + consumerIntervalMs: 20, + visibilityTimeoutMs: 60000, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 5, + defaultLimit: 5, + }, + { + name: "organization", + extractGroupId: (q) => (q.metadata.orgId as string) ?? "default", + getLimit: async () => 1, + defaultLimit: 1, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + + try { + queue.onMessage(async (ctx) => { + if (ctx.message.payload.value === "msg-0") { + await redis.hdel(keys.inflightDataKey(0), ctx.message.id); + } + processed.push(ctx.message.payload.value); + await ctx.complete(); }); - } - queue.start(); + for (let i = 0; i < 2; i++) { + await queue.enqueue({ + queueId: "tenant:t1:queue:q1", + tenantId: "t1", + metadata: { orgId: "org-1" }, + payload: { value: `msg-${i}` }, + }); + } - await vi.waitFor( - () => { - expect(processed).toHaveLength(2); - }, - { timeout: 10000 } - ); + queue.start(); - const held = await redis.scard(keys.concurrencyKey("tenant", "t1")); - expect(held).toBe(0); + await vi.waitFor( + () => { + expect(processed).toHaveLength(2); + }, + { timeout: 10000 } + ); - await redis.quit(); - await queue.close(); + expect(await redis.scard(keys.concurrencyKey("organization", "org-1"))).toBe(0); + expect(await redis.scard(keys.concurrencyKey("organization", "default"))).toBe(0); + } finally { + await redis.quit(); + await queue.close(); + } } ); });