{"version":3,"file":"utils.cjs","sources":["../../../../src/query/live/utils.ts"],"sourcesContent":["import { MultiSet } from '@tanstack/db-ivm'\nimport { UnsupportedRootScalarSelectError } from '../../errors.js'\nimport { normalizeOrderByPaths } from '../compiler/expressions.js'\nimport { buildQuery, getQueryIR } from '../builder/index.js'\nimport { collectCollectionSources, isExpressionLike } from '../ir.js'\nimport type { MultiSetArray, RootStreamBuilder } from '@tanstack/db-ivm'\nimport type { Collection } from '../../collection/index.js'\nimport type { ChangeMessage } from '../../types.js'\nimport type { InitialQueryBuilder, QueryBuilder } from '../builder/index.js'\nimport type { Context } from '../builder/types.js'\nimport type { OrderBy, QueryIR } from '../ir.js'\n\n/**\n * Helper function to extract collections from a compiled query.\n * Traverses the query IR to find all collection references.\n * Maps collections by their ID (not alias) as expected by the compiler.\n */\nexport function extractCollectionsFromQuery(\n  query: QueryIR,\n): Record<string, Collection<any, any, any>> {\n  const collections: Record<string, Collection<any, any, any>> = {}\n  for (const source of collectCollectionSources(query)) {\n    collections[source.collection.id] = source.collection\n  }\n  return collections\n}\n\nexport { collectCollectionSources as extractCollectionSources }\n\n/**\n * Helper function to extract the collection that is referenced in the query's FROM clause.\n * The FROM clause may refer directly to a collection or indirectly to a subquery.\n */\nexport function extractCollectionFromSource(\n  query: any,\n): Collection<any, any, any> {\n  const from = query.from\n\n  if (from.type === `collectionRef`) {\n    return from.collection\n  } else if (from.type === `queryRef`) {\n    // Recursively extract from subquery\n    return extractCollectionFromSource(from.query)\n  } else if (from.type === `unionFrom`) {\n    return extractCollectionFromSource({ from: from.sources[0] })\n  } else if (from.type === `unionAll`) {\n    return extractCollectionFromSource(from.queries[0])\n  }\n\n  throw new Error(\n    `Failed to extract collection. Invalid FROM clause: ${JSON.stringify(query)}`,\n  )\n}\n\n/**\n * Check if a value is a nested select object (plain object, not an expression)\n */\nfunction isNestedSelectObject(obj: any): boolean {\n  if (obj === null || typeof obj !== `object`) return false\n  if (isExpressionLike(obj)) return false\n  // Ref proxies from spread operations\n  if (obj.__refProxy) return false\n  return true\n}\n\n/**\n * Builds a query IR from a config object that contains either a query builder\n * function or a QueryBuilder instance.\n */\nexport function buildQueryFromConfig<TContext extends Context>(config: {\n  query:\n    | ((q: InitialQueryBuilder) => QueryBuilder<TContext>)\n    | QueryBuilder<TContext>\n  requireObjectResult?: boolean\n}): QueryIR {\n  // Build the query using the provided query builder function or instance\n  const query =\n    typeof config.query === `function`\n      ? buildQuery<TContext>(config.query)\n      : getQueryIR(config.query)\n\n  if (\n    config.requireObjectResult &&\n    query.select &&\n    !isNestedSelectObject(query.select)\n  ) {\n    throw new UnsupportedRootScalarSelectError()\n  }\n\n  return query\n}\n\n/**\n * Helper function to send changes to a D2 input stream.\n * Converts ChangeMessages to D2 MultiSet data and sends to the input.\n *\n * @returns The number of multiset entries sent\n */\nexport function sendChangesToInput(\n  input: RootStreamBuilder<unknown>,\n  changes: Iterable<ChangeMessage>,\n): number {\n  const multiSetArray: MultiSetArray<unknown> = []\n  for (const change of changes) {\n    const key = change.key\n    if (change.type === `insert`) {\n      multiSetArray.push([[key, change.value], 1])\n    } else if (change.type === `update`) {\n      multiSetArray.push([[key, change.previousValue], -1])\n      multiSetArray.push([[key, change.value], 1])\n    } else {\n      // change.type === `delete`\n      multiSetArray.push([[key, change.value], -1])\n    }\n  }\n\n  if (multiSetArray.length !== 0) {\n    input.sendData(new MultiSet(multiSetArray))\n  }\n\n  return multiSetArray.length\n}\n\n/** Splits updates into a delete of the old value and an insert of the new value */\nexport function* splitUpdates<\n  T extends object = Record<string, unknown>,\n  TKey extends string | number = string | number,\n>(\n  changes: Iterable<ChangeMessage<T, TKey>>,\n): Generator<ChangeMessage<T, TKey>> {\n  for (const change of changes) {\n    if (change.type === `update`) {\n      yield { type: `delete`, key: change.key, value: change.previousValue! }\n      yield { type: `insert`, key: change.key, value: change.value }\n    } else {\n      yield change\n    }\n  }\n}\n\n/** Keep each source key at one exact D2 contribution. */\nexport function reconcileChangesForD2<\n  T extends object,\n  TKey extends string | number,\n>(\n  changes: Array<ChangeMessage<T, TKey>>,\n  sentRows: Map<TKey, T>,\n): Array<ChangeMessage<T, TKey>> {\n  const reconciled: Array<ChangeMessage<T, TKey>> = []\n  for (const change of changes) {\n    const previousValue = sentRows.get(change.key)\n    if (change.type === `insert`) {\n      if (previousValue !== undefined) continue\n      sentRows.set(change.key, change.value)\n      reconciled.push(change)\n    } else if (change.type === `delete`) {\n      if (previousValue === undefined) continue\n      sentRows.delete(change.key)\n      reconciled.push({ ...change, value: previousValue })\n    } else {\n      sentRows.set(change.key, change.value)\n      reconciled.push(\n        previousValue === undefined\n          ? { type: `insert`, key: change.key, value: change.value }\n          : { ...change, previousValue },\n      )\n    }\n  }\n  return reconciled\n}\n\n/**\n * Compute orderBy/limit subscription hints for an alias.\n * Returns normalised orderBy and effective limit suitable for passing to\n * `subscribeChanges`, or `undefined` values when the query's orderBy cannot\n * be scoped to the given alias (e.g. cross-collection refs or aggregates).\n */\nexport function computeSubscriptionOrderByHints(\n  query: { orderBy?: OrderBy; limit?: number; offset?: number },\n  alias: string,\n): { orderBy: OrderBy | undefined; limit: number | undefined } {\n  const { orderBy, limit, offset } = query\n  const effectiveLimit =\n    limit !== undefined && offset !== undefined ? limit + offset : limit\n\n  const normalizedOrderBy = orderBy\n    ? normalizeOrderByPaths(orderBy, alias)\n    : undefined\n\n  // Only pass orderBy when it is scoped to this alias and uses simple refs,\n  // to avoid leaking cross-collection paths into backend-specific compilers.\n  const canPassOrderBy =\n    normalizedOrderBy?.every((clause) => {\n      const exp = clause.expression\n      if (exp.type !== `ref`) return false\n      const path = exp.path\n      return Array.isArray(path) && path.length === 1\n    }) ?? false\n\n  return {\n    orderBy: canPassOrderBy ? normalizedOrderBy : undefined,\n    limit: canPassOrderBy ? effectiveLimit : undefined,\n  }\n}\n"],"names":["collectCollectionSources","isExpressionLike","buildQuery","getQueryIR","UnsupportedRootScalarSelectError","MultiSet","normalizeOrderByPaths"],"mappings":";;;;;;;;AAiBO,SAAS,4BACd,OAC2C;AAC3C,QAAM,cAAyD,CAAA;AAC/D,aAAW,UAAUA,4BAAyB,KAAK,GAAG;AACpD,gBAAY,OAAO,WAAW,EAAE,IAAI,OAAO;AAAA,EAC7C;AACA,SAAO;AACT;AAQO,SAAS,4BACd,OAC2B;AAC3B,QAAM,OAAO,MAAM;AAEnB,MAAI,KAAK,SAAS,iBAAiB;AACjC,WAAO,KAAK;AAAA,EACd,WAAW,KAAK,SAAS,YAAY;AAEnC,WAAO,4BAA4B,KAAK,KAAK;AAAA,EAC/C,WAAW,KAAK,SAAS,aAAa;AACpC,WAAO,4BAA4B,EAAE,MAAM,KAAK,QAAQ,CAAC,GAAG;AAAA,EAC9D,WAAW,KAAK,SAAS,YAAY;AACnC,WAAO,4BAA4B,KAAK,QAAQ,CAAC,CAAC;AAAA,EACpD;AAEA,QAAM,IAAI;AAAA,IACR,sDAAsD,KAAK,UAAU,KAAK,CAAC;AAAA,EAAA;AAE/E;AAKA,SAAS,qBAAqB,KAAmB;AAC/C,MAAI,QAAQ,QAAQ,OAAO,QAAQ,SAAU,QAAO;AACpD,MAAIC,GAAAA,iBAAiB,GAAG,EAAG,QAAO;AAElC,MAAI,IAAI,WAAY,QAAO;AAC3B,SAAO;AACT;AAMO,SAAS,qBAA+C,QAKnD;AAEV,QAAM,QACJ,OAAO,OAAO,UAAU,aACpBC,iBAAqB,OAAO,KAAK,IACjCC,mBAAW,OAAO,KAAK;AAE7B,MACE,OAAO,uBACP,MAAM,UACN,CAAC,qBAAqB,MAAM,MAAM,GAClC;AACA,UAAM,IAAIC,OAAAA,iCAAA;AAAA,EACZ;AAEA,SAAO;AACT;AAQO,SAAS,mBACd,OACA,SACQ;AACR,QAAM,gBAAwC,CAAA;AAC9C,aAAW,UAAU,SAAS;AAC5B,UAAM,MAAM,OAAO;AACnB,QAAI,OAAO,SAAS,UAAU;AAC5B,oBAAc,KAAK,CAAC,CAAC,KAAK,OAAO,KAAK,GAAG,CAAC,CAAC;AAAA,IAC7C,WAAW,OAAO,SAAS,UAAU;AACnC,oBAAc,KAAK,CAAC,CAAC,KAAK,OAAO,aAAa,GAAG,EAAE,CAAC;AACpD,oBAAc,KAAK,CAAC,CAAC,KAAK,OAAO,KAAK,GAAG,CAAC,CAAC;AAAA,IAC7C,OAAO;AAEL,oBAAc,KAAK,CAAC,CAAC,KAAK,OAAO,KAAK,GAAG,EAAE,CAAC;AAAA,IAC9C;AAAA,EACF;AAEA,MAAI,cAAc,WAAW,GAAG;AAC9B,UAAM,SAAS,IAAIC,MAAAA,SAAS,aAAa,CAAC;AAAA,EAC5C;AAEA,SAAO,cAAc;AACvB;AAGO,UAAU,aAIf,SACmC;AACnC,aAAW,UAAU,SAAS;AAC5B,QAAI,OAAO,SAAS,UAAU;AAC5B,YAAM,EAAE,MAAM,UAAU,KAAK,OAAO,KAAK,OAAO,OAAO,cAAA;AACvD,YAAM,EAAE,MAAM,UAAU,KAAK,OAAO,KAAK,OAAO,OAAO,MAAA;AAAA,IACzD,OAAO;AACL,YAAM;AAAA,IACR;AAAA,EACF;AACF;AAGO,SAAS,sBAId,SACA,UAC+B;AAC/B,QAAM,aAA4C,CAAA;AAClD,aAAW,UAAU,SAAS;AAC5B,UAAM,gBAAgB,SAAS,IAAI,OAAO,GAAG;AAC7C,QAAI,OAAO,SAAS,UAAU;AAC5B,UAAI,kBAAkB,OAAW;AACjC,eAAS,IAAI,OAAO,KAAK,OAAO,KAAK;AACrC,iBAAW,KAAK,MAAM;AAAA,IACxB,WAAW,OAAO,SAAS,UAAU;AACnC,UAAI,kBAAkB,OAAW;AACjC,eAAS,OAAO,OAAO,GAAG;AAC1B,iBAAW,KAAK,EAAE,GAAG,QAAQ,OAAO,eAAe;AAAA,IACrD,OAAO;AACL,eAAS,IAAI,OAAO,KAAK,OAAO,KAAK;AACrC,iBAAW;AAAA,QACT,kBAAkB,SACd,EAAE,MAAM,UAAU,KAAK,OAAO,KAAK,OAAO,OAAO,MAAA,IACjD,EAAE,GAAG,QAAQ,cAAA;AAAA,MAAc;AAAA,IAEnC;AAAA,EACF;AACA,SAAO;AACT;AAQO,SAAS,gCACd,OACA,OAC6D;AAC7D,QAAM,EAAE,SAAS,OAAO,OAAA,IAAW;AACnC,QAAM,iBACJ,UAAU,UAAa,WAAW,SAAY,QAAQ,SAAS;AAEjE,QAAM,oBAAoB,UACtBC,YAAAA,sBAAsB,SAAS,KAAK,IACpC;AAIJ,QAAM,iBACJ,mBAAmB,MAAM,CAAC,WAAW;AACnC,UAAM,MAAM,OAAO;AACnB,QAAI,IAAI,SAAS,MAAO,QAAO;AAC/B,UAAM,OAAO,IAAI;AACjB,WAAO,MAAM,QAAQ,IAAI,KAAK,KAAK,WAAW;AAAA,EAChD,CAAC,KAAK;AAER,SAAO;AAAA,IACL,SAAS,iBAAiB,oBAAoB;AAAA,IAC9C,OAAO,iBAAiB,iBAAiB;AAAA,EAAA;AAE7C;;;;;;;;;"}