import { Mutex } from 'async-mutex'; /** * Logger interface for ConcurrencyManager */ export interface ConcurrencyLogger { info: (message: string, meta?: Record) => void; warn: (message: string, meta?: Record) => void; error: (message: string, meta?: Record) => void; } // Default console logger const consoleLogger: ConcurrencyLogger = { info: (message: string, meta?: Record) => meta ? console.info(message, meta) : console.info(message), warn: (message: string, meta?: Record) => meta ? console.warn(message, meta) : console.warn(message), error: (message: string, meta?: Record) => meta ? console.error(message, meta) : console.error(message) }; let concurrencyLogger: ConcurrencyLogger = consoleLogger; export const setConcurrencyLogger = (logger: ConcurrencyLogger): void => { concurrencyLogger = logger; }; /** * ConcurrencyManager - Manages concurrent request slots with mutex-based synchronization * * Prevents resource exhaustion by limiting the number of concurrent expensive operations. * Uses async-mutex for atomic slot acquisition/release. * * @example * ```typescript * const manager = new ConcurrencyManager(3); // Max 3 concurrent requests * * async function expensiveOperation() { * await manager.acquireSlot(); * try { * // ... expensive operation ... * } finally { * manager.releaseSlot(); * } * } * ``` */ export class ConcurrencyManager { private activeSynthesisRequests = 0; private readonly requestQueue: Array<{ callback: () => void; cancelled: boolean }> = []; private readonly slotMutex = new Mutex(); constructor(private readonly maxConcurrentRequests: number = 3) {} /** * Acquire a slot for concurrent execution * Blocks until a slot is available or timeout occurs (5 seconds) */ async acquireSlot(): Promise { return new Promise((resolve, reject) => { let callbackWrapper: { callback: () => void; cancelled: boolean } | null = null; let timedOut = false; // Add timeout to prevent indefinite hanging (5 second timeout) // CRITICAL: Mark callback as cancelled when timeout fires AND set timedOut flag const timeoutHandle = setTimeout(() => { timedOut = true; // Prevent mutex from queuing after timeout if (callbackWrapper) { callbackWrapper.cancelled = true; // Mark stale to prevent execution when dequeued } concurrencyLogger.warn('Slot acquisition timed out', { activeSynthesisRequests: this.activeSynthesisRequests, queueLength: this.requestQueue.length }); reject(new Error('Slot acquisition timed out after 5 seconds')); }, 5000); // Use async-mutex for proper atomic operations this.slotMutex.runExclusive(async () => { // Clear timeout since we're proceeding normally clearTimeout(timeoutHandle); // CRITICAL: Check if timeout already fired - if so, don't queue or increment if (timedOut) { concurrencyLogger.info('Mutex acquired after timeout - skipping queue operation'); return; } // Atomic check and increment within the mutex if (this.activeSynthesisRequests < this.maxConcurrentRequests) { this.activeSynthesisRequests++; resolve(); return; } // Create callback wrapper and queue it - resolve will be called by releaseSlot callbackWrapper = { cancelled: false, callback: () => { this.activeSynthesisRequests++; resolve(); } }; this.requestQueue.push(callbackWrapper); }).catch((error) => { // Clear timeout and ensure mutex releases properly even if error occurs clearTimeout(timeoutHandle); if (callbackWrapper) { callbackWrapper.cancelled = true; // Prevent execution if queued } concurrencyLogger.error('Error in slot acquisition', { error: error.message }); reject(error); }); }); } /** * Release a slot and process next queued request */ releaseSlot(): void { // Use async-mutex to ensure atomic decrement and queue processing this.slotMutex.runExclusive(async () => { this.activeSynthesisRequests--; // Process next queued request atomically, skipping cancelled ones // CRITICAL: Cancelled callbacks (from timeouts) must be skipped to prevent slot leaks let nextWrapper = this.requestQueue.shift(); while (nextWrapper && nextWrapper.cancelled) { concurrencyLogger.warn('Skipping cancelled callback from queue', { activeSynthesisRequests: this.activeSynthesisRequests, remainingQueueLength: this.requestQueue.length }); nextWrapper = this.requestQueue.shift(); // Get next non-cancelled callback } // CRITICAL: Double-check cancelled flag right before execution to handle // race condition where timeout fires between while-check and callback execution if (nextWrapper && !nextWrapper.cancelled) { nextWrapper.callback(); } else if (nextWrapper) { concurrencyLogger.warn('Callback cancelled between dequeue and execution', { activeSynthesisRequests: this.activeSynthesisRequests }); } }).catch((error) => { concurrencyLogger.error('Error in slot release', { error: error.message }); }); } /** * Get current number of active requests */ getActiveCount(): number { return this.activeSynthesisRequests; } /** * Get current queue length */ getQueueLength(): number { return this.requestQueue.length; } /** * Get max concurrent requests allowed */ getMaxConcurrent(): number { return this.maxConcurrentRequests; } /** * Execute a function with automatic slot management * * @example * ```typescript * const result = await manager.execute(async () => { * return await expensiveOperation(); * }); * ``` */ async execute(fn: () => Promise): Promise { await this.acquireSlot(); try { return await fn(); } finally { this.releaseSlot(); } } }