semaphore.ts
⎇
Raw
1type PromiseResolver = () => void;
2
3export 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