v2: live console (SSE+fast-poll), scheduled tasks, tag-group commands, CPU sparklines, dead-node alerts

This commit is contained in:
hermes
2026-09-23 19:40:36 +00:00
parent e3e5865419
commit 33dd524196
21 changed files with 9291 additions and 12 deletions

234
server.js
View File

@@ -961,3 +961,237 @@ server.listen(PORT, '0.0.0.0', () => {
}
console.log(`=======================================================`);
});
// ═══════════════════════════════════════════════════════════════════
// NEXUSOPS v2 ADDITIONS — live console, schedules, groups, alerts
// Patched 2026-09-23. Original server.js untouched below this block's
// insertion point; overrides registered here take effect after load.
// ═══════════════════════════════════════════════════════════════════
// ── Live console state ──
const liveConsoles = new Map(); // nodeId → { operatorCount, lastActivity }
const consoleBuffers = new Map(); // nodeId → [{ ts, text }] recent output lines
// ── Scheduled tasks ──
const schedules = []; // { id, name, command, tag, intervalSec, nextRun, lastRun, history: [], enabled }
const SCHEDULES_FILE = path.join(DATA_DIR, 'schedules.json');
const ALERTS_FILE = path.join(DATA_DIR, 'alerts.json');
const alertLog = [];
function loadSchedules() {
try {
const d = JSON.parse(fs.readFileSync(SCHEDULES_FILE, 'utf8') || '[]');
schedules.push(...d);
} catch (e) { /* first run */ }
}
function saveSchedules() {
_atomicWrite(SCHEDULES_FILE, JSON.stringify(schedules, null, 2));
}
function saveAlerts() {
_atomicWrite(ALERTS_FILE, JSON.stringify(alertLog.slice(-200), null, 2));
}
loadSchedules();
// ── Dead-node alerts ──
const ALERT_WEBHOOK = process.env.NEXUS_ALERT_WEBHOOK || '';
const OFFLINE_ALERT_AFTER_MS = 5 * 60 * 1000; // alert if dark > 5 min
const alertedOffline = new Set();
setInterval(() => {
const now = Date.now();
let changed = false;
nodes.forEach((node, id) => {
if (node.status === 'online' && now - node.lastHeartbeat > 20000) {
node.status = 'offline';
changed = true;
}
if (node.status === 'offline' && now - node.lastHeartbeat > OFFLINE_ALERT_AFTER_MS && !alertedOffline.has(id)) {
alertedOffline.add(id);
const entry = {
ts: now, nodeId: id, hostname: node.hostname,
ip: node.ip, type: 'node_offline',
message: `${node.hostname} (${node.ip}) offline > ${OFFLINE_ALERT_AFTER_MS / 60000} min`
};
alertLog.push(entry);
saveAlerts();
broadcastState();
if (ALERT_WEBHOOK) {
try {
const req = require('http');
const url = new URL(ALERT_WEBHOOK);
const data = JSON.stringify({ text: `⚠️ NexusOps: ${entry.message}` });
const r = req.request({ hostname: url.hostname, port: url.port || 80, path: url.pathname, method: 'POST', headers: { 'Content-Type': 'application/json', 'Content-Length': data.length } }, () => {});
r.on('error', () => {});
r.write(data); r.end();
} catch (e) { /* silent */ }
}
}
// re-arm when node comes back
if (node.status === 'online' && alertedOffline.has(id)) {
alertedOffline.delete(id);
}
});
}, 5000);
// ── Schedule runner (every 10s tick) ──
setInterval(() => {
const now = Date.now();
let ran = false;
schedules.forEach(s => {
if (!s.enabled) return;
if (now >= s.nextRun) {
s.lastRun = now;
s.nextRun = now + s.intervalSec * 1000;
// target nodes: by tag group, or all online
const targets = Array.from(nodes.values()).filter(n =>
n.status === 'online' && (!s.tag || s.tag === 'all' || (n.tags || []).includes(s.tag)));
targets.forEach(node => {
if (!commandQueues.has(node.id)) commandQueues.set(node.id, []);
const commandId = `sched-${s.id}-${now}`;
commandQueues.get(node.id).push({
id: commandId,
actionType: 'raw_command',
payload: { command: s.command },
command: s.command,
createdAt: now
});
commandHistory.push({
id: commandId, nodeId: node.id, hostname: node.hostname,
command: `[SCHED:${s.name}] ${s.command}`,
status: 'queued', createdAt: now, output: ''
});
});
s.history.push({ ts: now, targets: targets.length });
if (s.history.length > 50) s.history.shift();
ran = true;
}
});
if (ran) { saveSchedules(); broadcastState(); }
}, 10000);
// ═══ API: Schedules ═══
app.get('/api/schedules', (req, res) => res.json({ schedules }));
app.post('/api/schedules', (req, res) => {
const { name, command, tag, intervalSec } = req.body;
if (!command || !intervalSec || intervalSec < 15) {
return res.status(400).json({ error: 'command and intervalSec (>=15) required' });
}
const s = {
id: `sch-${Date.now()}`,
name: name || command.slice(0, 30),
command, tag: tag || 'all',
intervalSec: parseInt(intervalSec, 10),
nextRun: Date.now() + parseInt(intervalSec, 10) * 1000,
lastRun: 0, enabled: true, history: []
};
schedules.push(s);
saveSchedules(); broadcastState();
res.json({ success: true, schedule: s });
});
app.post('/api/schedules/:id/toggle', (req, res) => {
const s = schedules.find(x => x.id === req.params.id);
if (!s) return res.status(404).json({ error: 'not found' });
s.enabled = !s.enabled;
if (s.enabled) s.nextRun = Date.now() + s.intervalSec * 1000;
saveSchedules(); broadcastState();
res.json({ success: true, enabled: s.enabled });
});
app.delete('/api/schedules/:id', (req, res) => {
const i = schedules.findIndex(x => x.id === req.params.id);
if (i === -1) return res.status(404).json({ error: 'not found' });
schedules.splice(i, 1);
saveSchedules(); broadcastState();
res.json({ success: true });
});
// ═══ API: Alerts ═══
app.get('/api/alerts', (req, res) => res.json({ alerts: alertLog.slice(-100) }));
// ═══ API: Groups — list distinct tags across nodes ═══
app.get('/api/groups', (req, res) => {
const tagCounts = {};
nodes.forEach(n => (n.tags || []).forEach(t => { tagCounts[t] = (tagCounts[t] || 0) + 1; }));
res.json({ groups: Object.entries(tagCounts).map(([tag, count]) => ({ tag, count })) });
});
// ═══ API: Bulk command by group tag ═══
app.post('/api/groups/command', (req, res) => {
const { command, actionType, payload, tag } = req.body;
const targets = Array.from(nodes.values()).filter(n =>
n.status === 'online' && (!tag || tag === 'all' || (n.tags || []).includes(tag)));
if (targets.length === 0) return res.status(400).json({ error: 'No online nodes in group' });
const queuedIds = [];
targets.forEach(node => {
const commandId = `cmd-grp-${Date.now()}-${Math.random().toString(36).slice(2, 4)}`;
if (!commandQueues.has(node.id)) commandQueues.set(node.id, []);
commandQueues.get(node.id).push({
id: commandId, actionType: actionType || 'raw_command',
payload: payload || { command }, command: command || actionType, createdAt: Date.now()
});
commandHistory.push({
id: commandId, nodeId: node.id, hostname: node.hostname,
command: `[GROUP:${tag || 'all'}] ${command || actionType}`,
status: 'queued', createdAt: Date.now(), output: ''
});
queuedIds.push(commandId);
});
broadcastState();
res.json({ success: true, count: targets.length, commandIds: queuedIds });
});
// ═══ Live Console: SSE stream per node ═══
app.get('/api/nodes/:id/console', (req, res) => {
const nodeId = req.params.id;
if (!nodes.has(nodeId)) return res.status(404).json({ error: 'Node not found' });
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive',
'X-Accel-Buffering': 'no'
});
res.write(`data: ${JSON.stringify({ type: 'connected', nodeId })}\n\n`);
if (!consoleBuffers.has(nodeId)) consoleBuffers.set(nodeId, []);
const buf = consoleBuffers.get(nodeId);
// replay recent output
buf.forEach(l => res.write(`data: ${JSON.stringify({ type: 'output', text: l.text, ts: l.ts })}\n\n`));
// poll commandHistory for new output belonging to this node
let lastSeen = Date.now();
const interval = setInterval(() => {
// include both command outputs and syslog entries for this node
commandHistory.forEach(c => {
if (c.nodeId === nodeId && c.completedAt && c.completedAt > lastSeen && c.output) {
res.write(`data: ${JSON.stringify({ type: 'output', text: `$ ${c.command}\n${c.output}`, ts: c.completedAt })}\n\n`);
buf.push({ ts: c.completedAt, text: `$ ${c.command}\n${c.output}` });
lastSeen = c.completedAt;
}
});
if (buf.length > 200) buf.splice(0, buf.length - 200);
res.write(`: keepalive\n\n`);
}, 1500);
req.on('close', () => clearInterval(interval));
});
// Request the agent to fast-poll (1s) while console open, slow-poll (5s) after
app.post('/api/nodes/:id/fastpoll', (req, res) => {
const nodeId = req.params.id;
if (!nodes.has(nodeId)) return res.status(404).json({ error: 'Node not found' });
const node = nodes.get(nodeId);
node.heartbeatInterval = req.body && req.body.slow ? 5 : 1;
if (!commandQueues.has(nodeId)) commandQueues.set(nodeId, []);
commandQueues.get(nodeId).push({
id: `poll-${Date.now()}`,
actionType: 'set_heartbeat_rate',
payload: { interval: node.heartbeatInterval },
command: `set_heartbeat_rate ${node.heartbeatInterval}s`,
createdAt: Date.now()
});
broadcastState();
res.json({ success: true, interval: node.heartbeatInterval });
});
console.log('[v2] NexusOps v2 additions loaded: live console, schedules, groups, alerts');