diff --git a/src/persistence.ts b/src/persistence.ts new file mode 100644 index 0000000..2875f2e --- /dev/null +++ b/src/persistence.ts @@ -0,0 +1,149 @@ +import * as Y from 'yjs'; +import { docsDb, LOCAL_DELTA_STORE } from './sync'; +import { assert } from './util'; +import _ from 'lodash'; + +const COMPACTION_THRESHOLD = 200; + +export async function openDoc(id: string): Promise { + // Return existing doc if we still have it open + const previousState = idToState.get(id); + if (previousState) { + previousState.refCount += 1; + return previousState.ydoc; + } + + const state = await initialLoad(id); + idToState.set(id, state); + docToState.set(state.ydoc, state); + + return state.ydoc; +} + +export function closeDoc(doc: Y.Doc) { + const state = docToState.get(doc); + assert(state, 'Document is not already open'); + + state.refCount -= 1; + if (state.refCount <= 0) { + idToState.delete(state.id); + docToState.delete(doc); + doc.destroy(); + } +} + +const idToState = new Map(); +const docToState = new Map(); + +type State = { + // TODO: Move more of this into class + id: string; + deltaBuffer: Uint8Array[]; + numRecords: number; + refCount: number; + ydoc: Y.Doc; +}; + +async function initialLoad(id: string): Promise { + const ydoc = new Y.Doc(); + const records = await docsDb.getAll( + LOCAL_DELTA_STORE, + IDBKeyRange.bound([id, -Infinity], [id, Infinity]) + ); + Y.applyUpdate(ydoc, Y.mergeUpdates(records)); + return { + id, + deltaBuffer: [], + numRecords: records.length, + refCount: 1, + ydoc, + }; +} + +class BufferedWriter { + private deltaBuffer: Uint8Array[]; + private docId: string; + private numRecords: number; + private debouncedFlush: () => void; + private ydoc: Y.Doc; + + constructor(docId: string) { + this.deltaBuffer = []; + this.docId = docId; + this.numRecords = 0; + this.debouncedFlush = _.debounce( + () => { + this.flush(); + }, + 1_000, + { maxWait: 30_000 } + ); + this.ydoc = new Y.Doc(); + } + + public async connect(): Promise { + const records = await docsDb.getAll( + LOCAL_DELTA_STORE, + IDBKeyRange.bound([this.docId, -Infinity], [this.docId, Infinity]) + ); + this.numRecords = records.length; + Y.applyUpdate(this.ydoc, Y.mergeUpdates(records)); + this.ydoc.on('update', (blob) => this.addUpdate(blob)); + return this.ydoc; + } + + public addUpdate(update: Uint8Array) { + this.deltaBuffer.push(update); + this.debouncedFlush(); + } + + public async flush(): Promise { + if (this.deltaBuffer.length === 0) { + return; + } + + // Write delta + const delta = Y.mergeUpdates(this.deltaBuffer.splice(0)); + try { + await docsDb.put(LOCAL_DELTA_STORE, delta, [ + this.docId, + new Date().valueOf(), + ]); + this.numRecords += 1; + } catch (err) { + this.addUpdate(delta); + console.log(`Error flushing document ${this.docId}: ${err}`); + throw err; + } + + // Compact if necessary + try { + await this.compact(); + } catch (err) { + console.log(`Error compacting document ${this.docId}: ${err}`); + } + } + + public async compact() { + if (this.numRecords < COMPACTION_THRESHOLD) { + return; + } + const tx = docsDb.transaction(LOCAL_DELTA_STORE, 'readwrite'); + + // Make sure we're not missing anything by pulling in all changes. + // This could happen if a document is being edited in two different tabs. + const fullRange = IDBKeyRange.bound( + [this.docId, -Infinity], + [this.docId, Infinity] + ); + const records = await tx.store.getAll(fullRange); + Y.applyUpdate(this.ydoc, Y.mergeUpdates(records)); + + // Replace with a compacted blob + const compacted = Y.encodeStateAsUpdate(this.ydoc); + await tx.store.delete(fullRange); + await tx.store.put(compacted, [this.docId, new Date().valueOf()]); + tx.commit(); + this.numRecords = 1; + } +} diff --git a/src/sync.ts b/src/sync.ts index c6a10c8..5ce76d8 100644 --- a/src/sync.ts +++ b/src/sync.ts @@ -5,8 +5,8 @@ import { assert, generateId } from './util'; import { getDocsMap, getVaultMap, type DocMap } from './vault'; const DB_NAME = 'synced-docs'; -const LOCAL_DELTA_STORE = 'local-updates'; -const VAULT_STORE = 'vaults'; +export const LOCAL_DELTA_STORE = 'local-updates'; +export const VAULT_STORE = 'vaults'; type LocalDeltaKey = [string, number]; type VaultRecord = { id: string }; @@ -22,7 +22,7 @@ interface DocsDatabase extends DBSchema { }; } -const db = await openDB(DB_NAME, 3, { +export const docsDb = await openDB(DB_NAME, 3, { upgrade(db, oldVersion, newVersion) { console.log(`Running upgrade from ${oldVersion} to ${newVersion}`); if (oldVersion < 2) { @@ -59,7 +59,7 @@ export class LocalDocument { public static async load(id: string): Promise { // Build Y.Doc from records on disk - const updates: Uint8Array[] = await db.getAll( + const updates: Uint8Array[] = await docsDb.getAll( LOCAL_DELTA_STORE, IDBKeyRange.bound([id, -Infinity], [id, Infinity]) ); @@ -78,7 +78,7 @@ export class LocalDocument { if (batch.length > 0) { const update = Y.mergeUpdates(batch); const key: LocalDeltaKey = [this.id, new Date().valueOf()]; - await db.add(LOCAL_DELTA_STORE, update, key); + await docsDb.add(LOCAL_DELTA_STORE, update, key); this.totalUpdates += 1; } } @@ -90,13 +90,13 @@ export class LocalDocument { } async function getOrCreateVaultDoc() { - const records = await db.getAll(VAULT_STORE); + const records = await docsDb.getAll(VAULT_STORE); if (records.length > 1) { throw new Error(`Expected 1 vault record, got ${records.length}`); } else if (records.length === 0) { const key = generateId(); const value = { id: key }; - await db.add(VAULT_STORE, value, key); + await docsDb.add(VAULT_STORE, value, key); records.push(value); } @@ -135,7 +135,7 @@ if (docsMap.size === 0) { const docInfo = new Y.Map([['title', oldKey]]) as DocMap; docsMap.set(newId, docInfo); - const tx = db.transaction(LOCAL_DELTA_STORE, 'readwrite'); + const tx = docsDb.transaction(LOCAL_DELTA_STORE, 'readwrite'); const recordKeys = (await tx.store.getAllKeys( IDBKeyRange.bound([oldKey, -Infinity], [oldKey, Infinity]) )) as LocalDeltaKey[];