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