{"version":3,"sources":["../../src/internals/fetcher.ts"],"names":["RestfulClientError","FrameworkError","createURLParams","data","urlTokenParams","URLSearchParams","key","value","Object","entries","undefined","Array","isArray","forEach","v","append","String","isPlainObject","set","toString","RestfulClient","Serializable","emitter","Emitter","root","child","namespace","creator","input","stream","path","init","groupId","url","getUrl","options","method","headers","getHeaders","emit","emitterToGenerator","fetchEventSource","onopen","response","contentType","get","ok","includes","EventStreamContentType","context","err","text","isRetryable","status","onmessage","msg","event","onclose","onerror","then","catch","error","doNothing","finally","fetch","target","searchParams","json","extra","final","override","filter","isTruthy","Headers","assign","fromEntries","URL","baseUrl","extraPath","paths","pathname","endsWith","replace","createSnapshot","shallowCopy","loadSnapshot","snapshot"],"mappings":";;;;;;;;;;;;AAkBO,MAAMA,2BAA2BC,yBAAAA,CAAAA;EAlBxC;;;AAkBwD;AAGjD,SAASC,gBACdC,IAAAA,EAAyE;AAEzE,EAAA,MAAMC,cAAAA,GAAiB,IAAIC,eAAAA,EAAAA;AAC3B,EAAA,KAAA,MAAW,CAACC,GAAAA,EAAKC,KAAAA,KAAUC,MAAAA,CAAOC,OAAAA,CAAQN,IAAAA,CAAAA,EAAO;AAC/C,IAAA,IAAII,UAAUG,MAAAA,EAAW;AACvB,MAAA;AACF,IAAA;AAEA,IAAA,IAAIC,KAAAA,CAAMC,OAAAA,CAAQL,KAAAA,CAAAA,EAAQ;AACxBA,MAAAA,KAAAA,CAAMM,OAAAA,CAAQ,CAACC,CAAAA,KAAAA;AACb,QAAA,IAAIA,MAAMJ,MAAAA,EAAW;AACnBN,UAAAA,cAAAA,CAAeW,MAAAA,CAAOT,GAAAA,EAAKU,MAAAA,CAAOF,CAAAA,CAAAA,CAAAA;AACpC,QAAA;MACF,CAAA,CAAA;IACF,CAAA,MAAA,IAAWG,oBAAAA,CAAcV,KAAAA,CAAAA,EAAQ;AAC/BH,MAAAA,cAAAA,CAAec,IAAIZ,GAAAA,EAAKJ,eAAAA,CAAgBK,KAAAA,CAAAA,CAAOY,UAAQ,CAAA;IACzD,CAAA,MAAO;AACLf,MAAAA,cAAAA,CAAec,GAAAA,CAAIZ,GAAAA,EAAKU,MAAAA,CAAOT,KAAAA,CAAAA,CAAAA;AACjC,IAAA;AACF,EAAA;AACA,EAAA,OAAOH,cAAAA;AACT;AAtBgBF,MAAAA,CAAAA,eAAAA,EAAAA,iBAAAA,CAAAA;AAgDT,MAAMkB,sBAAwDC,6BAAAA,CAAAA;EArErE;;;;EAsEkBC,OAAAA,GAAUC,mBAAAA,CAAQC,KAAKC,KAAAA,CAA2B;IAChEC,SAAAA,EAAW;AAAC,MAAA,WAAA;AAAa,MAAA;;IACzBC,OAAAA,EAAS;GACX,CAAA;AAEA,EAAA,WAAA,CACYC,KAAAA,EAKV;AACA,IAAA,KAAA,EAAK,EAAA,KANKA,KAAAA,GAAAA,KAAAA;AAOZ,EAAA;EAEA,OAAOC,MAAAA,CACLC,MACAC,IAAAA,EACgD;AAChD,IAAA,MAAMT,OAAAA,GAAU,IAAA,CAAKA,OAAAA,CAAQG,KAAAA,CAAM;MACjCO,OAAAA,EAAS;KACX,CAAA;AAEA,IAAA,MAAMJ,KAAAA,GAAqB;AACzBK,MAAAA,GAAAA,EAAK,IAAA,CAAKC,MAAAA,CAAOJ,IAAAA,CAAAA,CAAMX,QAAAA,EAAQ;MAC/BgB,OAAAA,EAAS;QACPC,MAAAA,EAAQ,MAAA;QACR,GAAGL,IAAAA;AACHM,QAAAA,OAAAA,EAAS,MAAM,IAAA,CAAKC,UAAAA,CAAWP,IAAAA,EAAMM,OAAAA;AACvC;AACF,KAAA;AACA,IAAA,MAAMf,OAAAA,CAAQiB,KAAK,aAAA,EAAe;AAAEX,MAAAA;KAAM,CAAA;AAE1C,IAAA,OAAO,OAAOY,+BAAmB,OAAO,EAAED,MAAI,KAC5CE,qCAAAA,CAAiBb,MAAMK,GAAAA,EAAK;AAC1B,MAAA,GAAGL,KAAAA,CAAMO,OAAAA;AACT,MAAA,MAAMO,OAAOC,QAAAA,EAAQ;AACnB,QAAA,MAAMC,WAAAA,GAAcD,QAAAA,CAASN,OAAAA,CAAQQ,GAAAA,CAAI,cAAA,CAAA,IAAmB,EAAA;AAC5D,QAAA,IAAIF,QAAAA,CAASG,EAAAA,IAAMF,WAAAA,CAAYG,QAAAA,CAASC,2CAAAA,CAAAA,EAAyB;AAC/D,UAAA,MAAM1B,OAAAA,CAAQiB,KAAK,YAAA,EAAc;AAAEX,YAAAA;WAAM,CAAA;AACzC,UAAA;AACF,QAAA;AACA,QAAA,MAAM,IAAI5B,kBAAAA,CAAmB,mBAAA,EAAqB,EAAA,EAAI;UACpDiD,OAAAA,EAAS;AACPhB,YAAAA,GAAAA,EAAKU,QAAAA,CAASV,GAAAA;YACdiB,GAAAA,EAAK,MAAMP,SAASQ,IAAAA,EAAI;AACxBR,YAAAA;AACF,WAAA;AACAS,UAAAA,WAAAA,EAAaT,SAASU,MAAAA,IAAU,GAAA,IAAOV,SAASU,MAAAA,GAAS,GAAA,IAAOV,SAASU,MAAAA,KAAW;SACtF,CAAA;AACF,MAAA,CAAA;AACA,MAAA,MAAMC,UAAUC,GAAAA,EAAG;AACjB,QAAA,IAAIA,GAAAA,EAAKC,UAAU,OAAA,EAAS;AAC1B,UAAA,MAAM,IAAIxD,kBAAAA,CAAmB,CAAA,oCAAA,CAAA,EAAwC,EAAA,EAAI;YACvEiD,OAAAA,EAASM;WACX,CAAA;AACF,QAAA;AACA,QAAA,MAAMjC,OAAAA,CAAQiB,KAAK,eAAA,EAAiB;AAAEX,UAAAA,KAAAA;UAAOzB,IAAAA,EAAMoD;SAAI,CAAA;AACvDhB,QAAAA,IAAAA,CAAKgB,GAAAA,CAAAA;AACP,MAAA,CAAA;MACAE,OAAAA,GAAAA;AAAW,MAAA,CAAA;AACXC,MAAAA,OAAAA,CAAQR,GAAAA,EAAG;AACT,QAAA,MAAM,IAAIlD,mBAAmB,CAAA,oCAAA,CAAA,EAAwC;AAACkD,UAAAA;AAAI,SAAA,CAAA;AAC5E,MAAA;AACF,KAAA,CAAA,CACGS,IAAAA,CAAK,MAAMrC,OAAAA,CAAQiB,KAAK,eAAA,EAAiB;AAAEX,MAAAA;AAAM,KAAA,CAAA,CAAA,CACjDgC,KAAAA,CAAM,OAAOC,KAAAA,KAAAA;AACZ,MAAA,MAAMvC,OAAAA,CAAQiB,KAAK,aAAA,EAAe;AAAEX,QAAAA,KAAAA;AAAOiC,QAAAA;OAAM,CAAA,CAAGD,KAAAA,CAAME,gBAAAA,EAAAA,CAAAA;AAC1D,MAAA,MAAMD,KAAAA;AACR,IAAA,CAAA,CAAA,CACCE,OAAAA,CAAQ,MAAMzC,OAAAA,CAAQiB,KAAK,YAAA,EAAc;AAAEX,MAAAA;AAAM,KAAA,CAAA,CAAA,CAAA;AAExD,EAAA;EAEA,MAAMoC,KAAAA,CAAMlC,MAAeC,IAAAA,EAAyD;AAClF,IAAA,MAAMT,OAAAA,GAAU,IAAA,CAAKA,OAAAA,CAAQG,KAAAA,CAAM;MACjCO,OAAAA,EAAS;KACX,CAAA;AAEA,IAAA,MAAMiC,MAAAA,GAAS,IAAA,CAAK/B,MAAAA,CAAOJ,IAAAA,CAAAA;AAC3B,IAAA,IAAIC,MAAMmC,YAAAA,EAAc;AACtB,MAAA,KAAA,MAAW,CAAC5D,GAAAA,EAAKC,KAAAA,CAAAA,IAAUwB,KAAKmC,YAAAA,EAAc;AAC5CD,QAAAA,MAAAA,CAAOC,YAAAA,CAAahD,GAAAA,CAAIZ,GAAAA,EAAKC,KAAAA,CAAAA;AAC/B,MAAA;AACF,IAAA;AAEA,IAAA,MAAMqB,KAAAA,GAAoB;AACxBK,MAAAA,GAAAA,EAAKgC,OAAO9C,QAAAA,EAAQ;MACpBgB,OAAAA,EAAS;QACP,GAAGJ,IAAAA;AACHM,QAAAA,OAAAA,EAAS,MAAM,IAAA,CAAKC,UAAAA,CAAWP,IAAAA,EAAMM,OAAAA;AACvC;AACF,KAAA;AAEA,IAAA,MAAMf,OAAAA,CAAQiB,KAAK,YAAA,EAAc;AAAEX,MAAAA;KAAM,CAAA;AACzC,IAAA,IAAI;AACF,MAAA,MAAMe,WAAW,MAAMqB,KAAAA,CAAMpC,KAAAA,CAAMK,GAAAA,EAAKL,MAAMO,OAAO,CAAA;AAErD,MAAA,IAAI,CAACQ,SAASG,EAAAA,EAAI;AAChB,QAAA,MAAM,IAAI9C,kBAAAA,CAAmB,kBAAA,EAAoB,EAAA,EAAI;UACnDiD,OAAAA,EAAS;AACPhB,YAAAA,GAAAA,EAAKU,QAAAA,CAASV,GAAAA;YACd4B,KAAAA,EAAO,MAAMlB,SAASQ,IAAAA,EAAI;AAC1BR,YAAAA;AACF,WAAA;UACAS,WAAAA,EAAa;AAAC,YAAA,GAAA;AAAK,YAAA;YAAKL,QAAAA,CAASJ,QAAAA,CAASU,UAAU,GAAA;SACtD,CAAA;AACF,MAAA;AAEA,MAAA,MAAMlD,IAAAA,GAAO,MAAMwC,QAAAA,CAASwB,IAAAA,EAAI;AAChC,MAAA,MAAM7C,OAAAA,CAAQiB,KAAK,cAAA,EAAgB;AAAEI,QAAAA,QAAAA;AAAUxC,QAAAA,IAAAA;AAAMyB,QAAAA;OAAM,CAAA;AAC3D,MAAA,OAAOzB,IAAAA;AACT,IAAA,CAAA,CAAA,OAAS0D,KAAAA,EAAO;AACd,MAAA,MAAMvC,OAAAA,CAAQiB,KAAK,YAAA,EAAc;AAAEsB,QAAAA,KAAAA;AAAOjC,QAAAA;OAAa,CAAA;AACvD,MAAA,MAAMiC,KAAAA;IACR,CAAA,SAAA;AACE,MAAA,MAAMvC,OAAAA,CAAQiB,KAAK,WAAA,EAAa;AAAEX,QAAAA;OAAa,CAAA;AACjD,IAAA;AACF,EAAA;AAEA,EAAA,MAAgBU,WAAW8B,KAAAA,EAAsD;AAC/E,IAAA,MAAMC,QAAQ,EAAC;AACf,IAAA,KAAA,MAAWC,QAAAA,IAAY;MAAC,MAAM,IAAA,CAAK1C,MAAMS,OAAAA,IAAO;AAAM+B,MAAAA;AAAOG,KAAAA,CAAAA,MAAAA,CAAOC,eAAAA,CAAAA,EAAW;AAC7E,MAAA,IAAI5D,cAAAA,CAAQ0D,QAAAA,CAAAA,IAAaA,QAAAA,YAAoBG,OAAAA,EAAS;AACpDjE,QAAAA,MAAAA,CAAOkE,OAAOL,KAAAA,EAAO7D,MAAAA,CAAOmE,YAAYL,QAAAA,CAAS7D,OAAAA,EAAO,CAAA,CAAA;MAC1D,CAAA,MAAO;AACLD,QAAAA,MAAAA,CAAOkE,MAAAA,CAAOL,OAAOC,QAAAA,CAAAA;AACvB,MAAA;AACF,IAAA;AACA,IAAA,OAAOD,KAAAA;AACT,EAAA;AAEUnC,EAAAA,MAAAA,CAAOJ,IAAAA,EAAoB;AACnC,IAAA,MAAMG,GAAAA,GAAM,IAAI2C,GAAAA,CAAI,IAAA,CAAKhD,MAAMiD,OAAO,CAAA;AACtC,IAAA,MAAMC,SAAAA,GAAY,IAAA,CAAKlD,KAAAA,CAAMmD,KAAAA,CAAMjD,IAAAA,CAAAA,IAASA,IAAAA;AAC5C,IAAA,IAAIG,GAAAA,CAAI+C,QAAAA,CAASC,QAAAA,CAAS,GAAA,CAAA,EAAM;AAC9BhD,MAAAA,GAAAA,CAAI+C,QAAAA,IAAYF,SAAAA,CAAUI,OAAAA,CAAQ,KAAA,EAAO,EAAA,CAAA;IAC3C,CAAA,MAAO;AACLjD,MAAAA,GAAAA,CAAI+C,QAAAA,IAAYF,SAAAA;AAClB,IAAA;AACA,IAAA,OAAO7C,GAAAA;AACT,EAAA;EAEAkD,cAAAA,GAAiB;AACf,IAAA,OAAO;MACLvD,KAAAA,EAAOwD,qBAAAA,CAAY,KAAKxD,KAAK,CAAA;AAC7BN,MAAAA,OAAAA,EAAS,IAAA,CAAKA;AAChB,KAAA;AACF,EAAA;AAEA+D,EAAAA,YAAAA,CAAaC,QAAAA,EAAwD;AACnE9E,IAAAA,MAAAA,CAAOkE,MAAAA,CAAO,MAAMY,QAAAA,CAAAA;AACtB,EAAA;AACF","file":"fetcher.cjs","sourcesContent":["/**\n * Copyright 2025 © BeeAI a Series of LF Projects, LLC\n * SPDX-License-Identifier: Apache-2.0\n */\n\nimport { FrameworkError } from \"@/errors.js\";\nimport { Serializable } from \"@/internals/serializable.js\";\nimport {\n  EventSourceMessage,\n  EventStreamContentType,\n  fetchEventSource,\n} from \"@ai-zen/node-fetch-event-source\";\nimport { FetchEventSourceInit } from \"@ai-zen/node-fetch-event-source/lib/cjs/fetch.js\";\nimport { emitterToGenerator } from \"@/internals/helpers/promise.js\";\nimport { doNothing, isArray, isPlainObject, isTruthy } from \"remeda\";\nimport { Callback, Emitter } from \"@/emitter/emitter.js\";\nimport { shallowCopy } from \"@/serializer/utils.js\";\n\nexport class RestfulClientError extends FrameworkError {}\n\ntype URLParamType = string | number | boolean | null | undefined;\nexport function createURLParams(\n  data: Record<string, URLParamType | URLParamType[] | Record<string, any>>,\n) {\n  const urlTokenParams = new URLSearchParams();\n  for (const [key, value] of Object.entries(data)) {\n    if (value === undefined) {\n      continue;\n    }\n\n    if (Array.isArray(value)) {\n      value.forEach((v) => {\n        if (v !== undefined) {\n          urlTokenParams.append(key, String(v));\n        }\n      });\n    } else if (isPlainObject(value)) {\n      urlTokenParams.set(key, createURLParams(value).toString());\n    } else {\n      urlTokenParams.set(key, String(value));\n    }\n  }\n  return urlTokenParams;\n}\n\ninterface FetchInput {\n  url: string;\n  options: RequestInit;\n}\n\nexport interface StreamInput {\n  url: string;\n  options: FetchEventSourceInit;\n}\n\nexport interface RestfulClientEvents {\n  fetchStart: Callback<{ input: FetchInput }>;\n  fetchError: Callback<{ error: Error; input: FetchInput }>;\n  fetchSuccess: Callback<{ response: Response; data: any; input: FetchInput }>;\n  fetchDone: Callback<{ input: FetchInput }>;\n\n  streamStart: Callback<{ input: StreamInput }>;\n  streamOpen: Callback<{ input: StreamInput }>;\n  streamSuccess: Callback<{ input: StreamInput }>;\n  streamMessage: Callback<{ data: EventSourceMessage; input: StreamInput }>;\n  streamError: Callback<{ error: Error; input: StreamInput }>;\n  streamDone: Callback<{ input: StreamInput }>;\n}\n\nexport class RestfulClient<K extends Record<string, string>> extends Serializable {\n  public readonly emitter = Emitter.root.child<RestfulClientEvents>({\n    namespace: [\"internals\", \"restfulClient\"],\n    creator: this,\n  });\n\n  constructor(\n    protected input: {\n      baseUrl: string;\n      headers?: () => Promise<HeadersInit>;\n      paths: K;\n    },\n  ) {\n    super();\n  }\n\n  async *stream(\n    path: keyof K,\n    init: FetchEventSourceInit,\n  ): AsyncGenerator<EventSourceMessage, void, void> {\n    const emitter = this.emitter.child({\n      groupId: \"stream\",\n    });\n\n    const input: StreamInput = {\n      url: this.getUrl(path).toString(),\n      options: {\n        method: \"POST\",\n        ...init,\n        headers: await this.getHeaders(init?.headers),\n      },\n    };\n    await emitter.emit(\"streamStart\", { input });\n\n    return yield* emitterToGenerator(async ({ emit }) =>\n      fetchEventSource(input.url, {\n        ...input.options,\n        async onopen(response) {\n          const contentType = response.headers.get(\"content-type\") || \"\";\n          if (response.ok && contentType.includes(EventStreamContentType)) {\n            await emitter.emit(\"streamOpen\", { input });\n            return;\n          }\n          throw new RestfulClientError(\"Failed to stream!\", [], {\n            context: {\n              url: response.url,\n              err: await response.text(),\n              response,\n            },\n            isRetryable: response.status >= 400 && response.status < 500 && response.status !== 429,\n          });\n        },\n        async onmessage(msg) {\n          if (msg?.event === \"error\") {\n            throw new RestfulClientError(`Error during streaming has occurred.`, [], {\n              context: msg,\n            });\n          }\n          await emitter.emit(\"streamMessage\", { input, data: msg });\n          emit(msg);\n        },\n        onclose() {},\n        onerror(err) {\n          throw new RestfulClientError(`Error during streaming has occurred.`, [err]);\n        },\n      })\n        .then(() => emitter.emit(\"streamSuccess\", { input }))\n        .catch(async (error) => {\n          await emitter.emit(\"streamError\", { input, error }).catch(doNothing());\n          throw error;\n        })\n        .finally(() => emitter.emit(\"streamDone\", { input })),\n    );\n  }\n\n  async fetch(path: keyof K, init?: RequestInit & { searchParams?: URLSearchParams }) {\n    const emitter = this.emitter.child({\n      groupId: \"fetch\",\n    });\n\n    const target = this.getUrl(path);\n    if (init?.searchParams) {\n      for (const [key, value] of init.searchParams) {\n        target.searchParams.set(key, value);\n      }\n    }\n\n    const input: FetchInput = {\n      url: target.toString(),\n      options: {\n        ...init,\n        headers: await this.getHeaders(init?.headers),\n      },\n    };\n\n    await emitter.emit(\"fetchStart\", { input });\n    try {\n      const response = await fetch(input.url, input.options);\n\n      if (!response.ok) {\n        throw new RestfulClientError(\"Fetch has failed\", [], {\n          context: {\n            url: response.url,\n            error: await response.text(),\n            response,\n          },\n          isRetryable: [408, 503].includes(response.status ?? 500),\n        });\n      }\n\n      const data = await response.json();\n      await emitter.emit(\"fetchSuccess\", { response, data, input });\n      return data;\n    } catch (error) {\n      await emitter.emit(\"fetchError\", { error, input: input });\n      throw error;\n    } finally {\n      await emitter.emit(\"fetchDone\", { input: input });\n    }\n  }\n\n  protected async getHeaders(extra?: HeadersInit): Promise<Record<string, string>> {\n    const final = {};\n    for (const override of [await this.input.headers?.(), extra].filter(isTruthy)) {\n      if (isArray(override) || override instanceof Headers) {\n        Object.assign(final, Object.fromEntries(override.entries()));\n      } else {\n        Object.assign(final, override);\n      }\n    }\n    return final;\n  }\n\n  protected getUrl(path: keyof K): URL {\n    const url = new URL(this.input.baseUrl);\n    const extraPath = this.input.paths[path] ?? path;\n    if (url.pathname.endsWith(\"/\")) {\n      url.pathname += extraPath.replace(/^\\//, \"\");\n    } else {\n      url.pathname += extraPath;\n    }\n    return url;\n  }\n\n  createSnapshot() {\n    return {\n      input: shallowCopy(this.input),\n      emitter: this.emitter,\n    };\n  }\n\n  loadSnapshot(snapshot: ReturnType<typeof this.createSnapshot>): void {\n    Object.assign(this, snapshot);\n  }\n}\n"]}