{"version":3,"file":"streamedQuery.cjs","names":["addToEnd","addConsumeAwareSignal"],"sources":["../../src/streamedQuery.ts"],"sourcesContent":["import { addConsumeAwareSignal, addToEnd } from './utils'\nimport type {\n  OmitKeyof,\n  QueryFunction,\n  QueryFunctionContext,\n  QueryKey,\n} from './types'\n\ntype BaseStreamedQueryParams<TQueryFnData, TQueryKey extends QueryKey> = {\n  /** The function that returns an `AsyncIterable` to stream data from. */\n  streamFn: (\n    context: QueryFunctionContext<TQueryKey>,\n  ) => AsyncIterable<TQueryFnData> | Promise<AsyncIterable<TQueryFnData>>\n  /**\n   * Defines how refetches are handled.\n   * - `'reset'` (default): erases all data and puts the query back into `pending` state.\n   * - `'append'`: appends new data to the existing data.\n   * - `'replace'`: writes all data to the cache once the stream ends.\n   */\n  refetchMode?: 'append' | 'reset' | 'replace'\n}\n\ntype SimpleStreamedQueryParams<\n  TQueryFnData,\n  TQueryKey extends QueryKey,\n> = BaseStreamedQueryParams<TQueryFnData, TQueryKey> & {\n  reducer?: never\n  initialValue?: never\n}\n\ntype ReducibleStreamedQueryParams<\n  TQueryFnData,\n  TData,\n  TQueryKey extends QueryKey,\n> = BaseStreamedQueryParams<TQueryFnData, TQueryKey> & {\n  /**\n   * Reduces streamed chunks into the final data shape. Required whenever `TData` is not an\n   * array, since there is no default way to accumulate non-array chunks.\n   */\n  reducer: (acc: TData, chunk: TQueryFnData) => TData\n  /**\n   * The value used while the first chunk is being fetched, and returned if the stream yields no\n   * values. Required together with a custom `reducer`.\n   */\n  initialValue: TData\n}\n\ntype StreamedQueryParams<TQueryFnData, TData, TQueryKey extends QueryKey> =\n  | SimpleStreamedQueryParams<TQueryFnData, TQueryKey>\n  | ReducibleStreamedQueryParams<TQueryFnData, TData, TQueryKey>\n\n/**\n * This is a helper function to create a query function that streams data from an AsyncIterable.\n * Data will be an Array of all the chunks received.\n * The query will be in a 'pending' state until the first chunk of data is received, but will go to 'success' after that.\n * The query will stay in fetchStatus 'fetching' until the stream ends.\n * @param streamFn - The function that returns an AsyncIterable to stream data from.\n * @example\n * ```ts\n * await queryClient.query({\n *   queryKey: ['data'],\n *   queryFn: streamedQuery({\n *     streamFn: fetchDataInChunks,\n *   }),\n * })\n * ```\n */\nexport function streamedQuery<\n  TQueryFnData = unknown,\n  TData = Array<TQueryFnData>,\n  TQueryKey extends QueryKey = QueryKey,\n>({\n  streamFn,\n  refetchMode = 'reset',\n  reducer = (items, chunk) =>\n    addToEnd(items as Array<TQueryFnData>, chunk) as TData,\n  initialValue = [] as TData,\n}: StreamedQueryParams<TQueryFnData, TData, TQueryKey>): QueryFunction<\n  TData,\n  TQueryKey\n> {\n  return async (context) => {\n    const query = context.client\n      .getQueryCache()\n      .find({ queryKey: context.queryKey, exact: true })\n    const isRefetch = !!query && query.isFetched()\n    if (isRefetch && refetchMode === 'reset') {\n      query.setState({\n        ...query.resetState,\n        fetchStatus: 'fetching',\n      })\n    }\n\n    let result = initialValue\n\n    // eslint-disable-next-line @typescript-eslint/no-unnecessary-type-assertion\n    let cancelled: boolean = false as boolean\n    const streamFnContext = addConsumeAwareSignal<\n      OmitKeyof<typeof context, 'signal'>\n    >(\n      {\n        client: context.client,\n        meta: context.meta,\n        queryKey: context.queryKey,\n        pageParam: context.pageParam,\n        direction: context.direction,\n      },\n      () => context.signal,\n      () => (cancelled = true),\n    )\n\n    const stream = await streamFn(streamFnContext)\n\n    const isReplaceRefetch = isRefetch && refetchMode === 'replace'\n\n    for await (const chunk of stream) {\n      if (cancelled) {\n        break\n      }\n\n      if (isReplaceRefetch) {\n        // don't append to the cache directly when replace-refetching\n        result = reducer(result, chunk)\n      } else {\n        context.client.setQueryData<TData>(context.queryKey, (prev) =>\n          reducer(prev === undefined ? initialValue : prev, chunk),\n        )\n      }\n    }\n\n    // finalize result: replace-refetching needs to write to the cache\n    if (isReplaceRefetch && !cancelled) {\n      context.client.setQueryData<TData>(context.queryKey, result)\n    }\n\n    const data = context.client.getQueryData<TData>(context.queryKey)\n    return data === undefined ? initialValue : data\n  }\n}\n"],"mappings":";;;;;;;;;;;;;;;;;;;AAmEA,SAAgB,cAId,EACA,UACA,cAAc,SACd,WAAW,OAAO,UAChBA,cAAAA,SAAS,OAA8B,KAAK,GAC9C,eAAe,CAAC,KAIhB;CACA,OAAO,OAAO,YAAY;EACxB,MAAM,QAAQ,QAAQ,OACnB,cAAc,CAAC,CACf,KAAK;GAAE,UAAU,QAAQ;GAAU,OAAO;EAAK,CAAC;EACnD,MAAM,YAAY,CAAC,CAAC,SAAS,MAAM,UAAU;EAC7C,IAAI,aAAa,gBAAgB,SAC/B,MAAM,SAAS;GACb,GAAG,MAAM;GACT,aAAa;EACf,CAAC;EAGH,IAAI,SAAS;EAGb,IAAI,YAAqB;EAezB,MAAM,SAAS,MAAM,SAdGC,cAAAA,sBAGtB;GACE,QAAQ,QAAQ;GAChB,MAAM,QAAQ;GACd,UAAU,QAAQ;GAClB,WAAW,QAAQ;GACnB,WAAW,QAAQ;EACrB,SACM,QAAQ,cACP,YAAY,IAGuB,CAAC;EAE7C,MAAM,mBAAmB,aAAa,gBAAgB;EAEtD,WAAW,MAAM,SAAS,QAAQ;GAChC,IAAI,WACF;GAGF,IAAI,kBAEF,SAAS,QAAQ,QAAQ,KAAK;QAE9B,QAAQ,OAAO,aAAoB,QAAQ,WAAW,SACpD,QAAQ,SAAS,KAAA,IAAY,eAAe,MAAM,KAAK,CACzD;EAEJ;EAGA,IAAI,oBAAoB,CAAC,WACvB,QAAQ,OAAO,aAAoB,QAAQ,UAAU,MAAM;EAG7D,MAAM,OAAO,QAAQ,OAAO,aAAoB,QAAQ,QAAQ;EAChE,OAAO,SAAS,KAAA,IAAY,eAAe;CAC7C;AACF"}