import { 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 { 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>, ): 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 { 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 => { 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 { 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, ); } }