Something went wrong. Try again.
A local-first note taking app
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380import * as Automerge from '@automerge/automerge';
import { habitatAuth } from '../auth/habitatAuth';import { Libp2pAutomergeConnectionProvider, HabitatLibp2pNode,} from './automergeConnectionProvider';import { loadAutomergeDoc, resolveDocDisplayTitle, saveAutomergeDoc, type LoadDocResult, type AutomergeDocSchema,} from './automergeDoc';import { createHabitatLibp2pNode, dialRelayAndStartPeerDiscovery,} from './libp2pNode';import type { HabitatDoc, TypedRecord } from './types';
/** * Per-URI live document session for Automerge docs. Mirrors the Yjs * `docSession.ts` lifecycle but with Automerge's immutable doc model. * * Automerge docs are immutable — each change produces a new doc reference. * The session holds the current doc and notifies listeners when it changes. */
export type AutomergeDocSession = { uri: string; doc: Automerge.Doc<AutomergeDocSchema>; provider: Libp2pAutomergeConnectionProvider | null; node: HabitatLibp2pNode | null; ownerRecord: TypedRecord<HabitatDoc>; loadPromise: Promise<void>; getTitle: () => string; resync: () => Promise<void>; release: () => void; applyChange: (fn: (doc: AutomergeDocSchema) => void) => void; changeAt: (heads: string[], fn: (doc: AutomergeDocSchema) => void) => void; subscribe: (cb: () => void) => () => void;};
type RegistryEntry = { uri: string; doc: Automerge.Doc<AutomergeDocSchema>; loadPromise: Promise<LoadDocResult>; ownerRecord: TypedRecord<HabitatDoc> | null; node: HabitatLibp2pNode | null; provider: Libp2pAutomergeConnectionProvider | null; refCount: number; saveHandler: (() => void) | null; saveTimer: ReturnType<typeof setTimeout> | null; pendingSave: Promise<void> | null; teardownPromise: Promise<void> | null; /** Subscribers notified when `doc` reference changes. */ changeListeners: Set<() => void>;};
const registry = new Map<string, RegistryEntry>();
const SAVE_DEBOUNCE_MS = 1000;
function assignLoadPromise( entry: RegistryEntry, promise: Promise<LoadDocResult>,): void { entry.loadPromise = promise; promise.catch(() => undefined);}
let visibilityListenerInstalled = false;
function ensureVisibilityListener() { if (visibilityListenerInstalled) return; if (typeof document === 'undefined') return; visibilityListenerInstalled = true; document.addEventListener('visibilitychange', () => { if (document.visibilityState !== 'visible') return; for (const entry of registry.values()) { if (!entry.node) continue; void resyncEntry(entry); } });}
async function resyncEntry(entry: RegistryEntry) { try { const actor = Automerge.getActorId(entry.doc); const result = await loadAutomergeDoc( entry.uri, Automerge.init<AutomergeDocSchema>({ actor }), ); entry.doc = Automerge.merge(result.doc, entry.doc); entry.ownerRecord = result.ownerRecord; notifyChange(entry); } catch (e) { console.warn('[automerge] visibility resync failed', e); } if (entry.node) { await dialRelayAndStartPeerDiscovery(entry.uri, entry.node); }}
function scheduleSave(entry: RegistryEntry) { if (entry.saveTimer) clearTimeout(entry.saveTimer); entry.saveTimer = setTimeout(() => { entry.saveTimer = null; void runSave(entry); }, SAVE_DEBOUNCE_MS);}
async function runSave(entry: RegistryEntry) { const myDid = habitatAuth.getAuthInfo()?.did; if (!myDid || !entry.ownerRecord) return; const ownerRecord = entry.ownerRecord; entry.pendingSave = saveAutomergeDoc({ ownerRecord, doc: entry.doc, myDid, }) .then(() => { const name = resolveDocDisplayTitle(entry.doc, ownerRecord.value.name); entry.ownerRecord = { ...ownerRecord, value: { ...ownerRecord.value, name }, }; }) .catch((e) => { console.error('[automerge] save failed', e); }) .finally(() => { entry.pendingSave = null; }); await entry.pendingSave;}
function attachSaveHandler(entry: RegistryEntry) { if (entry.saveHandler) return; entry.saveHandler = () => { scheduleSave(entry); }; entry.changeListeners.add(entry.saveHandler);}
function detachSaveHandler(entry: RegistryEntry) { if (!entry.saveHandler) return; entry.changeListeners.delete(entry.saveHandler); entry.saveHandler = null;}
function notifyChange(entry: RegistryEntry): void { for (const listener of entry.changeListeners) { listener(); }}
export function acquireDocSession(uri: string): AutomergeDocSession { ensureVisibilityListener(); let entry = registry.get(uri); if (!entry) { const doc = Automerge.init<AutomergeDocSchema>(); entry = { uri, doc, loadPromise: Promise.resolve({} as LoadDocResult), ownerRecord: null, node: null, provider: null, refCount: 0, saveHandler: null, saveTimer: null, pendingSave: null, teardownPromise: null, changeListeners: new Set(), }; registry.set(uri, entry); assignLoadPromise(entry, initSession(entry)); } else if (entry.teardownPromise || !entry.node) { assignLoadPromise(entry, reinitSession(entry)); }
entry.refCount += 1; const captured = entry; return makeHandle(captured);}
// TODO(test): Add a regression test for the load/applyChange race that caused// "Attempting to change an out of date document". Repro: while loadAutomergeDoc// is awaiting the PDS fetch, a concurrent applyChange marks the in-flight doc// reference stale (Automerge progressDocument sets state.heads); the resumed// loadIncremental then throws. Fixed by loading into Automerge.init({ actor })// and reconciling via Automerge.merge(result.doc, entry.doc) after the await.// Same pattern guarded in reinitSession (the user's actual repro: switch away// then back) and resyncEntry. Test against a fake in-memory PDS behind// habitatAuth.fetch (not vi.mock of getPrivateRecord) so it stays decoupled// from the call sequence; gate one getRecord response to interleave applyChange,// then assert loadPromise resolves and the local edit survives the merge.// Deferred until the ./xrpc / habitatAuth.fetch layer settles.async function initSession(entry: RegistryEntry): Promise<LoadDocResult> { try { // Load into a fresh doc so concurrent applyChange calls can't make the // base reference stale during the async PDS fetch. const actor = Automerge.getActorId(entry.doc); const result = await loadAutomergeDoc( entry.uri, Automerge.init<AutomergeDocSchema>({ actor }), ); // Merge any in-flight local changes that happened during the load. entry.doc = Automerge.merge(result.doc, entry.doc); entry.ownerRecord = result.ownerRecord; attachSaveHandler(entry); await ensureLiveEditing(entry); notifyChange(entry); return result; } catch (e) { console.error('[automerge] initial load failed for', entry.uri, e); throw e; }}
async function reinitSession(entry: RegistryEntry): Promise<LoadDocResult> { if (entry.teardownPromise) { try { await entry.teardownPromise; } catch { // teardown errors are non-fatal } } try { const actor = Automerge.getActorId(entry.doc); const result = await loadAutomergeDoc( entry.uri, Automerge.init<AutomergeDocSchema>({ actor }), ); entry.doc = Automerge.merge(result.doc, entry.doc); entry.ownerRecord = result.ownerRecord; attachSaveHandler(entry); await ensureLiveEditing(entry); notifyChange(entry); return result; } catch (e) { console.error('[automerge] reload failed for', entry.uri, e); throw e; }}
async function ensureLiveEditing(entry: RegistryEntry) { if (entry.node) return; try { const node = await createHabitatLibp2pNode(); entry.node = node; entry.provider = new Libp2pAutomergeConnectionProvider( node, entry.doc, entry.uri, (newDoc) => { entry.doc = newDoc; notifyChange(entry); }, ); await dialRelayAndStartPeerDiscovery(entry.uri, node); } catch (e) { console.error( '[automerge] could not start libp2p; doc is read-write but offline', e, ); }}
function makeHandle(entry: RegistryEntry): AutomergeDocSession { let released = false; const loadPromise = entry.loadPromise.then(() => undefined); const handle: AutomergeDocSession = { uri: entry.uri, get doc() { return entry.doc; }, get provider() { return entry.provider; }, get node() { return entry.node; }, get ownerRecord() { return entry.ownerRecord!; }, loadPromise, getTitle() { if (!entry.ownerRecord) return ''; return resolveDocDisplayTitle(entry.doc, entry.ownerRecord.value.name); }, resync: () => resyncEntry(entry), applyChange: (fn) => { const newDoc = Automerge.change(entry.doc, fn); entry.doc = newDoc; entry.provider?.onLocalChange(newDoc); notifyChange(entry); }, changeAt: (heads, fn) => { const result = Automerge.changeAt(entry.doc, heads, fn); if (result.newHeads !== null) { entry.doc = result.newDoc; entry.provider?.onLocalChange(result.newDoc); notifyChange(entry); } }, subscribe: (cb) => { entry.changeListeners.add(cb); // Fire immediately if the doc is already loaded (ownerRecord is set). // This catches the case where the editor subscribes after the // initial load has already completed. if (entry.ownerRecord) { cb(); } return () => { entry.changeListeners.delete(cb); }; }, release: () => { if (released) return; released = true; entry.refCount -= 1; if (entry.refCount === 0) { entry.teardownPromise = teardownEntry(entry).finally(() => { entry.teardownPromise = null; }); } }, }; return handle;}
async function teardownEntry(entry: RegistryEntry) { if (entry.saveTimer) { clearTimeout(entry.saveTimer); entry.saveTimer = null; await runSave(entry).catch(() => undefined); } if (entry.pendingSave) { try { await entry.pendingSave; } catch { // already logged } } if (entry.provider) { try { entry.provider.destroy(); } catch { // best-effort } entry.provider = null; } if (entry.node) { try { await entry.node.stop(); } catch { // best-effort } entry.node = null; } detachSaveHandler(entry);}
/** Test-only: clear the registry between tests. */export function __resetDocSessionRegistryForTests() { for (const entry of registry.values()) { if (entry.saveTimer) clearTimeout(entry.saveTimer); try { entry.provider?.destroy(); } catch { // ignore } const stopped = entry.node?.stop(); if (stopped) void stopped.catch(() => undefined); } registry.clear();}