server-ws.js 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184
  1. // Simple server with WebSocket support
  2. // Run: npm install express ws && node server-ws.js
  3. const express = require('express');
  4. const { WebSocketServer } = require('ws');
  5. const crypto = require('crypto');
  6. const app = express();
  7. app.use(express.json());
  8. const PORT = process.env.PORT || 3000;
  9. // In-memory stores
  10. const events = [];
  11. const routers = new Map(); // router_id -> { last_seen, events_sent }
  12. // WebSocket clients
  13. const clients = new Set();
  14. // Auth middleware
  15. function auth(req, res, next) {
  16. const token = req.headers.authorization?.replace('Bearer ', '');
  17. if (token !== process.env.EVENT_TOKEN && token !== process.env.WS_TOKEN) {
  18. return res.status(401).json({ error: 'Unauthorized' });
  19. }
  20. next();
  21. }
  22. // ============================================
  23. // HTTP ENDPOINTS
  24. // ============================================
  25. // Health check
  26. app.get('/health', (req, res) => {
  27. res.json({
  28. status: 'ok',
  29. events_stored: events.length,
  30. routers_online: routers.size,
  31. ws_clients: clients.size,
  32. uptime: process.uptime()
  33. });
  34. });
  35. // Event webhook (HTTP fallback)
  36. app.post('/api/events', auth, (req, res) => {
  37. const event = {
  38. id: crypto.randomUUID(),
  39. ...req.body,
  40. received_at: new Date().toISOString(),
  41. connection: 'http'
  42. };
  43. // Store
  44. events.push(event);
  45. if (events.length > 10000) events.shift();
  46. // Track router
  47. const router_id = req.body.router_id;
  48. if (router_id) {
  49. routers.set(router_id, {
  50. last_seen: new Date(),
  51. last_event: event.event_type,
  52. events_sent: (routers.get(router_id)?.events_sent || 0) + 1
  53. });
  54. }
  55. console.log(`[${router_id}] ${event.event_type}`, event.payload);
  56. res.json({ success: true, event_id: event.id });
  57. });
  58. // Routers list
  59. app.get('/api/routers', (req, res) => {
  60. const router_list = [];
  61. for (const [id, data] of routers) {
  62. router_list.push({
  63. id,
  64. ...data,
  65. online: (Date.now() - data.last_seen.getTime()) < 60000 // 1 min
  66. });
  67. }
  68. res.json({ routers: router_list });
  69. });
  70. // Events query
  71. app.get('/api/events', (req, res) => {
  72. const { router, type, limit = 100 } = req.query;
  73. let filtered = events;
  74. if (router) filtered = filtered.filter(e => e.router_id === router);
  75. if (type) filtered = filtered.filter(e => e.event_type === type);
  76. res.json({
  77. events: filtered.slice(-parseInt(limit)),
  78. total: filtered.length
  79. });
  80. });
  81. // ============================================
  82. // WEBSOCKET SERVER
  83. // ============================================
  84. const server = require('http').createServer(app);
  85. const wss = new WebSocketServer({ server, path: '/ws' });
  86. wss.on('connection', (ws, req) => {
  87. const ip = req.socket.remoteAddress;
  88. let router_id = null;
  89. console.log(`Client connected: ${ip}`);
  90. clients.add(ws);
  91. ws.on('message', (data) => {
  92. try {
  93. const event = JSON.parse(data);
  94. router_id = event.router_id;
  95. // Store event
  96. events.push({
  97. ...event,
  98. id: crypto.randomUUID(),
  99. received_at: new Date().toISOString(),
  100. connection: 'websocket'
  101. });
  102. // Keep buffer size manageable
  103. if (events.length > 10000) events.shift();
  104. // Track router
  105. routers.set(router_id, {
  106. last_seen: new Date(),
  107. last_event: event.event_type,
  108. events_sent: (routers.get(router_id)?.events_sent || 0) + 1
  109. });
  110. console.log(`[WS ${router_id}] ${event.event_type}`, event.payload);
  111. // Echo back acknowledgment
  112. ws.send(JSON.stringify({ ack: true, event_id: event.id }));
  113. } catch (e) {
  114. console.error('WS parse error:', e.message);
  115. }
  116. });
  117. ws.on('close', () => {
  118. console.log(`Client disconnected: ${ip}, router: ${router_id}`);
  119. clients.delete(ws);
  120. });
  121. ws.on('error', (err) => {
  122. console.error(`WS error from ${ip}:`, err.message);
  123. });
  124. });
  125. // Broadcast to all clients (for real-time updates)
  126. function broadcast(type, data) {
  127. const msg = JSON.stringify({ type, data });
  128. for (const client of clients) {
  129. if (client.readyState === 1) { // OPEN
  130. client.send(msg);
  131. }
  132. }
  133. }
  134. // Start server
  135. server.listen(PORT, () => {
  136. console.log(`
  137. ╔═══════════════════════════════════════╗
  138. ║ 📡 client2server Central ║
  139. ║ HTTP: http://localhost:${PORT} ║
  140. ║ WS: ws://localhost:${PORT}/ws ║
  141. ╚═══════════════════════════════════════╝
  142. `);
  143. });
  144. // Graceful shutdown
  145. process.on('SIGINT', () => {
  146. console.log('\nShutting down...');
  147. wss.close();
  148. server.close();
  149. process.exit(0);
  150. });