wilooper's picture
initial commit
87cb242
Raw History Blame Contribute Delete
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';
}
}