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);
});