const express = require('express'); const http = require('http'); const WebSocket = require('ws'); const axios = require('axios'); const path = require('path'); const argv = require('yargs').argv; const app = express(); const server = http.createServer(app); const wss = new WebSocket.Server({ server }); let logCache = new Set(); app.use(express.static(path.join(__dirname, 'public'))); wss.on('connection', (ws) => { console.log('Client connected:', ws._socket.remoteAddress); ws.on('close', () => { console.log('Client disconnected:', ws._socket.remoteAddress); }); ws.on('error', (error) => { console.error('WebSocket error:', error); }); }); const debugMode = !!argv.debug; setInterval(async () => { console.log('Entering setInterval callback'); let logs; if (debugMode) { console.log('debugMode is true'); try { const fs = require('fs'); const data = fs.readFileSync('logsexample.out', 'utf8'); const logLines = data.split('\n'); const validLogs = []; logLines.forEach((line, index) => { if (line.trim() !== '') { try { const unescapedLine = line.replace(/\\"/g, '"'); const parsedLog = JSON.parse(unescapedLine); validLogs.push(parsedLog); } catch (parseError) { console.error(`Error parsing log line at index ${index}: ${parseError.message}`); } } }); logs = validLogs; } catch (error) { console.error('Error reading logsexample.out:', error); return; } } else { console.log('debugMode is false'); const query = 'query=_time:10s relay received'; console.log('Sending request to API'); try { const response = await axios.post('https://vmselect.riff.cc/select/logsql/query', query, { headers: { 'Content-Type': 'application/x-www-form-urlencoded' } }); console.log('Received response from API'); if (response.data.trim() === '') { console.log('Received empty response from API'); } else { console.log('Received non-empty response from API'); } const logLines = response.data.split('\n'); const validLogs = parseLogLines(response.data); logs = validLogs; } catch (error) { console.error('Error fetching logs:', error); return; } console.log('Fetched logs:', logs); if (Array.isArray(logs)) { console.log('logs is an array, processing...'); processLogs(logs); } else { console.error('Logs is not an array:', logs); } } }, 3000); server.listen(3000, () => { console.log('Server is listening on port 3000'); }); function parseLogLines(data) { const logLines = data.split('\n'); const validLogs = []; logLines.forEach((line, index) => { if (line.trim() !== '') { try { // Check if the line is already valid JSON try { const parsedLog = JSON.parse(line); validLogs.push(parsedLog); } catch { // If it's not valid JSON, apply transformations let jsonLine = line .replace(/\\"/g, '"') // Replace escaped quotes with regular quotes .replace(/(\w+)=/g, '"$1":') // Replace '=' with ':' to form valid JSON .replace(/,(\s*})/g, '$1') // Remove trailing commas before closing braces .replace(/_stream":"{(.+?)}"/g, (match, p1) => `_stream":{"${p1.replace(/=/g, '":"').replace(/,/g, '","')}"}`) // Transform _stream field into a JSON object .replace(/_msg":"(.+?)"/g, (match, p1) => `_msg":"${p1.replace(/([a-zA-Z0-9]+)=/g, '$1:').replace(/([a-zA-Z0-9]+): /g, '"$1": ').replace(/, /g, ',').replace(/([a-zA-Z0-9]+):/g, '"$1":')}"`); // Transform _msg field into a JSON object const parsedLog = JSON.parse(jsonLine); validLogs.push(parsedLog); } } catch (parseError) { console.error(`Error parsing log line at index ${index}: ${parseError.message}`); console.error(`Log line: ${line}`); const position = parseError.message.match(/position (\d+)/); if (position) { const pos = parseInt(position[1], 10); console.error(`Character at position ${pos}: ${line.charAt(pos)}`); } } } }); return validLogs; } function processLogs(logs) { logs.forEach((log, index) => { console.log(`Processing log at index ${index}`); const msgMatch = log._msg.match(/msg_hash=0x[0-9a-fA-F]+/); const timeMatch = log._msg.match(/receivedTime=\d+/); const nodeIdMatch = log.kubernetes_pod_name.match(/nodes-(\d+)/); if (msgMatch && timeMatch && nodeIdMatch) { const msg_hash = msgMatch[0].split('=')[1]; const receivedTime = timeMatch[0].split('=')[1]; const nodeId = parseInt(nodeIdMatch[1], 10); const logIdentifier = `${msg_hash}-${receivedTime}`; if (!logCache.has(logIdentifier)) { logCache.add(logIdentifier); const logData = { msg_hash, receivedTime, nodeId, newNode: !logCache.has(`node-${nodeId}`) }; logCache.add(`node-${nodeId}`); const eventData = { msg_hash: logData.msg_hash, receivedTime: logData.receivedTime, nodeId: logData.nodeId, newNode: logData.newNode }; wss.clients.forEach(client => { if (client.readyState === WebSocket.OPEN) { console.log(`Sending event to client ${client._socket.remoteAddress}:`, eventData); client.send(JSON.stringify(eventData)); } }); } else { console.log(`Duplicate log found: ${logIdentifier}`); } } else { console.warn(`Log at index ${index} did not match expected format`); } }); }