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