semaphore.ts
| 1 | type PromiseResolver = () => void; |
| 2 | |
| 3 | export class AsyncSemaphore { |
| 4 | private maxConcurrency: number; |
| 5 | private running: number; |
| 6 | private waitQueue: PromiseResolver[]; |
| 7 | |
| 8 | /** |
| 9 | * Creates a new AsyncSemaphore. |
| 10 | * @param maxConcurrency The maximum number of concurrent executions allowed. Must be a positive integer. |
| 11 | */ |
| 12 | constructor(maxConcurrency: number) { |
| 13 | if (typeof maxConcurrency !== "number" || !Number.isInteger(maxConcurrency) || maxConcurrency <= 0) { |
| 14 | throw new Error("AsyncSemaphore: maxConcurrency must be a positive integer."); |
| 15 | } |
| 16 | this.maxConcurrency = maxConcurrency; |
| 17 | this.running = 0; |
| 18 | this.waitQueue = []; |
| 19 | } |
| 20 | |
| 21 | setMaxConcurrency(newMax: number): void { |
| 22 | this.maxConcurrency = newMax; |
| 23 | // raising the limit has to wake waiters, otherwise the freed slots stay unused until a |
| 24 | // running task happens to finish. waking one counts as occupying a slot, hence running++ |
| 25 | while (this.waitQueue.length > 0 && this.running < this.maxConcurrency) { |
| 26 | const nextResolver = this.waitQueue.shift(); |
| 27 | this.running++; |
| 28 | if (nextResolver) queueMicrotask(nextResolver); |
| 29 | } |
| 30 | } |
| 31 | |
| 32 | /** |
| 33 | * Acquires a slot, waiting in the queue if none is free. |
| 34 | * @param signal aborts the wait, rejecting with an AbortError. Required to abort a *queued* |
| 35 | * waiter: it removes the resolver from the queue, so a later release() can't hand a slot to a |
| 36 | * promise nobody awaits (which would permanently reduce the effective concurrency). |
| 37 | */ |
| 38 | async acquire(signal?: AbortSignal): Promise<void> { |
| 39 | if (signal?.aborted) throw new DOMException("Semaphore acquire aborted", "AbortError"); |
| 40 | if (this.running < this.maxConcurrency) { |
| 41 | // slot acquired immediately |
| 42 | this.running++; |
| 43 | return; |
| 44 | } |
| 45 | |
| 46 | // No slots available, wait in the queue |
| 47 | return new Promise<void>((resolve, reject) => { |
| 48 | const resolver = () => { |
| 49 | signal?.removeEventListener("abort", onAbort); |
| 50 | resolve(); |
| 51 | }; |
| 52 | const onAbort = () => { |
| 53 | const index = this.waitQueue.indexOf(resolver); |
| 54 | if (index !== -1) this.waitQueue.splice(index, 1); |
| 55 | reject(new DOMException("Semaphore acquire aborted", "AbortError")); |
| 56 | }; |
| 57 | signal?.addEventListener("abort", onAbort, { once: true }); |
| 58 | this.waitQueue.push(resolver); |
| 59 | }); |
| 60 | } |
| 61 | |
| 62 | release(): void { |
| 63 | // If there are waiters, and we still have concurrency left (might have changed), wake one up |
| 64 | if (this.waitQueue.length > 0 && this.running <= this.maxConcurrency) { |
| 65 | const nextResolver = this.waitQueue.shift(); |
| 66 | if (nextResolver) { |
| 67 | // Defer resolution to avoid potential deep stacks |
| 68 | // and allow the current execution context to complete. |
| 69 | queueMicrotask(nextResolver); |
| 70 | } |
| 71 | } else if (this.running > 0) { |
| 72 | this.running--; |
| 73 | } |
| 74 | } |
| 75 | } |
| 76 |