{"version":3,"sources":["../../../src/internals/helpers/retryable.ts"],"names":["RunStrategy","THROW_IMMEDIATELY","SETTLE_ROUND","SETTLE_ALL","Retryable","ctx","createRandomHash","executor","onReset","onError","onRetry","config","maxRetries","Math","max","runGroup","strategy","inputs","Promise","all","map","input","get","controller","AbortController","results","allSettled","groupSignal","signal","undefined","catch","err","abort","throwIfAborted","result","value","runSequence","collect","R","values","asyncProperties","mapValues","attempt","executionId","Object","defineProperty","enumerable","isResolved","state","TaskState","RESOLVED","isRejected","REJECTED","_run","task","Task","assertAborted","lastError","pRetry","retries","factor","shouldRetry","e","FrameworkError","isRetryable","aborted","onFailedAttempt","meta","then","x","resolve","reject","resolvedValue","rejectedValue","PENDING","reset"],"mappings":";;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;AAmBO,MAAMA,WAAAA,GAAc;;;;EAIzBC,iBAAAA,EAAmB,mBAAA;;;;EAKnBC,YAAAA,EAAc,cAAA;;;;EAKdC,UAAAA,EAAY;AACd;AAiBO,MAAMC,SAAAA,CAAAA;EAnDb;;;AAoDW,EAAA,GAAA;AACT,EAAA,MAAA;AACA,EAAA,OAAA;AACA,EAAA,SAAA;AAEA,EAAA,WAAA,CAAYC,GAAAA,EAMT;AACD,IAAA,IAAA,CAAK,MAAA,GAAS,IAAA;AACd,IAAA,IAAA,CAAK,MAAMC,yBAAAA,EAAAA;AACX,IAAA,IAAA,CAAK,SAAA,GAAY;AACfC,MAAAA,QAAAA,EAAUF,GAAAA,CAAIE,QAAAA;AACdC,MAAAA,OAAAA,EAASH,GAAAA,CAAIG,OAAAA;AACbC,MAAAA,OAAAA,EAASJ,GAAAA,CAAII,OAAAA;AACbC,MAAAA,OAAAA,EAASL,GAAAA,CAAIK;AACf,KAAA;AACA,IAAA,IAAA,CAAK,OAAA,GAAU;AACb,MAAA,GAAGL,GAAAA,CAAIM,MAAAA;AACPC,MAAAA,UAAAA,EAAYC,KAAKC,GAAAA,CAAIT,GAAAA,CAAIM,MAAAA,EAAQC,UAAAA,IAAc,GAAG,CAAA;AACpD,KAAA;AACF,EAAA;EAEA,aAAaG,QAAAA,CACXC,UACAC,MAAAA,EACc;AACd,IAAA,IAAID,QAAAA,KAAahB,YAAYC,iBAAAA,EAAmB;AAC9C,MAAA,OAAO,MAAMiB,OAAAA,CAAQC,GAAAA,CAAIF,MAAAA,CAAOG,GAAAA,CAAI,CAACC,KAAAA,KAAUA,KAAAA,CAAMC,GAAAA,EAAG,CAAA,CAAA;AAC1D,IAAA;AAEA,IAAA,MAAMC,UAAAA,GAAa,IAAIC,eAAAA,EAAAA;AACvB,IAAA,MAAMC,OAAAA,GAAU,MAAMP,OAAAA,CAAQQ,UAAAA,CAC5BT,MAAAA,CAAOG,GAAAA,CAAI,CAACC,KAAAA,KACVA,KAAAA,CACGC,GAAAA,CAAIN,QAAAA,KAAahB,WAAAA,CAAYG,UAAAA,GAAa;AAAEwB,MAAAA,WAAAA,EAAaJ,UAAAA,CAAWK;AAAO,KAAA,GAAIC,MAAAA,CAAAA,CAC/EC,KAAAA,CAAM,CAACC,GAAAA,KAAAA;AACNR,MAAAA,UAAAA,CAAWS,MAAMD,GAAAA,CAAAA;AACjB,MAAA,MAAMA,GAAAA;AACR,IAAA,CAAA,CAAA,CAAA,CAAA;AAGNR,IAAAA,UAAAA,CAAWK,OAAOK,cAAAA,EAAc;AAChC,IAAA,OAAOR,OAAAA,CAAQL,GAAAA,CAAI,CAACc,MAAAA,KAAYA,OAAqCC,KAAK,CAAA;AAC5E,EAAA;AAEA,EAAA,cAAcC,YAAenB,MAAAA,EAAoD;AAC/E,IAAA,KAAA,MAAWI,SAASJ,MAAAA,EAAQ;AAC1B,MAAA,MAAM,MAAMI,MAAMC,GAAAA,EAAG;AACvB,IAAA;AACF,EAAA;AAEA,EAAA,aAAae,QAAWpB,MAAAA,EAA4C;AAElE,IAAA,MAAMC,OAAAA,CAAQC,GAAAA,CAAImB,YAAAA,CAAEC,MAAAA,CAAOtB,MAAAA,CAAAA,CAAQG,GAAAA,CAAI,CAACC,KAAAA,KAAUA,KAAAA,CAAMC,GAAAA,EAAG,CAAA,CAAA;AAG3D,IAAA,OAAO,MAAMkB,2BAAAA,CACXF,YAAAA,CAAEG,SAAAA,CAAUxB,MAAAA,EAAQ,CAACkB,KAAAA,KAAWA,KAAAA,CAAyBb,GAAAA,EAAG,CAAA,CAAA;AAIhE,EAAA;AAEA,EAAA,WAAA,CAAYoB,OAAAA,EAAe;AACzB,IAAA,MAAMrC,GAAAA,GAAwB;AAC5BqC,MAAAA,OAAAA;AACAC,MAAAA,WAAAA,EAAa,IAAA,CAAK,GAAA;AAClBf,MAAAA,MAAAA,EAAQ,KAAK,OAAA,CAAQA;AACvB,KAAA;AACAgB,IAAAA,MAAAA,CAAOC,cAAAA,CAAexC,KAAK,QAAA,EAAU;MACnCyC,UAAAA,EAAY;KACd,CAAA;AACA,IAAA,OAAOzC,GAAAA;AACT,EAAA;AAEA,EAAA,IAAI0C,UAAAA,GAAsB;AACxB,IAAA,OAAO,IAAA,CAAK,MAAA,EAAQC,KAAAA,KAAUC,0BAAAA,CAAUC,QAAAA;AAC1C,EAAA;AAEA,EAAA,IAAIC,UAAAA,GAAsB;AACxB,IAAA,OAAO,IAAA,CAAK,MAAA,EAAQH,KAAAA,KAAUC,0BAAAA,CAAUG,QAAAA;AAC1C,EAAA;AAEUC,EAAAA,IAAAA,CAAK1C,MAAAA,EAA6B;AAC1C,IAAA,MAAM2C,IAAAA,GAAO,IAAIC,qBAAAA,EAAAA;AAEjB,IAAA,MAAMC,gCAAgB,MAAA,CAAA,MAAA;AACpB,MAAA,IAAA,CAAK,OAAA,CAAQ5B,QAAQK,cAAAA,IAAAA;AACrBtB,MAAAA,MAAAA,EAAQgB,aAAaM,cAAAA,IAAAA;IACvB,CAAA,EAHsB,eAAA,CAAA;AAKtB,IAAA,IAAIwB,SAAAA,GAA0B,IAAA;AAC9BC,IAAAA,gBAAAA,CACE,OAAOhB,OAAAA,KAAAA;AACLc,MAAAA,aAAAA,EAAAA;AAEA,MAAA,MAAMnD,GAAAA,GAAM,IAAA,CAAK,WAAA,CAAYqC,OAAAA,CAAAA;AAC7B,MAAA,IAAIA,UAAU,CAAA,EAAG;AACf,QAAA,MAAM,IAAA,CAAK,SAAA,CAAUhC,OAAAA,GAAUL,GAAAA,EAAKoD,SAAAA,CAAAA;AACtC,MAAA;AACA,MAAA,OAAO,MAAM,IAAA,CAAK,SAAA,CAAUlD,QAAAA,CAASF,GAAAA,CAAAA;IACvC,CAAA,EACA;AACEsD,MAAAA,OAAAA,EAAS,KAAK,OAAA,CAAQ/C,UAAAA;AACtBgD,MAAAA,MAAAA,EAAQ,KAAK,OAAA,CAAQA,MAAAA;AACrBhC,MAAAA,MAAAA,EAAQ,KAAK,OAAA,CAAQA,MAAAA;AACrBiC,MAAAA,WAAAA,0BAAcC,CAAAA,KAAAA;AACZ,QAAA,IAAI,CAACC,yBAAAA,CAAeC,WAAAA,CAAYF,CAAAA,CAAAA,EAAI;AAClC,UAAA,OAAO,KAAA;AACT,QAAA;AACA,QAAA,OAAO,CAACnD,MAAAA,EAAQgB,WAAAA,EAAasC,WAAW,CAAC,IAAA,CAAK,QAAQrC,MAAAA,EAAQqC,OAAAA;MAChE,CAAA,EALa,aAAA,CAAA;MAMbC,eAAAA,kBAAiB,MAAA,CAAA,OAAOJ,GAAGK,IAAAA,KAAAA;AACzBV,QAAAA,SAAAA,GAAYK,CAAAA;AACZ,QAAA,MAAM,IAAA,CAAK,UAAUrD,OAAAA,GAAUqD,CAAAA,EAAG,KAAK,WAAA,CAAYK,IAAAA,CAAKzB,OAAO,CAAA,CAAA;AAC/D,QAAA,IAAI,CAACqB,yBAAAA,CAAeC,WAAAA,CAAYF,CAAAA,CAAAA,EAAI;AAClC,UAAA,MAAMA,CAAAA;AACR,QAAA;AACAN,QAAAA,aAAAA,EAAAA;MACF,CAAA,EAPiB,iBAAA;AAQnB,KAAA,CAAA,CAECY,IAAAA,CAAK,CAACC,CAAAA,KAAMf,KAAKgB,OAAAA,CAAQD,CAAAA,CAAAA,CAAAA,CACzBvC,MAAM,CAACuC,CAAAA,KAAMf,IAAAA,CAAKiB,MAAAA,CAAOF,CAAAA,CAAAA,CAAAA;AAE5B,IAAA,OAAOf,IAAAA;AACT,EAAA;AAEA,EAAA,MAAMhC,IAAIX,MAAAA,EAAyC;AACjD,IAAA,IAAI,KAAKoC,UAAAA,EAAY;AACnB,MAAA,OAAO,IAAA,CAAK,OAAQyB,aAAAA,EAAa;AACnC,IAAA;AACA,IAAA,IAAI,KAAKrB,UAAAA,EAAY;AACnB,MAAA,MAAM,IAAA,CAAK,QAAQsB,aAAAA,EAAAA;AACrB,IAAA;AAEA,IAAA,IAAI,KAAK,MAAA,EAAQzB,KAAAA,KAAUC,0BAAAA,CAAUyB,OAAAA,IAAW,CAAC/D,MAAAA,EAAQ;AACvD,MAAA,OAAO,IAAA,CAAK,MAAA;AACd,IAAA;AAEA,IAAA,IAAA,CAAK,MAAA,EAAQmB,QAAQ,MAAA;IAAO,CAAA,CAAA;AAC5B,IAAA,IAAA,CAAK,MAAA,GAAS,IAAA,CAAKuB,IAAAA,CAAK1C,MAAAA,CAAAA;AACxB,IAAA,OAAO,IAAA,CAAK,MAAA;AACd,EAAA;EAEAgE,KAAAA,GAAQ;AACN,IAAA,IAAA,CAAK,MAAA,EAAQ7C,QAAQ,MAAA;IAAO,CAAA,CAAA;AAC5B,IAAA,IAAA,CAAK,MAAA,GAAS,IAAA;AACd,IAAA,IAAA,CAAK,UAAUtB,OAAAA,IAAO;AACxB,EAAA;AACF","file":"retryable.cjs","sourcesContent":["/**\n * Copyright 2025 © BeeAI a Series of LF Projects, LLC\n * SPDX-License-Identifier: Apache-2.0\n */\n\nimport * as R from \"remeda\";\nimport { Task, TaskState } from \"promise-based-task\";\nimport { FrameworkError } from \"@/errors.js\";\nimport { EnumValue } from \"@/internals/types.js\";\nimport { asyncProperties } from \"@/internals/helpers/promise.js\";\nimport { createRandomHash } from \"@/internals/helpers/hash.js\";\nimport { pRetry } from \"@/internals/helpers/retry.js\";\n\nexport interface RetryableConfig {\n  maxRetries: number;\n  factor?: number;\n  signal?: AbortSignal;\n}\n\nexport const RunStrategy = {\n  /**\n   * Once a single Retryable throws, other retry ables get cancelled immediately.\n   */\n  THROW_IMMEDIATELY: \"THROW_IMMEDIATELY\",\n\n  /**\n   * Once a single Retryable throws, wait for other to completes, but prevent further retries.\n   */\n  SETTLE_ROUND: \"SETTLE_ROUND\",\n\n  /**\n   * Once a single Retryable throws, other Retryables remains to continue. Error is thrown by the end.\n   */\n  SETTLE_ALL: \"SETTLE_ALL\",\n} as const;\n\nexport interface RetryableRunConfig {\n  groupSignal: AbortSignal;\n}\n\nexport interface RetryableContext {\n  executionId: string;\n  attempt: number;\n  signal?: AbortSignal;\n}\n\nexport type RetryableHandler<T> = (ctx: RetryableContext) => Promise<T>;\nexport type ResetHandler = () => void;\nexport type ErrorHandler = (error: Error, ctx: RetryableContext) => void | Promise<void>;\nexport type RetryHandler = (ctx: RetryableContext, lastError: Error) => void | Promise<void>;\n\nexport class Retryable<T> {\n  readonly #id: string;\n  #value: Task<T, Error> | null;\n  #config: RetryableConfig;\n  #handlers;\n\n  constructor(ctx: {\n    executor: RetryableHandler<T>;\n    onReset?: ResetHandler;\n    onError?: ErrorHandler;\n    onRetry?: RetryHandler;\n    config?: Partial<RetryableConfig>;\n  }) {\n    this.#value = null;\n    this.#id = createRandomHash();\n    this.#handlers = {\n      executor: ctx.executor,\n      onReset: ctx.onReset,\n      onError: ctx.onError,\n      onRetry: ctx.onRetry,\n    } as const;\n    this.#config = {\n      ...ctx.config,\n      maxRetries: Math.max(ctx.config?.maxRetries || 0, 0),\n    };\n  }\n\n  static async runGroup<T>(\n    strategy: EnumValue<typeof RunStrategy>,\n    inputs: Retryable<T>[],\n  ): Promise<T[]> {\n    if (strategy === RunStrategy.THROW_IMMEDIATELY) {\n      return await Promise.all(inputs.map((input) => input.get()));\n    }\n\n    const controller = new AbortController();\n    const results = await Promise.allSettled(\n      inputs.map((input) =>\n        input\n          .get(strategy === RunStrategy.SETTLE_ALL ? { groupSignal: controller.signal } : undefined)\n          .catch((err) => {\n            controller.abort(err);\n            throw err;\n          }),\n      ),\n    );\n    controller.signal.throwIfAborted();\n    return results.map((result) => (result as PromiseFulfilledResult<T>).value!);\n  }\n\n  static async *runSequence<T>(inputs: readonly Retryable<T>[]): AsyncGenerator<T> {\n    for (const input of inputs) {\n      yield await input.get();\n    }\n  }\n\n  static async collect<T>(inputs: T & Record<string, Retryable<any>>) {\n    // Solve everything\n    await Promise.all(R.values(inputs).map((input) => input.get()));\n\n    // Obtain latest values\n    return await asyncProperties(\n      R.mapValues(inputs, (value) => (value as Retryable<any>).get()) as {\n        [K in keyof T]: Promise<T[K] extends Retryable<infer Q> ? Q : never>;\n      },\n    );\n  }\n\n  #getContext(attempt: number): RetryableContext {\n    const ctx: RetryableContext = {\n      attempt,\n      executionId: this.#id,\n      signal: this.#config.signal,\n    };\n    Object.defineProperty(ctx, \"signal\", {\n      enumerable: false,\n    });\n    return ctx;\n  }\n\n  get isResolved(): boolean {\n    return this.#value?.state === TaskState.RESOLVED;\n  }\n\n  get isRejected(): boolean {\n    return this.#value?.state === TaskState.REJECTED;\n  }\n\n  protected _run(config?: RetryableRunConfig) {\n    const task = new Task<T, Error>();\n\n    const assertAborted = () => {\n      this.#config.signal?.throwIfAborted?.();\n      config?.groupSignal?.throwIfAborted?.();\n    };\n\n    let lastError: Error | null = null;\n    pRetry(\n      async (attempt) => {\n        assertAborted();\n\n        const ctx = this.#getContext(attempt);\n        if (attempt > 1) {\n          await this.#handlers.onRetry?.(ctx, lastError!);\n        }\n        return await this.#handlers.executor(ctx);\n      },\n      {\n        retries: this.#config.maxRetries,\n        factor: this.#config.factor,\n        signal: this.#config.signal,\n        shouldRetry: (e) => {\n          if (!FrameworkError.isRetryable(e)) {\n            return false;\n          }\n          return !config?.groupSignal?.aborted && !this.#config.signal?.aborted;\n        },\n        onFailedAttempt: async (e, meta) => {\n          lastError = e;\n          await this.#handlers.onError?.(e, this.#getContext(meta.attempt));\n          if (!FrameworkError.isRetryable(e)) {\n            throw e;\n          }\n          assertAborted();\n        },\n      },\n    )\n      .then((x) => task.resolve(x))\n      .catch((x) => task.reject(x));\n\n    return task;\n  }\n\n  async get(config?: RetryableRunConfig): Promise<T> {\n    if (this.isResolved) {\n      return this.#value!.resolvedValue()!;\n    }\n    if (this.isRejected) {\n      throw this.#value?.rejectedValue();\n    }\n\n    if (this.#value?.state === TaskState.PENDING && !config) {\n      return this.#value;\n    }\n\n    this.#value?.catch?.(() => {});\n    this.#value = this._run(config);\n    return this.#value;\n  }\n\n  reset() {\n    this.#value?.catch?.(() => {});\n    this.#value = null;\n    this.#handlers.onReset?.();\n  }\n}\n"]}