From 8a517ab5e9ec70e540025d0202ea29ba0d261970 Mon Sep 17 00:00:00 2001 From: Ethan Graf Date: Mon, 25 May 2026 18:38:04 -0400 Subject: [PATCH] Fix p2p multiplayer connections --- src/habitat/connectionProvider.ts | 7 +-- src/habitat/libp2pNode.ts | 89 ++++++++++++++----------------- 2 files changed, 44 insertions(+), 52 deletions(-) diff --git a/src/habitat/connectionProvider.ts b/src/habitat/connectionProvider.ts index b3b74e3..53f9e4f 100644 --- a/src/habitat/connectionProvider.ts +++ b/src/habitat/connectionProvider.ts @@ -99,10 +99,12 @@ export class Libp2pConnectionProvider node.services.pubsub.addEventListener('message', (message) => { if (message.detail.topic !== this.topic) return; + console.log('[habitat:sync] received pubsub message'); const decoder = decoding.createDecoder(message.detail.data); const messageType = decoding.readVarUint(decoder); switch (messageType) { case messageTypes.sync: { + console.log('[habitat:sync] processing sync message'); const encoder = encoding.createEncoder(); encoding.writeVarUint(encoder, messageTypes.sync); syncProtocol.readSyncMessage(decoder, encoder, doc, this); @@ -127,14 +129,12 @@ export class Libp2pConnectionProvider }); // Send sync step 1 when a peer joins this specific topic's mesh. - // gossipsub:graft fires after the mesh is established so publish is - // guaranteed to have at least one recipient. The topic filter prevents - // spurious syncs for unrelated gossipsub-internal topics. node.services.pubsub.addEventListener('gossipsub:graft', (event) => { const { topic } = ( event as CustomEvent<{ peerId: unknown; topic: string }> ).detail; if (topic !== this.topic) return; + console.log('[habitat:sync] gossipsub:graft — sending sync step 1'); const encoder = encoding.createEncoder(); encoding.writeVarUint(encoder, messageTypes.sync); syncProtocol.writeSyncStep1(encoder, doc); @@ -167,6 +167,7 @@ export class Libp2pConnectionProvider doc.on('update', this.handleDocUpdate); this.awareness.on('update', this.handleAwarenessUpdate); + console.log('[habitat:sync] subscribing to topic', this.topic); node.services.pubsub.subscribe(this.topic); } diff --git a/src/habitat/libp2pNode.ts b/src/habitat/libp2pNode.ts index 5ce2c37..c425dca 100644 --- a/src/habitat/libp2pNode.ts +++ b/src/habitat/libp2pNode.ts @@ -16,15 +16,7 @@ import type { HabitatLibp2pNode } from './connectionProvider'; import { query } from './xrpc'; /** - * Build a libp2p node identical to upstream's docs app: - * - listen on `/p2p-circuit` (for circuit-relay-v2) and `/webrtc` (for inbound WebRTC). - * - transports: WebRTC, WebTransport, WebSockets, circuit-relay-v2. - * - noise + yamux + identify + gossipsub. - * - * The gossipsub config matches upstream: - * - `runOnLimitedConnection: true` so we can use the relay's limited connection slot. - * - `allowPublishToZeroTopicPeers: true` so a publish before mesh formation doesn't throw. - * - `emitSelf: false` so we don't re-process our own publishes (Y handles origin). + * Build a libp2p node identical to upstream's docs app. */ export async function createHabitatLibp2pNode(): Promise { const node = await createLibp2p({ @@ -56,28 +48,6 @@ export async function createHabitatLibp2pNode(): Promise { return node as unknown as HabitatLibp2pNode; } -/** - * Hide `node.components` from `for..in` enumeration without changing - * property-access semantics. Required because: - * - * - `node.components` is a `Proxy` whose `get` handler throws - * `MissingServiceError("$prop not set")` for any property name that - * isn't a registered libp2p service — including `$$typeof`. - * - React 19's dev-mode render logger (`addObjectToProperties` in - * `react-dom-client.development.js`) walks every prop value three - * levels deep with `for..in`, and at each step checks - * `value.$$typeof` to decide whether the object is a React element. - * - When the libp2p node is held (transitively) by a React prop (our - * `Libp2pConnectionProvider` → `node` → `components`), the 3rd-level - * walk hits `components.$$typeof`, the proxy throws, and React's - * commit phase crashes — surfacing as "second doc never loads after - * a switch". - * - * `for..in` skips non-enumerable own properties, so toggling enumerability - * is the minimal surgical fix. Libp2p's internal `this.components.xxx` - * accesses are unaffected because property access doesn't check the - * enumerability flag. - */ function hideComponentsFromIntrospection( node: Awaited>, ): void { @@ -91,28 +61,27 @@ function hideComponentsFromIntrospection( }); } -/** Default relay multiaddr derived from `HABITAT_RELAY_DOMAIN`. */ +/** Default relay multiaddr. */ export function relayMultiaddr() { return multiaddr(`/dns4/${HABITAT_RELAY_DOMAIN}/tcp/443/wss`); } /** - * Open the `/habitat/peer-discovery/1.0.0` stream to `relayPeerId`, hand the - * relay our serviceauth token for `HABITAT_DID`, then read newline-delimited - * peer-id strings and dial each one via a circuit address. libp2p upgrades - * those to direct WebRTC connections so gossipsub messages flow without - * passing through the relay. + * 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', { @@ -136,19 +105,31 @@ export async function startPeerDiscovery( const dialPeer = async (peerIdStr: string): Promise => { if (peerIdStr === node.peerId.toString()) return; if (peerIdStr === relayPeerId.toString()) return; - if (node.getConnections(peerIdFromString(peerIdStr)).length > 0) 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.log('[habitat] error dialing peer', peerIdStr, 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 }); @@ -157,36 +138,46 @@ export async function startPeerDiscovery( for (const line of lines) { const id = line.trim(); if (id) { + peersReceived++; + console.log('[habitat] peer discovery received', id); dialPeer(id).catch((e) => { - console.log('[habitat] error dialing peer', id, 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 stream ended', e); + console.warn('[habitat] peer discovery failed', e); } } /** - * Idempotent: returns immediately if we're already connected to the relay. - * Mirrors upstream so it can be called both at session open and on every - * `visibilitychange` to "visible" to recover after a sleep/wake cycle. + * Dial relay if not already connected, then start peer discovery. */ export async function dialRelayAndStartPeerDiscovery( uri: string, node: HabitatLibp2pNode, ): Promise { const relayAddr = relayMultiaddr(); - const connected = node + const openConn = node .getConnections() - .some((conn) => conn.remoteAddr.toString() === relayAddr.toString()); - if (connected) return; + .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); - const relayPeerId = conn.remotePeer; - void startPeerDiscovery(uri, relayPeerId, node); + 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', -- 2.51.2