type PromiseResolver = () => void; export class AsyncSemaphore { private maxConcurrency: number; private running: number; private waitQueue: PromiseResolver[]; /** * Creates a new AsyncSemaphore. * @param maxConcurrency The maximum number of concurrent executions allowed. Must be a positive integer. */ constructor(maxConcurrency: number) { if (typeof maxConcurrency !== "number" || !Number.isInteger(maxConcurrency) || maxConcurrency <= 0) { throw new Error("AsyncSemaphore: maxConcurrency must be a positive integer."); } this.maxConcurrency = maxConcurrency; this.running = 0; this.waitQueue = []; } setMaxConcurrency(newMax: number): void { this.maxConcurrency = newMax; } async acquire(): Promise { if (this.running < this.maxConcurrency) { // slot acquired immediately this.running++; return Promise.resolve(); } // No slots available, wait in the queue return new Promise((resolve) => { this.waitQueue.push(resolve); }); } release(): void { // If there are waiters, and we still have concurrency left (might have changed), wake one up if (this.waitQueue.length > 0 && this.running <= this.maxConcurrency) { const nextResolver = this.waitQueue.shift(); if (nextResolver) { // Defer resolution to avoid potential deep stacks // and allow the current execution context to complete. queueMicrotask(nextResolver); } } else if (this.running > 0) { this.running--; } } }