Initial commit: LogMaster with Viewer Collector Monitor modes
This commit is contained in:
497
dist-server/collector.js
Normal file
497
dist-server/collector.js
Normal file
@@ -0,0 +1,497 @@
|
||||
"use strict";
|
||||
var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) {
|
||||
if (k2 === undefined) k2 = k;
|
||||
var desc = Object.getOwnPropertyDescriptor(m, k);
|
||||
if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) {
|
||||
desc = { enumerable: true, get: function() { return m[k]; } };
|
||||
}
|
||||
Object.defineProperty(o, k2, desc);
|
||||
}) : (function(o, m, k, k2) {
|
||||
if (k2 === undefined) k2 = k;
|
||||
o[k2] = m[k];
|
||||
}));
|
||||
var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) {
|
||||
Object.defineProperty(o, "default", { enumerable: true, value: v });
|
||||
}) : function(o, v) {
|
||||
o["default"] = v;
|
||||
});
|
||||
var __importStar = (this && this.__importStar) || (function () {
|
||||
var ownKeys = function(o) {
|
||||
ownKeys = Object.getOwnPropertyNames || function (o) {
|
||||
var ar = [];
|
||||
for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k;
|
||||
return ar;
|
||||
};
|
||||
return ownKeys(o);
|
||||
};
|
||||
return function (mod) {
|
||||
if (mod && mod.__esModule) return mod;
|
||||
var result = {};
|
||||
if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]);
|
||||
__setModuleDefault(result, mod);
|
||||
return result;
|
||||
};
|
||||
})();
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
exports.LogCollector = void 0;
|
||||
const dgram = __importStar(require("dgram"));
|
||||
const fs = __importStar(require("fs"));
|
||||
const path = __importStar(require("path"));
|
||||
const zlib = __importStar(require("zlib"));
|
||||
class LogCollector {
|
||||
config;
|
||||
streams = new Map();
|
||||
buffers = new Map();
|
||||
fileSizes = new Map();
|
||||
fileCreatedAt = new Map();
|
||||
rotationTimer = null;
|
||||
udpServer = null;
|
||||
forwardBuffers = new Map();
|
||||
forwardTimers = new Map();
|
||||
constructor(config = {}) {
|
||||
this.config = {
|
||||
storageDir: config.storageDir || path.join(process.cwd(), 'collector-data'),
|
||||
maxLinesPerStream: config.maxLinesPerStream || 100000,
|
||||
maxStreams: config.maxStreams || 50,
|
||||
udpPort: config.udpPort || 5140,
|
||||
udpEnabled: config.udpEnabled !== false,
|
||||
rotation: config.rotation || {
|
||||
maxSizeBytes: 5 * 1024 * 1024, // 5MB default
|
||||
maxAgeMs: 10 * 60 * 1000, // 10 min default
|
||||
compression: 'gzip',
|
||||
fileExtension: '.log',
|
||||
},
|
||||
forwardTargets: config.forwardTargets || [],
|
||||
persistToDisk: config.persistToDisk !== false,
|
||||
};
|
||||
if (!fs.existsSync(this.config.storageDir)) {
|
||||
fs.mkdirSync(this.config.storageDir, { recursive: true });
|
||||
}
|
||||
this.loadExistingStreams();
|
||||
// Start rotation check timer
|
||||
if (this.config.rotation.maxAgeMs > 0) {
|
||||
this.rotationTimer = setInterval(() => this.checkTimeRotation(), Math.min(this.config.rotation.maxAgeMs / 2, 30000));
|
||||
}
|
||||
// Start UDP server
|
||||
if (this.config.udpEnabled) {
|
||||
this.startUdpServer();
|
||||
}
|
||||
}
|
||||
startUdpServer() {
|
||||
try {
|
||||
this.udpServer = dgram.createSocket('udp4');
|
||||
this.udpServer.on('message', (msg, rinfo) => {
|
||||
const message = msg.toString('utf-8').trim();
|
||||
if (!message)
|
||||
return;
|
||||
// Parse stream ID from message or use sender IP as stream
|
||||
// Format: <streamId>|<message> or just <message> (uses IP:port as stream)
|
||||
let streamId;
|
||||
let logLine;
|
||||
const pipeIdx = message.indexOf('|');
|
||||
if (pipeIdx > 0 && pipeIdx < 64 && !message.substring(0, pipeIdx).includes(' ')) {
|
||||
streamId = message.substring(0, pipeIdx);
|
||||
logLine = message.substring(pipeIdx + 1);
|
||||
}
|
||||
else {
|
||||
streamId = `udp-${rinfo.address}`;
|
||||
logLine = message;
|
||||
}
|
||||
this.ingest(streamId, [logLine], `UDP ${rinfo.address}:${rinfo.port}`);
|
||||
});
|
||||
this.udpServer.on('error', (err) => {
|
||||
console.error(`UDP server error: ${err.message}`);
|
||||
});
|
||||
this.udpServer.bind(this.config.udpPort, '0.0.0.0', () => {
|
||||
console.log(` UDP collector listening on port ${this.config.udpPort}`);
|
||||
});
|
||||
}
|
||||
catch (err) {
|
||||
console.error('Failed to start UDP server:', err.message);
|
||||
}
|
||||
}
|
||||
loadExistingStreams() {
|
||||
try {
|
||||
const metaFile = path.join(this.config.storageDir, '_streams.json');
|
||||
if (fs.existsSync(metaFile)) {
|
||||
const data = JSON.parse(fs.readFileSync(metaFile, 'utf-8'));
|
||||
for (const stream of data) {
|
||||
this.streams.set(stream.id, stream);
|
||||
const streamFile = this.getStreamFilePath(stream.id);
|
||||
if (fs.existsSync(streamFile)) {
|
||||
const content = fs.readFileSync(streamFile, 'utf-8');
|
||||
const lines = content.split('\n').filter(l => l.length > 0);
|
||||
this.buffers.set(stream.id, lines);
|
||||
const stat = fs.statSync(streamFile);
|
||||
this.fileSizes.set(stream.id, stat.size);
|
||||
this.fileCreatedAt.set(stream.id, stat.birthtimeMs || Date.now());
|
||||
}
|
||||
else {
|
||||
this.buffers.set(stream.id, []);
|
||||
this.fileSizes.set(stream.id, 0);
|
||||
this.fileCreatedAt.set(stream.id, Date.now());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (err) {
|
||||
console.error('Error loading collector streams:', err);
|
||||
}
|
||||
}
|
||||
saveMetadata() {
|
||||
const metaFile = path.join(this.config.storageDir, '_streams.json');
|
||||
fs.writeFileSync(metaFile, JSON.stringify(Array.from(this.streams.values()), null, 2));
|
||||
}
|
||||
getStreamFilePath(streamId) {
|
||||
return path.join(this.config.storageDir, `${streamId}${this.config.rotation.fileExtension}`);
|
||||
}
|
||||
getRotatedFilePath(streamId) {
|
||||
const ts = new Date().toISOString().replace(/[:.]/g, '-');
|
||||
const ext = this.config.rotation.fileExtension;
|
||||
const compExt = this.getCompressionExtension();
|
||||
return path.join(this.config.storageDir, `${streamId}_${ts}${ext}${compExt}`);
|
||||
}
|
||||
getCompressionExtension() {
|
||||
switch (this.config.rotation.compression) {
|
||||
case 'gzip': return '.gz';
|
||||
case 'bzip2': return '.bz2';
|
||||
case 'xz': return '.xz';
|
||||
case 'zstd': return '.zst';
|
||||
default: return '';
|
||||
}
|
||||
}
|
||||
saveStream(streamId) {
|
||||
if (!this.config.persistToDisk)
|
||||
return;
|
||||
const lines = this.buffers.get(streamId) || [];
|
||||
const streamFile = this.getStreamFilePath(streamId);
|
||||
const content = lines.join('\n') + (lines.length > 0 ? '\n' : '');
|
||||
fs.writeFileSync(streamFile, content);
|
||||
this.fileSizes.set(streamId, Buffer.byteLength(content, 'utf-8'));
|
||||
}
|
||||
checkSizeRotation(streamId) {
|
||||
if (this.config.rotation.maxSizeBytes <= 0)
|
||||
return;
|
||||
const size = this.fileSizes.get(streamId) || 0;
|
||||
if (size >= this.config.rotation.maxSizeBytes) {
|
||||
this.rotateStream(streamId);
|
||||
}
|
||||
}
|
||||
checkTimeRotation() {
|
||||
if (this.config.rotation.maxAgeMs <= 0)
|
||||
return;
|
||||
const now = Date.now();
|
||||
for (const [streamId, createdAt] of this.fileCreatedAt) {
|
||||
if (now - createdAt >= this.config.rotation.maxAgeMs) {
|
||||
const buffer = this.buffers.get(streamId);
|
||||
if (buffer && buffer.length > 0) {
|
||||
this.rotateStream(streamId);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
rotateStream(streamId) {
|
||||
const srcFile = this.getStreamFilePath(streamId);
|
||||
if (!fs.existsSync(srcFile))
|
||||
return;
|
||||
const content = fs.readFileSync(srcFile);
|
||||
if (content.length === 0)
|
||||
return;
|
||||
const destFile = this.getRotatedFilePath(streamId);
|
||||
// Compress and write
|
||||
let compressed;
|
||||
switch (this.config.rotation.compression) {
|
||||
case 'gzip':
|
||||
compressed = zlib.gzipSync(content);
|
||||
break;
|
||||
case 'bzip2':
|
||||
case 'xz':
|
||||
case 'zstd':
|
||||
// For these, just copy without compression (would need external tools)
|
||||
compressed = content;
|
||||
break;
|
||||
default:
|
||||
compressed = content;
|
||||
}
|
||||
fs.writeFileSync(destFile, compressed);
|
||||
console.log(` Rotated: ${streamId} -> ${path.basename(destFile)} (${compressed.length} bytes)`);
|
||||
// Reset current file
|
||||
fs.writeFileSync(srcFile, '');
|
||||
this.buffers.set(streamId, []);
|
||||
this.fileSizes.set(streamId, 0);
|
||||
this.fileCreatedAt.set(streamId, Date.now());
|
||||
const stream = this.streams.get(streamId);
|
||||
if (stream) {
|
||||
stream.lineCount = 0;
|
||||
stream.currentFileSize = 0;
|
||||
}
|
||||
this.saveMetadata();
|
||||
}
|
||||
// ---- Forwarding ----
|
||||
forwardLines(streamId, lines) {
|
||||
for (const target of this.config.forwardTargets) {
|
||||
const key = `${streamId}:${target.type}:${target.url || target.path || target.host}`;
|
||||
if (!this.forwardBuffers.has(key))
|
||||
this.forwardBuffers.set(key, []);
|
||||
const buf = this.forwardBuffers.get(key);
|
||||
buf.push(...lines);
|
||||
const batchSize = target.batchSize || 100;
|
||||
if (buf.length >= batchSize) {
|
||||
this.flushForward(key, target, streamId);
|
||||
}
|
||||
else if (!this.forwardTimers.has(key)) {
|
||||
const flushMs = target.flushIntervalMs || 5000;
|
||||
this.forwardTimers.set(key, setTimeout(() => {
|
||||
this.flushForward(key, target, streamId);
|
||||
this.forwardTimers.delete(key);
|
||||
}, flushMs));
|
||||
}
|
||||
}
|
||||
}
|
||||
flushForward(key, target, streamId) {
|
||||
const lines = this.forwardBuffers.get(key) || [];
|
||||
if (lines.length === 0)
|
||||
return;
|
||||
this.forwardBuffers.set(key, []);
|
||||
switch (target.type) {
|
||||
case 'http':
|
||||
this.forwardHttp(target, streamId, lines);
|
||||
break;
|
||||
case 'file':
|
||||
this.forwardFile(target, lines);
|
||||
break;
|
||||
case 'udp':
|
||||
this.forwardUdp(target, streamId, lines);
|
||||
break;
|
||||
case 'stream':
|
||||
this.forwardToStream(target, streamId, lines);
|
||||
break;
|
||||
case 'console':
|
||||
for (const l of lines)
|
||||
console.log(`[FWD:${streamId}] ${l}`);
|
||||
break;
|
||||
}
|
||||
}
|
||||
forwardHttp(target, streamId, lines) {
|
||||
if (!target.url)
|
||||
return;
|
||||
const payload = target.format === 'json'
|
||||
? JSON.stringify({ stream: streamId, lines, timestamp: new Date().toISOString() })
|
||||
: lines.join('\n');
|
||||
try {
|
||||
const url = new URL(target.url);
|
||||
const transport = url.protocol === 'https:' ? require('https') : require('http');
|
||||
const req = transport.request({
|
||||
hostname: url.hostname,
|
||||
port: url.port || (url.protocol === 'https:' ? 443 : 80),
|
||||
path: url.pathname,
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': target.format === 'json' ? 'application/json' : 'text/plain', 'X-Stream-Id': streamId },
|
||||
});
|
||||
req.on('error', (e) => console.error(`Forward HTTP error: ${e.message}`));
|
||||
req.write(payload);
|
||||
req.end();
|
||||
}
|
||||
catch (e) {
|
||||
console.error(`Forward HTTP error: ${e.message}`);
|
||||
}
|
||||
}
|
||||
forwardFile(target, lines) {
|
||||
if (!target.path)
|
||||
return;
|
||||
try {
|
||||
fs.appendFileSync(target.path, lines.join('\n') + '\n');
|
||||
}
|
||||
catch (e) {
|
||||
console.error(`Forward file error: ${e.message}`);
|
||||
}
|
||||
}
|
||||
forwardUdp(target, streamId, lines) {
|
||||
if (!target.host || !target.port)
|
||||
return;
|
||||
try {
|
||||
const client = dgram.createSocket('udp4');
|
||||
for (const line of lines) {
|
||||
const msg = Buffer.from(`${streamId}|${line}`);
|
||||
client.send(msg, target.port, target.host);
|
||||
}
|
||||
setTimeout(() => client.close(), 1000);
|
||||
}
|
||||
catch (e) {
|
||||
console.error(`Forward UDP error: ${e.message}`);
|
||||
}
|
||||
}
|
||||
forwardToStream(target, sourceStreamId, lines) {
|
||||
const targetId = target.targetStream;
|
||||
if (!targetId || targetId === sourceStreamId)
|
||||
return; // Prevent infinite loop
|
||||
try {
|
||||
// Prefix lines with source stream info
|
||||
const prefixed = lines.map(l => `[fwd:${sourceStreamId}] ${l}`);
|
||||
// Ingest directly into the target stream without triggering forwards again
|
||||
const stream = this.getOrCreateStream(targetId, `Forwarded from ${sourceStreamId}`);
|
||||
const buffer = this.buffers.get(targetId);
|
||||
for (const line of prefixed) {
|
||||
if (buffer.length >= this.config.maxLinesPerStream)
|
||||
buffer.shift();
|
||||
buffer.push(line);
|
||||
}
|
||||
stream.lineCount = buffer.length;
|
||||
stream.lastReceived = new Date().toISOString();
|
||||
if (this.config.persistToDisk)
|
||||
this.saveStream(targetId);
|
||||
this.saveMetadata();
|
||||
}
|
||||
catch (e) {
|
||||
console.error(`Forward stream error: ${e.message}`);
|
||||
}
|
||||
}
|
||||
// ---- Public API ----
|
||||
getOrCreateStream(streamId, name) {
|
||||
if (this.streams.has(streamId))
|
||||
return this.streams.get(streamId);
|
||||
if (this.streams.size >= this.config.maxStreams)
|
||||
throw new Error(`Max streams (${this.config.maxStreams}) reached`);
|
||||
const stream = {
|
||||
id: streamId, name: name || streamId, createdAt: new Date().toISOString(),
|
||||
lineCount: 0, lastReceived: null, currentFileSize: 0,
|
||||
};
|
||||
this.streams.set(streamId, stream);
|
||||
this.buffers.set(streamId, []);
|
||||
this.fileSizes.set(streamId, 0);
|
||||
this.fileCreatedAt.set(streamId, Date.now());
|
||||
this.saveMetadata();
|
||||
return stream;
|
||||
}
|
||||
ingest(streamId, lines, streamName) {
|
||||
const stream = this.getOrCreateStream(streamId, streamName);
|
||||
const buffer = this.buffers.get(streamId);
|
||||
let accepted = 0, dropped = 0;
|
||||
const validLines = [];
|
||||
for (const line of lines) {
|
||||
if (line.trim() === '')
|
||||
continue;
|
||||
if (buffer.length >= this.config.maxLinesPerStream) {
|
||||
buffer.shift();
|
||||
dropped++;
|
||||
}
|
||||
buffer.push(line);
|
||||
validLines.push(line);
|
||||
accepted++;
|
||||
}
|
||||
stream.lineCount = buffer.length;
|
||||
stream.lastReceived = new Date().toISOString();
|
||||
if (this.config.persistToDisk) {
|
||||
this.saveStream(streamId);
|
||||
stream.currentFileSize = this.fileSizes.get(streamId) || 0;
|
||||
this.checkSizeRotation(streamId);
|
||||
}
|
||||
// Forward to targets
|
||||
if (validLines.length > 0 && this.config.forwardTargets.length > 0) {
|
||||
this.forwardLines(streamId, validLines);
|
||||
}
|
||||
this.saveMetadata();
|
||||
return { accepted, dropped };
|
||||
}
|
||||
getStreams() { return Array.from(this.streams.values()); }
|
||||
readStream(streamId, tail) {
|
||||
const buffer = this.buffers.get(streamId);
|
||||
if (!buffer)
|
||||
throw new Error(`Stream not found: ${streamId}`);
|
||||
return tail && tail > 0 ? buffer.slice(-tail) : [...buffer];
|
||||
}
|
||||
deleteStream(streamId) {
|
||||
if (!this.streams.has(streamId))
|
||||
return false;
|
||||
this.streams.delete(streamId);
|
||||
this.buffers.delete(streamId);
|
||||
this.fileSizes.delete(streamId);
|
||||
this.fileCreatedAt.delete(streamId);
|
||||
const f = this.getStreamFilePath(streamId);
|
||||
if (fs.existsSync(f))
|
||||
fs.unlinkSync(f);
|
||||
this.saveMetadata();
|
||||
return true;
|
||||
}
|
||||
getConfig() { return { ...this.config }; }
|
||||
shutdown() {
|
||||
if (this.rotationTimer)
|
||||
clearInterval(this.rotationTimer);
|
||||
if (this.udpServer)
|
||||
this.udpServer.close();
|
||||
for (const timer of this.forwardTimers.values())
|
||||
clearTimeout(timer);
|
||||
}
|
||||
handleRequest(req, res, pathname) {
|
||||
const method = req.method || 'GET';
|
||||
// POST /collect/:streamId
|
||||
const collectMatch = pathname.match(/^\/collect\/([a-zA-Z0-9_\-\.]+)(\/json)?$/);
|
||||
if (collectMatch && method === 'POST') {
|
||||
const streamId = collectMatch[1];
|
||||
const isJson = !!collectMatch[2];
|
||||
let body = '';
|
||||
req.on('data', (chunk) => { body += chunk.toString(); });
|
||||
req.on('end', () => {
|
||||
try {
|
||||
let lines;
|
||||
if (isJson) {
|
||||
const parsed = JSON.parse(body);
|
||||
lines = Array.isArray(parsed) ? parsed.map((e) => typeof e === 'string' ? e : JSON.stringify(e)) : [JSON.stringify(parsed)];
|
||||
}
|
||||
else {
|
||||
lines = body.split('\n');
|
||||
}
|
||||
const streamName = req.headers['x-stream-name'];
|
||||
const result = this.ingest(streamId, lines, streamName);
|
||||
res.writeHead(200, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify({ ok: true, ...result }));
|
||||
}
|
||||
catch (err) {
|
||||
res.writeHead(400, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify({ error: err.message }));
|
||||
}
|
||||
});
|
||||
return true;
|
||||
}
|
||||
// GET /collector/streams
|
||||
if (pathname === '/collector/streams' && method === 'GET') {
|
||||
res.writeHead(200, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify(this.getStreams()));
|
||||
return true;
|
||||
}
|
||||
// GET /collector/config
|
||||
if (pathname === '/collector/config' && method === 'GET') {
|
||||
res.writeHead(200, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify(this.getConfig()));
|
||||
return true;
|
||||
}
|
||||
// GET /collector/stream/:id
|
||||
const readMatch = pathname.match(/^\/collector\/stream\/([a-zA-Z0-9_\-\.]+)$/);
|
||||
if (readMatch && method === 'GET') {
|
||||
const streamId = readMatch[1];
|
||||
try {
|
||||
const url = new URL(req.url || '/', `http://${req.headers.host}`);
|
||||
const tail = url.searchParams.get('tail');
|
||||
const lines = this.readStream(streamId, tail ? parseInt(tail, 10) : undefined);
|
||||
res.writeHead(200, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify({ streamId, lineCount: lines.length, lines }));
|
||||
}
|
||||
catch (err) {
|
||||
res.writeHead(404, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify({ error: err.message }));
|
||||
}
|
||||
return true;
|
||||
}
|
||||
// DELETE /collector/stream/:id
|
||||
const deleteMatch = pathname.match(/^\/collector\/stream\/([a-zA-Z0-9_\-\.]+)$/);
|
||||
if (deleteMatch && method === 'DELETE') {
|
||||
const deleted = this.deleteStream(deleteMatch[1]);
|
||||
res.writeHead(deleted ? 200 : 404, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify({ ok: deleted }));
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
}
|
||||
exports.LogCollector = LogCollector;
|
||||
//# sourceMappingURL=collector.js.map
|
||||
Reference in New Issue
Block a user