File size: 5,400 Bytes
87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 eccdcb1 87cb242 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 | import http from 'http';
import WebSocket from 'ws';
import jwt from 'jsonwebtoken';
import { CONFIG } from './config';
import { connectMongo } from './db/mongo';
import { presenceClient } from './redis/presenceClient';
import { messageClient } from './redis/messageClient';
import { onConnection } from './handlers/connectionHandler';
import { registry } from './services/connectionRegistry';
import { isRateLimited } from './middleware/rateLimiter';
import { logger } from './utils/logger';
async function bootstrap(): Promise<void> {
logger.info('Starting P2P Signal Server', { port: CONFIG.PORT, env: CONFIG.NODE_ENV });
await connectMongo();
await Promise.all([
waitRedisReady(presenceClient, 'Presence'),
waitRedisReady(messageClient, 'Message'),
]);
const httpServer = http.createServer((req, res) => {
const url = new URL(req.url ?? '/', `http://${req.headers.host}`);
// ββ CORS headers for browser fetch ββββββββββββββββββββββββββββββββββββββ
res.setHeader('Access-Control-Allow-Origin', '*');
res.setHeader('Access-Control-Allow-Methods', 'GET, POST, OPTIONS');
res.setHeader('Access-Control-Allow-Headers', 'Content-Type');
if (req.method === 'OPTIONS') { res.writeHead(204); res.end(); return; }
// ββ GET /health ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
if (req.method === 'GET' && url.pathname === '/health') {
res.writeHead(200, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ status: 'ok', connections: registry.size(), ts: Date.now() }));
return;
}
// ββ GET /token?uid=user1 β generates a real signed JWT for testing βββββββ
// Simple uid allowlist β only test UIDs work, not arbitrary strings
if (req.method === 'GET' && url.pathname === '/token') {
const uid = url.searchParams.get('uid') ?? '';
const allowed = ['user1', 'user2', 'user3', 'user4', 'user5'];
if (!allowed.includes(uid)) {
res.writeHead(400, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ error: `uid must be one of: ${allowed.join(', ')}` }));
return;
}
const token = jwt.sign({ uid }, CONFIG.JWT_SECRET, { expiresIn: '7d' });
res.writeHead(200, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ token, uid, note: 'Test token β do not use in production' }));
return;
}
// ββ Fallback βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
res.writeHead(200, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({
service: 'P2P Signal Server',
endpoints: {
health: 'GET /health',
token: 'GET /token?uid=user1 (test only)',
ws: 'WS /?token=<jwt>',
}
}));
});
const wss = new WebSocket.Server({
server: httpServer,
perMessageDeflate: false,
maxPayload: 65536,
handleProtocols: () => false,
});
wss.on('connection', async (ws: WebSocket, req: http.IncomingMessage) => {
if (isRateLimited(req)) { ws.close(1008, 'Rate limit exceeded'); return; }
const origin = req.headers.origin ?? '';
if (CONFIG.ALLOWED_ORIGINS[0] !== '*' && !CONFIG.ALLOWED_ORIGINS.includes(origin)) {
ws.close(1008, 'Origin not allowed');
return;
}
await onConnection(ws, req);
});
wss.on('error', (err) => logger.error('WSS error', { error: err.message }));
const ZOMBIE_TIMEOUT_MS = CONFIG.HEARTBEAT_INTERVAL_MS * 3;
setInterval(() => {
const pruned = registry.pruneZombies(ZOMBIE_TIMEOUT_MS);
if (pruned > 0) logger.warn('Pruned zombie sockets', { pruned });
}, ZOMBIE_TIMEOUT_MS);
setInterval(() => {
logger.info('Server metrics', {
connections: registry.size(),
memMB: Math.round(process.memoryUsage().rss / 1024 / 1024),
});
}, 60_000);
httpServer.listen(CONFIG.PORT, '0.0.0.0', () => {
logger.info(`Server listening on port ${CONFIG.PORT}`);
});
}
// eslint-disable-next-line @typescript-eslint/no-explicit-any
async function waitRedisReady(client: any, name: string): Promise<void> {
return new Promise((resolve, reject) => {
if (client.status === 'ready') { resolve(); return; }
client.once('ready', resolve);
client.once('error', (err: Error) => reject(new Error(`${name} Redis failed: ${err.message}`)));
setTimeout(() => reject(new Error(`${name} Redis connection timeout`)), 10_000);
});
}
process.on('unhandledRejection', (reason) => logger.error('Unhandled rejection', { reason: String(reason) }));
process.on('uncaughtException', (err) => logger.error('Uncaught exception', { error: err.message, stack: err.stack }));
async function shutdown(signal: string): Promise<void> {
logger.info(`Received ${signal}, shutting down`);
presenceClient.disconnect();
messageClient.disconnect();
process.exit(0);
}
process.on('SIGTERM', () => shutdown('SIGTERM'));
process.on('SIGINT', () => shutdown('SIGINT'));
bootstrap().catch((err) => {
logger.error('Bootstrap failed', { error: String(err) });
process.exit(1);
});
|