Something went wrong. Try again.
automerge-repo drop in replacement backed by Habitat "permissioned spaces" om your PDS
Something went wrong. Try again.
5.9 kB · 158 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159import { 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<unknown> /** 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<void> /** Flush any debounced pushes now and await them. */ flush(): Promise<void> /** 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<string, SourcePolicy>( 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<unknown> = load(save(options.initialDoc)) const timers = new Map<string, ReturnType<typeof setTimeout>>() const pending = new Set<Promise<void>>() 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<unknown>): 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<unknown>): 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<void> { 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<void> { 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() }, }}