/** * P2P Signal Client — Browser Reference Implementation * * Drop this into your frontend (React/Vue/Vanilla). * It handles: * - WS connection + JWT auth * - Heartbeat loop * - WebRTC offer/answer/ICE * - Auto-fallback to WS relay on ICE failure * - File upload via Supabase signed URL * - Offline message receipt on reconnect * * Usage: * const client = new P2PClient({ serverUrl, token, onMessage, onPeerOnline, onPeerOffline }); * await client.connect(); * await client.sendMessage({ toUid, type: 'text', content: 'Hello!' }); * await client.callPeer(peerUid); // initiates WebRTC */ export interface P2PClientOptions { serverUrl: string; // wss://your-space.hf.space token: string; // JWT from your auth system onMessage?: (msg: EnvelopeMessage) => void; onPeerOnline?: (uid: string) => void; onPeerOffline?: (uid: string) => void; onOfflineFlush?: (messages: EnvelopeMessage[]) => void; onP2PData?: (fromUid: string, data: unknown) => void; // P2P channel data heartbeatIntervalMs?: number; } export interface EnvelopeMessage { fromUid: string; toUid: string; type: 'text' | 'media'; content: string; mediaType?: string; sentAt: number; } interface IceServerConfig { urls: string; username?: string; credential?: string; } export class P2PClient { private ws: WebSocket | null = null; private myUid: string | null = null; private socketId: string | null = null; private iceServers: IceServerConfig[] = []; private heartbeatTimer: ReturnType | null = null; private reconnectTimer: ReturnType | null = null; private reconnectAttempts = 0; // uid → RTCPeerConnection private peerConnections = new Map(); // uid → RTCDataChannel (P2P) private dataChannels = new Map(); // uid → 'p2p' | 'relay' (fallback mode) private peerMode = new Map(); private opts: Required; constructor(opts: P2PClientOptions) { this.opts = { heartbeatIntervalMs: 30_000, onMessage: () => {}, onPeerOnline: () => {}, onPeerOffline: () => {}, onOfflineFlush: () => {}, onP2PData: () => {}, ...opts, }; } // ── Connection ───────────────────────────────────────────────────────────── connect(): Promise { return new Promise((resolve, reject) => { const url = `${this.opts.serverUrl}?token=${encodeURIComponent(this.opts.token)}`; this.ws = new WebSocket(url); this.ws.binaryType = 'arraybuffer'; this.ws.onopen = () => { this.reconnectAttempts = 0; this.startHeartbeat(); }; this.ws.onmessage = (evt) => { const msg = this.parseServer(evt.data as string); if (!msg) return; this.handleServer(msg, resolve, reject); }; this.ws.onclose = () => { this.stopHeartbeat(); this.scheduleReconnect(); }; this.ws.onerror = () => { reject(new Error('WebSocket connection failed')); }; }); } disconnect(): void { this.stopHeartbeat(); if (this.reconnectTimer) clearTimeout(this.reconnectTimer); this.ws?.close(); } // ── Heartbeat ────────────────────────────────────────────────────────────── private startHeartbeat(): void { this.heartbeatTimer = setInterval(() => { this.send('heartbeat', {}); }, this.opts.heartbeatIntervalMs); } private stopHeartbeat(): void { if (this.heartbeatTimer) { clearInterval(this.heartbeatTimer); this.heartbeatTimer = null; } } // ── Reconnect ────────────────────────────────────────────────────────────── private scheduleReconnect(): void { const delay = Math.min(1000 * 2 ** this.reconnectAttempts, 30_000); this.reconnectAttempts++; this.reconnectTimer = setTimeout(() => { this.connect().catch(() => {}); // Will schedule again on failure }, delay); } // ── Sending ──────────────────────────────────────────────────────────────── private send(type: string, payload: Record, requestId?: string): void { if (!this.ws || this.ws.readyState !== WebSocket.OPEN) return; this.ws.send(JSON.stringify({ type, payload, requestId })); } /** Send a text or media message to a peer */ sendMessage(opts: { toUid: string; type: 'text' | 'media'; content: string; mediaType?: string; }): void { // If P2P channel is open and healthy, use it const dc = this.dataChannels.get(opts.toUid); const mode = this.peerMode.get(opts.toUid); if (mode === 'p2p' && dc?.readyState === 'open') { dc.send(JSON.stringify({ ...opts, fromUid: this.myUid, sentAt: Date.now() })); return; } // Otherwise use WS signaling server (relay or direct delivery) this.send('message', { toUid: opts.toUid, type: opts.type, content: opts.content, mediaType: opts.mediaType ?? null, clientMsgId: crypto.randomUUID(), }); } // ── File upload ──────────────────────────────────────────────────────────── /** * Request signed upload URLs and upload files directly to Supabase. * Returns array of public URLs for use in sendMessage({ type: 'media' }). */ async uploadFiles(files: File[]): Promise { const requestId = crypto.randomUUID(); // Request signed URLs from server const urlsPromise = this.waitForRequestId<{ urls: Array<{ signedUrl: string; publicUrl: string; fileName: string }> }>( requestId ); this.send('get-upload-url', { files: files.map(f => ({ fileName: f.name, mimeType: f.type, sizeBytes: f.size, })), }, requestId); const { urls } = await urlsPromise; // Upload each file directly to Supabase (server sees 0 bytes) const publicUrls = await Promise.all( urls.map(async ({ signedUrl, publicUrl, fileName }, i) => { const file = files.find(f => f.name === fileName) ?? files[i]; await fetch(signedUrl, { method: 'PUT', headers: { 'Content-Type': file.type }, body: file, }); return publicUrl; }) ); return publicUrls; } // ── WebRTC ───────────────────────────────────────────────────────────────── /** Initiate a WebRTC connection to peerUid */ async callPeer(peerUid: string): Promise { const pc = this.createPeerConnection(peerUid); // Create data channel (caller side) const dc = pc.createDataChannel('chat', { ordered: true }); this.setupDataChannel(peerUid, dc); const offer = await pc.createOffer(); await pc.setLocalDescription(offer); this.send('offer', { toUid: peerUid, sdp: pc.localDescription, }); } private createPeerConnection(peerUid: string): RTCPeerConnection { // Close existing if any this.peerConnections.get(peerUid)?.close(); const pc = new RTCPeerConnection({ iceServers: this.iceServers }); this.peerConnections.set(peerUid, pc); this.peerMode.set(peerUid, 'p2p'); // ICE candidate forwarding pc.onicecandidate = (evt) => { if (evt.candidate) { this.send('ice-candidate', { toUid: peerUid, candidate: evt.candidate }); } }; // ── Fallback: if ICE fails, switch to WS relay ──────────────────────── pc.oniceconnectionstatechange = () => { if (pc.iceConnectionState === 'failed' || pc.iceConnectionState === 'disconnected') { console.warn(`[P2P] ICE failed for ${peerUid} — switching to WS relay`); this.peerMode.set(peerUid, 'relay'); // Close the dead connection pc.close(); this.peerConnections.delete(peerUid); this.dataChannels.delete(peerUid); } }; // Callee side: accept incoming data channel pc.ondatachannel = (evt) => { this.setupDataChannel(peerUid, evt.channel); }; return pc; } private setupDataChannel(peerUid: string, dc: RTCDataChannel): void { this.dataChannels.set(peerUid, dc); dc.onopen = () => { console.info(`[P2P] Data channel open with ${peerUid}`); this.peerMode.set(peerUid, 'p2p'); }; dc.onmessage = (evt) => { try { const data = JSON.parse(evt.data as string) as unknown; this.opts.onP2PData(peerUid, data); } catch { this.opts.onP2PData(peerUid, evt.data); } }; dc.onclose = () => { console.info(`[P2P] Data channel closed with ${peerUid}`); this.dataChannels.delete(peerUid); // Fall back to WS relay this.peerMode.set(peerUid, 'relay'); }; } // ── Server message handler ──────────────────────────────────────────────── private requestCallbacks = new Map void>(); private waitForRequestId(requestId: string): Promise { return new Promise((resolve, reject) => { const timer = setTimeout(() => { this.requestCallbacks.delete(requestId); reject(new Error(`Request ${requestId} timed out`)); }, 15_000); this.requestCallbacks.set(requestId, (payload) => { clearTimeout(timer); resolve(payload as T); }); }); } private parseServer(raw: string): Record | null { try { return JSON.parse(raw) as Record; } catch { return null; } } private async handleServer( msg: Record, onConnected: () => void, onConnectError: (e: Error) => void ): Promise { const type = msg.type as string; const payload = (msg.payload ?? {}) as Record; const requestId = msg.requestId as string | undefined; // Resolve pending request callbacks if (requestId && this.requestCallbacks.has(requestId)) { this.requestCallbacks.get(requestId)!(payload); this.requestCallbacks.delete(requestId); } switch (type) { case 'connected': this.myUid = payload.uid as string; this.socketId = payload.socketId as string; onConnected(); break; case 'ice-servers': this.iceServers = payload.iceServers as IceServerConfig[]; break; case 'message': this.opts.onMessage(payload as unknown as EnvelopeMessage); break; case 'offline-flush': { const messages = payload.messages as EnvelopeMessage[]; this.opts.onOfflineFlush(messages); break; } case 'peer-online': this.opts.onPeerOnline(payload.uid as string); break; case 'peer-offline': { const offlineUid = payload.uid as string; this.opts.onPeerOffline(offlineUid); // Switch to relay mode for this peer this.peerMode.set(offlineUid, 'relay'); break; } case 'offer': { const fromUid = payload.fromUid as string; const pc = this.createPeerConnection(fromUid); await pc.setRemoteDescription(payload.sdp as RTCSessionDescriptionInit); const answer = await pc.createAnswer(); await pc.setLocalDescription(answer); this.send('answer', { toUid: fromUid, sdp: pc.localDescription }); break; } case 'answer': { const fromUid = payload.fromUid as string; const pc = this.peerConnections.get(fromUid); if (pc) await pc.setRemoteDescription(payload.sdp as RTCSessionDescriptionInit); break; } case 'ice-candidate': { const fromUid = payload.fromUid as string; const pc = this.peerConnections.get(fromUid); if (pc && payload.candidate) { await pc.addIceCandidate(payload.candidate as RTCIceCandidateInit); } break; } case 'relay': // Message arrived via WS relay fallback this.opts.onMessage(payload as unknown as EnvelopeMessage); break; case 'error': console.error('[P2P Server Error]', payload.message); break; } } // ── Friend helpers ──────────────────────────────────────────────────────── sendFriendRequest(toUid: string): void { this.send('friend', { action: 'send-request', targetUid: toUid }); } acceptFriendRequest(fromUid: string): void { this.send('friend', { action: 'accept', targetUid: fromUid }); } rejectFriendRequest(fromUid: string): void { this.send('friend', { action: 'reject', targetUid: fromUid }); } removeFriend(uid: string): void { this.send('friend', { action: 'remove', targetUid: uid }); } async getFriends(): Promise> { const requestId = crypto.randomUUID(); const result = await this.waitForRequestId<{ friends: Array<{ uid: string; online: boolean }> }>(requestId); this.send('friend', { action: 'list' }, requestId); return result.friends; } // ── Getters ─────────────────────────────────────────────────────────────── get uid(): string | null { return this.myUid; } get connected(): boolean { return this.ws?.readyState === WebSocket.OPEN; } getPeerMode(peerUid: string): 'p2p' | 'relay' | 'none' { return this.peerMode.get(peerUid) ?? 'none'; } }