Something went wrong. Try again.
automerge-repo drop in replacement backed by Habitat "permissioned spaces" om your PDS
Something went wrong. Try again.
6.7 kB · 160 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161import { emptyChange, init, load, save, type Doc } from '@automerge/automerge'import { generateDocumentId, validateDocumentId, type DocumentId } from './document-id.js'import type { RootDocumentSync } from './habitat/index.js'import { LocalDocumentStore } from './local/document-store.js'import { MemoryLocalPersistence, type LocalPersistence } from './local/persistence.js'import type { LocalSyncChannel } from './local/sync-channel.js'import { createSession, type Session, type SourceBinding, type SourceProvider } from './sync/session.js'import type { OpenSessionOptions, SessionHost } from './sync/session-host.js'import { createEditorSource, type Handle } from './sync/sources/editorSource.js'import { createLocalSource } from './sync/sources/localSource.js'import { createPdsSource } from './sync/sources/pdsSource.js'
export interface MultiSourceSessionHostOptions { /** Local durable byte store backing this host (in-memory if omitted). */ readonly persistence?: LocalPersistence /** * The atproto oplog sync used to reconstruct documents and back each session's * PDS Source. When omitted, the host is local-only (no remote sync). */ readonly documentSync?: RootDocumentSync<unknown> /** The vault root document id, context for publishing to the PDS oplog. */ readonly rootDocumentId?: DocumentId /** Optional cross-tab channel; when present, sessions receive sibling-tab edits. */ readonly sync?: LocalSyncChannel /** * Additional Sources attached to every session (after the built-in local + PDS * Sources) — the injection point for a libp2p live transport or any custom * backend. Each provider is invoked per document. */ readonly sources?: readonly SourceProvider[]}
/** * A self-contained {@link SessionHost}: it manages documents purely as live sync * sessions, over the low-level primitives (a {@link LocalDocumentStore} plus * injected `documentSync`). `create`/`find` return a {@link Handle} (a DocHandle- * like adapter) backed by a session; there is no automerge-repo compatibility layer. * * Each session is backed by a **local Source** (on-device persistence + cross-tab) * and, when `documentSync`/`rootDocumentId` are configured, a **PDS Source** for * durable remote sync — neutral by default (catch up on open; publish only when * `pdsPolicy` opts in). Completely decoupled from `Repo`. */export class MultiSourceSessionHost implements SessionHost { readonly #store: LocalDocumentStore readonly #documentSync: RootDocumentSync<unknown> | undefined readonly #rootDocumentId: DocumentId | undefined readonly #sync: LocalSyncChannel | undefined readonly #sourceProviders: readonly SourceProvider[] readonly #sourceId = generateDocumentId() readonly #sessions = new Map<DocumentId, { session: Session; handle: Handle }>() readonly #pending = new Set<Promise<void>>() #writeQueue: Promise<void> = Promise.resolve()
constructor(options: MultiSourceSessionHostOptions = {}) { this.#store = new LocalDocumentStore(options.persistence ?? new MemoryLocalPersistence()) this.#documentSync = options.documentSync this.#rootDocumentId = options.rootDocumentId this.#sync = options.sync this.#sourceProviders = options.sources ?? [] }
create<T = unknown>(options: OpenSessionOptions = {}): Handle<T> { const documentId = generateDocumentId() const initialDoc = emptyChange(init<T>(), undefined) // Persist an independent copy of the initial doc (bytes captured now). const snapshot = save(initialDoc) this.#enqueue(() => this.#store.saveDoc(documentId, load(snapshot))) return this.#buildSession<T>(documentId, initialDoc, options).handle }
async find<T = unknown>(id: DocumentId, options: OpenSessionOptions = {}): Promise<Handle<T>> { const documentId = validateDocumentId(id) const existing = this.#sessions.get(documentId) if (existing) return existing.handle as Handle<T>
let doc = await this.#store.loadDoc<T>(documentId) if (!doc) { // Not local — reconstruct from the PDS oplog (throws if nothing published). const loaded = await this.#documentSync?.loadRemoteDocument?.({ documentId, store: this.#store }) doc = (loaded?.doc as Doc<T> | undefined) ?? null } if (!doc) throw new Error(`Document ${documentId} not found`)
const { session, handle } = this.#buildSession<T>(documentId, doc, options) await session.pull() await session.flush() return handle }
delete(id: DocumentId): void { const documentId = validateDocumentId(id) this.#sessions.get(documentId)?.session.dispose() this.#sessions.delete(documentId) this.#enqueue(() => this.#store.deleteDoc(documentId)) }
async listDocuments(): Promise<DocumentId[]> { return this.#store.listDocumentIds() }
async flush(): Promise<void> { for (const { session } of this.#sessions.values()) await session.flush() while (this.#pending.size > 0) await Promise.all([...this.#pending]) }
async close(): Promise<void> { for (const { session } of this.#sessions.values()) { await session.flush() session.dispose() } this.#sessions.clear() await this.flush() }
// --- internals ---
#buildSession<T>( documentId: DocumentId, initialDoc: Doc<T>, options: OpenSessionOptions, ): { session: Session; handle: Handle<T> } { // The editor is itself a Source (first in the list) and the client-facing Handle. const editor = createEditorSource<T>(documentId, initialDoc)
const sources: SourceBinding[] = [ { source: editor }, { source: createLocalSource({ store: this.#store, partition: this.#store.partition, sourceId: this.#sourceId, ...(this.#sync ? { sync: this.#sync } : {}), }), }, ]
if (this.#documentSync && this.#rootDocumentId) { sources.push({ source: createPdsSource({ documentSync: this.#documentSync, rootDocumentId: this.#rootDocumentId }), policy: options.pdsPolicy ?? { autoPush: false }, }) }
// Injected Sources (e.g. a libp2p live transport) attach after local + PDS. for (const provider of this.#sourceProviders) sources.push(provider(documentId))
const session = createSession({ documentId, initialDoc, sources }) this.#sessions.set(documentId, { session, handle: editor }) return { session, handle: editor } }
/** Serialize store writes and track them for flush(). */ #enqueue(task: () => Promise<void>): void { const run = this.#writeQueue.then(task, task).catch(() => {}) this.#writeQueue = run this.#pending.add(run) void run.finally(() => this.#pending.delete(run)) }}