import { 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 /** 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 | undefined readonly #rootDocumentId: DocumentId | undefined readonly #sync: LocalSyncChannel | undefined readonly #sourceProviders: readonly SourceProvider[] readonly #sourceId = generateDocumentId() readonly #sessions = new Map() readonly #pending = new Set>() #writeQueue: Promise = 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(options: OpenSessionOptions = {}): Handle { const documentId = generateDocumentId() const initialDoc = emptyChange(init(), 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(documentId, initialDoc, options).handle } async find(id: DocumentId, options: OpenSessionOptions = {}): Promise> { const documentId = validateDocumentId(id) const existing = this.#sessions.get(documentId) if (existing) return existing.handle as Handle let doc = await this.#store.loadDoc(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 | undefined) ?? null } if (!doc) throw new Error(`Document ${documentId} not found`) const { session, handle } = this.#buildSession(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 { return this.#store.listDocumentIds() } async flush(): Promise { for (const { session } of this.#sessions.values()) await session.flush() while (this.#pending.size > 0) await Promise.all([...this.#pending]) } async close(): Promise { for (const { session } of this.#sessions.values()) { await session.flush() session.dispose() } this.#sessions.clear() await this.flush() } // --- internals --- #buildSession( documentId: DocumentId, initialDoc: Doc, options: OpenSessionOptions, ): { session: Session; handle: Handle } { // The editor is itself a Source (first in the list) and the client-facing Handle. const editor = createEditorSource(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 { const run = this.#writeQueue.then(task, task).catch(() => {}) this.#writeQueue = run this.#pending.add(run) void run.finally(() => this.#pending.delete(run)) } }