{"version":3,"file":"OutboxManager.cjs","sources":["../../../src/outbox/OutboxManager.ts"],"sourcesContent":["import { withSpan } from '../telemetry/tracer'\nimport {\n  MissingTemporalConstructorError,\n  TransactionSerializer,\n} from './TransactionSerializer'\nimport type { OfflineTransaction, StorageAdapter } from '../types'\nimport type { Collection } from '@tanstack/db'\n\nexport class OutboxTransactionNotFoundError extends Error {}\n\nexport class OutboxManager {\n  private storage: StorageAdapter\n  private serializer: TransactionSerializer\n  private keyPrefix = `tx:`\n  private removalRevision = 0\n  private activeReads = new Set<{ revision: number }>()\n  private removedAtRevision = new Map<string, number>()\n\n  constructor(\n    storage: StorageAdapter,\n\n    collections: Record<string, Collection<any, any, any, any, any>>,\n  ) {\n    this.storage = storage\n    this.serializer = new TransactionSerializer(collections)\n  }\n\n  private getStorageKey(id: string): string {\n    return `${this.keyPrefix}${id}`\n  }\n\n  private recordRemovals(ids: Iterable<string>): void {\n    this.removalRevision++\n    if (this.activeReads.size === 0) return\n    for (const id of ids) this.removedAtRevision.set(id, this.removalRevision)\n    this.pruneRemovalHistory()\n  }\n\n  private pruneRemovalHistory(): void {\n    if (this.activeReads.size === 0) {\n      this.removedAtRevision.clear()\n      return\n    }\n    let oldestRead = Number.POSITIVE_INFINITY\n    for (const read of this.activeReads)\n      oldestRead = Math.min(oldestRead, read.revision)\n    for (const [id, revision] of this.removedAtRevision)\n      if (revision <= oldestRead) this.removedAtRevision.delete(id)\n  }\n\n  async add(transaction: OfflineTransaction): Promise<void> {\n    return withSpan(\n      `outbox.add`,\n      {\n        'transaction.id': transaction.id,\n        'transaction.mutationFnName': transaction.mutationFnName,\n        'transaction.keyCount': transaction.keys.length,\n      },\n      async () => {\n        const key = this.getStorageKey(transaction.id)\n        const serialized = this.serializer.serialize(transaction)\n        await this.storage.set(key, serialized)\n      },\n    )\n  }\n\n  async get(id: string): Promise<OfflineTransaction | null> {\n    return withSpan(`outbox.get`, { 'transaction.id': id }, async (span) => {\n      const key = this.getStorageKey(id)\n      const data = await this.storage.get(key)\n\n      if (!data) {\n        span.setAttribute(`result`, `not_found`)\n        return null\n      }\n\n      try {\n        const transaction = this.serializer.deserialize(data)\n        span.setAttribute(`result`, `found`)\n        return transaction\n      } catch (error) {\n        if (error instanceof MissingTemporalConstructorError) {\n          error.message = `transaction ${id}: ${error.message}`\n          throw error\n        }\n        console.warn(`Failed to deserialize transaction ${id}:`, error)\n        span.setAttribute(`result`, `deserialize_error`)\n        return null\n      }\n    })\n  }\n\n  async getAll(): Promise<Array<OfflineTransaction>> {\n    return this.withAll((transactions) => transactions)\n  }\n\n  async withAll<T>(\n    consume: (transactions: Array<OfflineTransaction>) => T,\n  ): Promise<T> {\n    const read = { revision: this.removalRevision }\n    this.activeReads.add(read)\n    try {\n      return await withSpan(`outbox.getAll`, {}, async (span) => {\n        const keys = await this.storage.keys()\n        const transactionKeys = keys.filter((key) =>\n          key.startsWith(this.keyPrefix),\n        )\n\n        span.setAttribute(`transactionCount`, transactionKeys.length)\n\n        const transactions: Array<OfflineTransaction> = []\n\n        for (const key of transactionKeys) {\n          const data = await this.storage.get(key)\n          if (data) {\n            try {\n              const transaction = this.serializer.deserialize(data)\n              transactions.push(transaction)\n            } catch (error) {\n              if (error instanceof MissingTemporalConstructorError) {\n                error.message = `transaction ${key.slice(this.keyPrefix.length)}: ${error.message}`\n                throw error\n              }\n              console.warn(\n                `Failed to deserialize transaction from key ${key}:`,\n                error,\n              )\n            }\n          }\n        }\n\n        const currentTransactions = transactions\n          .filter(\n            ({ id }) => (this.removedAtRevision.get(id) ?? 0) <= read.revision,\n          )\n          .sort((a, b) => a.createdAt.getTime() - b.createdAt.getTime())\n        // Keep the read registered until replay admission completes. A durable\n        // removal cannot land between filtering and this synchronous consumer.\n        return consume(currentTransactions)\n      })\n    } finally {\n      this.activeReads.delete(read)\n      this.pruneRemovalHistory()\n    }\n  }\n\n  async getByKeys(keys: Array<string>): Promise<Array<OfflineTransaction>> {\n    const allTransactions = await this.getAll()\n    const keySet = new Set(keys)\n\n    return allTransactions.filter((transaction) =>\n      transaction.keys.some((key) => keySet.has(key)),\n    )\n  }\n\n  async update(\n    id: string,\n    updates: Partial<OfflineTransaction>,\n  ): Promise<void> {\n    return withSpan(`outbox.update`, { 'transaction.id': id }, async () => {\n      const existing = await this.get(id)\n      if (!existing) {\n        throw new OutboxTransactionNotFoundError(`Transaction ${id} not found`)\n      }\n\n      const updated = { ...existing, ...updates }\n      await this.add(updated)\n    })\n  }\n\n  async remove(id: string): Promise<void> {\n    return withSpan(`outbox.remove`, { 'transaction.id': id }, async () => {\n      const key = this.getStorageKey(id)\n      await this.storage.delete(key)\n      this.recordRemovals([id])\n    })\n  }\n\n  async removeMany(ids: Array<string>): Promise<void> {\n    return withSpan(`outbox.removeMany`, { count: ids.length }, async () => {\n      await Promise.all(ids.map((id) => this.remove(id)))\n    })\n  }\n\n  async clear(): Promise<void> {\n    const keys = await this.storage.keys()\n    const transactionKeys = keys.filter((key) => key.startsWith(this.keyPrefix))\n\n    await this.removeMany(\n      transactionKeys.map((key) => key.slice(this.keyPrefix.length)),\n    )\n  }\n\n  async count(): Promise<number> {\n    const keys = await this.storage.keys()\n    return keys.filter((key) => key.startsWith(this.keyPrefix)).length\n  }\n}\n"],"names":["TransactionSerializer","withSpan","MissingTemporalConstructorError"],"mappings":";;;;AAQO,MAAM,uCAAuC,MAAM;AAAC;AAEpD,MAAM,cAAc;AAAA,EAQzB,YACE,SAEA,aACA;AATF,SAAQ,YAAY;AACpB,SAAQ,kBAAkB;AAC1B,SAAQ,kCAAkB,IAAA;AAC1B,SAAQ,wCAAwB,IAAA;AAO9B,SAAK,UAAU;AACf,SAAK,aAAa,IAAIA,sBAAAA,sBAAsB,WAAW;AAAA,EACzD;AAAA,EAEQ,cAAc,IAAoB;AACxC,WAAO,GAAG,KAAK,SAAS,GAAG,EAAE;AAAA,EAC/B;AAAA,EAEQ,eAAe,KAA6B;AAClD,SAAK;AACL,QAAI,KAAK,YAAY,SAAS,EAAG;AACjC,eAAW,MAAM,IAAK,MAAK,kBAAkB,IAAI,IAAI,KAAK,eAAe;AACzE,SAAK,oBAAA;AAAA,EACP;AAAA,EAEQ,sBAA4B;AAClC,QAAI,KAAK,YAAY,SAAS,GAAG;AAC/B,WAAK,kBAAkB,MAAA;AACvB;AAAA,IACF;AACA,QAAI,aAAa,OAAO;AACxB,eAAW,QAAQ,KAAK;AACtB,mBAAa,KAAK,IAAI,YAAY,KAAK,QAAQ;AACjD,eAAW,CAAC,IAAI,QAAQ,KAAK,KAAK;AAChC,UAAI,YAAY,WAAY,MAAK,kBAAkB,OAAO,EAAE;AAAA,EAChE;AAAA,EAEA,MAAM,IAAI,aAAgD;AACxD,WAAOC,OAAAA;AAAAA,MACL;AAAA,MACA;AAAA,QACE,kBAAkB,YAAY;AAAA,QAC9B,8BAA8B,YAAY;AAAA,QAC1C,wBAAwB,YAAY,KAAK;AAAA,MAAA;AAAA,MAE3C,YAAY;AACV,cAAM,MAAM,KAAK,cAAc,YAAY,EAAE;AAC7C,cAAM,aAAa,KAAK,WAAW,UAAU,WAAW;AACxD,cAAM,KAAK,QAAQ,IAAI,KAAK,UAAU;AAAA,MACxC;AAAA,IAAA;AAAA,EAEJ;AAAA,EAEA,MAAM,IAAI,IAAgD;AACxD,WAAOA,OAAAA,SAAS,cAAc,CAAuB,GAAG,OAAO,SAAS;AACtE,YAAM,MAAM,KAAK,cAAc,EAAE;AACjC,YAAM,OAAO,MAAM,KAAK,QAAQ,IAAI,GAAG;AAEvC,UAAI,CAAC,MAAM;AACT,aAAK,aAAa,UAAU,WAAW;AACvC,eAAO;AAAA,MACT;AAEA,UAAI;AACF,cAAM,cAAc,KAAK,WAAW,YAAY,IAAI;AACpD,aAAK,aAAa,UAAU,OAAO;AACnC,eAAO;AAAA,MACT,SAAS,OAAO;AACd,YAAI,iBAAiBC,sBAAAA,iCAAiC;AACpD,gBAAM,UAAU,eAAe,EAAE,KAAK,MAAM,OAAO;AACnD,gBAAM;AAAA,QACR;AACA,gBAAQ,KAAK,qCAAqC,EAAE,KAAK,KAAK;AAC9D,aAAK,aAAa,UAAU,mBAAmB;AAC/C,eAAO;AAAA,MACT;AAAA,IACF,CAAC;AAAA,EACH;AAAA,EAEA,MAAM,SAA6C;AACjD,WAAO,KAAK,QAAQ,CAAC,iBAAiB,YAAY;AAAA,EACpD;AAAA,EAEA,MAAM,QACJ,SACY;AACZ,UAAM,OAAO,EAAE,UAAU,KAAK,gBAAA;AAC9B,SAAK,YAAY,IAAI,IAAI;AACzB,QAAI;AACF,aAAO,MAAMD,OAAAA,SAAS,iBAAiB,CAAA,GAAI,OAAO,SAAS;AACzD,cAAM,OAAO,MAAM,KAAK,QAAQ,KAAA;AAChC,cAAM,kBAAkB,KAAK;AAAA,UAAO,CAAC,QACnC,IAAI,WAAW,KAAK,SAAS;AAAA,QAAA;AAG/B,aAAK,aAAa,oBAAoB,gBAAgB,MAAM;AAE5D,cAAM,eAA0C,CAAA;AAEhD,mBAAW,OAAO,iBAAiB;AACjC,gBAAM,OAAO,MAAM,KAAK,QAAQ,IAAI,GAAG;AACvC,cAAI,MAAM;AACR,gBAAI;AACF,oBAAM,cAAc,KAAK,WAAW,YAAY,IAAI;AACpD,2BAAa,KAAK,WAAW;AAAA,YAC/B,SAAS,OAAO;AACd,kBAAI,iBAAiBC,sBAAAA,iCAAiC;AACpD,sBAAM,UAAU,eAAe,IAAI,MAAM,KAAK,UAAU,MAAM,CAAC,KAAK,MAAM,OAAO;AACjF,sBAAM;AAAA,cACR;AACA,sBAAQ;AAAA,gBACN,8CAA8C,GAAG;AAAA,gBACjD;AAAA,cAAA;AAAA,YAEJ;AAAA,UACF;AAAA,QACF;AAEA,cAAM,sBAAsB,aACzB;AAAA,UACC,CAAC,EAAE,UAAU,KAAK,kBAAkB,IAAI,EAAE,KAAK,MAAM,KAAK;AAAA,QAAA,EAE3D,KAAK,CAAC,GAAG,MAAM,EAAE,UAAU,YAAY,EAAE,UAAU,QAAA,CAAS;AAG/D,eAAO,QAAQ,mBAAmB;AAAA,MACpC,CAAC;AAAA,IACH,UAAA;AACE,WAAK,YAAY,OAAO,IAAI;AAC5B,WAAK,oBAAA;AAAA,IACP;AAAA,EACF;AAAA,EAEA,MAAM,UAAU,MAAyD;AACvE,UAAM,kBAAkB,MAAM,KAAK,OAAA;AACnC,UAAM,SAAS,IAAI,IAAI,IAAI;AAE3B,WAAO,gBAAgB;AAAA,MAAO,CAAC,gBAC7B,YAAY,KAAK,KAAK,CAAC,QAAQ,OAAO,IAAI,GAAG,CAAC;AAAA,IAAA;AAAA,EAElD;AAAA,EAEA,MAAM,OACJ,IACA,SACe;AACf,WAAOD,OAAAA,SAAS,iBAAiB,CAAuB,GAAG,YAAY;AACrE,YAAM,WAAW,MAAM,KAAK,IAAI,EAAE;AAClC,UAAI,CAAC,UAAU;AACb,cAAM,IAAI,+BAA+B,eAAe,EAAE,YAAY;AAAA,MACxE;AAEA,YAAM,UAAU,EAAE,GAAG,UAAU,GAAG,QAAA;AAClC,YAAM,KAAK,IAAI,OAAO;AAAA,IACxB,CAAC;AAAA,EACH;AAAA,EAEA,MAAM,OAAO,IAA2B;AACtC,WAAOA,OAAAA,SAAS,iBAAiB,CAAuB,GAAG,YAAY;AACrE,YAAM,MAAM,KAAK,cAAc,EAAE;AACjC,YAAM,KAAK,QAAQ,OAAO,GAAG;AAC7B,WAAK,eAAe,CAAC,EAAE,CAAC;AAAA,IAC1B,CAAC;AAAA,EACH;AAAA,EAEA,MAAM,WAAW,KAAmC;AAClD,WAAOA,OAAAA,SAAS,qBAAqB,EAAE,OAAO,IAAI,OAAA,GAAU,YAAY;AACtE,YAAM,QAAQ,IAAI,IAAI,IAAI,CAAC,OAAO,KAAK,OAAO,EAAE,CAAC,CAAC;AAAA,IACpD,CAAC;AAAA,EACH;AAAA,EAEA,MAAM,QAAuB;AAC3B,UAAM,OAAO,MAAM,KAAK,QAAQ,KAAA;AAChC,UAAM,kBAAkB,KAAK,OAAO,CAAC,QAAQ,IAAI,WAAW,KAAK,SAAS,CAAC;AAE3E,UAAM,KAAK;AAAA,MACT,gBAAgB,IAAI,CAAC,QAAQ,IAAI,MAAM,KAAK,UAAU,MAAM,CAAC;AAAA,IAAA;AAAA,EAEjE;AAAA,EAEA,MAAM,QAAyB;AAC7B,UAAM,OAAO,MAAM,KAAK,QAAQ,KAAA;AAChC,WAAO,KAAK,OAAO,CAAC,QAAQ,IAAI,WAAW,KAAK,SAAS,CAAC,EAAE;AAAA,EAC9D;AACF;;;"}