import { getHeads, load, merge, save, type Doc } from '@automerge/automerge' import type { DocumentId } from '../document-id.js' import type { Source } from './source.js' import { defaultSourcePolicy, type SourcePolicy } from './policy.js' /** * A Source paired with the policy that governs it. Binding them together (rather * than a separate id-keyed policy map) makes it impossible to attach a policy to * a Source that isn't present, or to forget one. */ export interface SourceBinding { readonly source: Source /** Governs push/pull for this Source; omitted → {@link defaultSourcePolicy}. */ readonly policy?: SourcePolicy } /** * Produces a {@link SourceBinding} for a specific document — the injection point * for additional Sources (e.g. a libp2p live transport, or a custom backend) that * a {@link SessionHost} attaches to every session it opens. */ export type SourceProvider = (documentId: DocumentId) => SourceBinding export interface SessionOptions { /** The single document this session coordinates. */ readonly documentId: DocumentId /** The document's starting state (loaded locally, reconstructed, or a fresh doc). */ readonly initialDoc: Doc /** The Sources (each with its policy) that keep this document in sync. */ readonly sources: readonly SourceBinding[] } /** * A per-document sync coordinator. Its surface is *purely* coordination — the * document is read and edited through a {@link Handle}, which is itself one of the * Sources (see `createEditorSource`). The session knows nothing about clients or * transports; it only merges Sources into a thin canonical and fans changes out. */ export interface Session { readonly documentId: DocumentId /** Pull each Source into the canonical (load / catch-up), then fan out if advanced. */ pull(): Promise /** Flush any debounced pushes now and await them. */ flush(): Promise /** Stop coordinating: cancel timers and unsubscribe from all Sources. */ dispose(): void } function sameHeads(a: readonly string[], b: readonly string[]): boolean { return a.length === b.length && a.every((head, i) => head === b[i]) } /** * Holds a **thin canonical** doc (the merge point). A change entering from any * Source is merged into the canonical and, **only if the canonical's heads * advance**, pushed to the *other* Sources so they converge. * * Loops are prevented structurally: the session never pushes back to the Source * that delivered a change, and it propagates only when a merge advanced the * canonical. */ export function createSession(options: SessionOptions): Session { const { documentId } = options const sources = options.sources.map((binding) => binding.source) const policyById = new Map( options.sources.map((binding) => [binding.source.id, binding.policy ?? defaultSourcePolicy]), ) // Own an independent copy — never share a doc object with a Source. let canonical: Doc = load(save(options.initialDoc)) const timers = new Map>() const pending = new Set>() const unsubscribes: Array<() => void> = [] let disposed = false const policyFor = (sourceId: string): SourcePolicy => policyById.get(sourceId) ?? defaultSourcePolicy /** Merge an incoming doc into the canonical; return true iff heads advanced. */ function mergeIntoCanonical(incoming: Doc): boolean { const own = load(save(incoming)) const beforeHeads = getHeads(canonical) canonical = merge(canonical, own) return !sameHeads(beforeHeads, getHeads(canonical)) } function runPush(source: Source): void { const p = source.push(documentId, canonical).catch(() => {}) pending.add(p) void p.finally(() => pending.delete(p)) } function schedulePush(source: Source): void { if (disposed) return if (policyFor(source.id).autoPush === false) return // neutral: no automatic propagation const existing = timers.get(source.id) if (existing) clearTimeout(existing) const debounce = policyFor(source.id).pushDebounceMs ?? 0 if (debounce <= 0) { timers.delete(source.id) runPush(source) return } timers.set( source.id, setTimeout(() => { timers.delete(source.id) runPush(source) }, debounce), ) } /** A source delivered a change: merge it, and (structurally loop-safe) fan out to the rest. */ function onInbound(from: Source, id: DocumentId, doc: Doc): void { if (disposed) return if (id !== documentId) return // a multi-doc source (cross-tab) delivered another doc — ignore if (!mergeIntoCanonical(doc)) return // no advance → nothing to propagate (no loop) for (const source of sources) { if (source.id === from.id) continue // never push back to the deliverer (no loop) schedulePush(source) } } for (const source of sources) { unsubscribes.push(source.subscribe((id, doc) => onInbound(source, id, doc))) } return { documentId, async pull(): Promise { let advanced = false for (const source of sources) { if (policyFor(source.id).pullOnOpen === false) continue const pulled = await source.pull(documentId) if (pulled && mergeIntoCanonical(pulled)) advanced = true } if (advanced) for (const source of sources) schedulePush(source) }, async flush(): Promise { for (const [id, timer] of [...timers]) { clearTimeout(timer) timers.delete(id) const source = sources.find((s) => s.id === id) if (source) runPush(source) } while (pending.size > 0) await Promise.all([...pending]) }, dispose(): void { disposed = true for (const timer of timers.values()) clearTimeout(timer) timers.clear() for (const unsubscribe of unsubscribes) unsubscribe() }, } }