Download src/handlers/connectionHandler.ts from Collabos/chatapi: direct link, hf CLI and curl.
- Browser
- Download file 5.35 kB
-
https://huggingface.co/spaces/Collabos/chatapi/resolve/main/src/handlers/connectionHandler.ts
- Command line
-
hf download hf://spaces/Collabos/chatapi/src/handlers/connectionHandler.ts
-
curl -L -o connectionHandler.ts https://huggingface.co/spaces/Collabos/chatapi/resolve/main/src/handlers/connectionHandler.ts
5.35 kB
| 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<void> { | |
| // ββ 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<void> { | |
| 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<void> { | |
| 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); | |
| } | |