Something went wrong. Try again.
A local-first note taking app
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188import { gossipsub } from '@chainsafe/libp2p-gossipsub';import { noise } from '@chainsafe/libp2p-noise';import { yamux } from '@chainsafe/libp2p-yamux';import { circuitRelayTransport } from '@libp2p/circuit-relay-v2';import { identify } from '@libp2p/identify';import type { PeerId } from '@libp2p/interface';import { peerIdFromString } from '@libp2p/peer-id';import { webRTC } from '@libp2p/webrtc';import { webSockets } from '@libp2p/websockets';import { webTransport } from '@libp2p/webtransport';import { multiaddr } from '@multiformats/multiaddr';import { createLibp2p } from 'libp2p';
import { HABITAT_DID, HABITAT_RELAY_DOMAIN } from './config';import type { HabitatLibp2pNode } from './automergeConnectionProvider';import { query } from './xrpc';
/** * Build a libp2p node identical to upstream's docs app. */export async function createHabitatLibp2pNode(): Promise<HabitatLibp2pNode> { const node = await createLibp2p({ addresses: { listen: ['/p2p-circuit', '/webrtc'], }, transports: [ webRTC({ rtcConfiguration: { iceServers: [{ urls: ['stun:stun.l.google.com:19302'] }], }, }), webTransport(), webSockets(), circuitRelayTransport(), ], connectionEncrypters: [noise()], streamMuxers: [yamux()], services: { identify: identify(), pubsub: gossipsub({ runOnLimitedConnection: true, allowPublishToZeroTopicPeers: true, emitSelf: false, }), }, }); hideComponentsFromIntrospection(node); return node as unknown as HabitatLibp2pNode;}
function hideComponentsFromIntrospection( node: Awaited<ReturnType<typeof createLibp2p>>,): void { const descriptor = Object.getOwnPropertyDescriptor(node, 'components'); if (!descriptor || !descriptor.enumerable) return; Object.defineProperty(node, 'components', { value: descriptor.value, enumerable: false, writable: descriptor.writable ?? true, configurable: true, });}
/** Default relay multiaddr. */export function relayMultiaddr() { return multiaddr(`/dns4/${HABITAT_RELAY_DOMAIN}/tcp/443/wss`);}
/** * Open peer-discovery stream, send tokens, read peer IDs, dial them. */export async function startPeerDiscovery( uri: string, relayPeerId: PeerId, node: HabitatLibp2pNode,): Promise<void> { console.log('[habitat] starting peer discovery for', uri); try { const stream = await node.dialProtocol( relayPeerId, '/habitat/peer-discovery/1.0.0', ); console.log('[habitat] peer-discovery stream opened');
const { token: serviceAuthToken } = await query( 'com.atproto.server.getServiceAuth', { lxm: 'com.atproto.server.getServiceAuth', aud: HABITAT_DID, }, );
const encoder = new TextEncoder(); stream.sink( (async function* () { yield encoder.encode( JSON.stringify({ topic: uri, serviceauth_token: serviceAuthToken, }), ); })(), );
const dialPeer = async (peerIdStr: string): Promise<void> => { if (peerIdStr === node.peerId.toString()) return; if (peerIdStr === relayPeerId.toString()) return; // CRITICAL FIX: only skip if there is an OPEN connection. // Previously this checked .length > 0 which included closed/stale // connections, causing us to skip dialing when no usable path existed. const openConns = node .getConnections(peerIdFromString(peerIdStr)) .filter((c) => c.status === 'open'); if (openConns.length > 0) { console.log('[habitat] already connected to', peerIdStr); return; } const circuitAddr = multiaddr( `/p2p/${relayPeerId}/p2p-circuit/p2p/${peerIdStr}`, ); console.log('[habitat] dialing peer via relay', peerIdStr); try { await node.dial(circuitAddr); console.log('[habitat] successfully dialed peer', peerIdStr); } catch (e) { console.warn('[habitat] failed to dial peer', peerIdStr, e); } };
const decoder = new TextDecoder(); let buf = ''; let peersReceived = 0; for await (const chunk of stream.source) { const bytes = chunk instanceof Uint8Array ? chunk : chunk.subarray(); buf += decoder.decode(bytes, { stream: true }); const lines = buf.split('\n'); buf = lines.pop() ?? ''; for (const line of lines) { const id = line.trim(); if (id) { peersReceived++; console.log('[habitat] peer discovery received', id); dialPeer(id).catch((e) => { console.warn('[habitat] error dialing peer', id, e); }); } } } console.log( '[habitat] peer discovery stream closed; received', peersReceived, 'peers', ); } catch (e) { console.warn('[habitat] peer discovery failed', e); }}
/** * Dial relay if not already connected, then start peer discovery. */export async function dialRelayAndStartPeerDiscovery( uri: string, node: HabitatLibp2pNode,): Promise<void> { const relayAddr = relayMultiaddr(); const openConn = node .getConnections() .find((conn) => conn.remoteAddr.toString() === relayAddr.toString()); if (openConn && openConn.status === 'open') { console.log('[habitat] relay already connected, re-running discovery'); void startPeerDiscovery(uri, openConn.remotePeer, node); return; }
console.log('[habitat] dialing relay', relayAddr.toString()); try { const conn = await node.dial(relayAddr); console.log('[habitat] relay connected', conn.remotePeer.toString()); void startPeerDiscovery(uri, conn.remotePeer, node); } catch (e) { console.error( '[habitat] unable to dial relay; continuing without live collaboration', e, ); }}