semaphore.ts
⎇
Raw
1type 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
5export 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