chatapi / src /handlers /connectionHandler.ts
wilooper's picture
initial commit
87cb242
Raw History Blame Contribute Delete
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);
}