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 }
24
25 async acquire(): Promise<void> {
26 if (this.running < this.maxConcurrency) {
27 // slot acquired immediately
28 this.running++;
29 return Promise.resolve();
30 }
31
32 // No slots available, wait in the queue
33 return new Promise<void>((resolve) => {
34 this.waitQueue.push(resolve);
35 });
36 }
37
38 release(): void {
39 // If there are waiters, and we still have concurrency left (might have changed), wake one up
40 if (this.waitQueue.length > 0 && this.running <= this.maxConcurrency) {
41 const nextResolver = this.waitQueue.shift();
42 if (nextResolver) {
43 // Defer resolution to avoid potential deep stacks
44 // and allow the current execution context to complete.
45 queueMicrotask(nextResolver);
46 }
47 } else if (this.running > 0) {
48 this.running--;
49 }
50 }
51}
52