Something went wrong. Try again.
A local-first note taking app
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146import type { GossipsubEvents } from '@chainsafe/libp2p-gossipsub';import type { PubSub } from '@libp2p/interface';import { Libp2p } from 'libp2p';import * as Automerge from '@automerge/automerge';import * as decoding from 'lib0/decoding';import * as encoding from 'lib0/encoding';import { ObservableV2 } from 'lib0/observable';
import type { AutomergeDocSchema } from './automergeDoc';
/** * Wires Automerge sync messages through a libp2p gossipsub topic. * * Message framing: a single varuint prefix selects sync (0) vs * awareness (1), matching the Yjs path so both share the same * gossipsub transport layer. Awareness is stubbed as `null` for * future cursor support via Automerge ephemeral messages. */
const messageTypes = { sync: 0, awareness: 1,} as const;
type PubsubService = PubSub<GossipsubEvents>;type Node = Libp2p<{ pubsub: PubsubService }>;
/** * Stub surface for future awareness/cursor support. * Mirrors the shape of `CollabCaretProvider` from the Yjs path. */export type AutomergeCaretProvider = { _stub: true;};
export class Libp2pAutomergeConnectionProvider extends ObservableV2<{ 'connection-close': ( event: CloseEvent | null, provider: Libp2pAutomergeConnectionProvider, ) => void; status: (event: { status: 'connected' | 'disconnected' | 'connecting'; }) => void; 'connection-error': ( event: Event, provider: Libp2pAutomergeConnectionProvider, ) => void; sync: (state: boolean) => void;}> { /** Not enumerable — must not appear on React props / DevTools walks. */ #node: Node; topic: string;
/** Stubbed for future awareness/cursor support. */ awareness: null = null;
#syncState: Automerge.SyncState; #doc: Automerge.Doc<AutomergeDocSchema>; #onRemoteChange: (doc: Automerge.Doc<AutomergeDocSchema>) => void;
constructor( node: Node, doc: Automerge.Doc<AutomergeDocSchema>, topic: string, onRemoteChange: (doc: Automerge.Doc<AutomergeDocSchema>) => void, ) { super(); this.#node = node; this.#doc = doc; this.topic = topic; this.#onRemoteChange = onRemoteChange; this.#syncState = Automerge.initSyncState();
node.services.pubsub.addEventListener('message', (message) => { if (message.detail.topic !== this.topic) return; this.#handleMessage(message.detail.data); });
console.log('[automerge:sync] subscribing to topic', this.topic); node.services.pubsub.subscribe(this.topic);
// Send initial sync message when a peer joins this topic's mesh. node.services.pubsub.addEventListener('gossipsub:graft', (event) => { const { topic } = ( event as CustomEvent<{ peerId: unknown; topic: string }> ).detail; if (topic !== this.topic) return; this.#sendSyncMessage(); }); }
#handleMessage(data: Uint8Array): void { const decoder = decoding.createDecoder(data); const messageType = decoding.readVarUint(decoder); if (messageType !== messageTypes.sync) return;
const messageBytes = decoding.readVarUint8Array(decoder); const [newDoc, newSyncState] = Automerge.receiveSyncMessage( this.#doc, this.#syncState, messageBytes, ); this.#doc = newDoc; this.#syncState = newSyncState; this.#onRemoteChange(newDoc);
// Send reply if there's a response message. this.#sendSyncMessage(); }
/** Called by the docSession when the local doc has changed. */ onLocalChange(doc: Automerge.Doc<AutomergeDocSchema>): void { this.#doc = doc; this.#sendSyncMessage(); }
#sendSyncMessage(): void { const [syncState, message] = Automerge.generateSyncMessage( this.#doc, this.#syncState, ); this.#syncState = syncState; if (message) { const encoder = encoding.createEncoder(); encoding.writeVarUint(encoder, messageTypes.sync); encoding.writeVarUint8Array(encoder, message); this.#node.services.pubsub .publish(this.topic, encoding.toUint8Array(encoder)) .catch((err: Error) => { console.error('[pubsub] publish error', err); }); } }
destroy(): void { try { this.#node.services.pubsub.unsubscribe(this.topic); } catch { // node may already be stopped; tearing down anyway. } super.destroy(); }}
export type { Node as HabitatLibp2pNode };