import WebSocket from 'ws'; import { IncomingMessage } from 'http'; import { v4 as uuidv4 } from 'uuid'; import { verifyToken, extractToken } from '../utils/jwt'; import { safeSend, parseClientMessage, ClientMessage } from '../utils/protocol'; import { setOnline, setOffline, refreshPresence } from '../redis/presenceClient'; import { flushOfflineQueue } from '../redis/messageClient'; import { registry } from '../services/connectionRegistry'; import { buildIceServers } from '../config'; import { logger } from '../utils/logger'; import { handleSignaling } from './signalingHandler'; import { handleMessage } from './messageHandler'; import { handleFileRequest } from './fileHandler'; import { handleFriend } from './friendHandler'; export async function onConnection( ws: WebSocket, req: IncomingMessage ): Promise { // ── 1. Authenticate ────────────────────────────────────────────────────── let uid: string; try { const url = new URL(req.url ?? '/', `http://${req.headers.host}`); const token = extractToken( req.headers.authorization, url.searchParams.get('token') ?? undefined ); const payload = verifyToken(token); uid = payload.uid; } catch (err) { const msg = err instanceof Error ? err.message : 'Auth failed'; ws.send(JSON.stringify({ type: 'error', payload: { message: msg }, ts: Date.now() })); ws.terminate(); return; } const socketId = uuidv4(); // ── 2. Register in process registry + Redis Bucket 1 ──────────────────── registry.add(socketId, uid, ws); // Redis errors must NOT crash the connection handler try { await setOnline(uid, socketId); } catch (err) { logger.error('setOnline failed', { uid, error: String(err) }); } logger.info('Client connected', { uid, socketId, total: registry.size() }); // ── 3. Send welcome + ICE server config ───────────────────────────────── safeSend(ws, 'connected', { uid, socketId }); safeSend(ws, 'ice-servers', { iceServers: buildIceServers() }); // ── 4. Flush offline message queue (Bucket 2) ──────────────────────────── try { const pending = await flushOfflineQueue(uid); if (pending.length > 0) { safeSend(ws, 'offline-flush', { messages: pending, count: pending.length }); logger.info('Flushed offline queue', { uid, count: pending.length }); } } catch (err) { logger.error('Offline flush failed', { uid, error: String(err) }); } // ── 5. Message router ──────────────────────────────────────────────────── ws.on('message', (raw) => { const msg = parseClientMessage(raw); if (!msg) return; // Malformed — silently drop registry.heartbeat(socketId); // Any message counts as a heartbeat routeMessage(ws, uid, socketId, msg).catch((err: unknown) => { logger.error('Message handler error', { uid, type: msg.type, error: String(err) }); }); }); // ── 6. Cleanup on disconnect ───────────────────────────────────────────── ws.on('close', () => { registry.remove(socketId); setOffline(uid).catch((err: unknown) => logger.error('setOffline failed', { uid, error: String(err) }) ); logger.info('Client disconnected', { uid, socketId, total: registry.size() }); }); ws.on('error', (err) => { logger.error('WebSocket error', { uid, socketId, error: err.message }); ws.terminate(); }); } // ── Message router ──────────────────────────────────────────────────────────── async function routeMessage( ws: WebSocket, uid: string, socketId: string, msg: ClientMessage ): Promise { switch (msg.type) { case 'heartbeat': await refreshPresence(uid).catch(() => {}); // Best effort break; case 'offer': case 'answer': case 'ice-candidate': await handleSignaling(ws, uid, msg); break; case 'message': await handleMessage(ws, uid, msg); break; case 'get-upload-url': await handleFileRequest(ws, uid, msg); break; case 'relay': await handleRelay(ws, uid, msg); break; case 'friend': await handleFriend(ws, uid, msg); break; default: safeSend(ws, 'error', { message: `Unknown message type: ${msg.type}` }); } } // ── WebSocket relay fallback (when WebRTC ICE fails) ───────────────────────── async function handleRelay( _senderWs: WebSocket, fromUid: string, msg: ClientMessage ): Promise { const toUid = msg.payload.toUid as string | undefined; if (!toUid || typeof toUid !== 'string') return; const target = registry.getByUid(toUid); if (!target) return; safeSend(target.ws, 'relay', { fromUid, data: msg.payload.data, }, msg.requestId); }