Download client-reference/client.ts from Collabos/chatapi: direct link, hf CLI and curl.
- Browser
- Download file 14.4 kB
-
https://huggingface.co/spaces/Collabos/chatapi/resolve/main/client-reference/client.ts
- Command line
-
hf download hf://spaces/Collabos/chatapi/client-reference/client.ts
-
curl -L -o client.ts https://huggingface.co/spaces/Collabos/chatapi/resolve/main/client-reference/client.ts
14.4 kB
| /** | |
| * 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<typeof setInterval> | null = null; | |
| private reconnectTimer: ReturnType<typeof setTimeout> | null = null; | |
| private reconnectAttempts = 0; | |
| // uid β RTCPeerConnection | |
| private peerConnections = new Map<string, RTCPeerConnection>(); | |
| // uid β RTCDataChannel (P2P) | |
| private dataChannels = new Map<string, RTCDataChannel>(); | |
| // uid β 'p2p' | 'relay' (fallback mode) | |
| private peerMode = new Map<string, 'p2p' | 'relay'>(); | |
| private opts: Required<P2PClientOptions>; | |
| constructor(opts: P2PClientOptions) { | |
| this.opts = { | |
| heartbeatIntervalMs: 30_000, | |
| onMessage: () => {}, | |
| onPeerOnline: () => {}, | |
| onPeerOffline: () => {}, | |
| onOfflineFlush: () => {}, | |
| onP2PData: () => {}, | |
| ...opts, | |
| }; | |
| } | |
| // ββ Connection βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| connect(): Promise<void> { | |
| 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<string, unknown>, 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<string[]> { | |
| 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<void> { | |
| 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<string, (payload: unknown) => void>(); | |
| private waitForRequestId<T>(requestId: string): Promise<T> { | |
| 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<string, unknown> | null { | |
| try { return JSON.parse(raw) as Record<string, unknown>; } | |
| catch { return null; } | |
| } | |
| private async handleServer( | |
| msg: Record<string, unknown>, | |
| onConnected: () => void, | |
| onConnectError: (e: Error) => void | |
| ): Promise<void> { | |
| const type = msg.type as string; | |
| const payload = (msg.payload ?? {}) as Record<string, unknown>; | |
| 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<Array<{ uid: string; online: boolean }>> { | |
| 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'; | |
| } | |
| } | |