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; // raising the limit has to wake waiters, otherwise the freed slots stay unused until a // running task happens to finish. waking one counts as occupying a slot, hence running++ while (this.waitQueue.length > 0 && this.running < this.maxConcurrency) { const nextResolver = this.waitQueue.shift(); this.running++; if (nextResolver) queueMicrotask(nextResolver); } } /** * Acquires a slot, waiting in the queue if none is free. * @param signal aborts the wait, rejecting with an AbortError. Required to abort a *queued* * waiter: it removes the resolver from the queue, so a later release() can't hand a slot to a * promise nobody awaits (which would permanently reduce the effective concurrency). */ async acquire(signal?: AbortSignal): Promise { if (signal?.aborted) throw new DOMException("Semaphore acquire aborted", "AbortError"); if (this.running < this.maxConcurrency) { // slot acquired immediately this.running++; return; } // No slots available, wait in the queue return new Promise((resolve, reject) => { const resolver = () => { signal?.removeEventListener("abort", onAbort); resolve(); }; const onAbort = () => { const index = this.waitQueue.indexOf(resolver); if (index !== -1) this.waitQueue.splice(index, 1); reject(new DOMException("Semaphore acquire aborted", "AbortError")); }; signal?.addEventListener("abort", onAbort, { once: true }); this.waitQueue.push(resolver); }); } 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--; } } }