// Simple server with WebSocket support // Run: npm install express ws && node server-ws.js const express = require('express'); const { WebSocketServer } = require('ws'); const crypto = require('crypto'); const app = express(); app.use(express.json()); const PORT = process.env.PORT || 3000; // In-memory stores const events = []; const routers = new Map(); // router_id -> { last_seen, events_sent } // WebSocket clients const clients = new Set(); // Auth middleware function auth(req, res, next) { const token = req.headers.authorization?.replace('Bearer ', ''); if (token !== process.env.EVENT_TOKEN && token !== process.env.WS_TOKEN) { return res.status(401).json({ error: 'Unauthorized' }); } next(); } // ============================================ // HTTP ENDPOINTS // ============================================ // Health check app.get('/health', (req, res) => { res.json({ status: 'ok', events_stored: events.length, routers_online: routers.size, ws_clients: clients.size, uptime: process.uptime() }); }); // Event webhook (HTTP fallback) app.post('/api/events', auth, (req, res) => { const event = { id: crypto.randomUUID(), ...req.body, received_at: new Date().toISOString(), connection: 'http' }; // Store events.push(event); if (events.length > 10000) events.shift(); // Track router const router_id = req.body.router_id; if (router_id) { routers.set(router_id, { last_seen: new Date(), last_event: event.event_type, events_sent: (routers.get(router_id)?.events_sent || 0) + 1 }); } console.log(`[${router_id}] ${event.event_type}`, event.payload); res.json({ success: true, event_id: event.id }); }); // Routers list app.get('/api/routers', (req, res) => { const router_list = []; for (const [id, data] of routers) { router_list.push({ id, ...data, online: (Date.now() - data.last_seen.getTime()) < 60000 // 1 min }); } res.json({ routers: router_list }); }); // Events query app.get('/api/events', (req, res) => { const { router, type, limit = 100 } = req.query; let filtered = events; if (router) filtered = filtered.filter(e => e.router_id === router); if (type) filtered = filtered.filter(e => e.event_type === type); res.json({ events: filtered.slice(-parseInt(limit)), total: filtered.length }); }); // ============================================ // WEBSOCKET SERVER // ============================================ const server = require('http').createServer(app); const wss = new WebSocketServer({ server, path: '/ws' }); wss.on('connection', (ws, req) => { const ip = req.socket.remoteAddress; let router_id = null; console.log(`Client connected: ${ip}`); clients.add(ws); ws.on('message', (data) => { try { const event = JSON.parse(data); router_id = event.router_id; // Store event events.push({ ...event, id: crypto.randomUUID(), received_at: new Date().toISOString(), connection: 'websocket' }); // Keep buffer size manageable if (events.length > 10000) events.shift(); // Track router routers.set(router_id, { last_seen: new Date(), last_event: event.event_type, events_sent: (routers.get(router_id)?.events_sent || 0) + 1 }); console.log(`[WS ${router_id}] ${event.event_type}`, event.payload); // Echo back acknowledgment ws.send(JSON.stringify({ ack: true, event_id: event.id })); } catch (e) { console.error('WS parse error:', e.message); } }); ws.on('close', () => { console.log(`Client disconnected: ${ip}, router: ${router_id}`); clients.delete(ws); }); ws.on('error', (err) => { console.error(`WS error from ${ip}:`, err.message); }); }); // Broadcast to all clients (for real-time updates) function broadcast(type, data) { const msg = JSON.stringify({ type, data }); for (const client of clients) { if (client.readyState === 1) { // OPEN client.send(msg); } } } // Start server server.listen(PORT, () => { console.log(` ╔═══════════════════════════════════════╗ ║ 📡 client2server Central ║ ║ HTTP: http://localhost:${PORT} ║ ║ WS: ws://localhost:${PORT}/ws ║ ╚═══════════════════════════════════════╝ `); }); // Graceful shutdown process.on('SIGINT', () => { console.log('\nShutting down...'); wss.close(); server.close(); process.exit(0); });