import { getHeads, init, merge, type Doc } from '@automerge/automerge' import type { DocumentId } from '../../document-id.js' import type { Source } from '../source.js' type InboundListener = (documentId: DocumentId, doc: Doc) => void /** * A Fake {@link Source} with real in-memory state — for testing the Engine without * any transport. Accumulates pushed docs (via `merge`), serves them on `pull`, and * can simulate an inbound delivery from a peer via {@link deliver}. Assertions read * {@link current} / {@link pushes}. */ export class InMemorySource implements Source { readonly id: string readonly #docs = new Map>() readonly #listeners = new Set() /** One entry per `push`, for asserting call counts (e.g. debounce coalescing). */ readonly pushes: Array<{ documentId: DocumentId; heads: readonly string[] }> = [] constructor(id: string) { this.id = id } async push(documentId: DocumentId, doc: Doc): Promise { const next = merge(this.#docs.get(documentId) ?? init(), doc) this.#docs.set(documentId, next) this.pushes.push({ documentId, heads: getHeads(next) }) } async pull(documentId: DocumentId): Promise | null> { return this.#docs.get(documentId) ?? null } subscribe(onInbound: InboundListener): () => void { this.#listeners.add(onInbound) return () => this.#listeners.delete(onInbound) } // --- test helpers --- /** Simulate this source receiving a change from its transport and delivering it. */ deliver(documentId: DocumentId, doc: Doc): void { const next = merge(this.#docs.get(documentId) ?? init(), doc) this.#docs.set(documentId, next) for (const listener of this.#listeners) listener(documentId, next) } /** Seed a doc as if it already existed here (drives a later `pull`). */ seed(documentId: DocumentId, doc: Doc): void { this.#docs.set(documentId, doc) } current(documentId: DocumentId): Doc | null { return this.#docs.get(documentId) ?? null } }