{"version":3,"file":"subset-demand-controller.cjs","sources":["../../../../src/query/live/subset-demand-controller.ts"],"sourcesContent":["import { inArray } from '../builder/functions.js'\nimport { PropRef } from '../ir.js'\nimport { createValueIdentity } from '../equality-value-identity.js'\nimport type { ValueIdentity } from '../equality-value-identity.js'\nimport type { CollectionSubscription } from '../../collection/subscription.js'\nimport type {\n  LazyCollectionCallbacks,\n  LazyDemandPlan,\n} from '../compiler/joins.js'\nimport type { BasicExpression } from '../ir.js'\nimport type { LoadSubsetRequestResult } from '../../types.js'\n\ntype DemandSegment = {\n  keys: Map<string, unknown>\n  where: BasicExpression<boolean>\n  abortController: AbortController\n  ready: LoadSubsetRequestResult\n  state: `starting` | `pending` | `settled` | `failed`\n  replaces?: Array<DemandSegment>\n}\n\ntype DemandState = {\n  keys: Map<string, unknown>\n  segments: Array<DemandSegment>\n}\n\nexport type DemandUpdate = {\n  changed: boolean\n  empty: boolean\n  ready: Promise<Array<unknown>> | true\n}\n\n/**\n * Keeps lazy subset requests aligned with the current relation of demanded\n * keys. Growth loads only new keys. Churn replaces fragmented coverage after\n * the replacement applies, so prior coverage remains live in the meantime.\n */\nexport class SubsetDemandController {\n  private readonly states = new Map<string, DemandState>()\n  private readonly warnedPlans = new Set<string>()\n  private valueIdentity = createValueIdentity()\n\n  setDemand(\n    subscription: CollectionSubscription,\n    plan: LazyDemandPlan,\n    keys: Set<unknown>,\n  ): DemandUpdate {\n    const nextKeys = canonicalizeKeys(keys, this.valueIdentity)\n    const previous = this.states.get(plan.id)\n    const hasFailedCoverage = previous?.segments.some(\n      (segment) =>\n        segment.state === `failed` && intersects(segment.keys, nextKeys),\n    )\n    if (\n      previous &&\n      equalKeySets(previous.keys, nextKeys) &&\n      !hasFailedCoverage\n    ) {\n      return { changed: false, empty: nextKeys.size === 0, ready: true }\n    }\n\n    const segments = (previous?.segments ?? []).filter(\n      (segment) =>\n        segment.state !== `failed` &&\n        intersects(segment.keys, nextKeys) &&\n        (!segment.replaces || equalKeySets(segment.keys, nextKeys)),\n    )\n    const retired = (previous?.segments ?? []).filter(\n      (segment) => !segments.includes(segment),\n    )\n    const retryReplacement = previous?.segments.some(\n      (segment) =>\n        segment.replaces &&\n        segment.state === `failed` &&\n        equalKeySets(segment.keys, nextKeys),\n    )\n    const state: DemandState = { keys: nextKeys, segments }\n    this.states.set(plan.id, state)\n\n    for (const segment of retired) releaseSegment(subscription, segment)\n    if (this.states.get(plan.id) !== state) {\n      return { changed: false, empty: false, ready: true }\n    }\n\n    if (nextKeys.size === 0) {\n      this.states.delete(plan.id)\n      return { changed: true, empty: true, ready: true }\n    }\n\n    const coveredKeys = new Set(\n      segments.flatMap((segment) => [...segment.keys.keys()]),\n    )\n    const added = new Map(\n      [...nextKeys].filter(([key]) => !coveredKeys.has(key)),\n    )\n    const removed = previous\n      ? [...previous.keys.keys()].some((key) => !nextKeys.has(key))\n      : false\n    const replace =\n      retryReplacement || (removed && (added.size > 0 || segments.length > 1))\n    const requestedKeys = replace ? nextKeys : added\n    if (requestedKeys.size > 0) {\n      const segment = createSegment(plan, requestedKeys)\n      if (replace) segment.replaces = [...segments]\n      segments.push(segment)\n      const started = startSegment(subscription, segment, () =>\n        this.warnUnoptimized(plan),\n      )\n      if (this.states.get(plan.id) !== state) {\n        return { changed: false, empty: false, ready: true }\n      }\n      if (!started) throw new Error(`Subset demand snapshot did not start`)\n      if (replace) {\n        if (segment.ready instanceof Promise) {\n          void segment.ready.then(\n            () => this.finishReplacement(subscription, plan.id, segment),\n            () => {},\n          )\n        } else {\n          this.finishReplacement(subscription, plan.id, segment)\n        }\n      }\n    }\n\n    const pending = state.segments\n      .filter(\n        (segment) =>\n          segment.state === `pending` && intersects(segment.keys, nextKeys),\n      )\n      .map((segment) => segment.ready)\n      .filter((ready): ready is Promise<void> => ready instanceof Promise)\n    return {\n      changed: true,\n      empty: false,\n      ready: pending.length > 0 ? Promise.all(pending) : true,\n    }\n  }\n\n  clear(): void {\n    for (const state of this.states.values()) {\n      for (const segment of state.segments) segment.abortController.abort()\n    }\n    this.states.clear()\n    this.warnedPlans.clear()\n    this.valueIdentity = createValueIdentity()\n  }\n\n  hasPendingDemand(planId: string): boolean {\n    return (\n      this.states\n        .get(planId)\n        ?.segments.some((segment) => segment.state === `pending`) ?? false\n    )\n  }\n\n  private finishReplacement(\n    subscription: CollectionSubscription,\n    planId: string,\n    replacement: DemandSegment,\n  ): void {\n    const state = this.states.get(planId)\n    if (\n      !state?.segments.includes(replacement) ||\n      !equalKeySets(replacement.keys, state.keys)\n    ) {\n      releaseSegment(subscription, replacement)\n      return\n    }\n    const replaced = replacement.replaces ?? []\n    replacement.replaces = undefined\n    state.segments = [replacement]\n    for (const segment of replaced) releaseSegment(subscription, segment)\n  }\n\n  private warnUnoptimized(plan: LazyDemandPlan): void {\n    if (this.warnedPlans.has(plan.id)) return\n    this.warnedPlans.add(plan.id)\n    const path = plan.path.join(`.`)\n    console.warn(\n      `[TanStack DB]${plan.collectionId ? ` [${plan.collectionId}]` : ``} Join requires an index on \"${path}\" for efficient loading. ` +\n        `Falling back to scanning local data. ` +\n        `Consider creating an index on the collection with collection.createIndex((row) => row.${path}) ` +\n        `or enable auto-indexing with autoIndex: 'eager' and a defaultIndexType.`,\n    )\n  }\n}\n\n/** An ordered filter waits only for demand on its joined source. */\nexport function hasPendingJoinedWork(\n  joinedSourceId: string,\n  lazySourcesCallbacks: Record<string, LazyCollectionCallbacks>,\n  subscriptions: Record<string, CollectionSubscription>,\n  isPendingPlan: (planId: string) => boolean,\n): boolean {\n  return (\n    (lazySourcesCallbacks[joinedSourceId]?.plans?.some((plan) =>\n      isPendingPlan(plan.id),\n    ) ??\n      false) ||\n    subscriptions[joinedSourceId]?.status === `loadingSubset`\n  )\n}\n\nfunction canonicalizeKeys(\n  keys: Set<unknown>,\n  valueIdentity: ValueIdentity,\n): Map<string, unknown> {\n  return new Map(\n    [...keys].map((key) => [valueIdentity.serializeEquality(key), key]),\n  )\n}\n\nfunction equalKeySets(\n  left: Map<string, unknown>,\n  right: Map<string, unknown>,\n): boolean {\n  return (\n    left.size === right.size && [...left.keys()].every((key) => right.has(key))\n  )\n}\n\nfunction intersects(\n  left: Map<string, unknown>,\n  right: Map<string, unknown>,\n): boolean {\n  return [...left.keys()].some((key) => right.has(key))\n}\n\nfunction createSegment(\n  plan: LazyDemandPlan,\n  keys: Map<string, unknown>,\n): DemandSegment {\n  const where = inArray(new PropRef(plan.path), [...keys.values()])\n  return {\n    keys,\n    where,\n    abortController: new AbortController(),\n    ready: true,\n    state: `starting`,\n  }\n}\n\nfunction startSegment(\n  subscription: CollectionSubscription,\n  segment: DemandSegment,\n  onUnoptimized: () => void,\n): boolean {\n  let observed = false\n  const requested = subscription.requestSnapshot({\n    where: segment.where,\n    signal: segment.abortController.signal,\n    trackLoadSubsetPromise: false,\n    onUnoptimized,\n    onLoadSubsetResult: (result) => {\n      observed = true\n      segment.ready = result\n      segment.state = result instanceof Promise ? `pending` : `settled`\n    },\n  })\n  if (segment.ready instanceof Promise) {\n    void segment.ready.then(\n      () => {\n        segment.state = `settled`\n      },\n      () => {\n        segment.state = `failed`\n      },\n    )\n  }\n  return requested && observed\n}\n\nfunction releaseSegment(\n  subscription: CollectionSubscription,\n  segment: DemandSegment,\n): void {\n  segment.abortController.abort()\n  try {\n    subscription.releaseSnapshot(segment.where)\n  } catch {\n    // CollectionSubscription has already reported the one cleanup failure.\n  }\n}\n"],"names":["createValueIdentity","inArray","PropRef"],"mappings":";;;;;AAqCO,MAAM,uBAAuB;AAAA,EAA7B,cAAA;AACL,SAAiB,6BAAa,IAAA;AAC9B,SAAiB,kCAAkB,IAAA;AACnC,SAAQ,gBAAgBA,0CAAA;AAAA,EAAoB;AAAA,EAE5C,UACE,cACA,MACA,MACc;AACd,UAAM,WAAW,iBAAiB,MAAM,KAAK,aAAa;AAC1D,UAAM,WAAW,KAAK,OAAO,IAAI,KAAK,EAAE;AACxC,UAAM,oBAAoB,UAAU,SAAS;AAAA,MAC3C,CAAC,YACC,QAAQ,UAAU,YAAY,WAAW,QAAQ,MAAM,QAAQ;AAAA,IAAA;AAEnE,QACE,YACA,aAAa,SAAS,MAAM,QAAQ,KACpC,CAAC,mBACD;AACA,aAAO,EAAE,SAAS,OAAO,OAAO,SAAS,SAAS,GAAG,OAAO,KAAA;AAAA,IAC9D;AAEA,UAAM,YAAY,UAAU,YAAY,CAAA,GAAI;AAAA,MAC1C,CAAC,YACC,QAAQ,UAAU,YAClB,WAAW,QAAQ,MAAM,QAAQ,MAChC,CAAC,QAAQ,YAAY,aAAa,QAAQ,MAAM,QAAQ;AAAA,IAAA;AAE7D,UAAM,WAAW,UAAU,YAAY,CAAA,GAAI;AAAA,MACzC,CAAC,YAAY,CAAC,SAAS,SAAS,OAAO;AAAA,IAAA;AAEzC,UAAM,mBAAmB,UAAU,SAAS;AAAA,MAC1C,CAAC,YACC,QAAQ,YACR,QAAQ,UAAU,YAClB,aAAa,QAAQ,MAAM,QAAQ;AAAA,IAAA;AAEvC,UAAM,QAAqB,EAAE,MAAM,UAAU,SAAA;AAC7C,SAAK,OAAO,IAAI,KAAK,IAAI,KAAK;AAE9B,eAAW,WAAW,QAAS,gBAAe,cAAc,OAAO;AACnE,QAAI,KAAK,OAAO,IAAI,KAAK,EAAE,MAAM,OAAO;AACtC,aAAO,EAAE,SAAS,OAAO,OAAO,OAAO,OAAO,KAAA;AAAA,IAChD;AAEA,QAAI,SAAS,SAAS,GAAG;AACvB,WAAK,OAAO,OAAO,KAAK,EAAE;AAC1B,aAAO,EAAE,SAAS,MAAM,OAAO,MAAM,OAAO,KAAA;AAAA,IAC9C;AAEA,UAAM,cAAc,IAAI;AAAA,MACtB,SAAS,QAAQ,CAAC,YAAY,CAAC,GAAG,QAAQ,KAAK,MAAM,CAAC;AAAA,IAAA;AAExD,UAAM,QAAQ,IAAI;AAAA,MAChB,CAAC,GAAG,QAAQ,EAAE,OAAO,CAAC,CAAC,GAAG,MAAM,CAAC,YAAY,IAAI,GAAG,CAAC;AAAA,IAAA;AAEvD,UAAM,UAAU,WACZ,CAAC,GAAG,SAAS,KAAK,KAAA,CAAM,EAAE,KAAK,CAAC,QAAQ,CAAC,SAAS,IAAI,GAAG,CAAC,IAC1D;AACJ,UAAM,UACJ,oBAAqB,YAAY,MAAM,OAAO,KAAK,SAAS,SAAS;AACvE,UAAM,gBAAgB,UAAU,WAAW;AAC3C,QAAI,cAAc,OAAO,GAAG;AAC1B,YAAM,UAAU,cAAc,MAAM,aAAa;AACjD,UAAI,QAAS,SAAQ,WAAW,CAAC,GAAG,QAAQ;AAC5C,eAAS,KAAK,OAAO;AACrB,YAAM,UAAU;AAAA,QAAa;AAAA,QAAc;AAAA,QAAS,MAClD,KAAK,gBAAgB,IAAI;AAAA,MAAA;AAE3B,UAAI,KAAK,OAAO,IAAI,KAAK,EAAE,MAAM,OAAO;AACtC,eAAO,EAAE,SAAS,OAAO,OAAO,OAAO,OAAO,KAAA;AAAA,MAChD;AACA,UAAI,CAAC,QAAS,OAAM,IAAI,MAAM,sCAAsC;AACpE,UAAI,SAAS;AACX,YAAI,QAAQ,iBAAiB,SAAS;AACpC,eAAK,QAAQ,MAAM;AAAA,YACjB,MAAM,KAAK,kBAAkB,cAAc,KAAK,IAAI,OAAO;AAAA,YAC3D,MAAM;AAAA,YAAC;AAAA,UAAA;AAAA,QAEX,OAAO;AACL,eAAK,kBAAkB,cAAc,KAAK,IAAI,OAAO;AAAA,QACvD;AAAA,MACF;AAAA,IACF;AAEA,UAAM,UAAU,MAAM,SACnB;AAAA,MACC,CAAC,YACC,QAAQ,UAAU,aAAa,WAAW,QAAQ,MAAM,QAAQ;AAAA,IAAA,EAEnE,IAAI,CAAC,YAAY,QAAQ,KAAK,EAC9B,OAAO,CAAC,UAAkC,iBAAiB,OAAO;AACrE,WAAO;AAAA,MACL,SAAS;AAAA,MACT,OAAO;AAAA,MACP,OAAO,QAAQ,SAAS,IAAI,QAAQ,IAAI,OAAO,IAAI;AAAA,IAAA;AAAA,EAEvD;AAAA,EAEA,QAAc;AACZ,eAAW,SAAS,KAAK,OAAO,OAAA,GAAU;AACxC,iBAAW,WAAW,MAAM,SAAU,SAAQ,gBAAgB,MAAA;AAAA,IAChE;AACA,SAAK,OAAO,MAAA;AACZ,SAAK,YAAY,MAAA;AACjB,SAAK,gBAAgBA,0CAAA;AAAA,EACvB;AAAA,EAEA,iBAAiB,QAAyB;AACxC,WACE,KAAK,OACF,IAAI,MAAM,GACT,SAAS,KAAK,CAAC,YAAY,QAAQ,UAAU,SAAS,KAAK;AAAA,EAEnE;AAAA,EAEQ,kBACN,cACA,QACA,aACM;AACN,UAAM,QAAQ,KAAK,OAAO,IAAI,MAAM;AACpC,QACE,CAAC,OAAO,SAAS,SAAS,WAAW,KACrC,CAAC,aAAa,YAAY,MAAM,MAAM,IAAI,GAC1C;AACA,qBAAe,cAAc,WAAW;AACxC;AAAA,IACF;AACA,UAAM,WAAW,YAAY,YAAY,CAAA;AACzC,gBAAY,WAAW;AACvB,UAAM,WAAW,CAAC,WAAW;AAC7B,eAAW,WAAW,SAAU,gBAAe,cAAc,OAAO;AAAA,EACtE;AAAA,EAEQ,gBAAgB,MAA4B;AAClD,QAAI,KAAK,YAAY,IAAI,KAAK,EAAE,EAAG;AACnC,SAAK,YAAY,IAAI,KAAK,EAAE;AAC5B,UAAM,OAAO,KAAK,KAAK,KAAK,GAAG;AAC/B,YAAQ;AAAA,MACN,gBAAgB,KAAK,eAAe,KAAK,KAAK,YAAY,MAAM,EAAE,+BAA+B,IAAI,uJAEV,IAAI;AAAA,IAAA;AAAA,EAGnG;AACF;AAGO,SAAS,qBACd,gBACA,sBACA,eACA,eACS;AACT,UACG,qBAAqB,cAAc,GAAG,OAAO;AAAA,IAAK,CAAC,SAClD,cAAc,KAAK,EAAE;AAAA,EAAA,KAErB,UACF,cAAc,cAAc,GAAG,WAAW;AAE9C;AAEA,SAAS,iBACP,MACA,eACsB;AACtB,SAAO,IAAI;AAAA,IACT,CAAC,GAAG,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,cAAc,kBAAkB,GAAG,GAAG,GAAG,CAAC;AAAA,EAAA;AAEtE;AAEA,SAAS,aACP,MACA,OACS;AACT,SACE,KAAK,SAAS,MAAM,QAAQ,CAAC,GAAG,KAAK,KAAA,CAAM,EAAE,MAAM,CAAC,QAAQ,MAAM,IAAI,GAAG,CAAC;AAE9E;AAEA,SAAS,WACP,MACA,OACS;AACT,SAAO,CAAC,GAAG,KAAK,KAAA,CAAM,EAAE,KAAK,CAAC,QAAQ,MAAM,IAAI,GAAG,CAAC;AACtD;AAEA,SAAS,cACP,MACA,MACe;AACf,QAAM,QAAQC,UAAAA,QAAQ,IAAIC,GAAAA,QAAQ,KAAK,IAAI,GAAG,CAAC,GAAG,KAAK,OAAA,CAAQ,CAAC;AAChE,SAAO;AAAA,IACL;AAAA,IACA;AAAA,IACA,iBAAiB,IAAI,gBAAA;AAAA,IACrB,OAAO;AAAA,IACP,OAAO;AAAA,EAAA;AAEX;AAEA,SAAS,aACP,cACA,SACA,eACS;AACT,MAAI,WAAW;AACf,QAAM,YAAY,aAAa,gBAAgB;AAAA,IAC7C,OAAO,QAAQ;AAAA,IACf,QAAQ,QAAQ,gBAAgB;AAAA,IAChC,wBAAwB;AAAA,IACxB;AAAA,IACA,oBAAoB,CAAC,WAAW;AAC9B,iBAAW;AACX,cAAQ,QAAQ;AAChB,cAAQ,QAAQ,kBAAkB,UAAU,YAAY;AAAA,IAC1D;AAAA,EAAA,CACD;AACD,MAAI,QAAQ,iBAAiB,SAAS;AACpC,SAAK,QAAQ,MAAM;AAAA,MACjB,MAAM;AACJ,gBAAQ,QAAQ;AAAA,MAClB;AAAA,MACA,MAAM;AACJ,gBAAQ,QAAQ;AAAA,MAClB;AAAA,IAAA;AAAA,EAEJ;AACA,SAAO,aAAa;AACtB;AAEA,SAAS,eACP,cACA,SACM;AACN,UAAQ,gBAAgB,MAAA;AACxB,MAAI;AACF,iBAAa,gBAAgB,QAAQ,KAAK;AAAA,EAC5C,QAAQ;AAAA,EAER;AACF;;;"}