diff --git a/.changeset/fair-queue-concurrency-slot-leak.md b/.changeset/fair-queue-concurrency-slot-leak.md new file mode 100644 index 00000000000..41f41e7e3d2 --- /dev/null +++ b/.changeset/fair-queue-concurrency-slot-leak.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/redis-worker": patch +--- + +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 641899ccc8d..30fd0c33a3a 100644 --- a/packages/redis-worker/src/fair-queue/concurrency.ts +++ b/packages/redis-worker/src/fair-queue/concurrency.ts @@ -7,10 +7,20 @@ 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; groups: ConcurrencyGroupConfig[]; + logger?: { + debug: (message: string, context?: Record) => void; + error: (message: string, context?: Record) => void; + }; } /** @@ -26,12 +36,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(); } @@ -103,7 +118,37 @@ 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 { + 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)); + + 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 +172,73 @@ export class ConcurrencyManager { } } - await pipeline.exec(); + 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", + 1000 + ); + 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++; + + 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, + error: error instanceof Error ? error.message : String(error), + }); + } + } + } while (cursor !== "0"); + } + + return { scannedSets, removed }; } /** @@ -268,14 +379,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 @@ -288,6 +405,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 + `, + }); } } @@ -300,5 +448,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 e24de876d85..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"; @@ -25,6 +25,7 @@ import type { FairScheduler, QueueCooloffState, QueueDescriptor, + ReclaimedMessageInfo, SchedulerContext, StoredMessage, TenantQueues, @@ -85,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; @@ -109,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(); @@ -138,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; @@ -196,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), + }, }); } @@ -657,6 +665,10 @@ export class FairQueue { // Start reclaim loop for handling timed-out messages this.reclaimLoop = this.#runReclaimLoop(); + if (this.concurrencyManager && this.reconcileIntervalMs > 0) { + this.reconcileLoop = this.#runReconcileLoop(); + } + this.logger.info("FairQueue started", { consumerCount: this.consumerCount, shardCount: this.shardCount, @@ -675,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"); } @@ -1245,22 +1262,17 @@ 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 ?? {}, + }; + + await this.#releaseConcurrencySafely(descriptor, messageId, "completeMessage"); // Complete in visibility manager await this.visibilityManager.complete(messageId, queueId); - // Release concurrency - if (this.concurrencyManager && storedMessage) { - await this.concurrencyManager.release(descriptor, messageId); - } - // Update both old and new indexes, clean up caches if queue is empty const { queueEmpty } = await this.#updateAllIndexesAfterDequeue(queueId, descriptor.tenantId); if (queueEmpty) { @@ -1300,13 +1312,14 @@ 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 + await this.#releaseConcurrencySafely(descriptor, messageId, "releaseMessage"); // Release back to queue (visibility manager updates dispatch indexes atomically) // Dispatch shard is tenant-based, not queue-based @@ -1324,17 +1337,62 @@ 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, }); } + /** + * 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 + * 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.#releaseConcurrencySafely(descriptor, messageId, "failMessage:missing-record"); + } + /** * Mark a message as failed. This will trigger retry logic if configured, * or move the message to the dead letter queue. @@ -1353,6 +1411,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; } @@ -1364,6 +1423,7 @@ export class FairQueue { messageId, queueId, }); + await this.#releaseOrphanedConcurrency(messageId, queueId); return; } @@ -1411,6 +1471,8 @@ export class FairQueue { attempt: storedMessage.attempt + 1, }; + 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. const tenantQueueIndexKey = this.keys.tenantQueueIndexKey(descriptor.tenantId); @@ -1427,11 +1489,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", { @@ -1445,13 +1502,11 @@ export class FairQueue { } } + // Release concurrency + await this.#releaseConcurrencySafely(descriptor, storedMessage.id, "deadLetter"); + // Move to DLQ await this.#moveToDeadLetterQueue(storedMessage, error?.message); - - // Release concurrency - if (this.concurrencyManager) { - await this.concurrencyManager.release(descriptor, storedMessage.id); - } } async #moveToDeadLetterQueue( @@ -1525,55 +1580,129 @@ export class FairQueue { } } + /** + * 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 { + if (!this.concurrencyManager || messages.length === 0) { + return; + } + + const descriptorFor = (message: ReclaimedMessageInfo): QueueDescriptor => ({ + 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), + }); + } + + for (const message of messages) { + await this.#releaseConcurrencySafely(descriptorFor(message), message.messageId, "reclaim"); + } + } + 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, - }; + 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; + } 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) { + this.logger.info("Reclaimed timed-out messages", { count: totalReclaimed }); + } + } + + /** + * 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, }); - 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), - }); - } + 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), + }); } - - // Dispatch indexes are updated atomically by the releaseMessage Lua script - // inside reclaimTimedOut, so no separate index update needed here. } + } catch (error) { + if (isAbortError(error)) { + this.logger.debug("Concurrency reconcile loop aborted"); + return; + } + throw error; + } + } - totalReclaimed += reclaimedMessages.length; + /** + * 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; } - if (totalReclaimed > 0) { - this.logger.info("Reclaimed timed-out messages", { count: totalReclaimed }); + 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), + }); } } 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..7f669cb4106 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,88 @@ 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 { + 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(orphans.sort()); + expect((await redis.smembers(concurrencyKey)).sort()).toEqual(["active-1", "active-2"]); + } 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 5ae1f390f9f..edf5b447c7a 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,825 @@ 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); + + 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}` }, + }); + } + + 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 groups from the cached descriptor 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, + }); + + 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(); + }); + + 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}` }, + }); + } + + queue.start(); + + await vi.waitFor( + () => { + expect(processed).toHaveLength(2); + }, + { timeout: 10000 } + ); + + 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(); + } + } + ); + + redisTest( + "should requeue a reclaimed message even when freeing its slot fails", + { timeout: 25000 }, + 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:release-fails"; + const concurrencyKey = keys.concurrencyKey("tenant", "t1"); + let starts = 0; + + try { + queue.onMessage(async (ctx) => { + starts++; + if (starts === 1) { + await new Promise((resolve) => setTimeout(resolve, 3000)); + return; + } + await ctx.complete(); + }); + + await queue.enqueue({ + queueId, + tenantId: "t1", + payload: { value: "msg-0" }, + }); + + queue.start(); + + await vi.waitFor( + async () => { + expect(await redis.zcard(keys.inflightKey(0))).toBe(1); + }, + { timeout: 5000 } + ); + + await redis.del(concurrencyKey); + await redis.set(concurrencyKey, "not-a-set"); + + 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); + + 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 } + ); + + 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 () => { + expect(await queue.getDeadLetterQueueLength("t1")).toBe(1); + }, + { timeout: 5000 } + ); + } finally { + await redis.del(concurrencyKey); + await redis.quit(); + await queue.close(); + } + } + ); + + 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: 30000, + 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"; + const inflightKey = keys.inflightKey(0); + + let unblockHandler: (() => void) | undefined; + const handlerBlocked = new Promise((resolve) => { + unblockHandler = resolve; + }); + + try { + queue.onMessage(async () => { + await handlerBlocked; + }); + + 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); + expect(await redis.zcard(inflightKey)).toBe(1); + }, + { timeout: 5000 } + ); + + 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(); + } + } + ); + + 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(); + } + } + ); + }); + + 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(); + } + } + ); + + 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/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/types.ts b/packages/redis-worker/src/fair-queue/types.ts index 3bc56b599fa..1bb3e6576c4 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 ): Promise { const inflightKey = this.keys.inflightKey(shardId); const inflightDataKey = this.keys.inflightDataKey(shardId); @@ -415,30 +419,68 @@ 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 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 (!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 { - // 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 - } - } - // 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; @@ -457,34 +499,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", { @@ -707,7 +727,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 @@ -717,14 +739,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 @@ -771,15 +794,20 @@ 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 + -- a slot in every future reclaim scan. + redis.call('ZREM', inflightKey, member) end end