Files

136 lines
4.1 KiB
JavaScript

const cron = require('node-cron');
const netify = require('../netify');
const { Summary, AppStat, ProtocolStat, DeviceStat, Flow, Threat } = require('../models/Schemas');
const SITE_UUID = process.env.NETIFY_SITE_UUID || process.env.BACKONE_SITE_UUID;
let isRunning = false;
async function runPoll() {
if (isRunning) return;
isRunning = true;
const timestamp = new Date();
console.log(`[Mongo-Ingestion] Started polling at ${timestamp.toISOString()}`);
try {
const agents = await netify.fetchAgents();
const agentList = agents && agents.length > 0 ? agents.map(a => a.uuid) : [null]; // null for global
for (const agentUuid of agentList) {
console.log(`[Mongo-Ingestion] Fetching data for Agent: ${agentUuid || 'Global'}`);
// 1. Summary
const summary = await netify.fetchBandwidthSummary(1440, agentUuid);
if (summary) {
await new Summary({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
...summary
}).save();
}
// 2. Apps
const apps = await netify.fetchTopApps(1440, 200, agentUuid); // high limit for data lake
if (apps && apps.length > 0) {
const appDocs = apps.map(app => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
app_label: app.application?.label || 'Unknown',
download: app.download || 0,
upload: app.upload || 0,
flows: app.flows || 0
}));
await AppStat.insertMany(appDocs);
}
const devices = await netify.fetchDiscoveredDevices(1440, 500, agentUuid);
if (devices && devices.length > 0) {
const devDocs = devices.map(d => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
ip_address: d.ip_address,
mac_address: d.mac_address,
device_label: d.device_label,
device_type: d.device_type,
os_label: d.os_label,
manufacturer: d.manufacturer,
download: d.download || 0,
upload: d.upload || 0,
flows: d.flows || 0,
last_seen: d.last_seen
})).filter(d => d.ip_address); // Ensure ip_address exists to avoid validation error
if (devDocs.length > 0) {
await DeviceStat.insertMany(devDocs);
}
}
// 4. Flows
const flows = await netify.fetchFlows(500, agentUuid);
if (flows && flows.length > 0) {
const flowDocs = flows.map(f => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
flow_id: f.flow_id,
src_ip: f.src_ip,
src_mac: f.src_mac,
dst_ip: f.dst_ip,
dst_port: f.dst_port,
protocol: f.protocol,
app_label: f.app_label,
domain: f.domain,
download: f.download || 0,
upload: f.upload || 0,
first_seen: f.first_seen,
last_seen: f.last_seen
})).filter(f => f.src_ip);
if (flowDocs.length > 0) {
await Flow.insertMany(flowDocs);
}
}
// 5. Threats
const threats = await netify.fetchCyberThreats(agentUuid);
if (threats && threats.length > 0) {
const threatDocs = threats.map(t => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
threat_type: t.threat_type || 'Unknown Threat',
severity: t.severity || 'Medium',
src_ip: t.src_ip,
dst_ip: t.dst_ip,
dst_port: t.dst_port,
protocol: t.protocol,
description: t.description,
event_at: t.event_at || new Date().toISOString()
}));
if (threatDocs.length > 0) {
await Threat.insertMany(threatDocs);
}
}
}
} catch (error) {
console.error('[Mongo-Ingestion] Error during polling:', error);
} finally {
isRunning = false;
}
}
function startScheduler() {
// Run every 5 minutes
cron.schedule('*/5 * * * *', () => {
runPoll();
});
console.log('[Mongo-Ingestion] Scheduler started (every 5 minutes)');
// Initial run
runPoll();
}
module.exports = { startScheduler };