added /token
Browse files- README.md +191 -8
- src/index.ts +47 -64
README.md
CHANGED
|
@@ -1,11 +1,194 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
---
|
| 2 |
-
|
| 3 |
-
|
| 4 |
-
|
| 5 |
-
|
| 6 |
-
|
| 7 |
-
|
| 8 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 9 |
---
|
| 10 |
|
| 11 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
# P2P Signal Server
|
| 2 |
+
|
| 3 |
+
Real-time WebRTC signaling server for 10,000 concurrent users.
|
| 4 |
+
Runs on **Hugging Face Spaces Free Tier** (Docker, port 7860).
|
| 5 |
+
|
| 6 |
+
---
|
| 7 |
+
|
| 8 |
+
## Architecture
|
| 9 |
+
|
| 10 |
+
```
|
| 11 |
+
Client A ββWSβββΊ HF Space (Signal Server)
|
| 12 |
+
β
|
| 13 |
+
ββ Redis Bucket 1 (Presence) uidβsocketId, online/offline
|
| 14 |
+
ββ Redis Bucket 2 (Msg Queue) offline_msgs:{uid} list
|
| 15 |
+
ββ MongoDB Atlas user profiles, friend lists
|
| 16 |
+
ββ Supabase Storage media files (server sees 0 bytes)
|
| 17 |
+
|
| 18 |
+
Client A ββββββββββββββββββββββββββββββββΊ Client B (WebRTC P2P after handshake)
|
| 19 |
+
βββ WS relay fallback if ICE fails ββββββββββΊ
|
| 20 |
+
```
|
| 21 |
+
|
| 22 |
+
---
|
| 23 |
+
|
| 24 |
+
## Quick Start
|
| 25 |
+
|
| 26 |
+
### 1. Clone & install
|
| 27 |
+
|
| 28 |
+
```bash
|
| 29 |
+
git clone https://github.com/you/p2p-signal-server
|
| 30 |
+
cd p2p-signal-server
|
| 31 |
+
npm install
|
| 32 |
+
cp .env.example .env
|
| 33 |
+
# Fill in all values in .env
|
| 34 |
+
npm run dev
|
| 35 |
+
```
|
| 36 |
+
|
| 37 |
+
### 2. Build & run (Docker)
|
| 38 |
+
|
| 39 |
+
```bash
|
| 40 |
+
docker build -t p2p-signal .
|
| 41 |
+
docker run --env-file .env -p 7860:7860 p2p-signal
|
| 42 |
+
```
|
| 43 |
+
|
| 44 |
---
|
| 45 |
+
|
| 46 |
+
## Required Secrets
|
| 47 |
+
|
| 48 |
+
Add all of these to **HF Spaces β Settings β Repository Secrets**.
|
| 49 |
+
|
| 50 |
+
| Variable | How to get it |
|
| 51 |
+
|---|---|
|
| 52 |
+
| `JWT_SECRET` | `node -e "console.log(require('crypto').randomBytes(64).toString('hex'))"` |
|
| 53 |
+
| `REDIS_PRESENCE_URL` | [Upstash](https://upstash.com) β Create Redis DB #1 β Copy Redis URL |
|
| 54 |
+
| `REDIS_MESSAGE_URL` | Upstash β Create Redis DB **#2** (separate!) β Copy Redis URL |
|
| 55 |
+
| `MONGO_URI` | [MongoDB Atlas](https://cloud.mongodb.com) β M0 Free β Connect β Driver URL |
|
| 56 |
+
| `SUPABASE_URL` | [Supabase](https://supabase.com) β Project Settings β API β Project URL |
|
| 57 |
+
| `SUPABASE_SERVICE_KEY` | Supabase β Project Settings β API β `service_role` key |
|
| 58 |
+
| `SUPABASE_BUCKET` | Create a bucket named `media` in Supabase Storage, set to **public** |
|
| 59 |
+
| `TURN_URL` | [metered.ca/tools/openrelay](https://www.metered.ca/tools/openrelay/) free TURN |
|
| 60 |
+
| `TURN_USERNAME` | From the same TURN provider |
|
| 61 |
+
| `TURN_CREDENTIAL` | From the same TURN provider |
|
| 62 |
+
| `ALLOWED_ORIGINS` | Your frontend domain(s), comma-separated e.g. `https://myapp.vercel.app` |
|
| 63 |
+
|
| 64 |
+
Optional (have defaults):
|
| 65 |
+
|
| 66 |
+
| Variable | Default | Notes |
|
| 67 |
+
|---|---|---|
|
| 68 |
+
| `PORT` | `7860` | Must stay 7860 for HF Spaces |
|
| 69 |
+
| `STUN_SERVERS` | Google STUN | Comma-separated stun: URIs |
|
| 70 |
+
| `PRESENCE_TTL_S` | `60` | Redis TTL for online status |
|
| 71 |
+
| `OFFLINE_MSG_TTL_S` | `604800` | 7 days β how long to hold queued msgs |
|
| 72 |
+
| `HEARTBEAT_INTERVAL_MS` | `30000` | Client sends ping every 30s |
|
| 73 |
+
| `MAX_FILE_SIZE_MB` | `99` | Max single file size |
|
| 74 |
+
| `MAX_FILE_COUNT` | `10` | Max files per upload request |
|
| 75 |
+
|
| 76 |
---
|
| 77 |
|
| 78 |
+
## HF Spaces Deploy
|
| 79 |
+
|
| 80 |
+
1. Create a new Space β **Docker** SDK
|
| 81 |
+
2. Push this repo to the Space's git remote
|
| 82 |
+
3. Add all secrets in Space settings
|
| 83 |
+
4. The Space will auto-build via the Dockerfile
|
| 84 |
+
|
| 85 |
+
Health check: `GET /health` β `{ status: 'ok', connections: N, ts: epoch }`
|
| 86 |
+
|
| 87 |
+
---
|
| 88 |
+
|
| 89 |
+
## WebSocket Protocol
|
| 90 |
+
|
| 91 |
+
Connect: `wss://your-space.hf.space?token=<JWT>`
|
| 92 |
+
|
| 93 |
+
### Client β Server
|
| 94 |
+
|
| 95 |
+
```jsonc
|
| 96 |
+
// Heartbeat (every 30s)
|
| 97 |
+
{ "type": "heartbeat", "payload": {} }
|
| 98 |
+
|
| 99 |
+
// Send a text message
|
| 100 |
+
{ "type": "message", "payload": { "toUid": "bob", "type": "text", "content": "Hello" } }
|
| 101 |
+
|
| 102 |
+
// Send a media message (after uploading to Supabase)
|
| 103 |
+
{ "type": "message", "payload": { "toUid": "bob", "type": "media", "content": "https://...", "mediaType": "image/jpeg" } }
|
| 104 |
+
|
| 105 |
+
// Request signed upload URLs
|
| 106 |
+
{ "type": "get-upload-url", "requestId": "abc123",
|
| 107 |
+
"payload": { "files": [{ "fileName": "photo.jpg", "mimeType": "image/jpeg", "sizeBytes": 1048576 }] } }
|
| 108 |
+
|
| 109 |
+
// WebRTC offer
|
| 110 |
+
{ "type": "offer", "payload": { "toUid": "bob", "sdp": { ...RTCSessionDescription } } }
|
| 111 |
+
|
| 112 |
+
// WebRTC answer
|
| 113 |
+
{ "type": "answer", "payload": { "toUid": "alice", "sdp": { ...RTCSessionDescription } } }
|
| 114 |
+
|
| 115 |
+
// ICE candidate
|
| 116 |
+
{ "type": "ice-candidate", "payload": { "toUid": "bob", "candidate": { ...RTCIceCandidate } } }
|
| 117 |
+
|
| 118 |
+
// WS relay fallback (when ICE fails)
|
| 119 |
+
{ "type": "relay", "payload": { "toUid": "bob", "data": { ...anything } } }
|
| 120 |
+
|
| 121 |
+
// Friend actions
|
| 122 |
+
{ "type": "friend", "payload": { "action": "send-request", "targetUid": "bob" } }
|
| 123 |
+
{ "type": "friend", "payload": { "action": "accept", "targetUid": "alice" } }
|
| 124 |
+
{ "type": "friend", "payload": { "action": "reject", "targetUid": "alice" } }
|
| 125 |
+
{ "type": "friend", "payload": { "action": "remove", "targetUid": "bob" } }
|
| 126 |
+
{ "type": "friend", "payload": { "action": "list" } }
|
| 127 |
+
{ "type": "friend", "payload": { "action": "status", "targetUid": "bob" } }
|
| 128 |
+
```
|
| 129 |
+
|
| 130 |
+
### Server β Client
|
| 131 |
+
|
| 132 |
+
```jsonc
|
| 133 |
+
{ "type": "connected", "payload": { "uid": "...", "socketId": "..." } }
|
| 134 |
+
{ "type": "ice-servers", "payload": { "iceServers": [...] } }
|
| 135 |
+
{ "type": "offline-flush", "payload": { "messages": [...], "count": N } }
|
| 136 |
+
{ "type": "message", "payload": { "fromUid": "...", "content": "...", ... } }
|
| 137 |
+
{ "type": "peer-online", "payload": { "uid": "bob" } }
|
| 138 |
+
{ "type": "peer-offline", "payload": { "uid": "bob" } }
|
| 139 |
+
{ "type": "upload-url", "payload": { "urls": [{ "signedUrl", "objectPath", "publicUrl" }] } }
|
| 140 |
+
{ "type": "offer", "payload": { "fromUid": "...", "sdp": {...} } }
|
| 141 |
+
{ "type": "answer", "payload": { "fromUid": "...", "sdp": {...} } }
|
| 142 |
+
{ "type": "ice-candidate", "payload": { "fromUid": "...", "candidate": {...} } }
|
| 143 |
+
{ "type": "relay", "payload": { "fromUid": "...", "data": {...} } }
|
| 144 |
+
{ "type": "error", "payload": { "message": "..." } }
|
| 145 |
+
```
|
| 146 |
+
|
| 147 |
+
---
|
| 148 |
+
|
| 149 |
+
## File Upload Flow
|
| 150 |
+
|
| 151 |
+
```
|
| 152 |
+
1. Client β server: { type: 'get-upload-url', payload: { files: [...] } }
|
| 153 |
+
2. Server β Supabase: createSignedUploadUrl (no bytes transferred)
|
| 154 |
+
3. Server β client: { type: 'upload-url', payload: { urls: [...] } }
|
| 155 |
+
4. Client β Supabase: PUT file directly (server sees 0 bytes, HF RAM unaffected)
|
| 156 |
+
5. Client β server: { type: 'message', payload: { type: 'media', content: publicUrl } }
|
| 157 |
+
6. Server β recipient: deliver or queue offline
|
| 158 |
+
```
|
| 159 |
+
|
| 160 |
+
---
|
| 161 |
+
|
| 162 |
+
## File Structure
|
| 163 |
+
|
| 164 |
+
```
|
| 165 |
+
src/
|
| 166 |
+
βββ index.ts # Entry point, HTTP+WS server
|
| 167 |
+
βββ config.ts # All env vars, ICE server builder
|
| 168 |
+
βββ db/
|
| 169 |
+
β βββ mongo.ts # MongoDB connection, User + Message models
|
| 170 |
+
βββ redis/
|
| 171 |
+
β βββ presenceClient.ts # Bucket 1: uidβsocketId, online/offline
|
| 172 |
+
β βββ messageClient.ts # Bucket 2: offline message queue
|
| 173 |
+
βββ services/
|
| 174 |
+
β βββ connectionRegistry.ts # In-process socket map + zombie pruner
|
| 175 |
+
β βββ supabaseService.ts # Signed URL generation, file validation
|
| 176 |
+
β βββ friendService.ts # Friend request/accept/remove logic
|
| 177 |
+
βββ handlers/
|
| 178 |
+
β βββ connectionHandler.ts # WS connect/disconnect, message router
|
| 179 |
+
β βββ signalingHandler.ts # WebRTC SDP/ICE forwarding
|
| 180 |
+
β βββ messageHandler.ts # Text/media delivery + offline queue
|
| 181 |
+
β βββ fileHandler.ts # Upload URL generation
|
| 182 |
+
β βββ friendHandler.ts # Friend CRUD over WS
|
| 183 |
+
βββ middleware/
|
| 184 |
+
β βββ rateLimiter.ts # Token bucket, 20 conn/IP/60s
|
| 185 |
+
βββ models/
|
| 186 |
+
β βββ friend.ts # FriendRequest mongoose model
|
| 187 |
+
βββ utils/
|
| 188 |
+
βββ jwt.ts # Token verification (uid always from JWT)
|
| 189 |
+
βββ logger.ts # Structured JSON logger
|
| 190 |
+
βββ protocol.ts # WS message types + safe serialization
|
| 191 |
+
|
| 192 |
+
client-reference/
|
| 193 |
+
βββ client.ts # Browser client (copy into your frontend)
|
| 194 |
+
```
|
src/index.ts
CHANGED
|
@@ -1,5 +1,6 @@
|
|
| 1 |
import http from 'http';
|
| 2 |
import WebSocket from 'ws';
|
|
|
|
| 3 |
import { CONFIG } from './config';
|
| 4 |
import { connectMongo } from './db/mongo';
|
| 5 |
import { presenceClient } from './redis/presenceClient';
|
|
@@ -9,67 +10,71 @@ import { registry } from './services/connectionRegistry';
|
|
| 9 |
import { isRateLimited } from './middleware/rateLimiter';
|
| 10 |
import { logger } from './utils/logger';
|
| 11 |
|
| 12 |
-
// ββ Graceful startup βββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 13 |
-
|
| 14 |
async function bootstrap(): Promise<void> {
|
| 15 |
-
logger.info('Starting P2P Signal Server', {
|
| 16 |
-
port: CONFIG.PORT,
|
| 17 |
-
env: CONFIG.NODE_ENV,
|
| 18 |
-
});
|
| 19 |
|
| 20 |
-
// Connect dependencies β fail fast if they are unreachable
|
| 21 |
await connectMongo();
|
| 22 |
-
|
| 23 |
-
// Redis clients connect automatically on instantiation;
|
| 24 |
-
// wait for 'ready' event before accepting traffic
|
| 25 |
await Promise.all([
|
| 26 |
waitRedisReady(presenceClient, 'Presence'),
|
| 27 |
waitRedisReady(messageClient, 'Message'),
|
| 28 |
]);
|
| 29 |
|
| 30 |
-
// ββ HTTP health endpoint βββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 31 |
-
// HF Spaces probes / to check the container is alive
|
| 32 |
const httpServer = http.createServer((req, res) => {
|
| 33 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 34 |
res.writeHead(200, { 'Content-Type': 'application/json' });
|
| 35 |
-
res.end(
|
| 36 |
-
JSON.stringify({
|
| 37 |
-
status: 'ok',
|
| 38 |
-
connections: registry.size(),
|
| 39 |
-
ts: Date.now(),
|
| 40 |
-
})
|
| 41 |
-
);
|
| 42 |
return;
|
| 43 |
}
|
| 44 |
-
|
| 45 |
-
|
| 46 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 47 |
});
|
| 48 |
|
| 49 |
-
// ββ WebSocket server βββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 50 |
const wss = new WebSocket.Server({
|
| 51 |
server: httpServer,
|
| 52 |
-
// No per-message deflate β compresses well but costs CPU on HF free tier
|
| 53 |
perMessageDeflate: false,
|
| 54 |
-
// Limit frame size to 64 KB to protect against memory exhaustion
|
| 55 |
maxPayload: 65536,
|
| 56 |
-
// Manual upgrade handling for rate limiting
|
| 57 |
handleProtocols: () => false,
|
| 58 |
});
|
| 59 |
|
| 60 |
wss.on('connection', async (ws: WebSocket, req: http.IncomingMessage) => {
|
| 61 |
-
|
| 62 |
-
if (isRateLimited(req)) {
|
| 63 |
-
ws.close(1008, 'Rate limit exceeded');
|
| 64 |
-
return;
|
| 65 |
-
}
|
| 66 |
|
| 67 |
-
// CORS check β only allow configured origins
|
| 68 |
const origin = req.headers.origin ?? '';
|
| 69 |
-
if (
|
| 70 |
-
CONFIG.ALLOWED_ORIGINS[0] !== '*' &&
|
| 71 |
-
!CONFIG.ALLOWED_ORIGINS.includes(origin)
|
| 72 |
-
) {
|
| 73 |
ws.close(1008, 'Origin not allowed');
|
| 74 |
return;
|
| 75 |
}
|
|
@@ -79,15 +84,12 @@ async function bootstrap(): Promise<void> {
|
|
| 79 |
|
| 80 |
wss.on('error', (err) => logger.error('WSS error', { error: err.message }));
|
| 81 |
|
| 82 |
-
// ββ Zombie socket pruner βββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 83 |
-
// Kill sockets that missed 3Γ the heartbeat window (90s default)
|
| 84 |
const ZOMBIE_TIMEOUT_MS = CONFIG.HEARTBEAT_INTERVAL_MS * 3;
|
| 85 |
setInterval(() => {
|
| 86 |
const pruned = registry.pruneZombies(ZOMBIE_TIMEOUT_MS);
|
| 87 |
if (pruned > 0) logger.warn('Pruned zombie sockets', { pruned });
|
| 88 |
}, ZOMBIE_TIMEOUT_MS);
|
| 89 |
|
| 90 |
-
// ββ Metrics logger βββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 91 |
setInterval(() => {
|
| 92 |
logger.info('Server metrics', {
|
| 93 |
connections: registry.size(),
|
|
@@ -95,52 +97,33 @@ async function bootstrap(): Promise<void> {
|
|
| 95 |
});
|
| 96 |
}, 60_000);
|
| 97 |
|
| 98 |
-
// ββ Listen βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 99 |
httpServer.listen(CONFIG.PORT, '0.0.0.0', () => {
|
| 100 |
logger.info(`Server listening on port ${CONFIG.PORT}`);
|
| 101 |
});
|
| 102 |
}
|
| 103 |
|
| 104 |
-
// ββ Helpers ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 105 |
-
|
| 106 |
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
| 107 |
async function waitRedisReady(client: any, name: string): Promise<void> {
|
| 108 |
return new Promise((resolve, reject) => {
|
| 109 |
if (client.status === 'ready') { resolve(); return; }
|
| 110 |
client.once('ready', resolve);
|
| 111 |
-
client.once('error', (err: Error) =>
|
| 112 |
-
reject(new Error(`${name} Redis failed: ${err.message}`))
|
| 113 |
-
);
|
| 114 |
-
// Timeout after 10s
|
| 115 |
setTimeout(() => reject(new Error(`${name} Redis connection timeout`)), 10_000);
|
| 116 |
});
|
| 117 |
}
|
| 118 |
|
| 119 |
-
|
| 120 |
-
|
| 121 |
-
|
| 122 |
-
process.on('unhandledRejection', (reason) => {
|
| 123 |
-
logger.error('Unhandled rejection', { reason: String(reason) });
|
| 124 |
-
});
|
| 125 |
-
|
| 126 |
-
process.on('uncaughtException', (err) => {
|
| 127 |
-
logger.error('Uncaught exception', { error: err.message, stack: err.stack });
|
| 128 |
-
// Do NOT exit β HF container restart is slow
|
| 129 |
-
});
|
| 130 |
-
|
| 131 |
-
// ββ Graceful shutdown ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 132 |
|
| 133 |
async function shutdown(signal: string): Promise<void> {
|
| 134 |
-
logger.info(`Received ${signal}, shutting down
|
| 135 |
presenceClient.disconnect();
|
| 136 |
messageClient.disconnect();
|
| 137 |
process.exit(0);
|
| 138 |
}
|
| 139 |
|
| 140 |
process.on('SIGTERM', () => shutdown('SIGTERM'));
|
| 141 |
-
process.on('SIGINT',
|
| 142 |
-
|
| 143 |
-
// ββ Boot βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 144 |
|
| 145 |
bootstrap().catch((err) => {
|
| 146 |
logger.error('Bootstrap failed', { error: String(err) });
|
|
|
|
| 1 |
import http from 'http';
|
| 2 |
import WebSocket from 'ws';
|
| 3 |
+
import jwt from 'jsonwebtoken';
|
| 4 |
import { CONFIG } from './config';
|
| 5 |
import { connectMongo } from './db/mongo';
|
| 6 |
import { presenceClient } from './redis/presenceClient';
|
|
|
|
| 10 |
import { isRateLimited } from './middleware/rateLimiter';
|
| 11 |
import { logger } from './utils/logger';
|
| 12 |
|
|
|
|
|
|
|
| 13 |
async function bootstrap(): Promise<void> {
|
| 14 |
+
logger.info('Starting P2P Signal Server', { port: CONFIG.PORT, env: CONFIG.NODE_ENV });
|
|
|
|
|
|
|
|
|
|
| 15 |
|
|
|
|
| 16 |
await connectMongo();
|
|
|
|
|
|
|
|
|
|
| 17 |
await Promise.all([
|
| 18 |
waitRedisReady(presenceClient, 'Presence'),
|
| 19 |
waitRedisReady(messageClient, 'Message'),
|
| 20 |
]);
|
| 21 |
|
|
|
|
|
|
|
| 22 |
const httpServer = http.createServer((req, res) => {
|
| 23 |
+
const url = new URL(req.url ?? '/', `http://${req.headers.host}`);
|
| 24 |
+
|
| 25 |
+
// ββ CORS headers for browser fetch ββββββββββββββββββββββββββββββββββββββ
|
| 26 |
+
res.setHeader('Access-Control-Allow-Origin', '*');
|
| 27 |
+
res.setHeader('Access-Control-Allow-Methods', 'GET, POST, OPTIONS');
|
| 28 |
+
res.setHeader('Access-Control-Allow-Headers', 'Content-Type');
|
| 29 |
+
if (req.method === 'OPTIONS') { res.writeHead(204); res.end(); return; }
|
| 30 |
+
|
| 31 |
+
// ββ GET /health ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 32 |
+
if (req.method === 'GET' && url.pathname === '/health') {
|
| 33 |
res.writeHead(200, { 'Content-Type': 'application/json' });
|
| 34 |
+
res.end(JSON.stringify({ status: 'ok', connections: registry.size(), ts: Date.now() }));
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 35 |
return;
|
| 36 |
}
|
| 37 |
+
|
| 38 |
+
// ββ GET /token?uid=user1 β generates a real signed JWT for testing βββββββ
|
| 39 |
+
// Simple uid allowlist β only test UIDs work, not arbitrary strings
|
| 40 |
+
if (req.method === 'GET' && url.pathname === '/token') {
|
| 41 |
+
const uid = url.searchParams.get('uid') ?? '';
|
| 42 |
+
const allowed = ['user1', 'user2', 'user3', 'user4', 'user5'];
|
| 43 |
+
if (!allowed.includes(uid)) {
|
| 44 |
+
res.writeHead(400, { 'Content-Type': 'application/json' });
|
| 45 |
+
res.end(JSON.stringify({ error: `uid must be one of: ${allowed.join(', ')}` }));
|
| 46 |
+
return;
|
| 47 |
+
}
|
| 48 |
+
const token = jwt.sign({ uid }, CONFIG.JWT_SECRET, { expiresIn: '7d' });
|
| 49 |
+
res.writeHead(200, { 'Content-Type': 'application/json' });
|
| 50 |
+
res.end(JSON.stringify({ token, uid, note: 'Test token β do not use in production' }));
|
| 51 |
+
return;
|
| 52 |
+
}
|
| 53 |
+
|
| 54 |
+
// ββ Fallback βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 55 |
+
res.writeHead(200, { 'Content-Type': 'application/json' });
|
| 56 |
+
res.end(JSON.stringify({
|
| 57 |
+
service: 'P2P Signal Server',
|
| 58 |
+
endpoints: {
|
| 59 |
+
health: 'GET /health',
|
| 60 |
+
token: 'GET /token?uid=user1 (test only)',
|
| 61 |
+
ws: 'WS /?token=<jwt>',
|
| 62 |
+
}
|
| 63 |
+
}));
|
| 64 |
});
|
| 65 |
|
|
|
|
| 66 |
const wss = new WebSocket.Server({
|
| 67 |
server: httpServer,
|
|
|
|
| 68 |
perMessageDeflate: false,
|
|
|
|
| 69 |
maxPayload: 65536,
|
|
|
|
| 70 |
handleProtocols: () => false,
|
| 71 |
});
|
| 72 |
|
| 73 |
wss.on('connection', async (ws: WebSocket, req: http.IncomingMessage) => {
|
| 74 |
+
if (isRateLimited(req)) { ws.close(1008, 'Rate limit exceeded'); return; }
|
|
|
|
|
|
|
|
|
|
|
|
|
| 75 |
|
|
|
|
| 76 |
const origin = req.headers.origin ?? '';
|
| 77 |
+
if (CONFIG.ALLOWED_ORIGINS[0] !== '*' && !CONFIG.ALLOWED_ORIGINS.includes(origin)) {
|
|
|
|
|
|
|
|
|
|
| 78 |
ws.close(1008, 'Origin not allowed');
|
| 79 |
return;
|
| 80 |
}
|
|
|
|
| 84 |
|
| 85 |
wss.on('error', (err) => logger.error('WSS error', { error: err.message }));
|
| 86 |
|
|
|
|
|
|
|
| 87 |
const ZOMBIE_TIMEOUT_MS = CONFIG.HEARTBEAT_INTERVAL_MS * 3;
|
| 88 |
setInterval(() => {
|
| 89 |
const pruned = registry.pruneZombies(ZOMBIE_TIMEOUT_MS);
|
| 90 |
if (pruned > 0) logger.warn('Pruned zombie sockets', { pruned });
|
| 91 |
}, ZOMBIE_TIMEOUT_MS);
|
| 92 |
|
|
|
|
| 93 |
setInterval(() => {
|
| 94 |
logger.info('Server metrics', {
|
| 95 |
connections: registry.size(),
|
|
|
|
| 97 |
});
|
| 98 |
}, 60_000);
|
| 99 |
|
|
|
|
| 100 |
httpServer.listen(CONFIG.PORT, '0.0.0.0', () => {
|
| 101 |
logger.info(`Server listening on port ${CONFIG.PORT}`);
|
| 102 |
});
|
| 103 |
}
|
| 104 |
|
|
|
|
|
|
|
| 105 |
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
| 106 |
async function waitRedisReady(client: any, name: string): Promise<void> {
|
| 107 |
return new Promise((resolve, reject) => {
|
| 108 |
if (client.status === 'ready') { resolve(); return; }
|
| 109 |
client.once('ready', resolve);
|
| 110 |
+
client.once('error', (err: Error) => reject(new Error(`${name} Redis failed: ${err.message}`)));
|
|
|
|
|
|
|
|
|
|
| 111 |
setTimeout(() => reject(new Error(`${name} Redis connection timeout`)), 10_000);
|
| 112 |
});
|
| 113 |
}
|
| 114 |
|
| 115 |
+
process.on('unhandledRejection', (reason) => logger.error('Unhandled rejection', { reason: String(reason) }));
|
| 116 |
+
process.on('uncaughtException', (err) => logger.error('Uncaught exception', { error: err.message, stack: err.stack }));
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 117 |
|
| 118 |
async function shutdown(signal: string): Promise<void> {
|
| 119 |
+
logger.info(`Received ${signal}, shutting down`);
|
| 120 |
presenceClient.disconnect();
|
| 121 |
messageClient.disconnect();
|
| 122 |
process.exit(0);
|
| 123 |
}
|
| 124 |
|
| 125 |
process.on('SIGTERM', () => shutdown('SIGTERM'));
|
| 126 |
+
process.on('SIGINT', () => shutdown('SIGINT'));
|
|
|
|
|
|
|
| 127 |
|
| 128 |
bootstrap().catch((err) => {
|
| 129 |
logger.error('Bootstrap failed', { error: String(err) });
|