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