{"version":3,"file":"compiled-query.cjs","sources":["../../../src/query/compiled-query.ts"],"sourcesContent":["import { D2, MultiSet, output } from \"@electric-sql/d2mini\"\nimport { createCollection } from \"../collection.js\"\nimport { compileQueryPipeline } from \"./pipeline-compiler.js\"\nimport type { Collection } from \"../collection.js\"\nimport type { ChangeMessage, SyncConfig } from \"../types.js\"\nimport type {\n  IStreamBuilder,\n  MultiSetArray,\n  RootStreamBuilder,\n} from \"@electric-sql/d2mini\"\nimport type { QueryBuilder, ResultsFromContext } from \"./query-builder.js\"\nimport type { Context, Schema } from \"./types.js\"\n\nexport function compileQuery<TContext extends Context<Schema>>(\n  queryBuilder: QueryBuilder<TContext>\n) {\n  return new CompiledQuery<\n    ResultsFromContext<TContext> & { _key?: string | number }\n  >(queryBuilder)\n}\n\nexport class CompiledQuery<TResults extends object = Record<string, unknown>> {\n  private graph: D2\n  private inputs: Record<string, RootStreamBuilder<any>>\n  private inputCollections: Record<string, Collection<any>>\n  private resultCollection: Collection<TResults>\n  public state: `compiled` | `running` | `stopped` = `compiled`\n  private unsubscribeCallbacks: Array<() => void> = []\n\n  constructor(queryBuilder: QueryBuilder<Context<Schema>>) {\n    const query = queryBuilder._query\n    const collections = query.collections\n\n    if (!collections) {\n      throw new Error(`No collections provided`)\n    }\n\n    this.inputCollections = collections\n\n    const graph = new D2()\n    const inputs = Object.fromEntries(\n      Object.entries(collections).map(([key]) => [key, graph.newInput<any>()])\n    )\n\n    const sync: SyncConfig<TResults>[`sync`] = ({ begin, write, commit }) => {\n      compileQueryPipeline<IStreamBuilder<[unknown, TResults]>>(\n        query,\n        inputs\n      ).pipe(\n        output((data) => {\n          begin()\n          data\n            .getInner()\n            .reduce((acc, [[key, value], multiplicity]) => {\n              const changes = acc.get(key) || {\n                deletes: 0,\n                inserts: 0,\n                value,\n              }\n              if (multiplicity < 0) {\n                changes.deletes += Math.abs(multiplicity)\n              } else if (multiplicity > 0) {\n                changes.inserts += multiplicity\n                changes.value = value\n              }\n              acc.set(key, changes)\n              return acc\n            }, new Map<unknown, { deletes: number; inserts: number; value: TResults }>())\n            .forEach((changes, rawKey) => {\n              const { deletes, inserts, value } = changes\n              const valueWithKey = { ...value, _key: rawKey }\n              if (inserts && !deletes) {\n                write({\n                  value: valueWithKey,\n                  type: `insert`,\n                })\n              } else if (inserts >= deletes) {\n                write({\n                  value: valueWithKey,\n                  type: `update`,\n                })\n              } else if (deletes > 0) {\n                write({\n                  value: valueWithKey,\n                  type: `delete`,\n                })\n              }\n            })\n          commit()\n        })\n      )\n      graph.finalize()\n    }\n\n    this.graph = graph\n    this.inputs = inputs\n    this.resultCollection = createCollection<TResults>({\n      id: crypto.randomUUID(), // TODO: remove when we don't require any more\n      getKey: (val: unknown) => {\n        return (val as any)._key\n      },\n      sync: {\n        sync,\n      },\n    })\n  }\n\n  get results() {\n    return this.resultCollection\n  }\n\n  private sendChangesToInput(\n    inputKey: string,\n    changes: Array<ChangeMessage>,\n    getKey: (item: ChangeMessage[`value`]) => any\n  ) {\n    const input = this.inputs[inputKey]!\n    const multiSetArray: MultiSetArray<unknown> = []\n    for (const change of changes) {\n      const key = getKey(change.value)\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    input.sendData(new MultiSet(multiSetArray))\n  }\n\n  private runGraph() {\n    this.graph.run()\n  }\n\n  start() {\n    if (this.state === `running`) {\n      throw new Error(`Query is already running`)\n    } else if (this.state === `stopped`) {\n      throw new Error(`Query is stopped`)\n    }\n\n    // Send initial state\n    Object.entries(this.inputCollections).forEach(([key, collection]) => {\n      this.sendChangesToInput(\n        key,\n        collection.currentStateAsChanges(),\n        collection.config.getKey\n      )\n    })\n    this.runGraph()\n\n    // Subscribe to changes\n    Object.entries(this.inputCollections).forEach(([key, collection]) => {\n      const unsubscribe = collection.subscribeChanges((changes) => {\n        this.sendChangesToInput(key, changes, collection.config.getKey)\n        this.runGraph()\n      })\n\n      this.unsubscribeCallbacks.push(unsubscribe)\n    })\n\n    this.state = `running`\n    return () => {\n      this.stop()\n    }\n  }\n\n  stop() {\n    this.unsubscribeCallbacks.forEach((unsubscribe) => unsubscribe())\n    this.unsubscribeCallbacks = []\n    this.state = `stopped`\n  }\n}\n"],"names":["D2","compileQueryPipeline","output","createCollection","MultiSet","collection"],"mappings":";;;;;AAaO,SAAS,aACd,cACA;AACO,SAAA,IAAI,cAET,YAAY;AAChB;AAEO,MAAM,cAAiE;AAAA,EAQ5E,YAAY,cAA6C;AAHzD,SAAO,QAA4C;AACnD,SAAQ,uBAA0C,CAAC;AAGjD,UAAM,QAAQ,aAAa;AAC3B,UAAM,cAAc,MAAM;AAE1B,QAAI,CAAC,aAAa;AACV,YAAA,IAAI,MAAM,yBAAyB;AAAA,IAAA;AAG3C,SAAK,mBAAmB;AAElB,UAAA,QAAQ,IAAIA,UAAG;AACrB,UAAM,SAAS,OAAO;AAAA,MACpB,OAAO,QAAQ,WAAW,EAAE,IAAI,CAAC,CAAC,GAAG,MAAM,CAAC,KAAK,MAAM,SAAA,CAAe,CAAC;AAAA,IACzE;AAEA,UAAM,OAAqC,CAAC,EAAE,OAAO,OAAO,aAAa;AACvEC,uBAAA;AAAA,QACE;AAAA,QACA;AAAA,MAAA,EACA;AAAA,QACAC,OAAA,OAAO,CAAC,SAAS;AACT,gBAAA;AAEH,eAAA,WACA,OAAO,CAAC,KAAK,CAAC,CAAC,KAAK,KAAK,GAAG,YAAY,MAAM;AAC7C,kBAAM,UAAU,IAAI,IAAI,GAAG,KAAK;AAAA,cAC9B,SAAS;AAAA,cACT,SAAS;AAAA,cACT;AAAA,YACF;AACA,gBAAI,eAAe,GAAG;AACZ,sBAAA,WAAW,KAAK,IAAI,YAAY;AAAA,YAAA,WAC/B,eAAe,GAAG;AAC3B,sBAAQ,WAAW;AACnB,sBAAQ,QAAQ;AAAA,YAAA;AAEd,gBAAA,IAAI,KAAK,OAAO;AACb,mBAAA;AAAA,UAAA,uBACF,IAAoE,CAAC,EAC3E,QAAQ,CAAC,SAAS,WAAW;AAC5B,kBAAM,EAAE,SAAS,SAAS,MAAU,IAAA;AACpC,kBAAM,eAAe,EAAE,GAAG,OAAO,MAAM,OAAO;AAC1C,gBAAA,WAAW,CAAC,SAAS;AACjB,oBAAA;AAAA,gBACJ,OAAO;AAAA,gBACP,MAAM;AAAA,cAAA,CACP;AAAA,YAAA,WACQ,WAAW,SAAS;AACvB,oBAAA;AAAA,gBACJ,OAAO;AAAA,gBACP,MAAM;AAAA,cAAA,CACP;AAAA,YAAA,WACQ,UAAU,GAAG;AAChB,oBAAA;AAAA,gBACJ,OAAO;AAAA,gBACP,MAAM;AAAA,cAAA,CACP;AAAA,YAAA;AAAA,UACH,CACD;AACI,iBAAA;AAAA,QACR,CAAA;AAAA,MACH;AACA,YAAM,SAAS;AAAA,IACjB;AAEA,SAAK,QAAQ;AACb,SAAK,SAAS;AACd,SAAK,mBAAmBC,4BAA2B;AAAA,MACjD,IAAI,OAAO,WAAW;AAAA;AAAA,MACtB,QAAQ,CAAC,QAAiB;AACxB,eAAQ,IAAY;AAAA,MACtB;AAAA,MACA,MAAM;AAAA,QACJ;AAAA,MAAA;AAAA,IACF,CACD;AAAA,EAAA;AAAA,EAGH,IAAI,UAAU;AACZ,WAAO,KAAK;AAAA,EAAA;AAAA,EAGN,mBACN,UACA,SACA,QACA;AACM,UAAA,QAAQ,KAAK,OAAO,QAAQ;AAClC,UAAM,gBAAwC,CAAC;AAC/C,eAAW,UAAU,SAAS;AACtB,YAAA,MAAM,OAAO,OAAO,KAAK;AAC3B,UAAA,OAAO,SAAS,UAAU;AACd,sBAAA,KAAK,CAAC,CAAC,KAAK,OAAO,KAAK,GAAG,CAAC,CAAC;AAAA,MAC7C,WAAW,OAAO,SAAS,UAAU;AACrB,sBAAA,KAAK,CAAC,CAAC,KAAK,OAAO,aAAa,GAAG,EAAE,CAAC;AACtC,sBAAA,KAAK,CAAC,CAAC,KAAK,OAAO,KAAK,GAAG,CAAC,CAAC;AAAA,MAAA,OACtC;AAES,sBAAA,KAAK,CAAC,CAAC,KAAK,OAAO,KAAK,GAAG,EAAE,CAAC;AAAA,MAAA;AAAA,IAC9C;AAEF,UAAM,SAAS,IAAIC,OAAS,SAAA,aAAa,CAAC;AAAA,EAAA;AAAA,EAGpC,WAAW;AACjB,SAAK,MAAM,IAAI;AAAA,EAAA;AAAA,EAGjB,QAAQ;AACF,QAAA,KAAK,UAAU,WAAW;AACtB,YAAA,IAAI,MAAM,0BAA0B;AAAA,IAC5C,WAAW,KAAK,UAAU,WAAW;AAC7B,YAAA,IAAI,MAAM,kBAAkB;AAAA,IAAA;AAI7B,WAAA,QAAQ,KAAK,gBAAgB,EAAE,QAAQ,CAAC,CAAC,KAAKC,WAAU,MAAM;AAC9D,WAAA;AAAA,QACH;AAAA,QACAA,YAAW,sBAAsB;AAAA,QACjCA,YAAW,OAAO;AAAA,MACpB;AAAA,IAAA,CACD;AACD,SAAK,SAAS;AAGP,WAAA,QAAQ,KAAK,gBAAgB,EAAE,QAAQ,CAAC,CAAC,KAAKA,WAAU,MAAM;AACnE,YAAM,cAAcA,YAAW,iBAAiB,CAAC,YAAY;AAC3D,aAAK,mBAAmB,KAAK,SAASA,YAAW,OAAO,MAAM;AAC9D,aAAK,SAAS;AAAA,MAAA,CACf;AAEI,WAAA,qBAAqB,KAAK,WAAW;AAAA,IAAA,CAC3C;AAED,SAAK,QAAQ;AACb,WAAO,MAAM;AACX,WAAK,KAAK;AAAA,IACZ;AAAA,EAAA;AAAA,EAGF,OAAO;AACL,SAAK,qBAAqB,QAAQ,CAAC,gBAAgB,aAAa;AAChE,SAAK,uBAAuB,CAAC;AAC7B,SAAK,QAAQ;AAAA,EAAA;AAEjB;;;"}