{"version":3,"file":"index.cjs","names":[],"sources":["../src/index.ts"],"sourcesContent":["// Shared by all fast-path run() calls to avoid allocating a new promise per call.\nconst RESOLVED_PROMISE = Promise.resolve();\n\nexport class PromisePool<T = unknown> {\n  private readonly promises: Set<Promise<T>>;\n  private readonly resumeFunctions: Array<(() => void) | undefined>;\n  private resumeFunctionIndex: number;\n  private reservedPromiseCount: number;\n  private _concurrency: number;\n  private _queuedPromiseCount: number;\n\n  constructor(concurrency = 10) {\n    this._concurrency = concurrency;\n    this._queuedPromiseCount = 0;\n    this.promises = new Set();\n    this.resumeFunctions = [];\n    this.resumeFunctionIndex = 0;\n    this.reservedPromiseCount = 0;\n  }\n\n  get workingPromiseCount(): number {\n    return this.promises.size;\n  }\n\n  get concurrency(): number {\n    return this._concurrency;\n  }\n\n  set concurrency(concurrency: number) {\n    const oldConcurrency = this._concurrency;\n    this._concurrency = concurrency;\n    if (concurrency > oldConcurrency) {\n      this.resume();\n    }\n  }\n\n  get queuedPromiseCount(): number {\n    return this._queuedPromiseCount;\n  }\n\n  promiseAll(): Promise<T[]> {\n    return Promise.all(this.promises);\n  }\n\n  promiseAllSettled(): Promise<PromiseSettledResult<T>[]> {\n    return Promise.allSettled(this.promises);\n  }\n\n  run(startPromise: () => Promise<T>): Promise<void> {\n    this._queuedPromiseCount++;\n    // Start the task synchronously when the pool has capacity to avoid the\n    // promise allocations and microtask hops of the queued (slow) path.\n    if (this.tryReserveCapacity()) {\n      try {\n        void this.startReservedTask(startPromise);\n      } catch (error) {\n        return Promise.reject(error);\n      }\n      return RESOLVED_PROMISE;\n    }\n    return this.runQueued(startPromise);\n  }\n\n  runAndWaitForReturnValue<R extends T>(startPromise: () => Promise<R>): Promise<R> {\n    this._queuedPromiseCount++;\n    if (this.tryReserveCapacity()) {\n      try {\n        return this.startReservedTask(startPromise);\n      } catch (error) {\n        return Promise.reject(error);\n      }\n    }\n    return this.runQueuedAndWaitForReturnValue(startPromise);\n  }\n\n  private async runQueued(startPromise: () => Promise<T>): Promise<void> {\n    while (!this.tryReserveCapacity()) {\n      await this.waitForResume();\n      if (this.tryKeepReservedSlot()) break;\n    }\n    // The started promise is intentionally discarded since run() resolves once the task starts.\n    void this.startReservedTask(startPromise);\n  }\n\n  private async runQueuedAndWaitForReturnValue<R extends T>(startPromise: () => Promise<R>): Promise<R> {\n    while (!this.tryReserveCapacity()) {\n      await this.waitForResume();\n      if (this.tryKeepReservedSlot()) break;\n    }\n    // The async function awaits the returned promise, so the caller receives the task's resolved value.\n    return this.startReservedTask(startPromise);\n  }\n\n  /**\n   * Starts the given task using a slot the caller has already reserved.\n   * The reservation is converted into a working promise, or released if `startPromise` throws.\n   */\n  private startReservedTask<R extends T>(startPromise: () => Promise<R>): Promise<R> {\n    let promise: Promise<R>;\n    try {\n      // Use `.then()` with two handlers instead of `.finally()` since `.finally()` allocates extra internal promises.\n      promise = startPromise().then(\n        (value) => {\n          this.completeTask(promise);\n          return value;\n        },\n        (error: unknown) => {\n          this.completeTask(promise);\n          throw error;\n        }\n      );\n    } catch (error) {\n      this._queuedPromiseCount--;\n      this.reservedPromiseCount--;\n      this.resume();\n      throw error;\n    }\n    this.reservedPromiseCount--;\n    this.promises.add(promise);\n    return promise;\n  }\n\n  private tryReserveCapacity(): boolean {\n    if (this.promises.size + this.reservedPromiseCount >= this._concurrency) {\n      return false;\n    }\n    this.reservedPromiseCount++;\n    return true;\n  }\n\n  // This is not an async function so that awaiting it costs no extra promise beyond the waiter itself.\n  private waitForResume(): Promise<void> {\n    return new Promise<void>((resolve) => {\n      this.resumeFunctions.push(resolve);\n    });\n  }\n\n  /**\n   * Decides whether the slot reserved by resume() can be kept, releasing it otherwise.\n   * The check happens after a microtask hop so that a concurrency reduction applied\n   * right after resume() can still cancel stale reservations.\n   */\n  private tryKeepReservedSlot(): boolean {\n    if (this.promises.size + this.reservedPromiseCount <= this._concurrency) {\n      return true;\n    }\n    this.reservedPromiseCount--;\n    return false;\n  }\n\n  private completeTask(promise: Promise<T>): void {\n    this._queuedPromiseCount--;\n    this.promises.delete(promise);\n    this.resume();\n  }\n\n  private resume(): void {\n    while (this.hasCapacity()) {\n      const resumeFunctionIndex = this.resumeFunctionIndex;\n      const resumeFunction = this.resumeFunctions[resumeFunctionIndex];\n      if (!resumeFunction) break;\n\n      this.resumeFunctionIndex++;\n      this.resumeFunctions[resumeFunctionIndex] = undefined;\n      this.reservedPromiseCount++;\n      resumeFunction();\n    }\n\n    if (this.resumeFunctionIndex === this.resumeFunctions.length) {\n      this.resumeFunctions.length = 0;\n      this.resumeFunctionIndex = 0;\n    } else if (this.resumeFunctionIndex >= 256 && this.resumeFunctionIndex >= this.resumeFunctions.length / 2) {\n      this.resumeFunctions.splice(0, this.resumeFunctionIndex);\n      this.resumeFunctionIndex = 0;\n    }\n  }\n\n  private hasCapacity(): boolean {\n    return this.promises.size + this.reservedPromiseCount < this._concurrency;\n  }\n}\n"],"mappings":"gFACA,MAAM,EAAmB,QAAQ,QAAQ,EAEzC,IAAa,EAAb,KAAsC,CACpC,SACA,gBACA,oBACA,qBACA,aACA,oBAEA,YAAY,EAAc,GAAI,CAC5B,KAAK,aAAe,EACpB,KAAK,oBAAsB,EAC3B,KAAK,SAAW,IAAI,IACpB,KAAK,gBAAkB,CAAC,EACxB,KAAK,oBAAsB,EAC3B,KAAK,qBAAuB,CAC9B,CAEA,IAAI,qBAA8B,CAChC,OAAO,KAAK,SAAS,IACvB,CAEA,IAAI,aAAsB,CACxB,OAAO,KAAK,YACd,CAEA,IAAI,YAAY,EAAqB,CACnC,IAAM,EAAiB,KAAK,aAC5B,KAAK,aAAe,EAChB,EAAc,GAChB,KAAK,OAAO,CAEhB,CAEA,IAAI,oBAA6B,CAC/B,OAAO,KAAK,mBACd,CAEA,YAA2B,CACzB,OAAO,QAAQ,IAAI,KAAK,QAAQ,CAClC,CAEA,mBAAwD,CACtD,OAAO,QAAQ,WAAW,KAAK,QAAQ,CACzC,CAEA,IAAI,EAA+C,CAIjD,GAHA,KAAK,sBAGD,KAAK,mBAAmB,EAAG,CAC7B,GAAI,CACF,KAAU,kBAAkB,CAAY,CAC1C,OAAS,EAAO,CACd,OAAO,QAAQ,OAAO,CAAK,CAC7B,CACA,OAAO,CACT,CACA,OAAO,KAAK,UAAU,CAAY,CACpC,CAEA,yBAAsC,EAA4C,CAEhF,GADA,KAAK,sBACD,KAAK,mBAAmB,EAC1B,GAAI,CACF,OAAO,KAAK,kBAAkB,CAAY,CAC5C,OAAS,EAAO,CACd,OAAO,QAAQ,OAAO,CAAK,CAC7B,CAEF,OAAO,KAAK,+BAA+B,CAAY,CACzD,CAEA,MAAc,UAAU,EAA+C,CACrE,KAAO,CAAC,KAAK,mBAAmB,IAC9B,MAAM,KAAK,cAAc,EACrB,MAAK,oBAAoB,KAG/B,KAAU,kBAAkB,CAAY,CAC1C,CAEA,MAAc,+BAA4C,EAA4C,CACpG,KAAO,CAAC,KAAK,mBAAmB,IAC9B,MAAM,KAAK,cAAc,EACrB,MAAK,oBAAoB,KAG/B,OAAO,KAAK,kBAAkB,CAAY,CAC5C,CAMA,kBAAuC,EAA4C,CACjF,IAAI,EACJ,GAAI,CAEF,EAAU,EAAa,CAAC,CAAC,KACtB,IACC,KAAK,aAAa,CAAO,EAClB,GAER,GAAmB,CAElB,MADA,KAAK,aAAa,CAAO,EACnB,CACR,CACF,CACF,OAAS,EAAO,CAId,KAHA,MAAK,sBACL,KAAK,uBACL,KAAK,OAAO,EACN,CACR,CAGA,MAFA,MAAK,uBACL,KAAK,SAAS,IAAI,CAAO,EAClB,CACT,CAEA,oBAAsC,CAKpC,OAJI,KAAK,SAAS,KAAO,KAAK,sBAAwB,KAAK,aAClD,IAET,KAAK,uBACE,GACT,CAGA,eAAuC,CACrC,OAAO,IAAI,QAAe,GAAY,CACpC,KAAK,gBAAgB,KAAK,CAAO,CACnC,CAAC,CACH,CAOA,qBAAuC,CAKrC,OAJI,KAAK,SAAS,KAAO,KAAK,sBAAwB,KAAK,aAClD,IAET,KAAK,uBACE,GACT,CAEA,aAAqB,EAA2B,CAC9C,KAAK,sBACL,KAAK,SAAS,OAAO,CAAO,EAC5B,KAAK,OAAO,CACd,CAEA,QAAuB,CACrB,KAAO,KAAK,YAAY,GAAG,CACzB,IAAM,EAAsB,KAAK,oBAC3B,EAAiB,KAAK,gBAAgB,GAC5C,GAAI,CAAC,EAAgB,MAErB,KAAK,sBACL,KAAK,gBAAgB,GAAuB,IAAA,GAC5C,KAAK,uBACL,EAAe,CACjB,CAEI,KAAK,sBAAwB,KAAK,gBAAgB,QACpD,KAAK,gBAAgB,OAAS,EAC9B,KAAK,oBAAsB,GAClB,KAAK,qBAAuB,KAAO,KAAK,qBAAuB,KAAK,gBAAgB,OAAS,IACtG,KAAK,gBAAgB,OAAO,EAAG,KAAK,mBAAmB,EACvD,KAAK,oBAAsB,EAE/B,CAEA,aAA+B,CAC7B,OAAO,KAAK,SAAS,KAAO,KAAK,qBAAuB,KAAK,YAC/D,CACF"}