From 4df6cda09dbe42868bd2f38bfb21b6993319a917 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 8 Aug 2026 21:34:07 +0100 Subject: [PATCH 01/12] 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 02/12] 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(); + } } ); }); From b2ebf18abc314a633e515594c8ca5ec26dbf38e3 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 11 Aug 2026 12:41:21 +0100 Subject: [PATCH 03/12] fix(redis-worker): close the remaining concurrency slot leaks in fair queue Four paths could strand a slot with nothing left to reclaim it. Reclaim released the slot only after requeuing, and only for messages that survived the requeue, so a message whose requeue threw kept its slot forever. It now releases before the message becomes claimable, which also stops a redelivery having its fresh reservation deleted by the previous holder. The dead-letter branch completed the message before releasing, so anything throwing in between stranded the slot. Release now happens first. failMessage returned early when the stored message was missing or unparseable without releasing, completing or requeuing, leaving the slot held. The requeue script returned without removing the in-flight entry when the payload was gone, so the member was rescanned on every reclaim tick forever. Enough of them fill the scan window and starve a whole shard's reclaim. reclaimTimedOut takes an optional pre-requeue hook; no existing caller changes. --- packages/redis-worker/src/fair-queue/index.ts | 112 ++++++++------ .../src/fair-queue/tests/fairQueue.test.ts | 141 ++++++++++++++++++ .../src/fair-queue/tests/visibility.test.ts | 65 ++++++++ .../redis-worker/src/fair-queue/visibility.ts | 20 ++- 4 files changed, 294 insertions(+), 44 deletions(-) diff --git a/packages/redis-worker/src/fair-queue/index.ts b/packages/redis-worker/src/fair-queue/index.ts index d8a2cd00523..4ab6fd62569 100644 --- a/packages/redis-worker/src/fair-queue/index.ts +++ b/packages/redis-worker/src/fair-queue/index.ts @@ -25,6 +25,7 @@ import type { FairScheduler, QueueCooloffState, QueueDescriptor, + ReclaimedMessageInfo, SchedulerContext, StoredMessage, TenantQueues, @@ -1331,6 +1332,26 @@ export class FairQueue { }); } + /** + * Release a concurrency slot for a message we can no longer describe, deriving the + * group from the queue id alone. Used on paths that bail out before the stored + * message is available, where the slot would otherwise be held with nothing left + * to reclaim it. + */ + async #releaseOrphanedConcurrency(messageId: string, queueId: string): Promise { + if (!this.concurrencyManager) { + return; + } + + const descriptor: QueueDescriptor = this.queueDescriptorCache.get(queueId) ?? { + id: queueId, + tenantId: this.keys.extractTenantId(queueId), + metadata: {}, + }; + + await this.concurrencyManager.release(descriptor, messageId); + } + /** * Mark a message as failed. This will trigger retry logic if configured, * or move the message to the dead letter queue. @@ -1349,6 +1370,7 @@ export class FairQueue { const dataJson = await this.redis.hget(inflightDataKey, messageId); if (!dataJson) { this.logger.error("Cannot fail message: not found in in-flight data", { messageId, queueId }); + await this.#releaseOrphanedConcurrency(messageId, queueId); return; } @@ -1360,6 +1382,7 @@ export class FairQueue { messageId, queueId, }); + await this.#releaseOrphanedConcurrency(messageId, queueId); return; } @@ -1441,13 +1464,13 @@ export class FairQueue { } } - // Move to DLQ - await this.#moveToDeadLetterQueue(storedMessage, error?.message); - // Release concurrency if (this.concurrencyManager) { await this.concurrencyManager.release(descriptor, storedMessage.id); } + + // Move to DLQ + await this.#moveToDeadLetterQueue(storedMessage, error?.message); } async #moveToDeadLetterQueue( @@ -1521,49 +1544,54 @@ export class FairQueue { } } - async #reclaimTimedOutMessages(): Promise { - let totalReclaimed = 0; + /** + * Free a reclaimed message's concurrency slot. Runs before the message is put back + * on the queue, so another consumer cannot reserve the same messageId and have its + * fresh reservation deleted by this release. Failures are logged rather than thrown + * so one bad message cannot stop the rest of the shard being reclaimed. + */ + async #releaseReclaimedConcurrency(message: ReclaimedMessageInfo): Promise { + if (!this.concurrencyManager) { + return; + } - for (let shardId = 0; shardId < this.shardCount; shardId++) { - const reclaimedMessages = await this.visibilityManager.reclaimTimedOut(shardId, (queueId) => { - const tenantId = this.keys.extractTenantId(queueId); - const dispatchShardId = this.tenantDispatch.getShardForTenant(tenantId); - return { - queueKey: this.keys.queueKey(queueId), - queueItemsKey: this.keys.queueItemsKey(queueId), - tenantQueueIndexKey: this.keys.tenantQueueIndexKey(tenantId), - dispatchKey: this.keys.dispatchKey(dispatchShardId), - tenantId, - }; + try { + await this.concurrencyManager.release( + { + id: message.queueId, + tenantId: message.tenantId, + metadata: message.metadata ?? {}, + }, + message.messageId + ); + } catch (error) { + this.logger.error("Failed to release concurrency for reclaimed message", { + messageId: message.messageId, + queueId: message.queueId, + error: error instanceof Error ? error.message : String(error), }); + } + } - if (reclaimedMessages.length > 0) { - // Release concurrency for all reclaimed messages in a single batch - // This is critical: when a message times out, its concurrency slot must be freed - // so the message can be processed again when it's re-claimed from the queue - if (this.concurrencyManager) { - try { - await this.concurrencyManager.releaseBatch( - reclaimedMessages.map((msg) => ({ - queue: { - id: msg.queueId, - tenantId: msg.tenantId, - metadata: msg.metadata ?? {}, - }, - messageId: msg.messageId, - })) - ); - } catch (error) { - this.logger.error("Failed to release concurrency for reclaimed messages", { - count: reclaimedMessages.length, - error: error instanceof Error ? error.message : String(error), - }); - } - } + async #reclaimTimedOutMessages(): Promise { + let totalReclaimed = 0; - // Dispatch indexes are updated atomically by the releaseMessage Lua script - // inside reclaimTimedOut, so no separate index update needed here. - } + for (let shardId = 0; shardId < this.shardCount; shardId++) { + const reclaimedMessages = await this.visibilityManager.reclaimTimedOut( + shardId, + (queueId) => { + const tenantId = this.keys.extractTenantId(queueId); + const dispatchShardId = this.tenantDispatch.getShardForTenant(tenantId); + return { + queueKey: this.keys.queueKey(queueId), + queueItemsKey: this.keys.queueItemsKey(queueId), + tenantQueueIndexKey: this.keys.tenantQueueIndexKey(tenantId), + dispatchKey: this.keys.dispatchKey(dispatchShardId), + tenantId, + }; + }, + this.#releaseReclaimedConcurrency.bind(this) + ); totalReclaimed += reclaimedMessages.length; } 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 f0ca5ef758f..cd1492ce43e 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -1516,5 +1516,146 @@ describe("FairQueue", () => { } } ); + + redisTest( + "should release the concurrency slot even when the reclaim requeue fails", + { timeout: 15000 }, + async ({ redisOptions }) => { + 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: 200, + reclaimIntervalMs: 100, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 1, + defaultLimit: 1, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + const queueId = "tenant:t1:queue:reclaim-fail"; + + try { + queue.onMessage(async () => { + await new Promise((resolve) => setTimeout(resolve, 10000)); + }); + + await queue.enqueue({ + queueId, + tenantId: "t1", + payload: { value: "msg-0" }, + }); + + queue.start(); + + await vi.waitFor( + async () => { + expect(await redis.scard(keys.concurrencyKey("tenant", "t1"))).toBe(1); + }, + { timeout: 5000 } + ); + + await redis.del(keys.queueKey(queueId)); + await redis.set(keys.queueKey(queueId), "not-a-zset"); + + await vi.waitFor( + async () => { + expect(await redis.scard(keys.concurrencyKey("tenant", "t1"))).toBe(0); + }, + { timeout: 8000 } + ); + } finally { + await redis.del(keys.queueKey(queueId)); + await redis.quit(); + await queue.close(); + } + } + ); + + redisTest( + "should release the concurrency slot when failMessage cannot read the message", + { timeout: 15000 }, + async ({ redisOptions }) => { + const started: 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 () => 5, + defaultLimit: 5, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + + try { + queue.onMessage(async (ctx) => { + started.push(ctx.message.payload.value); + await redis.hdel(keys.inflightDataKey(0), ctx.message.id); + await ctx.fail(new Error("boom")); + }); + + await queue.enqueue({ + queueId: "tenant:t1:queue:q1", + tenantId: "t1", + payload: { value: "msg-0" }, + }); + + queue.start(); + + await vi.waitFor( + () => { + expect(started).toHaveLength(1); + }, + { timeout: 10000 } + ); + + await vi.waitFor( + async () => { + expect(await redis.scard(keys.concurrencyKey("tenant", "t1"))).toBe(0); + }, + { timeout: 5000 } + ); + } finally { + await redis.quit(); + await queue.close(); + } + } + ); }); }); diff --git a/packages/redis-worker/src/fair-queue/tests/visibility.test.ts b/packages/redis-worker/src/fair-queue/tests/visibility.test.ts index e20b6671075..83d7e7f3e88 100644 --- a/packages/redis-worker/src/fair-queue/tests/visibility.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/visibility.test.ts @@ -912,5 +912,70 @@ describe("VisibilityManager", () => { await redis.quit(); } ); + + redisTest( + "should drop a dangling in-flight entry whose payload is gone", + { timeout: 10000 }, + async ({ redisOptions }) => { + keys = new DefaultFairQueueKeyProducer({ prefix: "test" }); + + const manager = new VisibilityManager({ + redis: redisOptions, + keys, + shardCount: 1, + defaultTimeoutMs: 100, + }); + + const redis = createRedisClient(redisOptions); + const queueId = "tenant:t1:queue:dangling"; + const queueKey = keys.queueKey(queueId); + const queueItemsKey = keys.queueItemsKey(queueId); + const dispatchKey = keys.dispatchKey(0); + const inflightKey = keys.inflightKey(0); + const inflightDataKey = keys.inflightDataKey(0); + + try { + const messageId = "dangling-msg"; + const storedMessage = { + id: messageId, + queueId, + tenantId: "t1", + payload: { id: 1, value: "test" }, + timestamp: Date.now() - 1000, + attempt: 1, + }; + + await redis.zadd(queueKey, storedMessage.timestamp, messageId); + await redis.hset(queueItemsKey, messageId, JSON.stringify(storedMessage)); + + const claimResult = await manager.claim( + queueId, + queueKey, + queueItemsKey, + "consumer-1", + 100 + ); + expect(claimResult.claimed).toBe(true); + + await redis.hdel(inflightDataKey, messageId); + expect(await redis.zcard(inflightKey)).toBe(1); + + await new Promise((resolve) => setTimeout(resolve, 150)); + + await manager.reclaimTimedOut(0, (qId) => ({ + queueKey: keys.queueKey(qId), + queueItemsKey: keys.queueItemsKey(qId), + tenantQueueIndexKey: keys.tenantQueueIndexKey(keys.extractTenantId(qId)), + dispatchKey, + tenantId: keys.extractTenantId(qId), + })); + + expect(await redis.zcard(inflightKey)).toBe(0); + } finally { + await manager.close(); + await redis.quit(); + } + } + ); }); }); diff --git a/packages/redis-worker/src/fair-queue/visibility.ts b/packages/redis-worker/src/fair-queue/visibility.ts index 8ecc3d9927a..5713c853378 100644 --- a/packages/redis-worker/src/fair-queue/visibility.ts +++ b/packages/redis-worker/src/fair-queue/visibility.ts @@ -398,7 +398,8 @@ export class VisibilityManager { tenantQueueIndexKey: string; dispatchKey: string; tenantId: string; - } + }, + onBeforeRequeue?: (message: ReclaimedMessageInfo) => Promise ): Promise { const inflightKey = this.keys.inflightKey(shardId); const inflightDataKey = this.keys.inflightDataKey(shardId); @@ -439,6 +440,15 @@ export class VisibilityManager { } } + if (onBeforeRequeue) { + await onBeforeRequeue({ + messageId, + queueId, + tenantId: storedMessage?.tenantId ?? tenantId, + metadata: storedMessage?.metadata, + }); + } + // Re-add to queue with original timestamp to preserve priority // Fall back to now if we can't get the original timestamp const score = storedMessage?.timestamp ?? now; @@ -707,7 +717,9 @@ local tenantId = ARGV[6] -- Get message data from in-flight local payload = redis.call('HGET', inflightDataKey, messageId) if not payload then - -- Message not in in-flight or already released + -- Message not in in-flight or already released. Drop the dangling in-flight + -- entry so it stops consuming a slot in every future reclaim scan. + redis.call('ZREM', inflightKey, member) return 0 end @@ -780,6 +792,10 @@ for i = 0, numMessages - 1 do redis.call('HSET', queueItemsKey, messageId, payload) releasedCount = releasedCount + 1 + else + -- Dangling in-flight entry with no payload: drop it so it stops consuming + -- a slot in every future reclaim scan. + redis.call('ZREM', inflightKey, member) end end From c6524ca51592d7ef50ac11d95a46fefdf1c22110 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 11 Aug 2026 12:51:43 +0100 Subject: [PATCH 04/12] fix(redis-worker): stop losing dead-lettered messages when the DLQ write fails The message was removed from in-flight before the dead-letter entry was written, and the write went through a pipeline whose per-command errors were never inspected. A failed write therefore lost the message silently: gone from in-flight, absent from the dead-letter queue, invisible to the reclaim loop. Write the entry first and only complete the message once it lands. On failure the message stays in-flight, so the reclaim loop picks it up and it is retried rather than dropped. A persistently failing write now loops visibly instead of discarding work. --- packages/redis-worker/src/fair-queue/index.ts | 22 +++++- .../src/fair-queue/tests/fairQueue.test.ts | 75 +++++++++++++++++++ 2 files changed, 93 insertions(+), 4 deletions(-) diff --git a/packages/redis-worker/src/fair-queue/index.ts b/packages/redis-worker/src/fair-queue/index.ts index 4ab6fd62569..989dd776ae4 100644 --- a/packages/redis-worker/src/fair-queue/index.ts +++ b/packages/redis-worker/src/fair-queue/index.ts @@ -1498,14 +1498,28 @@ export class FairQueue { originalTimestamp: storedMessage.timestamp, }; - // Complete in visibility manager - await this.visibilityManager.complete(storedMessage.id, storedMessage.queueId); - // Add to DLQ const pipeline = this.redis.pipeline(); pipeline.zadd(dlqKey, dlqMessage.deadLetteredAt, storedMessage.id); pipeline.hset(dlqDataKey, storedMessage.id, JSON.stringify(dlqMessage)); - await pipeline.exec(); + const dlqResults = await pipeline.exec(); + + const dlqErrors = (dlqResults ?? []) + .map(([error]) => error) + .filter((error): error is Error => Boolean(error)); + + if (dlqErrors.length > 0) { + this.logger.error("Failed to write message to DLQ, leaving it in-flight to be reclaimed", { + messageId: storedMessage.id, + queueId: storedMessage.queueId, + tenantId: storedMessage.tenantId, + errors: dlqErrors.map((error) => error.message), + }); + return; + } + + // Complete in visibility manager + await this.visibilityManager.complete(storedMessage.id, storedMessage.queueId); this.telemetry.recordDLQ(); 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 cd1492ce43e..e84e2af3324 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -1517,6 +1517,81 @@ describe("FairQueue", () => { } ); + redisTest( + "should keep a message in-flight when the DLQ write fails", + { timeout: 15000 }, + async ({ redisOptions }) => { + const started: 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, + reclaimIntervalMs: 60000, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 5, + defaultLimit: 5, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + const dlqKey = keys.deadLetterQueueKey("t1"); + + try { + await redis.set(dlqKey, "not-a-zset"); + + queue.onMessage(async (ctx) => { + started.push(ctx.message.payload.value); + await ctx.fail(new Error("boom")); + }); + + await queue.enqueue({ + queueId: "tenant:t1:queue:q1", + tenantId: "t1", + payload: { value: "msg-0" }, + }); + + queue.start(); + + await vi.waitFor( + () => { + expect(started).toHaveLength(1); + }, + { timeout: 10000 } + ); + + await vi.waitFor( + async () => { + expect(await redis.scard(keys.concurrencyKey("tenant", "t1"))).toBe(0); + }, + { timeout: 5000 } + ); + + expect(await redis.zcard(keys.inflightKey(0))).toBe(1); + } finally { + await redis.del(dlqKey); + await redis.quit(); + await queue.close(); + } + } + ); + redisTest( "should release the concurrency slot even when the reclaim requeue fails", { timeout: 15000 }, From 3f8ed8b05d8eae6f0814e49310879a6dc95255d8 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 11 Aug 2026 13:13:26 +0100 Subject: [PATCH 05/12] fix(redis-worker,run-engine): close the remaining ways a batch item is stranded or run twice Releasing a slot could silently do nothing. The release pipeline never inspected its per-command errors, and ioredis resolves a pipeline whose commands failed, so a failed SREM reported success and the caller went on to destroy the in-flight record. Release now throws on a failed command. Reclaim then acted on that false success: it released per message, swallowed any error, and requeued regardless, which is exactly how a slot ends up held with the message gone. It now frees the whole timed-out batch in one pipeline before any of them is requeued, and a failure aborts the requeue so the messages stay in-flight for the next tick. That also restores the single round trip per shard. Requeuing removed the message from in-flight before writing it back to the queue. Redis does not roll back a script that fails partway, so a failed queue write lost the message outright. The order is now write-then-remove. Batch items were never heartbeated, so any item slower than the visibility timeout was reclaimed and handed to a second consumer while the first was still running it, executing the item twice. Items are now heartbeated for as long as their callback runs, and the visibility timeout is configurable so this is testable. --- .../run-engine/src/batch-queue/index.ts | 45 ++++++- .../src/batch-queue/tests/index.test.ts | 52 ++++++++ .../run-engine/src/batch-queue/types.ts | 6 + .../src/fair-queue/concurrency.ts | 28 +++- packages/redis-worker/src/fair-queue/index.ts | 89 ++++++------- .../src/fair-queue/tests/fairQueue.test.ts | 58 +++++---- .../redis-worker/src/fair-queue/visibility.ts | 120 ++++++++++-------- 7 files changed, 259 insertions(+), 139 deletions(-) diff --git a/internal-packages/run-engine/src/batch-queue/index.ts b/internal-packages/run-engine/src/batch-queue/index.ts index 10354e77c52..b2c7757a901 100644 --- a/internal-packages/run-engine/src/batch-queue/index.ts +++ b/internal-packages/run-engine/src/batch-queue/index.ts @@ -59,6 +59,9 @@ const ENV_CONCURRENCY_KEY_PREFIX = "batch:env_concurrency"; // then all messages are routed to this queue for BatchQueue's own consumer loop. const BATCH_WORKER_QUEUE_ID = "batch-worker-queue"; +/** How long a claimed batch item stays invisible before the reclaim loop takes it back. */ +const BATCH_ITEM_VISIBILITY_TIMEOUT_MS = 60_000; + export class BatchQueue { private fairQueue: FairQueue; private workerQueueManager: WorkerQueueManager; @@ -67,6 +70,8 @@ export class BatchQueue { private tracer?: Tracer; private concurrencyRedis: Redis; private defaultConcurrency: number; + private heartbeatIntervalMs: number; + private visibilityTimeoutMs: number; private maxAttempts: number; private processItemCallback?: ProcessBatchItemCallback; @@ -95,6 +100,8 @@ export class BatchQueue { this.logger = options.logger ?? new Logger("BatchQueue", options.logLevel ?? "info"); this.tracer = options.tracer; this.defaultConcurrency = options.defaultConcurrency ?? 10; + this.visibilityTimeoutMs = options.visibilityTimeoutMs ?? BATCH_ITEM_VISIBILITY_TIMEOUT_MS; + this.heartbeatIntervalMs = Math.max(50, Math.floor(this.visibilityTimeoutMs / 3)); this.maxAttempts = options.retry?.maxAttempts ?? 1; this.abortController = new AbortController(); this.workerQueueBlockingTimeoutSeconds = options.workerQueueBlockingTimeoutSeconds ?? 10; @@ -154,7 +161,7 @@ export class BatchQueue { shardCount: options.shardCount ?? 1, consumerCount: options.consumerCount, consumerIntervalMs: options.consumerIntervalMs, - visibilityTimeoutMs: 60_000, // 1 minute for batch item processing + visibilityTimeoutMs: this.visibilityTimeoutMs, startConsumers: false, // We control when to start cooloff: { enabled: false, @@ -752,6 +759,29 @@ export class BatchQueue { // Private - Message Handling // ============================================================================ + /** + * Keep extending a message's visibility deadline while its callback runs, and return a + * function that stops doing so. Without this an item slower than the visibility timeout + * is redelivered while still being processed, and the original consumer's completion + * then destroys the redelivery's in-flight record, silently dropping the item and + * leaving the batch short of its expected count forever. + */ + #startHeartbeat(messageId: string, queueId: string): () => void { + const interval = setInterval(() => { + this.fairQueue.heartbeatMessage(messageId, queueId).catch((error) => { + this.logger.debug("Batch item heartbeat failed", { + messageId, + queueId, + error: error instanceof Error ? error.message : String(error), + }); + }); + }, this.heartbeatIntervalMs); + + interval.unref?.(); + + return () => clearInterval(interval); + } + async #handleMessage(consumerId: string, messageId: string, queueId: string): Promise { // Get message data from FairQueue's in-flight storage const storedMessage = await this.fairQueue.getMessageData(messageId, queueId); @@ -820,9 +850,10 @@ export class BatchQueue { let processedCount: number; try { - const result = await this.#startSpan( - "BatchQueue.processItemCallback", - async (innerSpan) => { + const stopHeartbeat = this.#startHeartbeat(messageId, queueId); + let result: Awaited>; + try { + result = await this.#startSpan("BatchQueue.processItemCallback", async (innerSpan) => { innerSpan?.setAttributes({ "batch.id": batchId, "batch.itemIndex": itemIndex, @@ -837,8 +868,10 @@ export class BatchQueue { attempt, isFinalAttempt, }); - } - ); + }); + } finally { + stopHeartbeat(); + } if (result.success) { span?.setAttribute("batch.result", "success"); diff --git a/internal-packages/run-engine/src/batch-queue/tests/index.test.ts b/internal-packages/run-engine/src/batch-queue/tests/index.test.ts index 56386b32ab7..07a822eb606 100644 --- a/internal-packages/run-engine/src/batch-queue/tests/index.test.ts +++ b/internal-packages/run-engine/src/batch-queue/tests/index.test.ts @@ -953,4 +953,56 @@ describe("BatchQueue", () => { } ); }); + + describe("visibility heartbeat", () => { + redisTest( + "should not redeliver an item that takes longer than the visibility timeout", + { timeout: 60_000 }, + async ({ redisContainer }) => { + const queue = new BatchQueue({ + redis: { + host: redisContainer.getHost(), + port: redisContainer.getPort(), + keyPrefix: "test:", + }, + drr: { quantum: 5, maxDeficit: 50 }, + consumerCount: 2, + consumerIntervalMs: 50, + visibilityTimeoutMs: 1_000, + startConsumers: false, + }); + + const invocations: number[] = []; + + try { + queue.onProcessItem(async ({ itemIndex }) => { + const isFirst = invocations.length === 0; + invocations.push(itemIndex); + if (isFirst) { + await new Promise((resolve) => setTimeout(resolve, 9_000)); + } + return { success: true, runId: `run-${itemIndex}` }; + }); + + await queue.initializeBatch(createInitOptions("batch-hb", "env-hb", 1)); + await enqueueItems(queue, "batch-hb", "env-hb", createBatchItems(1)); + + queue.start(); + + await vi.waitFor( + () => { + expect(invocations.length).toBeGreaterThanOrEqual(1); + }, + { timeout: 10_000 } + ); + + await new Promise((resolve) => setTimeout(resolve, 14_000)); + + expect(invocations).toEqual([0]); + } finally { + await queue.close(); + } + } + ); + }); }); diff --git a/internal-packages/run-engine/src/batch-queue/types.ts b/internal-packages/run-engine/src/batch-queue/types.ts index 9aa085e286a..11ca1d8451a 100644 --- a/internal-packages/run-engine/src/batch-queue/types.ts +++ b/internal-packages/run-engine/src/batch-queue/types.ts @@ -214,6 +214,12 @@ export type BatchQueueOptions = { * Items wait in queue until capacity frees up. */ defaultConcurrency?: number; + /** + * How long a claimed item stays invisible before the reclaim loop takes it back. + * The item is heartbeated for as long as its callback runs, so this only bites when + * a consumer stops making progress. Defaults to 60s. + */ + visibilityTimeoutMs?: number; /** * Optional global rate limiter to limit processing across all consumers. * When configured, limits the max items/second processed globally. diff --git a/packages/redis-worker/src/fair-queue/concurrency.ts b/packages/redis-worker/src/fair-queue/concurrency.ts index 641899ccc8d..9d69cc342c2 100644 --- a/packages/redis-worker/src/fair-queue/concurrency.ts +++ b/packages/redis-worker/src/fair-queue/concurrency.ts @@ -103,7 +103,31 @@ export class ConcurrencyManager { pipeline.srem(key, messageId); } - await pipeline.exec(); + this.#assertPipelineSucceeded(await pipeline.exec(), 1); + } + + /** + * Throw if any command in a released pipeline failed. ioredis resolves `exec()` even when + * individual commands error, so an unchecked pipeline reports success while leaving the + * slot held, which strands it permanently once the caller drops the in-flight record. + */ + #assertPipelineSucceeded( + results: Array<[Error | null, unknown]> | null, + messageCount: number + ): void { + const errors = (results ?? []) + .map(([error]) => error) + .filter((error): error is Error => Boolean(error)); + + if (errors.length > 0) { + throw new Error( + `Failed to release ${errors.length} of ${ + results?.length ?? 0 + } concurrency slot commands across ${messageCount} message(s): ${errors + .map((error) => error.message) + .join("; ")}` + ); + } } /** @@ -127,7 +151,7 @@ export class ConcurrencyManager { } } - await pipeline.exec(); + this.#assertPipelineSucceeded(await pipeline.exec(), messages.length); } /** diff --git a/packages/redis-worker/src/fair-queue/index.ts b/packages/redis-worker/src/fair-queue/index.ts index 989dd776ae4..09bd9fb5ca0 100644 --- a/packages/redis-worker/src/fair-queue/index.ts +++ b/packages/redis-worker/src/fair-queue/index.ts @@ -1498,28 +1498,14 @@ export class FairQueue { originalTimestamp: storedMessage.timestamp, }; + // Complete in visibility manager + await this.visibilityManager.complete(storedMessage.id, storedMessage.queueId); + // Add to DLQ const pipeline = this.redis.pipeline(); pipeline.zadd(dlqKey, dlqMessage.deadLetteredAt, storedMessage.id); pipeline.hset(dlqDataKey, storedMessage.id, JSON.stringify(dlqMessage)); - const dlqResults = await pipeline.exec(); - - const dlqErrors = (dlqResults ?? []) - .map(([error]) => error) - .filter((error): error is Error => Boolean(error)); - - if (dlqErrors.length > 0) { - this.logger.error("Failed to write message to DLQ, leaving it in-flight to be reclaimed", { - messageId: storedMessage.id, - queueId: storedMessage.queueId, - tenantId: storedMessage.tenantId, - errors: dlqErrors.map((error) => error.message), - }); - return; - } - - // Complete in visibility manager - await this.visibilityManager.complete(storedMessage.id, storedMessage.queueId); + await pipeline.exec(); this.telemetry.recordDLQ(); @@ -1559,55 +1545,56 @@ export class FairQueue { } /** - * Free a reclaimed message's concurrency slot. Runs before the message is put back - * on the queue, so another consumer cannot reserve the same messageId and have its - * fresh reservation deleted by this release. Failures are logged rather than thrown - * so one bad message cannot stop the rest of the shard being reclaimed. + * Free every timed-out message's concurrency slot in one pipeline, before any of them + * is put back on the queue. Throwing here aborts the requeue for the whole batch, which + * leaves the messages in-flight for the next reclaim tick. That is the safe direction: + * requeuing a message whose slot is still held is what strands the slot permanently. */ - async #releaseReclaimedConcurrency(message: ReclaimedMessageInfo): Promise { - if (!this.concurrencyManager) { + async #releaseReclaimedConcurrency(messages: ReclaimedMessageInfo[]): Promise { + if (!this.concurrencyManager || messages.length === 0) { return; } - try { - await this.concurrencyManager.release( - { + await this.concurrencyManager.releaseBatch( + messages.map((message) => ({ + queue: { id: message.queueId, tenantId: message.tenantId, metadata: message.metadata ?? {}, }, - message.messageId - ); - } catch (error) { - this.logger.error("Failed to release concurrency for reclaimed message", { messageId: message.messageId, - queueId: message.queueId, - error: error instanceof Error ? error.message : String(error), - }); - } + })) + ); } async #reclaimTimedOutMessages(): Promise { let totalReclaimed = 0; for (let shardId = 0; shardId < this.shardCount; shardId++) { - const reclaimedMessages = await this.visibilityManager.reclaimTimedOut( - shardId, - (queueId) => { - const tenantId = this.keys.extractTenantId(queueId); - const dispatchShardId = this.tenantDispatch.getShardForTenant(tenantId); - return { - queueKey: this.keys.queueKey(queueId), - queueItemsKey: this.keys.queueItemsKey(queueId), - tenantQueueIndexKey: this.keys.tenantQueueIndexKey(tenantId), - dispatchKey: this.keys.dispatchKey(dispatchShardId), - tenantId, - }; - }, - this.#releaseReclaimedConcurrency.bind(this) - ); + try { + const reclaimedMessages = await this.visibilityManager.reclaimTimedOut( + shardId, + (queueId) => { + const tenantId = this.keys.extractTenantId(queueId); + const dispatchShardId = this.tenantDispatch.getShardForTenant(tenantId); + return { + queueKey: this.keys.queueKey(queueId), + queueItemsKey: this.keys.queueItemsKey(queueId), + tenantQueueIndexKey: this.keys.tenantQueueIndexKey(tenantId), + dispatchKey: this.keys.dispatchKey(dispatchShardId), + tenantId, + }; + }, + this.#releaseReclaimedConcurrency.bind(this) + ); - totalReclaimed += reclaimedMessages.length; + totalReclaimed += reclaimedMessages.length; + } catch (error) { + this.logger.error("Failed to reclaim shard, leaving messages in-flight for the next tick", { + shardId, + error: error instanceof Error ? error.message : String(error), + }); + } } if (totalReclaimed > 0) { 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 e84e2af3324..e6255a885dc 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -1442,7 +1442,7 @@ describe("FairQueue", () => { ); redisTest( - "should release metadata-derived concurrency groups when the in-flight record is gone", + "should release metadata-derived groups from the cached descriptor when the in-flight record is gone", { timeout: 15000 }, async ({ redisOptions }) => { const processed: string[] = []; @@ -1518,10 +1518,9 @@ describe("FairQueue", () => { ); redisTest( - "should keep a message in-flight when the DLQ write fails", - { timeout: 15000 }, + "should not requeue a reclaimed message when freeing its slot fails", + { timeout: 20000 }, async ({ redisOptions }) => { - const started: string[] = []; keys = new DefaultFairQueueKeyProducer({ prefix: "test" }); const scheduler = new DRRScheduler({ @@ -1537,8 +1536,8 @@ describe("FairQueue", () => { shardCount: 1, consumerCount: 1, consumerIntervalMs: 20, - visibilityTimeoutMs: 60000, - reclaimIntervalMs: 60000, + visibilityTimeoutMs: 500, + reclaimIntervalMs: 100, concurrencyGroups: [ { name: "tenant", @@ -1551,41 +1550,38 @@ describe("FairQueue", () => { }); const redis = createRedisClient(redisOptions); - const dlqKey = keys.deadLetterQueueKey("t1"); + const queueId = "tenant:t1:queue:release-fails"; + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); try { - await redis.set(dlqKey, "not-a-zset"); - - queue.onMessage(async (ctx) => { - started.push(ctx.message.payload.value); - await ctx.fail(new Error("boom")); + queue.onMessage(async () => { + await new Promise((resolve) => setTimeout(resolve, 15000)); }); await queue.enqueue({ - queueId: "tenant:t1:queue:q1", + queueId, tenantId: "t1", payload: { value: "msg-0" }, }); queue.start(); - await vi.waitFor( - () => { - expect(started).toHaveLength(1); - }, - { timeout: 10000 } - ); - await vi.waitFor( async () => { - expect(await redis.scard(keys.concurrencyKey("tenant", "t1"))).toBe(0); + expect(await redis.zcard(keys.inflightKey(0))).toBe(1); }, { timeout: 5000 } ); + await redis.del(concurrencyKey); + await redis.set(concurrencyKey, "not-a-set"); + + await new Promise((resolve) => setTimeout(resolve, 2000)); + expect(await redis.zcard(keys.inflightKey(0))).toBe(1); + expect(await redis.zcard(keys.queueKey(queueId))).toBe(0); } finally { - await redis.del(dlqKey); + await redis.del(concurrencyKey); await redis.quit(); await queue.close(); } @@ -1611,7 +1607,7 @@ describe("FairQueue", () => { shardCount: 1, consumerCount: 1, consumerIntervalMs: 20, - visibilityTimeoutMs: 200, + visibilityTimeoutMs: 30000, reclaimIntervalMs: 100, concurrencyGroups: [ { @@ -1626,10 +1622,16 @@ describe("FairQueue", () => { const redis = createRedisClient(redisOptions); const queueId = "tenant:t1:queue:reclaim-fail"; + const inflightKey = keys.inflightKey(0); + + let unblockHandler: (() => void) | undefined; + const handlerBlocked = new Promise((resolve) => { + unblockHandler = resolve; + }); try { queue.onMessage(async () => { - await new Promise((resolve) => setTimeout(resolve, 10000)); + await handlerBlocked; }); await queue.enqueue({ @@ -1643,6 +1645,7 @@ describe("FairQueue", () => { await vi.waitFor( async () => { expect(await redis.scard(keys.concurrencyKey("tenant", "t1"))).toBe(1); + expect(await redis.zcard(inflightKey)).toBe(1); }, { timeout: 5000 } ); @@ -1650,13 +1653,20 @@ describe("FairQueue", () => { await redis.del(keys.queueKey(queueId)); await redis.set(keys.queueKey(queueId), "not-a-zset"); + const [member] = await redis.zrange(inflightKey, 0, 0); + expect(member).toBeDefined(); + await redis.zadd(inflightKey, Date.now() - 60000, member!); + await vi.waitFor( async () => { expect(await redis.scard(keys.concurrencyKey("tenant", "t1"))).toBe(0); }, { timeout: 8000 } ); + + expect(await redis.zcard(inflightKey)).toBe(1); } finally { + unblockHandler?.(); await redis.del(keys.queueKey(queueId)); await redis.quit(); await queue.close(); diff --git a/packages/redis-worker/src/fair-queue/visibility.ts b/packages/redis-worker/src/fair-queue/visibility.ts index 5713c853378..14f559ca607 100644 --- a/packages/redis-worker/src/fair-queue/visibility.ts +++ b/packages/redis-worker/src/fair-queue/visibility.ts @@ -399,7 +399,7 @@ export class VisibilityManager { dispatchKey: string; tenantId: string; }, - onBeforeRequeue?: (message: ReclaimedMessageInfo) => Promise + onBeforeRequeue?: (messages: ReclaimedMessageInfo[]) => Promise ): Promise { const inflightKey = this.keys.inflightKey(shardId); const inflightDataKey = this.keys.inflightDataKey(shardId); @@ -416,39 +416,67 @@ export class VisibilityManager { 100 // Process in batches ); - const reclaimedMessages: ReclaimedMessageInfo[] = []; + const candidates: Array<{ + member: string; + deadlineScore: string; + info: ReclaimedMessageInfo; + storedMessage: StoredMessage | null; + }> = []; for (let i = 0; i < timedOut.length; i += 2) { const member = timedOut[i]; - const _deadlineScore = timedOut[i + 1]; // This is the visibility deadline, not the original timestamp - if (!member || !_deadlineScore) { + const deadlineScore = timedOut[i + 1]; + if (!member || !deadlineScore) { continue; } const { messageId, queueId } = this.#parseMember(member); - const { queueKey, queueItemsKey, tenantQueueIndexKey, dispatchKey, tenantId } = - getQueueKeys(queueId); - try { - // Get message data BEFORE releasing so we can extract tenantId for concurrency release - const dataJson = await this.redis.hget(inflightDataKey, messageId); - let storedMessage: StoredMessage | null = null; - if (dataJson) { - try { - storedMessage = JSON.parse(dataJson); - } catch { - // Ignore parse error, proceed with reclaim - } + const dataJson = await this.redis.hget(inflightDataKey, messageId); + let storedMessage: StoredMessage | null = null; + if (dataJson) { + try { + storedMessage = JSON.parse(dataJson); + } catch { + // Ignore parse error, proceed with reclaim } + } - if (onBeforeRequeue) { - await onBeforeRequeue({ - messageId, - queueId, - tenantId: storedMessage?.tenantId ?? tenantId, - metadata: storedMessage?.metadata, - }); - } + if (!storedMessage) { + this.logger.error("Missing or corrupted message data during reclaim, using fallback", { + messageId, + queueId, + }); + } + + candidates.push({ + member, + deadlineScore, + storedMessage, + info: { + messageId, + queueId, + tenantId: storedMessage?.tenantId ?? this.keys.extractTenantId(queueId), + metadata: storedMessage?.metadata ?? {}, + }, + }); + } + + if (candidates.length === 0) { + return []; + } + + if (onBeforeRequeue) { + await onBeforeRequeue(candidates.map((candidate) => candidate.info)); + } + const reclaimedMessages: ReclaimedMessageInfo[] = []; + + for (const { member, deadlineScore, storedMessage, info } of candidates) { + const { messageId, queueId } = info; + const { queueKey, queueItemsKey, tenantQueueIndexKey, dispatchKey, tenantId } = + getQueueKeys(queueId); + + try { // Re-add to queue with original timestamp to preserve priority // Fall back to now if we can't get the original timestamp const score = storedMessage?.timestamp ?? now; @@ -467,34 +495,12 @@ export class VisibilityManager { tenantId ); - // Track reclaimed message for concurrency release - // Always add to reclaimedMessages to avoid concurrency leaks - if (storedMessage) { - reclaimedMessages.push({ - messageId, - queueId, - tenantId: storedMessage.tenantId, - metadata: storedMessage.metadata, - }); - } else { - // Fallback: extract tenantId from queueId when message data is missing or corrupted - // This ensures concurrency is released even if we can't get the full metadata - this.logger.error("Missing or corrupted message data during reclaim, using fallback", { - messageId, - queueId, - }); - reclaimedMessages.push({ - messageId, - queueId, - tenantId: this.keys.extractTenantId(queueId), - metadata: {}, - }); - } + reclaimedMessages.push(info); this.logger.debug("Reclaimed timed-out message", { messageId, queueId, - deadline: _deadlineScore, + deadline: deadlineScore, }); } catch (error) { this.logger.error("Failed to reclaim message", { @@ -729,14 +735,15 @@ if updatedData and updatedData ~= "" then payload = updatedData end +-- Add back to queue before dropping the in-flight record. Lua does not roll back on +-- error, so removing first would lose the message outright if the queue write failed. +redis.call('ZADD', queueKey, score, messageId) +redis.call('HSET', queueItemsKey, messageId, payload) + -- Remove from in-flight redis.call('ZREM', inflightKey, member) redis.call('HDEL', inflightDataKey, messageId) --- Add back to queue -redis.call('ZADD', queueKey, score, messageId) -redis.call('HSET', queueItemsKey, messageId, payload) - -- Update tenant queue index (Level 2) with queue's oldest message local oldest = redis.call('ZRANGE', queueKey, 0, 0, 'WITHSCORES') if #oldest >= 2 then @@ -783,14 +790,15 @@ for i = 0, numMessages - 1 do -- Get message data from in-flight local payload = redis.call('HGET', inflightDataKey, messageId) if payload then + -- Add back to queue before dropping the in-flight record, so a failed queue write + -- cannot lose the message: Lua does not roll back already-applied commands. + redis.call('ZADD', queueKey, score, messageId) + redis.call('HSET', queueItemsKey, messageId, payload) + -- Remove from in-flight redis.call('ZREM', inflightKey, member) redis.call('HDEL', inflightDataKey, messageId) - -- Add back to queue - redis.call('ZADD', queueKey, score, messageId) - redis.call('HSET', queueItemsKey, messageId, payload) - releasedCount = releasedCount + 1 else -- Dangling in-flight entry with no payload: drop it so it stops consuming From a2f12f2b5d983ead4dbc29a345fa58cf83dae229 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 11 Aug 2026 13:51:34 +0100 Subject: [PATCH 06/12] fix(redis-worker,run-engine): give the item heartbeat real margin and stop one bad key stalling a shard Each beat extended the deadline by the same amount as the gap between beats, so timer drift and round-trip latency left the message briefly past its deadline on every cycle, and a reclaim scan landing in that window still handed a running item to a second consumer. Beats now extend by a full visibility timeout while the tick stays at a third of it. Reclaim released the whole timed-out batch in one pipeline and aborted every requeue if any command failed, so a single unusable concurrency key stalled reclaim for every other tenant in that shard. A failed batch now falls back to releasing message by message and only the messages whose slot could not be freed are held back. --- .../run-engine/src/batch-queue/index.ts | 1 + packages/redis-worker/src/fair-queue/index.ts | 48 ++++++++++++++----- .../src/fair-queue/tests/fairQueue.test.ts | 10 ++++ .../redis-worker/src/fair-queue/visibility.ts | 18 +++++-- 4 files changed, 61 insertions(+), 16 deletions(-) diff --git a/internal-packages/run-engine/src/batch-queue/index.ts b/internal-packages/run-engine/src/batch-queue/index.ts index b2c7757a901..d6ec9110698 100644 --- a/internal-packages/run-engine/src/batch-queue/index.ts +++ b/internal-packages/run-engine/src/batch-queue/index.ts @@ -162,6 +162,7 @@ export class BatchQueue { consumerCount: options.consumerCount, consumerIntervalMs: options.consumerIntervalMs, visibilityTimeoutMs: this.visibilityTimeoutMs, + heartbeatIntervalMs: this.visibilityTimeoutMs, startConsumers: false, // We control when to start cooloff: { enabled: false, diff --git a/packages/redis-worker/src/fair-queue/index.ts b/packages/redis-worker/src/fair-queue/index.ts index 09bd9fb5ca0..749259479f6 100644 --- a/packages/redis-worker/src/fair-queue/index.ts +++ b/packages/redis-worker/src/fair-queue/index.ts @@ -1550,21 +1550,45 @@ export class FairQueue { * leaves the messages in-flight for the next reclaim tick. That is the safe direction: * requeuing a message whose slot is still held is what strands the slot permanently. */ - async #releaseReclaimedConcurrency(messages: ReclaimedMessageInfo[]): Promise { + async #releaseReclaimedConcurrency(messages: ReclaimedMessageInfo[]): Promise { if (!this.concurrencyManager || messages.length === 0) { - return; + return []; } - await this.concurrencyManager.releaseBatch( - messages.map((message) => ({ - queue: { - id: message.queueId, - tenantId: message.tenantId, - metadata: message.metadata ?? {}, - }, - messageId: message.messageId, - })) - ); + const descriptorFor = (message: ReclaimedMessageInfo) => ({ + id: message.queueId, + tenantId: message.tenantId, + metadata: message.metadata ?? {}, + }); + + try { + await this.concurrencyManager.releaseBatch( + messages.map((message) => ({ queue: descriptorFor(message), messageId: message.messageId })) + ); + return []; + } catch (error) { + this.logger.error("Batch concurrency release failed, retrying message by message", { + count: messages.length, + error: error instanceof Error ? error.message : String(error), + }); + } + + const failed: string[] = []; + + for (const message of messages) { + try { + await this.concurrencyManager.release(descriptorFor(message), message.messageId); + } catch (error) { + failed.push(message.messageId); + this.logger.error("Failed to release concurrency for reclaimed message", { + messageId: message.messageId, + queueId: message.queueId, + error: error instanceof Error ? error.message : String(error), + }); + } + } + + return failed; } async #reclaimTimedOutMessages(): Promise { 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 e6255a885dc..cf08c04c0e6 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -1580,6 +1580,16 @@ describe("FairQueue", () => { expect(await redis.zcard(keys.inflightKey(0))).toBe(1); expect(await redis.zcard(keys.queueKey(queueId))).toBe(0); + + await redis.del(concurrencyKey); + + await vi.waitFor( + async () => { + expect(await redis.zcard(keys.queueKey(queueId))).toBe(1); + expect(await redis.zcard(keys.inflightKey(0))).toBe(0); + }, + { timeout: 8000 } + ); } finally { await redis.del(concurrencyKey); await redis.quit(); diff --git a/packages/redis-worker/src/fair-queue/visibility.ts b/packages/redis-worker/src/fair-queue/visibility.ts index 14f559ca607..65d633a5fc1 100644 --- a/packages/redis-worker/src/fair-queue/visibility.ts +++ b/packages/redis-worker/src/fair-queue/visibility.ts @@ -399,7 +399,7 @@ export class VisibilityManager { dispatchKey: string; tenantId: string; }, - onBeforeRequeue?: (messages: ReclaimedMessageInfo[]) => Promise + onBeforeRequeue?: (messages: ReclaimedMessageInfo[]) => Promise ): Promise { const inflightKey = this.keys.inflightKey(shardId); const inflightDataKey = this.keys.inflightDataKey(shardId); @@ -465,14 +465,24 @@ export class VisibilityManager { return []; } - if (onBeforeRequeue) { - await onBeforeRequeue(candidates.map((candidate) => candidate.info)); - } + const notReleased = new Set( + (onBeforeRequeue + ? await onBeforeRequeue(candidates.map((candidate) => candidate.info)) + : undefined) ?? [] + ); const reclaimedMessages: ReclaimedMessageInfo[] = []; for (const { member, deadlineScore, storedMessage, info } of candidates) { const { messageId, queueId } = info; + + if (notReleased.has(messageId)) { + this.logger.error("Skipping requeue, concurrency slot was not released", { + messageId, + queueId, + }); + continue; + } const { queueKey, queueItemsKey, tenantQueueIndexKey, dispatchKey, tenantId } = getQueueKeys(queueId); From ff348c4822eb6d4489d571a2f6ed474c8bf1bfd9 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 11 Aug 2026 13:56:26 +0100 Subject: [PATCH 07/12] fix(run-engine,redis-worker): drop a batch item result once its lease is lost The heartbeat already learns when another consumer has taken the item over: the in-flight member is gone, so the extend returns false. That signal was discarded, leaving the original consumer free to finish and complete over the new owner's in-flight record. It now stops and discards its result instead, so the consumer that actually owns the item is the one that finishes it. Also treat a discarded release pipeline as a failure rather than a success, since a null result means none of the commands ran. --- .../run-engine/src/batch-queue/index.ts | 42 ++++++++++++++----- .../src/fair-queue/concurrency.ts | 8 +++- 2 files changed, 39 insertions(+), 11 deletions(-) diff --git a/internal-packages/run-engine/src/batch-queue/index.ts b/internal-packages/run-engine/src/batch-queue/index.ts index d6ec9110698..f2ba0654643 100644 --- a/internal-packages/run-engine/src/batch-queue/index.ts +++ b/internal-packages/run-engine/src/batch-queue/index.ts @@ -767,20 +767,32 @@ export class BatchQueue { * then destroys the redelivery's in-flight record, silently dropping the item and * leaving the batch short of its expected count forever. */ - #startHeartbeat(messageId: string, queueId: string): () => void { + #startHeartbeat( + messageId: string, + queueId: string + ): { stop: () => void; lostLease: () => boolean } { + let lostLease = false; + const interval = setInterval(() => { - this.fairQueue.heartbeatMessage(messageId, queueId).catch((error) => { - this.logger.debug("Batch item heartbeat failed", { - messageId, - queueId, - error: error instanceof Error ? error.message : String(error), + this.fairQueue + .heartbeatMessage(messageId, queueId) + .then((stillOwned) => { + if (!stillOwned) { + lostLease = true; + } + }) + .catch((error) => { + this.logger.debug("Batch item heartbeat failed", { + messageId, + queueId, + error: error instanceof Error ? error.message : String(error), + }); }); - }); }, this.heartbeatIntervalMs); interval.unref?.(); - return () => clearInterval(interval); + return { stop: () => clearInterval(interval), lostLease: () => lostLease }; } async #handleMessage(consumerId: string, messageId: string, queueId: string): Promise { @@ -851,7 +863,7 @@ export class BatchQueue { let processedCount: number; try { - const stopHeartbeat = this.#startHeartbeat(messageId, queueId); + const heartbeat = this.#startHeartbeat(messageId, queueId); let result: Awaited>; try { result = await this.#startSpan("BatchQueue.processItemCallback", async (innerSpan) => { @@ -871,7 +883,17 @@ export class BatchQueue { }); }); } finally { - stopHeartbeat(); + heartbeat.stop(); + } + + if (heartbeat.lostLease()) { + this.logger.warn("Discarding batch item result, another consumer now owns it", { + batchId, + itemIndex, + messageId, + attempt, + }); + return; } if (result.success) { diff --git a/packages/redis-worker/src/fair-queue/concurrency.ts b/packages/redis-worker/src/fair-queue/concurrency.ts index 9d69cc342c2..a06592b644a 100644 --- a/packages/redis-worker/src/fair-queue/concurrency.ts +++ b/packages/redis-worker/src/fair-queue/concurrency.ts @@ -115,7 +115,13 @@ export class ConcurrencyManager { results: Array<[Error | null, unknown]> | null, messageCount: number ): void { - const errors = (results ?? []) + if (results === null) { + throw new Error( + `Concurrency release pipeline for ${messageCount} message(s) was discarded without executing` + ); + } + + const errors = results .map(([error]) => error) .filter((error): error is Error => Boolean(error)); From 1b2491520ed4433ab10a0efbd3285f01b2a0d2da Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 11 Aug 2026 14:34:14 +0100 Subject: [PATCH 08/12] docs(run-engine): state what the lease signal actually guarantees The in-flight member is keyed only by message and queue id, so once another consumer re-claims a reclaimed item the member exists again and an extend from the previous consumer succeeds. The signal therefore only catches the window where the item is back on the queue and unclaimed, which is narrower than the comment and the changeset claimed. --- .changeset/fair-queue-concurrency-slot-leak.md | 4 +++- .../run-engine/src/batch-queue/index.ts | 13 ++++++++----- 2 files changed, 11 insertions(+), 6 deletions(-) diff --git a/.changeset/fair-queue-concurrency-slot-leak.md b/.changeset/fair-queue-concurrency-slot-leak.md index 4ccced1c9c9..aa325118421 100644 --- a/.changeset/fair-queue-concurrency-slot-leak.md +++ b/.changeset/fair-queue-concurrency-slot-leak.md @@ -2,4 +2,6 @@ "@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. +Fair queue consumers no longer leak the concurrency slots that gate a tenant's throughput. Slots were held by messages that had already finished, were never reclaimed, and once enough of them accumulated every queue belonging to that tenant stopped being served. Slots are now freed on the paths that previously skipped them, freed before the record needed to recover them is discarded, and released before a reclaimed message goes back on the queue. A failed release is now surfaced instead of being silently treated as success. + +Concurrency groups keyed on queue metadata rather than the tenant can still resolve to the wrong group when a consumer completes a message it did not enqueue, so this does not yet cover that case. diff --git a/internal-packages/run-engine/src/batch-queue/index.ts b/internal-packages/run-engine/src/batch-queue/index.ts index f2ba0654643..bbf39dafce0 100644 --- a/internal-packages/run-engine/src/batch-queue/index.ts +++ b/internal-packages/run-engine/src/batch-queue/index.ts @@ -761,11 +761,14 @@ export class BatchQueue { // ============================================================================ /** - * Keep extending a message's visibility deadline while its callback runs, and return a - * function that stops doing so. Without this an item slower than the visibility timeout - * is redelivered while still being processed, and the original consumer's completion - * then destroys the redelivery's in-flight record, silently dropping the item and - * leaving the batch short of its expected count forever. + * Keep extending a message's visibility deadline while its callback runs, so an item + * slower than the visibility timeout is not redelivered and executed a second time. + * + * `lostLease` reports that an extend found no in-flight entry, which means the item was + * reclaimed and is now back on the queue. It is a best-effort signal, not a fence: the + * in-flight member is keyed only by message and queue id, so once another consumer + * re-claims the item the member exists again and an extend from this consumer succeeds. + * Distinguishing owners would need a per-claim token in the member. */ #startHeartbeat( messageId: string, From 01d06660054ae64c6b3f6bf99ddfaa9dce5e98da Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 11 Aug 2026 16:12:46 +0100 Subject: [PATCH 09/12] chore(run-engine): move batch item heartbeating to its own branch The heartbeat fixes a separate pre-existing bug (items slower than the visibility timeout are redelivered and executed twice) and shares no files with the concurrency slot fixes, so it ships on its own. --- .../run-engine/src/batch-queue/index.ts | 71 ++----------------- .../src/batch-queue/tests/index.test.ts | 52 -------------- .../run-engine/src/batch-queue/types.ts | 6 -- 3 files changed, 6 insertions(+), 123 deletions(-) diff --git a/internal-packages/run-engine/src/batch-queue/index.ts b/internal-packages/run-engine/src/batch-queue/index.ts index bbf39dafce0..10354e77c52 100644 --- a/internal-packages/run-engine/src/batch-queue/index.ts +++ b/internal-packages/run-engine/src/batch-queue/index.ts @@ -59,9 +59,6 @@ const ENV_CONCURRENCY_KEY_PREFIX = "batch:env_concurrency"; // then all messages are routed to this queue for BatchQueue's own consumer loop. const BATCH_WORKER_QUEUE_ID = "batch-worker-queue"; -/** How long a claimed batch item stays invisible before the reclaim loop takes it back. */ -const BATCH_ITEM_VISIBILITY_TIMEOUT_MS = 60_000; - export class BatchQueue { private fairQueue: FairQueue; private workerQueueManager: WorkerQueueManager; @@ -70,8 +67,6 @@ export class BatchQueue { private tracer?: Tracer; private concurrencyRedis: Redis; private defaultConcurrency: number; - private heartbeatIntervalMs: number; - private visibilityTimeoutMs: number; private maxAttempts: number; private processItemCallback?: ProcessBatchItemCallback; @@ -100,8 +95,6 @@ export class BatchQueue { this.logger = options.logger ?? new Logger("BatchQueue", options.logLevel ?? "info"); this.tracer = options.tracer; this.defaultConcurrency = options.defaultConcurrency ?? 10; - this.visibilityTimeoutMs = options.visibilityTimeoutMs ?? BATCH_ITEM_VISIBILITY_TIMEOUT_MS; - this.heartbeatIntervalMs = Math.max(50, Math.floor(this.visibilityTimeoutMs / 3)); this.maxAttempts = options.retry?.maxAttempts ?? 1; this.abortController = new AbortController(); this.workerQueueBlockingTimeoutSeconds = options.workerQueueBlockingTimeoutSeconds ?? 10; @@ -161,8 +154,7 @@ export class BatchQueue { shardCount: options.shardCount ?? 1, consumerCount: options.consumerCount, consumerIntervalMs: options.consumerIntervalMs, - visibilityTimeoutMs: this.visibilityTimeoutMs, - heartbeatIntervalMs: this.visibilityTimeoutMs, + visibilityTimeoutMs: 60_000, // 1 minute for batch item processing startConsumers: false, // We control when to start cooloff: { enabled: false, @@ -760,44 +752,6 @@ export class BatchQueue { // Private - Message Handling // ============================================================================ - /** - * Keep extending a message's visibility deadline while its callback runs, so an item - * slower than the visibility timeout is not redelivered and executed a second time. - * - * `lostLease` reports that an extend found no in-flight entry, which means the item was - * reclaimed and is now back on the queue. It is a best-effort signal, not a fence: the - * in-flight member is keyed only by message and queue id, so once another consumer - * re-claims the item the member exists again and an extend from this consumer succeeds. - * Distinguishing owners would need a per-claim token in the member. - */ - #startHeartbeat( - messageId: string, - queueId: string - ): { stop: () => void; lostLease: () => boolean } { - let lostLease = false; - - const interval = setInterval(() => { - this.fairQueue - .heartbeatMessage(messageId, queueId) - .then((stillOwned) => { - if (!stillOwned) { - lostLease = true; - } - }) - .catch((error) => { - this.logger.debug("Batch item heartbeat failed", { - messageId, - queueId, - error: error instanceof Error ? error.message : String(error), - }); - }); - }, this.heartbeatIntervalMs); - - interval.unref?.(); - - return { stop: () => clearInterval(interval), lostLease: () => lostLease }; - } - async #handleMessage(consumerId: string, messageId: string, queueId: string): Promise { // Get message data from FairQueue's in-flight storage const storedMessage = await this.fairQueue.getMessageData(messageId, queueId); @@ -866,10 +820,9 @@ export class BatchQueue { let processedCount: number; try { - const heartbeat = this.#startHeartbeat(messageId, queueId); - let result: Awaited>; - try { - result = await this.#startSpan("BatchQueue.processItemCallback", async (innerSpan) => { + const result = await this.#startSpan( + "BatchQueue.processItemCallback", + async (innerSpan) => { innerSpan?.setAttributes({ "batch.id": batchId, "batch.itemIndex": itemIndex, @@ -884,20 +837,8 @@ export class BatchQueue { attempt, isFinalAttempt, }); - }); - } finally { - heartbeat.stop(); - } - - if (heartbeat.lostLease()) { - this.logger.warn("Discarding batch item result, another consumer now owns it", { - batchId, - itemIndex, - messageId, - attempt, - }); - return; - } + } + ); if (result.success) { span?.setAttribute("batch.result", "success"); diff --git a/internal-packages/run-engine/src/batch-queue/tests/index.test.ts b/internal-packages/run-engine/src/batch-queue/tests/index.test.ts index 07a822eb606..56386b32ab7 100644 --- a/internal-packages/run-engine/src/batch-queue/tests/index.test.ts +++ b/internal-packages/run-engine/src/batch-queue/tests/index.test.ts @@ -953,56 +953,4 @@ describe("BatchQueue", () => { } ); }); - - describe("visibility heartbeat", () => { - redisTest( - "should not redeliver an item that takes longer than the visibility timeout", - { timeout: 60_000 }, - async ({ redisContainer }) => { - const queue = new BatchQueue({ - redis: { - host: redisContainer.getHost(), - port: redisContainer.getPort(), - keyPrefix: "test:", - }, - drr: { quantum: 5, maxDeficit: 50 }, - consumerCount: 2, - consumerIntervalMs: 50, - visibilityTimeoutMs: 1_000, - startConsumers: false, - }); - - const invocations: number[] = []; - - try { - queue.onProcessItem(async ({ itemIndex }) => { - const isFirst = invocations.length === 0; - invocations.push(itemIndex); - if (isFirst) { - await new Promise((resolve) => setTimeout(resolve, 9_000)); - } - return { success: true, runId: `run-${itemIndex}` }; - }); - - await queue.initializeBatch(createInitOptions("batch-hb", "env-hb", 1)); - await enqueueItems(queue, "batch-hb", "env-hb", createBatchItems(1)); - - queue.start(); - - await vi.waitFor( - () => { - expect(invocations.length).toBeGreaterThanOrEqual(1); - }, - { timeout: 10_000 } - ); - - await new Promise((resolve) => setTimeout(resolve, 14_000)); - - expect(invocations).toEqual([0]); - } finally { - await queue.close(); - } - } - ); - }); }); diff --git a/internal-packages/run-engine/src/batch-queue/types.ts b/internal-packages/run-engine/src/batch-queue/types.ts index 11ca1d8451a..9aa085e286a 100644 --- a/internal-packages/run-engine/src/batch-queue/types.ts +++ b/internal-packages/run-engine/src/batch-queue/types.ts @@ -214,12 +214,6 @@ export type BatchQueueOptions = { * Items wait in queue until capacity frees up. */ defaultConcurrency?: number; - /** - * How long a claimed item stays invisible before the reclaim loop takes it back. - * The item is heartbeated for as long as its callback runs, so this only bites when - * a consumer stops making progress. Defaults to 60s. - */ - visibilityTimeoutMs?: number; /** * Optional global rate limiter to limit processing across all consumers. * When configured, limits the max items/second processed globally. From 72947fe22689272cb1411b3700ca0d3eab518b00 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Tue, 11 Aug 2026 17:38:46 +0100 Subject: [PATCH 10/12] test(redis-worker): assert reclaim recovery on a stable signal The test waited for the message to be queued and absent from in-flight at the same instant, a state that only lasts about one dispatch interval before the message is re-claimed. It now waits for the stuck in-flight entry's deadline to move forward instead, which is monotonic once reclaim has run. --- .../redis-worker/src/fair-queue/tests/fairQueue.test.ts | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) 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 cf08c04c0e6..5fbf2299fab 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -1581,12 +1581,15 @@ describe("FairQueue", () => { expect(await redis.zcard(keys.inflightKey(0))).toBe(1); expect(await redis.zcard(keys.queueKey(queueId))).toBe(0); + const [stuckMember] = await redis.zrange(keys.inflightKey(0), 0, 0); + const stuckDeadline = Number(await redis.zscore(keys.inflightKey(0), stuckMember!)); + await redis.del(concurrencyKey); await vi.waitFor( async () => { - expect(await redis.zcard(keys.queueKey(queueId))).toBe(1); - expect(await redis.zcard(keys.inflightKey(0))).toBe(0); + const score = await redis.zscore(keys.inflightKey(0), stuckMember!); + expect(score === null || Number(score) > stuckDeadline).toBe(true); }, { timeout: 8000 } ); From e7c3c86307dac25eb3ecfb1329312205a9faefee Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Wed, 12 Aug 2026 13:25:47 +0100 Subject: [PATCH 11/12] fix(redis-worker): heal leaked concurrency slots instead of blocking on release failures A failed slot release previously aborted its caller: completeMessage left the message in flight to be re-delivered as a duplicate execution, the retry path lost the attempt increment, and the reclaim path held messages in flight indefinitely. Release is now best-effort on every path: the primary state transition proceeds and the failure is logged. The resulting leaks are recoverable in two ways. reserve re-admits a message that already holds its own slot, since re-admission does not increase concurrency. A reconcile loop removes any slot member with no in-flight record; the check-and-remove is atomic and sound because a message is always registered in flight before its slot is reserved. --- .../fair-queue-concurrency-slot-leak.md | 4 +- .../src/fair-queue/concurrency.ts | 125 +++++- packages/redis-worker/src/fair-queue/index.ts | 149 +++++-- .../src/fair-queue/tests/concurrency.test.ts | 83 ++++ .../src/fair-queue/tests/fairQueue.test.ts | 385 +++++++++++++++++- packages/redis-worker/src/fair-queue/types.ts | 2 + .../redis-worker/src/fair-queue/visibility.ts | 22 +- 7 files changed, 699 insertions(+), 71 deletions(-) diff --git a/.changeset/fair-queue-concurrency-slot-leak.md b/.changeset/fair-queue-concurrency-slot-leak.md index aa325118421..bfdfa312c1a 100644 --- a/.changeset/fair-queue-concurrency-slot-leak.md +++ b/.changeset/fair-queue-concurrency-slot-leak.md @@ -2,6 +2,4 @@ "@trigger.dev/redis-worker": patch --- -Fair queue consumers no longer leak the concurrency slots that gate a tenant's throughput. Slots were held by messages that had already finished, were never reclaimed, and once enough of them accumulated every queue belonging to that tenant stopped being served. Slots are now freed on the paths that previously skipped them, freed before the record needed to recover them is discarded, and released before a reclaimed message goes back on the queue. A failed release is now surfaced instead of being silently treated as success. - -Concurrency groups keyed on queue metadata rather than the tenant can still resolve to the wrong group when a consumer completes a message it did not enqueue, so this does not yet cover that case. +Fair queue consumers no longer leak the concurrency slots that gate a tenant's throughput, and leaked slots now heal themselves. Slots are freed on the completion, retry, dead-letter, and reclaim paths that previously skipped them, and a failed release never blocks the message's own state transition, so a Redis error can no longer turn into a duplicate execution or a lost retry. A periodic reconcile loop removes any slot whose message is no longer in flight, and a message that still holds its own slot from an earlier failed release is re-admitted instead of being blocked by it. diff --git a/packages/redis-worker/src/fair-queue/concurrency.ts b/packages/redis-worker/src/fair-queue/concurrency.ts index a06592b644a..45f449f3190 100644 --- a/packages/redis-worker/src/fair-queue/concurrency.ts +++ b/packages/redis-worker/src/fair-queue/concurrency.ts @@ -11,6 +11,10 @@ export interface ConcurrencyManagerOptions { redis: RedisOptions; keys: FairQueueKeyProducer; groups: ConcurrencyGroupConfig[]; + logger?: { + debug: (message: string, context?: Record) => void; + error: (message: string, context?: Record) => void; + }; } /** @@ -26,12 +30,17 @@ export class ConcurrencyManager { private keys: FairQueueKeyProducer; private groups: ConcurrencyGroupConfig[]; private groupsByName: Map; + private logger: NonNullable; constructor(private options: ConcurrencyManagerOptions) { this.redis = createRedisClient(options.redis); this.keys = options.keys; this.groups = options.groups; this.groupsByName = new Map(options.groups.map((g) => [g.name, g])); + this.logger = options.logger ?? { + debug: () => {}, + error: () => {}, + }; this.#registerCommands(); } @@ -160,6 +169,69 @@ export class ConcurrencyManager { this.#assertPipelineSucceeded(await pipeline.exec(), messages.length); } + /** + * Remove concurrency set members that no longer correspond to an in-flight message, + * healing slots leaked by failed releases. Scans every set of every group and, per + * member, atomically removes it unless the message id appears in one of the given + * in-flight data hashes. Sound because a message is registered in-flight before its + * slot is reserved, so at the moment of the atomic check a member with no in-flight + * record can only be a leak; if the message is about to be re-claimed, reserve simply + * re-adds the member. + * + * @param inflightDataKeys - The in-flight data hash keys for every shard + * @returns The message ids that were removed, and how many sets were checked + */ + async sweepOrphanedSlots( + inflightDataKeys: string[] + ): Promise<{ scannedSets: number; removed: string[] }> { + const keyPrefix = this.options.redis.keyPrefix ?? ""; + let scannedSets = 0; + const removed: string[] = []; + + for (const group of this.groups) { + const pattern = `${keyPrefix}${this.keys.concurrencyKey(group.name, "*")}`; + let cursor = "0"; + + do { + const [nextCursor, foundKeys] = await this.redis.scan( + cursor, + "MATCH", + pattern, + "COUNT", + 100 + ); + cursor = nextCursor; + + for (const fullKey of foundKeys) { + const key = + keyPrefix && fullKey.startsWith(keyPrefix) ? fullKey.slice(keyPrefix.length) : fullKey; + + try { + const members = await this.redis.smembers(key); + if (members.length === 0) { + continue; + } + scannedSets++; + + const removedIds = await this.redis.removeOrphanedConcurrencySlots( + 1 + inflightDataKeys.length, + [key, ...inflightDataKeys], + ...members + ); + removed.push(...removedIds); + } catch (error) { + this.logger.error("Failed to sweep concurrency set, skipping it", { + key, + error: error instanceof Error ? error.message : String(error), + }); + } + } + } while (cursor !== "0"); + } + + return { scannedSets, removed }; + } + /** * Get current concurrency for a specific group. */ @@ -298,14 +370,20 @@ export class ConcurrencyManager { local numGroups = #KEYS local messageId = ARGV[1] --- Check all groups first +-- Check all groups first. A message that is already a member of a group's set passes +-- that group's check: re-admitting it does not increase concurrency (SADD is a no-op), +-- and counting its own leftover slot against it would let a message whose earlier +-- release failed block its own retry forever. for i = 1, numGroups do local key = KEYS[i] local limit = tonumber(ARGV[1 + i]) -- Limits start at ARGV[2] - local current = redis.call('SCARD', key) - - if current >= limit then - return 0 -- At capacity + + if redis.call('SISMEMBER', key, messageId) == 0 then + local current = redis.call('SCARD', key) + + if current >= limit then + return 0 -- At capacity + end end end @@ -318,6 +396,37 @@ end return 1 `, }); + + // Atomic orphan sweep for one concurrency set + // KEYS[1]: concurrency set key + // KEYS[2..n]: in-flight data hash keys for every shard + // ARGV: candidate messageIds (a snapshot of the set's members) + this.redis.defineCommand("removeOrphanedConcurrencySlots", { + lua: ` +local concurrencyKey = KEYS[1] +local removedIds = {} + +for i = 1, #ARGV do + local messageId = ARGV[i] + local inflight = false + + for j = 2, #KEYS do + if redis.call('HEXISTS', KEYS[j], messageId) == 1 then + inflight = true + break + end + end + + if not inflight then + if redis.call('SREM', concurrencyKey, messageId) == 1 then + table.insert(removedIds, messageId) + end + end +end + +return removedIds + `, + }); } } @@ -330,5 +439,11 @@ declare module "@internal/redis" { messageId: string, ...limits: string[] ): Promise; + + removeOrphanedConcurrencySlots( + numKeys: number, + keys: string[], + ...messageIds: string[] + ): Promise; } } diff --git a/packages/redis-worker/src/fair-queue/index.ts b/packages/redis-worker/src/fair-queue/index.ts index 749259479f6..41e9c74fec7 100644 --- a/packages/redis-worker/src/fair-queue/index.ts +++ b/packages/redis-worker/src/fair-queue/index.ts @@ -86,6 +86,7 @@ export class FairQueue { private visibilityTimeoutMs: number; private heartbeatIntervalMs: number; private reclaimIntervalMs: number; + private reconcileIntervalMs: number; private workerQueueResolver: (message: StoredMessage>) => string; private batchClaimSize: number; @@ -110,6 +111,7 @@ export class FairQueue { private abortController: AbortController; private masterQueueConsumerLoops: Promise[] = []; private reclaimLoop?: Promise; + private reconcileLoop?: Promise; // Queue descriptor cache for message processing private queueDescriptorCache = new Map(); @@ -139,6 +141,7 @@ export class FairQueue { this.visibilityTimeoutMs = options.visibilityTimeoutMs ?? 30_000; this.heartbeatIntervalMs = options.heartbeatIntervalMs ?? this.visibilityTimeoutMs / 3; this.reclaimIntervalMs = options.reclaimIntervalMs ?? 5_000; + this.reconcileIntervalMs = options.reconcileIntervalMs ?? 60_000; // Worker queue resolver (required) this.workerQueueResolver = options.workerQueue.resolveWorkerQueue; @@ -197,6 +200,10 @@ export class FairQueue { redis: options.redis, keys: options.keys, groups: options.concurrencyGroups, + logger: { + debug: (msg, ctx) => this.logger.debug(msg, ctx), + error: (msg, ctx) => this.logger.error(msg, ctx), + }, }); } @@ -658,6 +665,10 @@ export class FairQueue { // Start reclaim loop for handling timed-out messages this.reclaimLoop = this.#runReclaimLoop(); + if (this.concurrencyManager) { + this.reconcileLoop = this.#runReconcileLoop(); + } + this.logger.info("FairQueue started", { consumerCount: this.consumerCount, shardCount: this.shardCount, @@ -676,10 +687,15 @@ export class FairQueue { this.isRunning = false; this.abortController.abort(); - await Promise.allSettled([...this.masterQueueConsumerLoops, this.reclaimLoop]); + await Promise.allSettled([ + ...this.masterQueueConsumerLoops, + this.reclaimLoop, + this.reconcileLoop, + ]); this.masterQueueConsumerLoops = []; this.reclaimLoop = undefined; + this.reconcileLoop = undefined; this.logger.info("FairQueue stopped"); } @@ -1252,10 +1268,7 @@ export class FairQueue { metadata: storedMessage?.metadata ?? {}, }; - // Release concurrency - if (this.concurrencyManager) { - await this.concurrencyManager.release(descriptor, messageId); - } + await this.#releaseConcurrencySafely(descriptor, messageId, "completeMessage"); // Complete in visibility manager await this.visibilityManager.complete(messageId, queueId); @@ -1306,9 +1319,7 @@ export class FairQueue { }; // Release concurrency - if (this.concurrencyManager) { - await this.concurrencyManager.release(descriptor, messageId); - } + await this.#releaseConcurrencySafely(descriptor, messageId, "releaseMessage"); // Release back to queue (visibility manager updates dispatch indexes atomically) // Dispatch shard is tenant-based, not queue-based @@ -1332,6 +1343,36 @@ export class FairQueue { }); } + /** + * Release a message's concurrency slots without letting a failure escape. Release is + * best-effort cleanup: the primary state transition (complete, retry, DLQ, requeue) + * must proceed even when the slot cannot be freed, because aborting the transition + * re-delivers or strands the message, while a leaked slot is recoverable — reserve() + * re-admits a message that already holds its own slot, and the reconcile loop removes + * members that are no longer in flight. + */ + async #releaseConcurrencySafely( + descriptor: QueueDescriptor, + messageId: string, + context: string + ): Promise { + if (!this.concurrencyManager) { + return; + } + + try { + await this.concurrencyManager.release(descriptor, messageId); + } catch (error) { + this.logger.error("Failed to release concurrency slot, leaving it for the reconciler", { + context, + messageId, + queueId: descriptor.id, + tenantId: descriptor.tenantId, + error: error instanceof Error ? error.message : String(error), + }); + } + } + /** * Release a concurrency slot for a message we can no longer describe, deriving the * group from the queue id alone. Used on paths that bail out before the stored @@ -1349,7 +1390,7 @@ export class FairQueue { metadata: {}, }; - await this.concurrencyManager.release(descriptor, messageId); + await this.#releaseConcurrencySafely(descriptor, messageId, "failMessage:missing-record"); } /** @@ -1430,10 +1471,7 @@ export class FairQueue { attempt: storedMessage.attempt + 1, }; - // Release concurrency - if (this.concurrencyManager) { - await this.concurrencyManager.release(descriptor, storedMessage.id); - } + await this.#releaseConcurrencySafely(descriptor, storedMessage.id, "retry"); // Release with delay, passing the updated message data so the Lua script // atomically writes the incremented attempt count when re-queuing. @@ -1465,9 +1503,7 @@ export class FairQueue { } // Release concurrency - if (this.concurrencyManager) { - await this.concurrencyManager.release(descriptor, storedMessage.id); - } + await this.#releaseConcurrencySafely(descriptor, storedMessage.id, "deadLetter"); // Move to DLQ await this.#moveToDeadLetterQueue(storedMessage, error?.message); @@ -1545,17 +1581,18 @@ export class FairQueue { } /** - * Free every timed-out message's concurrency slot in one pipeline, before any of them - * is put back on the queue. Throwing here aborts the requeue for the whole batch, which - * leaves the messages in-flight for the next reclaim tick. That is the safe direction: - * requeuing a message whose slot is still held is what strands the slot permanently. + * Free every timed-out message's concurrency slot in one pipeline before the batch is + * requeued. A failed release never blocks the requeue: holding the message in-flight + * would stall it behind a persistent release failure, while the leaked slot self-heals + * when the message is re-admitted (reserve is idempotent for a message that already + * holds its own slot) or via the reconcile loop. */ - async #releaseReclaimedConcurrency(messages: ReclaimedMessageInfo[]): Promise { + async #releaseReclaimedConcurrency(messages: ReclaimedMessageInfo[]): Promise { if (!this.concurrencyManager || messages.length === 0) { - return []; + return; } - const descriptorFor = (message: ReclaimedMessageInfo) => ({ + const descriptorFor = (message: ReclaimedMessageInfo): QueueDescriptor => ({ id: message.queueId, tenantId: message.tenantId, metadata: message.metadata ?? {}, @@ -1565,7 +1602,7 @@ export class FairQueue { await this.concurrencyManager.releaseBatch( messages.map((message) => ({ queue: descriptorFor(message), messageId: message.messageId })) ); - return []; + return; } catch (error) { this.logger.error("Batch concurrency release failed, retrying message by message", { count: messages.length, @@ -1573,22 +1610,9 @@ export class FairQueue { }); } - const failed: string[] = []; - for (const message of messages) { - try { - await this.concurrencyManager.release(descriptorFor(message), message.messageId); - } catch (error) { - failed.push(message.messageId); - this.logger.error("Failed to release concurrency for reclaimed message", { - messageId: message.messageId, - queueId: message.queueId, - error: error instanceof Error ? error.message : String(error), - }); - } + await this.#releaseConcurrencySafely(descriptorFor(message), message.messageId, "reclaim"); } - - return failed; } async #reclaimTimedOutMessages(): Promise { @@ -1626,6 +1650,55 @@ export class FairQueue { } } + async #runReconcileLoop(): Promise { + try { + for await (const _ of setInterval(this.reconcileIntervalMs, null, { + signal: this.abortController.signal, + })) { + try { + await this.#reconcileConcurrency(); + } catch (error) { + this.logger.error("Concurrency reconcile loop error", { + error: error instanceof Error ? error.message : String(error), + }); + } + } + } catch (error) { + if (isAbortError(error)) { + this.logger.debug("Concurrency reconcile loop aborted"); + return; + } + throw error; + } + } + + /** + * Heal orphaned concurrency slots: remove any concurrency set member whose message is + * not registered in-flight on any shard. This is the recovery path for slots leaked by + * failed releases, which otherwise count against a tenant's limit forever and can + * wedge every queue the tenant owns. + */ + async #reconcileConcurrency(): Promise { + if (!this.concurrencyManager) { + return; + } + + const inflightDataKeys = Array.from({ length: this.shardCount }, (_, shardId) => + this.keys.inflightDataKey(shardId) + ); + + const { scannedSets, removed } = + await this.concurrencyManager.sweepOrphanedSlots(inflightDataKeys); + + if (removed.length > 0) { + this.logger.info("Reconciled orphaned concurrency slots", { + count: removed.length, + scannedSets, + messageIds: removed.slice(0, 50), + }); + } + } + // ============================================================================ // Private - Cooloff State // ============================================================================ diff --git a/packages/redis-worker/src/fair-queue/tests/concurrency.test.ts b/packages/redis-worker/src/fair-queue/tests/concurrency.test.ts index 7137ea8ac8b..26e37e5e89d 100644 --- a/packages/redis-worker/src/fair-queue/tests/concurrency.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/concurrency.test.ts @@ -1,5 +1,6 @@ import { describe, expect } from "vitest"; import { redisTest } from "@internal/testcontainers"; +import { createRedisClient } from "@internal/redis"; import { ConcurrencyManager } from "../concurrency.js"; import { DefaultFairQueueKeyProducer } from "../keyProducer.js"; import type { FairQueueKeyProducer, QueueDescriptor } from "../types.js"; @@ -736,4 +737,86 @@ describe("ConcurrencyManager", () => { } ); }); + + describe("reserve idempotency", () => { + redisTest( + "should re-admit a message that already holds its own slot at capacity", + { timeout: 10000 }, + async ({ redisOptions }) => { + keys = new DefaultFairQueueKeyProducer({ prefix: "test" }); + + const manager = new ConcurrencyManager({ + redis: redisOptions, + keys, + groups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 1, + defaultLimit: 1, + }, + ], + }); + + const redis = createRedisClient(redisOptions); + const queue: QueueDescriptor = { id: "queue-1", tenantId: "t1", metadata: {} }; + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + + try { + await redis.sadd(concurrencyKey, "m1"); + + expect(await manager.reserve(queue, "m1")).toBe(true); + expect(await redis.scard(concurrencyKey)).toBe(1); + + expect(await manager.reserve(queue, "m2")).toBe(false); + } finally { + await redis.del(concurrencyKey); + await redis.quit(); + await manager.close(); + } + } + ); + }); + + describe("orphaned slot sweep", () => { + redisTest( + "should remove only members without an in-flight record", + { timeout: 10000 }, + async ({ redisOptions }) => { + keys = new DefaultFairQueueKeyProducer({ prefix: "test" }); + + const manager = new ConcurrencyManager({ + redis: redisOptions, + keys, + groups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 5, + defaultLimit: 5, + }, + ], + }); + + const redis = createRedisClient(redisOptions); + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + const inflightDataKey = keys.inflightDataKey(0); + + try { + await redis.sadd(concurrencyKey, "orphan-1", "orphan-2", "active-1"); + await redis.hset(inflightDataKey, "active-1", "{}"); + + const result = await manager.sweepOrphanedSlots([inflightDataKey]); + + expect(result.removed.sort()).toEqual(["orphan-1", "orphan-2"]); + expect(await redis.smembers(concurrencyKey)).toEqual(["active-1"]); + } finally { + await redis.del(concurrencyKey); + await redis.del(inflightDataKey); + await redis.quit(); + await manager.close(); + } + } + ); + }); }); 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 5fbf2299fab..b02d667d71b 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -1518,8 +1518,8 @@ describe("FairQueue", () => { ); redisTest( - "should not requeue a reclaimed message when freeing its slot fails", - { timeout: 20000 }, + "should requeue a reclaimed message even when freeing its slot fails", + { timeout: 25000 }, async ({ redisOptions }) => { keys = new DefaultFairQueueKeyProducer({ prefix: "test" }); @@ -1538,6 +1538,7 @@ describe("FairQueue", () => { consumerIntervalMs: 20, visibilityTimeoutMs: 500, reclaimIntervalMs: 100, + reconcileIntervalMs: 600_000, concurrencyGroups: [ { name: "tenant", @@ -1552,10 +1553,16 @@ describe("FairQueue", () => { const redis = createRedisClient(redisOptions); const queueId = "tenant:t1:queue:release-fails"; const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + let starts = 0; try { - queue.onMessage(async () => { - await new Promise((resolve) => setTimeout(resolve, 15000)); + queue.onMessage(async (ctx) => { + starts++; + if (starts === 1) { + await new Promise((resolve) => setTimeout(resolve, 3000)); + return; + } + await ctx.complete(); }); await queue.enqueue({ @@ -1576,22 +1583,215 @@ describe("FairQueue", () => { await redis.del(concurrencyKey); await redis.set(concurrencyKey, "not-a-set"); - await new Promise((resolve) => setTimeout(resolve, 2000)); + await vi.waitFor( + async () => { + expect(await redis.zcard(keys.queueKey(queueId))).toBe(1); + expect(await redis.zcard(keys.inflightKey(0))).toBe(0); + }, + { timeout: 5000 } + ); + + await redis.del(concurrencyKey); - expect(await redis.zcard(keys.inflightKey(0))).toBe(1); + await vi.waitFor( + () => { + expect(starts).toBe(2); + }, + { timeout: 10000 } + ); + + await vi.waitFor( + async () => { + expect(await redis.zcard(keys.inflightKey(0))).toBe(0); + expect(await redis.zcard(keys.queueKey(queueId))).toBe(0); + expect(await redis.scard(concurrencyKey)).toBe(0); + }, + { timeout: 5000 } + ); + } finally { + await redis.del(concurrencyKey); + await redis.quit(); + await queue.close(); + } + } + ); + + redisTest( + "should complete a message even when freeing its slot fails, never re-delivering it", + { timeout: 20000 }, + async ({ redisOptions }) => { + 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: 500, + reclaimIntervalMs: 100, + reconcileIntervalMs: 600_000, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 5, + defaultLimit: 5, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + const queueId = "tenant:t1:queue:complete-release-fails"; + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + let executions = 0; + + try { + queue.onMessage(async (ctx) => { + executions++; + await redis.del(concurrencyKey); + await redis.set(concurrencyKey, "not-a-set"); + await ctx.complete(); + }); + + await queue.enqueue({ + queueId, + tenantId: "t1", + payload: { value: "msg-0" }, + }); + + queue.start(); + + await vi.waitFor( + () => { + expect(executions).toBe(1); + }, + { timeout: 5000 } + ); + + await vi.waitFor( + async () => { + expect(await redis.zcard(keys.inflightKey(0))).toBe(0); + }, + { timeout: 5000 } + ); + + await redis.del(concurrencyKey); + + await new Promise((resolve) => setTimeout(resolve, 1500)); + + expect(executions).toBe(1); + expect(await redis.zcard(keys.inflightKey(0))).toBe(0); expect(await redis.zcard(keys.queueKey(queueId))).toBe(0); + } finally { + await redis.del(concurrencyKey); + await redis.quit(); + await queue.close(); + } + } + ); + + redisTest( + "should schedule a retry with an advanced attempt even when freeing its slot fails", + { timeout: 20000 }, + async ({ redisOptions }) => { + keys = new DefaultFairQueueKeyProducer({ prefix: "test" }); + const attempts: number[] = []; + + 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, + reclaimIntervalMs: 100, + reconcileIntervalMs: 600_000, + retry: { + strategy: new FixedDelayRetry({ maxAttempts: 2, delayMs: 50 }), + deadLetterQueue: true, + }, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 5, + defaultLimit: 5, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + const queueId = "tenant:t1:queue:retry-release-fails"; + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + + try { + queue.onMessage(async (ctx) => { + attempts.push(ctx.message.attempt); + if (ctx.message.attempt === 1) { + await redis.del(concurrencyKey); + await redis.set(concurrencyKey, "not-a-set"); + } + await ctx.fail(new Error("boom")); + }); + + await queue.enqueue({ + queueId, + tenantId: "t1", + payload: { value: "msg-0" }, + messageId: "retry-msg", + }); + + queue.start(); + + await vi.waitFor( + () => { + expect(attempts).toEqual([1]); + }, + { timeout: 5000 } + ); + + await vi.waitFor( + async () => { + const stored = await redis.hget(keys.queueItemsKey(queueId), "retry-msg"); + expect(stored).not.toBeNull(); + expect(JSON.parse(stored!).attempt).toBe(2); + }, + { timeout: 5000 } + ); - const [stuckMember] = await redis.zrange(keys.inflightKey(0), 0, 0); - const stuckDeadline = Number(await redis.zscore(keys.inflightKey(0), stuckMember!)); + expect(await redis.zcard(keys.inflightKey(0))).toBe(0); await redis.del(concurrencyKey); + await vi.waitFor( + () => { + expect(attempts).toEqual([1, 2]); + }, + { timeout: 5000 } + ); + await vi.waitFor( async () => { - const score = await redis.zscore(keys.inflightKey(0), stuckMember!); - expect(score === null || Number(score) > stuckDeadline).toBe(true); + expect(await queue.getDeadLetterQueueLength("t1")).toBe(1); }, - { timeout: 8000 } + { timeout: 5000 } ); } finally { await redis.del(concurrencyKey); @@ -1756,4 +1956,167 @@ describe("FairQueue", () => { } ); }); + + describe("concurrency slot reconciliation", () => { + redisTest( + "should free slots held by messages that are no longer in flight", + { timeout: 15000 }, + async ({ redisOptions }) => { + 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, + reclaimIntervalMs: 100, + reconcileIntervalMs: 200, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 1, + defaultLimit: 1, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + const queueId = "tenant:t1:queue:orphaned-slot"; + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + const processed: string[] = []; + + try { + await redis.sadd(concurrencyKey, "ghost-message"); + + queue.onMessage(async (ctx) => { + processed.push(ctx.message.payload.value); + await ctx.complete(); + }); + + await queue.enqueue({ + queueId, + tenantId: "t1", + payload: { value: "msg-0" }, + }); + + queue.start(); + + await vi.waitFor( + () => { + expect(processed).toEqual(["msg-0"]); + }, + { timeout: 8000 } + ); + + expect(await redis.sismember(concurrencyKey, "ghost-message")).toBe(0); + + await vi.waitFor( + async () => { + expect(await redis.scard(concurrencyKey)).toBe(0); + }, + { timeout: 5000 } + ); + } finally { + await redis.del(concurrencyKey); + await redis.quit(); + await queue.close(); + } + } + ); + + redisTest( + "should leave the slot of an in-flight message alone", + { timeout: 15000 }, + async ({ redisOptions }) => { + 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, + reclaimIntervalMs: 100, + reconcileIntervalMs: 100, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 1, + defaultLimit: 1, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + const queueId = "tenant:t1:queue:active-slot"; + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + let executions = 0; + + try { + queue.onMessage(async (ctx) => { + executions++; + await new Promise((resolve) => setTimeout(resolve, 1500)); + await ctx.complete(); + }); + + await queue.enqueue({ + queueId, + tenantId: "t1", + payload: { value: "msg-0" }, + messageId: "active-msg", + }); + + queue.start(); + + await vi.waitFor( + async () => { + expect(await redis.zcard(keys.inflightKey(0))).toBe(1); + }, + { timeout: 5000 } + ); + + await new Promise((resolve) => setTimeout(resolve, 1000)); + + expect(await redis.sismember(concurrencyKey, "active-msg")).toBe(1); + expect(await redis.zcard(keys.inflightKey(0))).toBe(1); + + await vi.waitFor( + async () => { + expect(executions).toBe(1); + expect(await redis.scard(concurrencyKey)).toBe(0); + expect(await redis.zcard(keys.inflightKey(0))).toBe(0); + }, + { timeout: 5000 } + ); + + expect(executions).toBe(1); + } finally { + await redis.del(concurrencyKey); + await redis.quit(); + await queue.close(); + } + } + ); + }); }); diff --git a/packages/redis-worker/src/fair-queue/types.ts b/packages/redis-worker/src/fair-queue/types.ts index 3bc56b599fa..db762d0ded8 100644 --- a/packages/redis-worker/src/fair-queue/types.ts +++ b/packages/redis-worker/src/fair-queue/types.ts @@ -413,6 +413,8 @@ export interface FairQueueOptions Promise + onBeforeRequeue?: (messages: ReclaimedMessageInfo[]) => Promise ): Promise { const inflightKey = this.keys.inflightKey(shardId); const inflightDataKey = this.keys.inflightDataKey(shardId); @@ -465,24 +468,15 @@ export class VisibilityManager { return []; } - const notReleased = new Set( - (onBeforeRequeue - ? await onBeforeRequeue(candidates.map((candidate) => candidate.info)) - : undefined) ?? [] - ); + if (onBeforeRequeue) { + await onBeforeRequeue(candidates.map((candidate) => candidate.info)); + } const reclaimedMessages: ReclaimedMessageInfo[] = []; for (const { member, deadlineScore, storedMessage, info } of candidates) { const { messageId, queueId } = info; - if (notReleased.has(messageId)) { - this.logger.error("Skipping requeue, concurrency slot was not released", { - messageId, - queueId, - }); - continue; - } const { queueKey, queueItemsKey, tenantQueueIndexKey, dispatchKey, tenantId } = getQueueKeys(queueId); From 21eb0d330f0ca1139aae9d860911c4e0ad836079 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Wed, 12 Aug 2026 13:43:31 +0100 Subject: [PATCH 12/12] fix(redis-worker): bound the reconcile sweep cost and write the release note for users The sweep now scans with a larger COUNT, checks set members in fixed-size chunks so one script call cannot hold up Redis on a set full of leaked members, jitters its first run so instances sharing a Redis do not all sweep at once, and is disabled entirely by a zero interval instead of spinning. The changeset now leads with the user-visible fix instead of the mechanism. --- .../fair-queue-concurrency-slot-leak.md | 2 +- .../src/fair-queue/concurrency.ts | 23 ++++-- packages/redis-worker/src/fair-queue/index.ts | 11 ++- .../src/fair-queue/tests/concurrency.test.ts | 8 ++- .../src/fair-queue/tests/fairQueue.test.ts | 72 +++++++++++++++++++ packages/redis-worker/src/fair-queue/types.ts | 2 +- 6 files changed, 104 insertions(+), 14 deletions(-) diff --git a/.changeset/fair-queue-concurrency-slot-leak.md b/.changeset/fair-queue-concurrency-slot-leak.md index bfdfa312c1a..41f41e7e3d2 100644 --- a/.changeset/fair-queue-concurrency-slot-leak.md +++ b/.changeset/fair-queue-concurrency-slot-leak.md @@ -2,4 +2,4 @@ "@trigger.dev/redis-worker": patch --- -Fair queue consumers no longer leak the concurrency slots that gate a tenant's throughput, and leaked slots now heal themselves. Slots are freed on the completion, retry, dead-letter, and reclaim paths that previously skipped them, and a failed release never blocks the message's own state transition, so a Redis error can no longer turn into a duplicate execution or a lost retry. A periodic reconcile loop removes any slot whose message is no longer in flight, and a message that still holds its own slot from an earlier failed release is re-admitted instead of being blocked by it. +Fair queue tenants can no longer get permanently stuck behind leaked concurrency slots. Slots are now freed on every path that finishes a message, a failed release no longer causes a message to run twice or lose its retry, and a background sweep frees any slot that does leak, so a tenant's queues recover on their own instead of needing manual cleanup. diff --git a/packages/redis-worker/src/fair-queue/concurrency.ts b/packages/redis-worker/src/fair-queue/concurrency.ts index 45f449f3190..30fd0c33a3a 100644 --- a/packages/redis-worker/src/fair-queue/concurrency.ts +++ b/packages/redis-worker/src/fair-queue/concurrency.ts @@ -7,6 +7,12 @@ import type { QueueDescriptor, } from "./types.js"; +/** + * Caps how many set members a single sweep script invocation checks, bounding the time the + * atomic Lua script can hold up Redis when a set has accumulated many leaked members. + */ +const SWEEP_MEMBER_CHUNK_SIZE = 500; + export interface ConcurrencyManagerOptions { redis: RedisOptions; keys: FairQueueKeyProducer; @@ -198,7 +204,7 @@ export class ConcurrencyManager { "MATCH", pattern, "COUNT", - 100 + 1000 ); cursor = nextCursor; @@ -213,12 +219,15 @@ export class ConcurrencyManager { } scannedSets++; - const removedIds = await this.redis.removeOrphanedConcurrencySlots( - 1 + inflightDataKeys.length, - [key, ...inflightDataKeys], - ...members - ); - removed.push(...removedIds); + for (let i = 0; i < members.length; i += SWEEP_MEMBER_CHUNK_SIZE) { + const chunk = members.slice(i, i + SWEEP_MEMBER_CHUNK_SIZE); + const removedIds = await this.redis.removeOrphanedConcurrencySlots( + 1 + inflightDataKeys.length, + [key, ...inflightDataKeys], + ...chunk + ); + removed.push(...removedIds); + } } catch (error) { this.logger.error("Failed to sweep concurrency set, skipping it", { key, diff --git a/packages/redis-worker/src/fair-queue/index.ts b/packages/redis-worker/src/fair-queue/index.ts index 41e9c74fec7..826ac73edcb 100644 --- a/packages/redis-worker/src/fair-queue/index.ts +++ b/packages/redis-worker/src/fair-queue/index.ts @@ -2,7 +2,7 @@ import { createRedisClient, type Redis } from "@internal/redis"; import { SpanKind, type Span } from "@internal/tracing"; import { Logger } from "@trigger.dev/core/logger"; import { nanoid } from "nanoid"; -import { setInterval } from "node:timers/promises"; +import { setInterval, setTimeout as delay } from "node:timers/promises"; import { type z } from "zod"; import { isAbortError } from "../utils.js"; import { ConcurrencyManager } from "./concurrency.js"; @@ -665,7 +665,7 @@ export class FairQueue { // Start reclaim loop for handling timed-out messages this.reclaimLoop = this.#runReclaimLoop(); - if (this.concurrencyManager) { + if (this.concurrencyManager && this.reconcileIntervalMs > 0) { this.reconcileLoop = this.#runReconcileLoop(); } @@ -1650,8 +1650,15 @@ export class FairQueue { } } + /** + * The initial delay is jittered so instances sharing a Redis do not all sweep at once. + */ async #runReconcileLoop(): Promise { try { + await delay(Math.random() * this.reconcileIntervalMs, null, { + signal: this.abortController.signal, + }); + for await (const _ of setInterval(this.reconcileIntervalMs, null, { signal: this.abortController.signal, })) { diff --git a/packages/redis-worker/src/fair-queue/tests/concurrency.test.ts b/packages/redis-worker/src/fair-queue/tests/concurrency.test.ts index 26e37e5e89d..7f669cb4106 100644 --- a/packages/redis-worker/src/fair-queue/tests/concurrency.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/concurrency.test.ts @@ -803,13 +803,15 @@ describe("ConcurrencyManager", () => { const inflightDataKey = keys.inflightDataKey(0); try { - await redis.sadd(concurrencyKey, "orphan-1", "orphan-2", "active-1"); + const orphans = Array.from({ length: 1200 }, (_, i) => `orphan-${i}`); + await redis.sadd(concurrencyKey, ...orphans, "active-1", "active-2"); await redis.hset(inflightDataKey, "active-1", "{}"); + await redis.hset(inflightDataKey, "active-2", "{}"); const result = await manager.sweepOrphanedSlots([inflightDataKey]); - expect(result.removed.sort()).toEqual(["orphan-1", "orphan-2"]); - expect(await redis.smembers(concurrencyKey)).toEqual(["active-1"]); + expect(result.removed.sort()).toEqual(orphans.sort()); + expect((await redis.smembers(concurrencyKey)).sort()).toEqual(["active-1", "active-2"]); } finally { await redis.del(concurrencyKey); await redis.del(inflightDataKey); 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 b02d667d71b..edf5b447c7a 100644 --- a/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts +++ b/packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts @@ -2118,5 +2118,77 @@ describe("FairQueue", () => { } } ); + + redisTest( + "should not sweep when the reconcile interval is zero", + { timeout: 10000 }, + async ({ redisOptions }) => { + 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, + reclaimIntervalMs: 100, + reconcileIntervalMs: 0, + concurrencyGroups: [ + { + name: "tenant", + extractGroupId: (q) => q.tenantId, + getLimit: async () => 5, + defaultLimit: 5, + }, + ], + startConsumers: false, + }); + + const redis = createRedisClient(redisOptions); + const queueId = "tenant:t1:queue:sweep-disabled"; + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + const processed: string[] = []; + + try { + await redis.sadd(concurrencyKey, "ghost-message"); + + queue.onMessage(async (ctx) => { + processed.push(ctx.message.payload.value); + await ctx.complete(); + }); + + await queue.enqueue({ + queueId, + tenantId: "t1", + payload: { value: "msg-0" }, + }); + + queue.start(); + + await vi.waitFor( + () => { + expect(processed).toEqual(["msg-0"]); + }, + { timeout: 5000 } + ); + + await new Promise((resolve) => setTimeout(resolve, 1000)); + + expect(await redis.sismember(concurrencyKey, "ghost-message")).toBe(1); + } finally { + await redis.del(concurrencyKey); + await redis.quit(); + await queue.close(); + } + } + ); }); }); diff --git a/packages/redis-worker/src/fair-queue/types.ts b/packages/redis-worker/src/fair-queue/types.ts index db762d0ded8..1bb3e6576c4 100644 --- a/packages/redis-worker/src/fair-queue/types.ts +++ b/packages/redis-worker/src/fair-queue/types.ts @@ -413,7 +413,7 @@ export interface FairQueueOptions