feat: complete platform upgrade - user guide overhaul, unified export dropdowns, 24h agent tracking, HelpTrigger implementation on all pages
This commit is contained in:
1 parent
040ef4d45d
commit
64b5e87341
215 files changed
+7886
-4129
No files matched your search
@@ -0,0 +1,63 @@
|
||||
const path = require('path');
|
||||
require('dotenv').config({ path: path.join(__dirname, '..', '.env.local') });
|
||||
const mongoose = require('mongoose');
|
||||
|
||||
const MONGODB_URI = process.env.MONGODB_URI || 'mongodb://127.0.0.1:27017/backone_dpi';
|
||||
const SIAB = '6681452d_9cae_4ff4_8ae8_0d504774265e';
|
||||
|
||||
mongoose.connect(MONGODB_URI).then(async () => {
|
||||
const db = mongoose.connection.db;
|
||||
const since24h = new Date(Date.now() - 24 * 3600000);
|
||||
const since7d = new Date(Date.now() - 7 * 24 * 3600000);
|
||||
|
||||
// Count site-level summary docs
|
||||
const count24h = await db.collection('summaries').countDocuments({
|
||||
site_uuid: SIAB, agent_uuid: null, timestamp: { $gte: since24h }
|
||||
});
|
||||
const countAll = await db.collection('summaries').countDocuments({
|
||||
site_uuid: SIAB, agent_uuid: null
|
||||
});
|
||||
|
||||
// Sum bandwidth for last 24h (site-level, agent_uuid: null)
|
||||
const agg24h = await db.collection('summaries').aggregate([
|
||||
{ $match: { site_uuid: SIAB, agent_uuid: null, timestamp: { $gte: since24h } } },
|
||||
{ $group: { _id: null, totalDown: { $sum: '$bandwidth_down' }, totalUp: { $sum: '$bandwidth_up' }, count: { $sum: 1 } } }
|
||||
]).toArray();
|
||||
|
||||
// Sum bandwidth ALL time (site-level)
|
||||
const aggAll = await db.collection('summaries').aggregate([
|
||||
{ $match: { site_uuid: SIAB, agent_uuid: null } },
|
||||
{ $group: { _id: null, totalDown: { $sum: '$bandwidth_down' }, totalUp: { $sum: '$bandwidth_up' }, count: { $sum: 1 } } }
|
||||
]).toArray();
|
||||
|
||||
// Oldest and newest
|
||||
const oldest = await db.collection('summaries').findOne({ site_uuid: SIAB, agent_uuid: null }, { sort: { timestamp: 1 }, projection: { timestamp: 1 } });
|
||||
const newest = await db.collection('summaries').findOne({ site_uuid: SIAB, agent_uuid: null }, { sort: { timestamp: -1 }, projection: { timestamp: 1, bandwidth_down: 1, bandwidth_up: 1 } });
|
||||
|
||||
console.log('\n=== MongoDB Summary Check (SIAB site) ===');
|
||||
console.log('Total site-level docs:', countAll);
|
||||
console.log('Site-level docs in last 24h:', count24h);
|
||||
console.log('\nBandwidth SUM (last 24h):');
|
||||
console.log(' Down:', agg24h[0] ? (agg24h[0].totalDown / 1024 / 1024).toFixed(2) + ' MB' : '0 MB');
|
||||
console.log(' Up :', agg24h[0] ? (agg24h[0].totalUp / 1024 / 1024).toFixed(2) + ' MB' : '0 MB');
|
||||
console.log('\nBandwidth SUM (ALL time):');
|
||||
console.log(' Down:', aggAll[0] ? (aggAll[0].totalDown / 1024 / 1024).toFixed(2) + ' MB' : '0 MB');
|
||||
console.log(' Up :', aggAll[0] ? (aggAll[0].totalUp / 1024 / 1024).toFixed(2) + ' MB' : '0 MB');
|
||||
console.log('\nOldest entry :', oldest?.timestamp);
|
||||
console.log('Newest entry :', newest?.timestamp);
|
||||
console.log('Latest bandwidth_down per 5min:', newest ? (newest.bandwidth_down / 1024 / 1024).toFixed(4) + ' MB' : 'N/A');
|
||||
console.log('Latest bandwidth_up per 5min :', newest ? (newest.bandwidth_up / 1024 / 1024).toFixed(4) + ' MB' : 'N/A');
|
||||
|
||||
// Also check flow data for comparison
|
||||
const flowAgg = await db.collection('flows').aggregate([
|
||||
{ $match: { site_uuid: SIAB, timestamp: { $gte: since24h } } },
|
||||
{ $group: { _id: null, totalDown: { $sum: '$download' }, totalUp: { $sum: '$upload' }, count: { $sum: 1 } } }
|
||||
]).toArray();
|
||||
console.log('\nFlow-level bandwidth (last 24h from flows collection):');
|
||||
console.log(' Down:', flowAgg[0] ? (flowAgg[0].totalDown / 1024 / 1024).toFixed(2) + ' MB' : '0 MB');
|
||||
console.log(' Up :', flowAgg[0] ? (flowAgg[0].totalUp / 1024 / 1024).toFixed(2) + ' MB' : '0 MB');
|
||||
console.log(' Flow count:', flowAgg[0]?.count || 0);
|
||||
console.log('=====================================\n');
|
||||
|
||||
process.exit(0);
|
||||
}).catch(e => { console.error(e.message); process.exit(1); });
|
||||
+38
-37
@@ -10,14 +10,10 @@ require('dotenv').config({ path: path.join(__dirname, '..', '.env.local') });
|
||||
|
||||
const netify = require('./netifyClient');
|
||||
const { collectSecondaryTelemetry } = require('./collectorHelper');
|
||||
const {
|
||||
collectDevicesAndApps,
|
||||
collectFlows,
|
||||
collectThreats,
|
||||
collectEvents
|
||||
} = require('./collectorHelperDpi2');
|
||||
const { collectDevicesAndApps, collectFlows } = require('./collectorHelperDpi2');
|
||||
const { collectThreats, collectEvents } = require('./collectorHelperDpi3');
|
||||
|
||||
const { Summary, AppStat } = require('./models/Schemas');
|
||||
const { Summary, AppStat, AgentRegistry } = require('./models/Schemas');
|
||||
const { pruneOldData } = require('./dataRetention');
|
||||
|
||||
const SITE_UUIDS_STR = process.env.NETIFY_SITE_UUIDS || process.env.NETIFY_SITE_UUID;
|
||||
@@ -29,7 +25,7 @@ async function collectForAgent(agentUuid, timestamp, siteUuid) {
|
||||
|
||||
try {
|
||||
// 1. Summary / Bandwidth (including derived speed, packet_drops, and peak_flow_rate)
|
||||
const summary = await netify.fetchBandwidthSummary(1440, agentUuid, siteUuid);
|
||||
const summary = await netify.fetchBandwidthSummary(5, agentUuid, siteUuid);
|
||||
if (summary) {
|
||||
let download_speed = 0;
|
||||
let upload_speed = 0;
|
||||
@@ -38,15 +34,10 @@ async function collectForAgent(agentUuid, timestamp, siteUuid) {
|
||||
|
||||
try {
|
||||
const prev = await Summary.findOne({ agent_uuid: agentUuid }).sort({ timestamp: -1 }).lean();
|
||||
if (prev && prev.timestamp) {
|
||||
const timeDiffSec = (timestamp.getTime() - new Date(prev.timestamp).getTime()) / 1000;
|
||||
if (timeDiffSec > 0) {
|
||||
const bytesDiffDown = Math.max(0, summary.bandwidth_down - (prev.bandwidth_down || 0));
|
||||
const bytesDiffUp = Math.max(0, summary.bandwidth_up - (prev.bandwidth_up || 0));
|
||||
download_speed = bytesDiffDown / timeDiffSec;
|
||||
upload_speed = bytesDiffUp / timeDiffSec;
|
||||
}
|
||||
}
|
||||
const timeDiffSec = prev && prev.timestamp ? (timestamp.getTime() - new Date(prev.timestamp).getTime()) / 1000 : 300;
|
||||
const activeTimeDiff = timeDiffSec > 0 ? timeDiffSec : 300;
|
||||
download_speed = summary.bandwidth_down / activeTimeDiff;
|
||||
upload_speed = summary.bandwidth_up / activeTimeDiff;
|
||||
} catch (err) {
|
||||
console.error('[Collector] Error calculating summary speeds:', err.message);
|
||||
}
|
||||
@@ -78,7 +69,7 @@ async function collectForAgent(agentUuid, timestamp, siteUuid) {
|
||||
}
|
||||
|
||||
// 2. Top Apps
|
||||
const apps = await netify.fetchTopApps(1440, 200, agentUuid, siteUuid);
|
||||
const apps = await netify.fetchTopApps(5, 200, agentUuid, siteUuid);
|
||||
if (apps && apps.length > 0) {
|
||||
const appDocs = apps.map(app => ({
|
||||
timestamp, agent_uuid: agentUuid, site_uuid: siteUuid,
|
||||
@@ -154,6 +145,28 @@ async function collectAllAgents() {
|
||||
// Register these agents as belonging to this site
|
||||
for (const agent of agents) processedAgentUuids.add(agent.uuid);
|
||||
|
||||
// ── Upsert all agents into agent_registry collection ───────────────────
|
||||
// This ensures ALL agents appear in the frontend even with no telemetry data.
|
||||
await Promise.allSettled(agents.map(a =>
|
||||
AgentRegistry.findOneAndUpdate(
|
||||
{ uuid: a.uuid },
|
||||
{
|
||||
$set: {
|
||||
uuid: a.uuid,
|
||||
serial: a.serial || a.uuid,
|
||||
label: a.label,
|
||||
site_uuid: siteUuid,
|
||||
provisioned: a.provisioned ?? true,
|
||||
activated: a.activated ?? false,
|
||||
last_seen_at: a.last_seen_at ?? null,
|
||||
netify_id: a.id ?? null,
|
||||
}
|
||||
},
|
||||
{ upsert: true, new: true }
|
||||
)
|
||||
));
|
||||
console.log(`[Collector] ✓ ${agents.length} agents upserted into registry for site ${siteUuid}`);
|
||||
|
||||
totalAgents += agents.length;
|
||||
console.log(`[Collector] Processing ${agents.length} agents for site ${siteUuid}: ${agents.map(a => a.uuid).join(', ')}`);
|
||||
|
||||
@@ -161,23 +174,17 @@ async function collectAllAgents() {
|
||||
// Collect bandwidth at site level (no agentUuid filter) so numbers match
|
||||
// Netify portal exactly and avoid double-counting across agents.
|
||||
try {
|
||||
console.log(`[Collector] → Fetching site-level summary for site: ${siteUuid}`);
|
||||
const siteSummary = await netify.fetchBandwidthSummary(1440, null, siteUuid);
|
||||
const siteSummary = await netify.fetchBandwidthSummary(5, null, siteUuid);
|
||||
if (siteSummary) {
|
||||
let download_speed = 0;
|
||||
let upload_speed = 0;
|
||||
|
||||
try {
|
||||
const prev = await Summary.findOne({ agent_uuid: null, site_uuid: siteUuid }).sort({ timestamp: -1 }).lean();
|
||||
if (prev && prev.timestamp) {
|
||||
const timeDiffSec = (timestamp.getTime() - new Date(prev.timestamp).getTime()) / 1000;
|
||||
if (timeDiffSec > 0) {
|
||||
const bytesDiffDown = Math.max(0, siteSummary.bandwidth_down - (prev.bandwidth_down || 0));
|
||||
const bytesDiffUp = Math.max(0, siteSummary.bandwidth_up - (prev.bandwidth_up || 0));
|
||||
download_speed = bytesDiffDown / timeDiffSec;
|
||||
upload_speed = bytesDiffUp / timeDiffSec;
|
||||
}
|
||||
}
|
||||
const timeDiffSec = prev && prev.timestamp ? (timestamp.getTime() - new Date(prev.timestamp).getTime()) / 1000 : 300;
|
||||
const activeTimeDiff = timeDiffSec > 0 ? timeDiffSec : 300;
|
||||
download_speed = siteSummary.bandwidth_down / activeTimeDiff;
|
||||
upload_speed = siteSummary.bandwidth_up / activeTimeDiff;
|
||||
} catch (err) {
|
||||
console.error('[Collector] Error calculating site summary speeds:', err.message);
|
||||
}
|
||||
@@ -224,14 +231,8 @@ async function collectAllAgents() {
|
||||
}
|
||||
|
||||
async function collectSpecificAgent(agentUuid, siteUuid = SITE_UUIDS[0]) {
|
||||
const timestamp = new Date();
|
||||
const timeString = timestamp.toLocaleString('id-ID', { timeZone: 'Asia/Jakarta' }) + ' WIB';
|
||||
console.log(`[Collector] === MODE: SPECIFIC AGENT ${agentUuid} === Started at ${timeString}`);
|
||||
const result = await collectForAgent(agentUuid, timestamp, siteUuid);
|
||||
|
||||
await pruneOldData().catch(err => console.error('[Collector] [Retention] error:', err.message));
|
||||
|
||||
return { success: result.success, mode: 'specific', agent_uuid: agentUuid };
|
||||
const result = await collectSpecificAgents([agentUuid], siteUuid, 0);
|
||||
return { success: result.successful > 0, mode: 'specific', agent_uuid: agentUuid };
|
||||
}
|
||||
|
||||
async function collectSpecificAgents(agentUuids, siteUuid = SITE_UUIDS[0], delayMs = 5000) {
|
||||
|
||||
@@ -20,7 +20,7 @@ const {
|
||||
async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify, label) {
|
||||
try {
|
||||
// 2b. App Categories
|
||||
const categories = await netify.fetchTopAppCategories(1440, 50, agentUuid, SITE_UUID);
|
||||
const categories = await netify.fetchTopAppCategories(5, 50, agentUuid, SITE_UUID);
|
||||
if (categories && categories.length > 0) {
|
||||
const catDocs = categories.map(c => ({
|
||||
timestamp,
|
||||
@@ -36,7 +36,7 @@ async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify
|
||||
}
|
||||
|
||||
// 2c. TLS Versions
|
||||
const tlsVersions = await netify.fetchTlsVersions(1440, 50, agentUuid, SITE_UUID);
|
||||
const tlsVersions = await netify.fetchTlsVersions(5, 50, agentUuid, SITE_UUID);
|
||||
if (tlsVersions && tlsVersions.length > 0) {
|
||||
const tvDocs = tlsVersions.map(v => ({
|
||||
timestamp,
|
||||
@@ -52,7 +52,7 @@ async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify
|
||||
}
|
||||
|
||||
// 2d. TLS Ciphers
|
||||
const tlsCiphers = await netify.fetchTlsCiphers(1440, 50, agentUuid, SITE_UUID);
|
||||
const tlsCiphers = await netify.fetchTlsCiphers(5, 50, agentUuid, SITE_UUID);
|
||||
if (tlsCiphers && tlsCiphers.length > 0) {
|
||||
const tcDocs = tlsCiphers.map(c => ({
|
||||
timestamp,
|
||||
@@ -68,7 +68,7 @@ async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify
|
||||
}
|
||||
|
||||
// 2e. TLS Security
|
||||
const tlsSecurity = await netify.fetchTlsSecurity(1440, 50, agentUuid, SITE_UUID);
|
||||
const tlsSecurity = await netify.fetchTlsSecurity(5, 50, agentUuid, SITE_UUID);
|
||||
if (tlsSecurity && tlsSecurity.length > 0) {
|
||||
const tsDocs = tlsSecurity.map(s => ({
|
||||
timestamp,
|
||||
@@ -84,7 +84,7 @@ async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify
|
||||
}
|
||||
|
||||
// 2f. Top Countries
|
||||
const countries = await netify.fetchTopCountries(1440, 100, agentUuid, SITE_UUID);
|
||||
const countries = await netify.fetchTopCountries(5, 100, agentUuid, SITE_UUID);
|
||||
if (countries && countries.length > 0) {
|
||||
const coDocs = countries.map(c => ({
|
||||
timestamp,
|
||||
@@ -101,7 +101,7 @@ async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify
|
||||
}
|
||||
|
||||
// 2g. Top Protocols
|
||||
const protocols = await netify.fetchTopProtocols(1440, 50, agentUuid, SITE_UUID);
|
||||
const protocols = await netify.fetchTopProtocols(5, 50, agentUuid, SITE_UUID);
|
||||
if (protocols && protocols.length > 0) {
|
||||
const protoDocs = protocols.map(p => ({
|
||||
timestamp,
|
||||
@@ -117,7 +117,7 @@ async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify
|
||||
}
|
||||
|
||||
// 2h. SNI Hostnames
|
||||
const snis = await netify.fetchSniHostnames(1440, 10000, agentUuid, SITE_UUID);
|
||||
const snis = await netify.fetchSniHostnames(5, 10000, agentUuid, SITE_UUID);
|
||||
if (snis && snis.length > 0) {
|
||||
const sniDocs = snis.map(s => ({
|
||||
timestamp, agent_uuid: agentUuid, site_uuid: SITE_UUID,
|
||||
|
||||
@@ -14,7 +14,7 @@ const {
|
||||
} = require('./deviceResolver');
|
||||
|
||||
async function collectDevicesAndApps(agentUuid, timestamp, SITE_UUID, netify, label) {
|
||||
const devices = await netify.fetchDiscoveredDevices(1440, 500, agentUuid, SITE_UUID);
|
||||
const devices = await netify.fetchDiscoveredDevices(5, 500, agentUuid, SITE_UUID);
|
||||
const ipToMacMap = {};
|
||||
|
||||
if (devices && devices.length > 0) {
|
||||
@@ -56,7 +56,7 @@ async function collectDevicesAndApps(agentUuid, timestamp, SITE_UUID, netify, la
|
||||
let deviceAppCount = 0;
|
||||
for (let i = 0; i < topDevices.length; i += 5) {
|
||||
const batch = topDevices.slice(i, i + 5);
|
||||
const results = await Promise.allSettled(batch.map(d => netify.fetchDeviceApps(d.ip_address, 1440, 50, agentUuid, SITE_UUID)));
|
||||
const results = await Promise.allSettled(batch.map(d => netify.fetchDeviceApps(d.ip_address, 5, 50, agentUuid, SITE_UUID)));
|
||||
const appDocs = [];
|
||||
results.forEach((res, idx) => {
|
||||
if (res.status === 'fulfilled' && Array.isArray(res.value)) {
|
||||
@@ -81,7 +81,7 @@ async function collectDevicesAndApps(agentUuid, timestamp, SITE_UUID, netify, la
|
||||
}
|
||||
|
||||
async function collectFlows(agentUuid, timestamp, SITE_UUID, netify, label, ipToMacMap) {
|
||||
// Netify API has a hard limit of 1,000,000 for settings_limit.
|
||||
// Netify API has a hard limit of 1,000,000 for settings_limit. Use 1000000 as default per rule.
|
||||
const flowLimit = parseInt(process.env.PROXY_FLOW_LIMIT || '1000000');
|
||||
const flows = await netify.fetchFlows(flowLimit, agentUuid, SITE_UUID);
|
||||
if (flows && flows.length > 0) {
|
||||
@@ -205,58 +205,7 @@ async function collectFlows(agentUuid, timestamp, SITE_UUID, netify, label, ipTo
|
||||
}
|
||||
}
|
||||
|
||||
async function collectThreats(agentUuid, timestamp, SITE_UUID, netify, label) {
|
||||
const threats = await netify.fetchCyberThreats(agentUuid, SITE_UUID);
|
||||
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(),
|
||||
}));
|
||||
await Threat.insertMany(threatDocs);
|
||||
console.log(`[Collector] ✓ ${threatDocs.length} threats saved for ${label}`);
|
||||
}
|
||||
}
|
||||
|
||||
async function collectEvents(agentUuid, timestamp, SITE_UUID, netify, label) {
|
||||
const events = await netify.fetchEvents(100, agentUuid, SITE_UUID);
|
||||
if (events && events.length > 0) {
|
||||
const eventIds = events.map(e => e.event_id).filter(id => id !== null);
|
||||
const existing = await Event.find({ site_uuid: SITE_UUID, event_id: { $in: eventIds } }).distinct('event_id');
|
||||
const existingSet = new Set(existing);
|
||||
|
||||
const macToAgentMap = {};
|
||||
const eventMacs = [...new Set(events.map(e => e.mac_address).filter(Boolean))];
|
||||
if (eventMacs.length > 0) {
|
||||
const storedDevices = await DeviceStat.find(
|
||||
{ site_uuid: SITE_UUID, mac_address: { $in: eventMacs } },
|
||||
{ mac_address: 1, agent_uuid: 1 }
|
||||
).lean();
|
||||
for (const d of storedDevices) {
|
||||
if (d.mac_address && d.agent_uuid) macToAgentMap[d.mac_address] = d.agent_uuid;
|
||||
}
|
||||
}
|
||||
|
||||
const eventDocs = events.filter(e => e.event_id === null || !existingSet.has(e.event_id)).map(e => {
|
||||
const resolvedAgentUuid = (e.mac_address && macToAgentMap[e.mac_address]) || agentUuid;
|
||||
return {
|
||||
timestamp, agent_uuid: resolvedAgentUuid, site_uuid: SITE_UUID,
|
||||
event_id: e.event_id, event_type: e.event_type, severity: e.severity,
|
||||
description: e.description, category_label: e.category_label,
|
||||
ip_address: e.ip_address, mac_address: e.mac_address, event_at: e.event_at,
|
||||
};
|
||||
});
|
||||
if (eventDocs.length > 0) {
|
||||
await Event.insertMany(eventDocs);
|
||||
console.log(`[Collector] ✓ ${eventDocs.length} new events saved for ${label}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
collectDevicesAndApps,
|
||||
collectFlows,
|
||||
collectThreats,
|
||||
collectEvents
|
||||
collectFlows
|
||||
};
|
||||
@@ -0,0 +1,61 @@
|
||||
// proxy/collectorHelperDpi3.js
|
||||
// ─────────────────────────────────────────────────────────────────────────────
|
||||
// Supplementary Telemetry collection steps for threats and events.
|
||||
// Split from collectorHelperDpi2.js to satisfy the 256-line file size limit.
|
||||
// ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
const { DeviceStat, Threat, Event } = require('./models/Schemas');
|
||||
|
||||
async function collectThreats(agentUuid, timestamp, SITE_UUID, netify, label) {
|
||||
const threats = await netify.fetchCyberThreats(agentUuid, SITE_UUID);
|
||||
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(),
|
||||
}));
|
||||
await Threat.insertMany(threatDocs);
|
||||
console.log(`[Collector] ✓ ${threatDocs.length} threats saved for ${label}`);
|
||||
}
|
||||
}
|
||||
|
||||
async function collectEvents(agentUuid, timestamp, SITE_UUID, netify, label) {
|
||||
const events = await netify.fetchEvents(100, agentUuid, SITE_UUID);
|
||||
if (events && events.length > 0) {
|
||||
const eventIds = events.map(e => e.event_id).filter(id => id !== null);
|
||||
const existing = await Event.find({ site_uuid: SITE_UUID, event_id: { $in: eventIds } }).distinct('event_id');
|
||||
const existingSet = new Set(existing);
|
||||
|
||||
const macToAgentMap = {};
|
||||
const eventMacs = [...new Set(events.map(e => e.mac_address).filter(Boolean))];
|
||||
if (eventMacs.length > 0) {
|
||||
const storedDevices = await DeviceStat.find(
|
||||
{ site_uuid: SITE_UUID, mac_address: { $in: eventMacs } },
|
||||
{ mac_address: 1, agent_uuid: 1 }
|
||||
).lean();
|
||||
for (const d of storedDevices) {
|
||||
if (d.mac_address && d.agent_uuid) macToAgentMap[d.mac_address] = d.agent_uuid;
|
||||
}
|
||||
}
|
||||
|
||||
const eventDocs = events.filter(e => e.event_id === null || !existingSet.has(e.event_id)).map(e => {
|
||||
const resolvedAgentUuid = (e.mac_address && macToAgentMap[e.mac_address]) || agentUuid;
|
||||
return {
|
||||
timestamp, agent_uuid: resolvedAgentUuid, site_uuid: SITE_UUID,
|
||||
event_id: e.event_id, event_type: e.event_type, severity: e.severity,
|
||||
description: e.description, category_label: e.category_label,
|
||||
ip_address: e.ip_address, mac_address: e.mac_address, event_at: e.event_at,
|
||||
};
|
||||
});
|
||||
if (eventDocs.length > 0) {
|
||||
await Event.insertMany(eventDocs);
|
||||
console.log(`[Collector] ✓ ${eventDocs.length} new events saved for ${label}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
collectThreats,
|
||||
collectEvents
|
||||
};
|
||||
@@ -11,9 +11,9 @@ const Schemas = require('./models/Schemas');
|
||||
* Prune all time-series documents older than 7 days.
|
||||
*/
|
||||
async function pruneOldData() {
|
||||
const sevenDaysAgo = new Date(Date.now() - 7 * 24 * 60 * 60 * 1000);
|
||||
const timeString = sevenDaysAgo.toLocaleString('id-ID', { timeZone: 'Asia/Jakarta' }) + ' WIB';
|
||||
console.log(`[Collector] [Retention] Checking for telemetry data older than 7 days (before ${timeString})...`);
|
||||
const thirtyDaysAgo = new Date(Date.now() - 30 * 24 * 60 * 60 * 1000);
|
||||
const timeString = thirtyDaysAgo.toLocaleString('id-ID', { timeZone: 'Asia/Jakarta' }) + ' WIB';
|
||||
console.log(`[Collector] [Retention] Checking for telemetry data older than 30 days (before ${timeString})...`);
|
||||
|
||||
// Prune from all collections in Schemas except CustomDeviceLabel
|
||||
const collections = Object.keys(Schemas).filter(name => name !== 'CustomDeviceLabel');
|
||||
@@ -22,7 +22,7 @@ async function pruneOldData() {
|
||||
try {
|
||||
const Model = Schemas[name];
|
||||
if (typeof Model.deleteMany === 'function') {
|
||||
const res = await Model.deleteMany({ timestamp: { $lt: sevenDaysAgo } });
|
||||
const res = await Model.deleteMany({ timestamp: { $lt: thirtyDaysAgo } });
|
||||
if (res.deletedCount > 0) {
|
||||
console.log(`[Collector] [Retention] ✓ Cleaned up ${res.deletedCount} old records from ${name}`);
|
||||
}
|
||||
|
||||
+21
-6
@@ -139,6 +139,7 @@ AppStatSchema.index({ agent_uuid: 1, timestamp: -1, download: -1 });
|
||||
DeviceStatSchema.index({ agent_uuid: 1, ip_address: 1 }, { unique: true });
|
||||
FlowSchema.index({ agent_uuid: 1, timestamp: -1 });
|
||||
FlowSchema.index({ agent_uuid: 1, flow_id: 1 });
|
||||
FlowSchema.index({ site_uuid: 1, timestamp: -1 });
|
||||
FlowSchema.index({ site_uuid: 1, app_label: 1, timestamp: -1 });
|
||||
ThreatSchema.index({ agent_uuid: 1, timestamp: -1 });
|
||||
AppCategoryStatSchema.index({ agent_uuid: 1, timestamp: -1 });
|
||||
@@ -167,16 +168,30 @@ DeviceAppStatSchema.index({ site_uuid: 1, app_label: 1, timestamp: -1 });
|
||||
const telemetrySchemas = require('./SchemasTelemetry');
|
||||
const auxSchemas = require('./SchemasAux');
|
||||
|
||||
// ─── Agent Registry (all agents registered in Netify, regardless of activity) ──
|
||||
// Upserted every collector cycle. Source of truth for the agents list page.
|
||||
const AgentRegistrySchema = new mongoose.Schema({
|
||||
uuid: { type: String, required: true, unique: true, index: true },
|
||||
serial: { type: String },
|
||||
label: { type: String },
|
||||
site_uuid: { type: String, index: true },
|
||||
provisioned: { type: Boolean, default: false },
|
||||
activated: { type: Boolean, default: false },
|
||||
last_seen_at: { type: mongoose.Schema.Types.Mixed },
|
||||
netify_id: { type: Number },
|
||||
}, { ...baseOptions, collection: 'agent_registry' });
|
||||
|
||||
module.exports = {
|
||||
Summary: mongoose.model('Summary', SummarySchema),
|
||||
AppStat: mongoose.model('AppStat', AppStatSchema),
|
||||
ProtocolStat:mongoose.model('ProtocolStat',ProtocolStatSchema),
|
||||
DeviceStat: mongoose.model('DeviceStat', DeviceStatSchema),
|
||||
Summary: mongoose.model('Summary', SummarySchema),
|
||||
AppStat: mongoose.model('AppStat', AppStatSchema),
|
||||
ProtocolStat: mongoose.model('ProtocolStat', ProtocolStatSchema),
|
||||
DeviceStat: mongoose.model('DeviceStat', DeviceStatSchema),
|
||||
DeviceAppStat: mongoose.model('DeviceAppStat', DeviceAppStatSchema),
|
||||
Flow: mongoose.model('Flow', FlowSchema),
|
||||
Threat: mongoose.model('Threat', ThreatSchema),
|
||||
Flow: mongoose.model('Flow', FlowSchema),
|
||||
Threat: mongoose.model('Threat', ThreatSchema),
|
||||
AppCategoryStat: mongoose.model('AppCategoryStat', AppCategoryStatSchema),
|
||||
Event: mongoose.model('Event', EventSchema),
|
||||
AgentRegistry: mongoose.model('AgentRegistry', AgentRegistrySchema),
|
||||
...auxSchemas,
|
||||
...telemetrySchemas,
|
||||
};
|
||||
|
||||
+16
-227
@@ -3,157 +3,11 @@
|
||||
// DPI API wrapper for the BackOne Proxy Server targeting original Netify API.
|
||||
// ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
const { netifyFetch, BASE_URL, agentMap } = require('./netifyClientCore');
|
||||
const { netifyFetch, BASE_URL } = require('./netifyClientCore');
|
||||
const telemetry = require('./netifyTelemetry');
|
||||
const stats = require('./netifyClientStats');
|
||||
|
||||
const PORT_SERVICE_MAP = {
|
||||
80: 'HTTP', 443: 'HTTPS / TLS', 8080: 'HTTP Alt', 8443: 'HTTPS Alt',
|
||||
53: 'DNS', 5353: 'mDNS', 853: 'DNS-over-TLS',
|
||||
25: 'SMTP', 587: 'SMTP TLS', 465: 'SMTPS', 110: 'POP3', 143: 'IMAP',
|
||||
22: 'SSH', 23: 'Telnet', 3389: 'RDP', 5900: 'VNC',
|
||||
21: 'FTP', 20: 'FTP Data', 989: 'FTPS', 990: 'FTPS Control',
|
||||
3306: 'MySQL', 5432: 'PostgreSQL', 6379: 'Redis', 27017: 'MongoDB',
|
||||
1194: 'OpenVPN', 51820: 'WireGuard', 500: 'IPSec IKE', 4500: 'IPSec NAT-T',
|
||||
67: 'DHCP', 68: 'DHCP Client', 123: 'NTP',
|
||||
6881: 'BitTorrent', 6882: 'BitTorrent', 6883: 'BitTorrent',
|
||||
9993: 'ZeroTier VPN',
|
||||
};
|
||||
|
||||
async function fetchAgents(siteUuid = null) {
|
||||
const data = await netifyFetch('/data/stats/top/agent/download', { filter_interval: 43200, settings_limit: 100 }, null, siteUuid);
|
||||
if (!data || !Array.isArray(data)) return [];
|
||||
const list = data.map(r => ({
|
||||
id: r.agent?.id,
|
||||
uuid: r.agent?.uuid,
|
||||
label: r.agent?.label,
|
||||
})).filter(a => a.uuid);
|
||||
|
||||
for (const a of list) {
|
||||
if (a.uuid && a.id) agentMap[a.uuid] = a.id;
|
||||
}
|
||||
|
||||
// Secondary validation: if this site already has data in MongoDB, only return agents
|
||||
// that have at least one summary record for THIS site_uuid. This prevents the
|
||||
// org-level stats endpoint from cross-contaminating agents across sites.
|
||||
if (siteUuid) {
|
||||
try {
|
||||
const mongoose = require('mongoose');
|
||||
if (mongoose.connection.readyState === 1) {
|
||||
const db = mongoose.connection.db;
|
||||
const knownAgents = await db.collection('summaries').distinct('agent_uuid', {
|
||||
site_uuid: siteUuid,
|
||||
agent_uuid: { $ne: null },
|
||||
});
|
||||
|
||||
if (knownAgents.length > 0) {
|
||||
const knownSet = new Set(knownAgents);
|
||||
const validated = list.filter(a => knownSet.has(a.uuid));
|
||||
// If MongoDB cross-check yields results, use the validated list.
|
||||
// On first boot (no DB data yet), fall through and use the full API list.
|
||||
if (validated.length > 0) return validated;
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
console.warn('[Collector] fetchAgents DB cross-check failed:', err.message);
|
||||
}
|
||||
}
|
||||
|
||||
return list;
|
||||
}
|
||||
|
||||
|
||||
async function fetchBandwidthSummary(interval = 1440, agentUuid = null, siteUuid = null) {
|
||||
const [dlData, ulData, flowsData] = await Promise.all([
|
||||
netifyFetch('/data/stats/top/local_ip/download', { filter_interval: interval, settings_limit: 500 }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/local_ip/upload', { filter_interval: interval, settings_limit: 500 }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/local_ip/flow_count', { filter_interval: interval, settings_limit: 500 }, agentUuid, siteUuid),
|
||||
]);
|
||||
const bandwidth_down = dlData?.reduce((s, r) => s + (r.download ?? 0), 0) ?? 0;
|
||||
const bandwidth_up = ulData?.reduce((s, r) => s + (r.upload ?? 0), 0) ?? 0;
|
||||
const total_devices = dlData?.length ?? 0;
|
||||
const active_flows = flowsData?.reduce((s, r) => s + (r.flows ?? r.flow_count ?? 0), 0) ?? 0;
|
||||
return { bandwidth_down, bandwidth_up, total_devices, active_flows, total_threats: 0 };
|
||||
}
|
||||
|
||||
async function fetchTopApps(interval = 1440, limit = 200, agentUuid = null, siteUuid = null) {
|
||||
const [dlData, ulData] = await Promise.all([
|
||||
netifyFetch('/data/stats/top/application/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/application/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
|
||||
]);
|
||||
if (!dlData) return null;
|
||||
const ulMap = {};
|
||||
if (ulData) {
|
||||
for (const r of ulData) {
|
||||
const id = r.application?.id;
|
||||
if (id) ulMap[id] = r.upload ?? 0;
|
||||
}
|
||||
}
|
||||
return dlData.map(r => ({
|
||||
application: {
|
||||
id: r.application?.id ?? null,
|
||||
label: r.application?.label ?? 'Unknown',
|
||||
tag: r.application?.tag ?? null,
|
||||
},
|
||||
download: r.download ?? 0,
|
||||
upload: ulMap[r.application?.id] ?? 0,
|
||||
flows: r.flows ?? 0,
|
||||
}));
|
||||
}
|
||||
|
||||
async function fetchDiscoveredDevices(interval = 1440, limit = 500, agentUuid = null, siteUuid = null) {
|
||||
const [dlData, ulData] = await Promise.all([
|
||||
netifyFetch('/data/stats/top/local_ip/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/local_ip/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
|
||||
]);
|
||||
if (!dlData) return null;
|
||||
const ulMap = {};
|
||||
if (ulData) {
|
||||
for (const r of ulData) {
|
||||
const ip = r.local_ip?.address ?? String(r.local_ip);
|
||||
ulMap[ip] = r.upload ?? 0;
|
||||
}
|
||||
}
|
||||
return dlData.map(r => {
|
||||
const ip = r.local_ip?.address ?? String(r.local_ip ?? '');
|
||||
return {
|
||||
ip_address: ip,
|
||||
mac_address: r.local_mac ?? null,
|
||||
device_label: r.device_label ?? ip,
|
||||
device_type: r.device_type ?? null,
|
||||
os_label: r.os_label ?? null,
|
||||
manufacturer: r.manufacturer ?? null,
|
||||
download: r.download ?? 0,
|
||||
upload: ulMap[ip] ?? 0,
|
||||
flows: r.flows ?? 0,
|
||||
last_seen: r.last_seen_at?.date ?? null,
|
||||
};
|
||||
}).filter(d => d.ip_address);
|
||||
}
|
||||
|
||||
async function fetchDeviceApps(ipAddress, interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
|
||||
const ipFilter = JSON.stringify([ipAddress]);
|
||||
const [dlData, ulData] = await Promise.all([
|
||||
netifyFetch('/data/stats/top/application/download', { filter_interval: interval, settings_limit: limit, filter_local_ips: ipFilter }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/application/upload', { filter_interval: interval, settings_limit: limit, filter_local_ips: ipFilter }, agentUuid, siteUuid),
|
||||
]);
|
||||
if (!dlData || !Array.isArray(dlData)) return [];
|
||||
const ulMap = {};
|
||||
if (ulData && Array.isArray(ulData)) {
|
||||
for (const r of ulData) {
|
||||
const id = r.application?.id;
|
||||
if (id) ulMap[id] = r.upload ?? 0;
|
||||
}
|
||||
}
|
||||
return dlData.map(r => ({
|
||||
app_label: r.application?.label ?? 'Unknown',
|
||||
app_id: r.application?.id ?? null,
|
||||
download: r.download ?? 0,
|
||||
upload: ulMap[r.application?.id] ?? 0,
|
||||
flows: r.flows ?? 0,
|
||||
})).filter(a => a.download > 0 || a.upload > 0);
|
||||
}
|
||||
|
||||
async function fetchFlows(limit = 10000, agentUuid = null, siteUuid = null) {
|
||||
async function fetchFlows(limit = 1000000, agentUuid = null, siteUuid = null) {
|
||||
const [raw, sniRaw] = await Promise.all([
|
||||
netifyFetch('/data/flows', { settings_limit: limit }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/tls_sni/download', { filter_interval: 1440, settings_limit: 50 }, agentUuid, siteUuid),
|
||||
@@ -163,12 +17,14 @@ async function fetchFlows(limit = 10000, agentUuid = null, siteUuid = null) {
|
||||
if (sniRaw && Array.isArray(sniRaw)) {
|
||||
for (const r of sniRaw) {
|
||||
const sni = typeof r.tls_sni === 'object' ? r.tls_sni?.label : r.tls_sni;
|
||||
if (sni && typeof sni === 'string' && sni.trim() !== '') sniList.push(sni.replace(/^\*\./, '').trim());
|
||||
if (sni && typeof sni === 'string' && sni.trim() !== '') {
|
||||
sniList.push(sni.replace(/^\*\./, '').trim());
|
||||
}
|
||||
}
|
||||
}
|
||||
return raw.map(r => {
|
||||
const port = r.remote_port ?? null;
|
||||
const portService = port ? (PORT_SERVICE_MAP[port] ?? `Port ${port}`) : null;
|
||||
const portService = port ? (stats.PORT_SERVICE_MAP[port] ?? `Port ${port}`) : null;
|
||||
const appLabel = r.application?.label || portService;
|
||||
const domain = r.tls_sni || r.dns_hostname || r.hostname || null;
|
||||
return {
|
||||
@@ -188,84 +44,17 @@ async function fetchFlows(limit = 10000, agentUuid = null, siteUuid = null) {
|
||||
}).filter(f => f.src_ip);
|
||||
}
|
||||
|
||||
async function fetchCyberThreats(agentUuid = null, siteUuid = null) {
|
||||
const ipRepData = await netifyFetch('/data/stats/top/remote_ip/download', { filter_interval: 1440, settings_limit: 50 }, agentUuid, siteUuid);
|
||||
if (!ipRepData || !Array.isArray(ipRepData)) return [];
|
||||
const SUSPICIOUS_PORTS = new Set([23, 4444, 1337, 6667, 31337, 12345, 54321, 4899, 5554, 9999]);
|
||||
const threats = [];
|
||||
for (const r of ipRepData) {
|
||||
const ip = r.remote_ip?.address ?? null;
|
||||
const port = r.remote_port ?? 0;
|
||||
if (ip && SUSPICIOUS_PORTS.has(port)) {
|
||||
threats.push({
|
||||
threat_type: `Suspicious Port ${port}`,
|
||||
severity: 'High',
|
||||
src_ip: null,
|
||||
dst_ip: ip,
|
||||
dst_port: port,
|
||||
protocol: r.ip_protocol?.label ?? null,
|
||||
description: `Suspicious outbound connection to ${ip}:${port}`,
|
||||
event_at: new Date().toISOString(),
|
||||
});
|
||||
}
|
||||
}
|
||||
return threats;
|
||||
}
|
||||
|
||||
async function fetchEvents(limit = 100, agentUuid = null, siteUuid = null) {
|
||||
const raw = await netifyFetch('/event/events', { settings_limit: limit }, agentUuid, siteUuid);
|
||||
if (!raw || !Array.isArray(raw)) return [];
|
||||
return raw.map(r => {
|
||||
let msg = r.label || '';
|
||||
if (r.description) {
|
||||
try {
|
||||
const descObj = JSON.parse(r.description);
|
||||
msg = descObj.default || r.label || '';
|
||||
if (descObj.tags) {
|
||||
for (const k in descObj.tags) {
|
||||
const tagVal = descObj.tags[k];
|
||||
const val = Array.isArray(tagVal) ? (tagVal[0] === 'Unknown' && tagVal[1] ? tagVal[1] : tagVal[0]) : tagVal;
|
||||
msg = msg.replace(`{{ ${k} }}`, val).replace(`{{${k}}}`, val);
|
||||
}
|
||||
}
|
||||
} catch (e) {
|
||||
msg = r.description;
|
||||
}
|
||||
}
|
||||
let sevLabel = 'Info';
|
||||
if (r.severity >= 30) sevLabel = 'Critical';
|
||||
else if (r.severity >= 20) sevLabel = 'High';
|
||||
else if (r.severity >= 10) sevLabel = 'Warning';
|
||||
let srcIp = null;
|
||||
if (r.description) {
|
||||
try {
|
||||
const descObj = JSON.parse(r.description);
|
||||
srcIp = descObj.tags?.device_ip || descObj.tags?.ip || null;
|
||||
} catch {}
|
||||
}
|
||||
return {
|
||||
event_id: r.id || null,
|
||||
event_type: r.basename || 'unknown',
|
||||
severity: sevLabel,
|
||||
description: msg,
|
||||
category_label: r.category?.label || 'Intelligence',
|
||||
ip_address: srcIp,
|
||||
mac_address: r.additional?.device?.mac?.address || null,
|
||||
event_at: r.created_at?.date ? new Date(r.created_at.date) : new Date()
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
fetchAgents,
|
||||
fetchBandwidthSummary,
|
||||
fetchTopApps,
|
||||
fetchDiscoveredDevices,
|
||||
fetchDeviceApps,
|
||||
fetchFlows,
|
||||
fetchCyberThreats,
|
||||
fetchEvents,
|
||||
fetchAgents: stats.fetchAgents,
|
||||
fetchBandwidthSummary: stats.fetchBandwidthSummary,
|
||||
fetchTopApps: stats.fetchTopApps,
|
||||
fetchDiscoveredDevices: stats.fetchDiscoveredDevices,
|
||||
fetchDeviceApps: stats.fetchDeviceApps,
|
||||
fetchCyberThreats: stats.fetchCyberThreats,
|
||||
fetchEvents: stats.fetchEvents,
|
||||
syncApplicationDictionary: stats.syncApplicationDictionary,
|
||||
BASE_URL,
|
||||
PORT_SERVICE_MAP,
|
||||
PORT_SERVICE_MAP: stats.PORT_SERVICE_MAP,
|
||||
...telemetry,
|
||||
};
|
||||
@@ -0,0 +1,238 @@
|
||||
// proxy/netifyClientStats.js
|
||||
// ─────────────────────────────────────────────────────────────────────────────
|
||||
// Supplementary fetchers split from netifyClient.js to satisfy the 256-line limit.
|
||||
// ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
const { netifyFetch, agentMap } = require('./netifyClientCore');
|
||||
const PORT_SERVICE_MAP = {
|
||||
80: 'HTTP', 443: 'HTTPS / TLS', 8080: 'HTTP Alt', 8443: 'HTTPS Alt', 53: 'DNS', 5353: 'mDNS', 853: 'DNS-over-TLS', 25: 'SMTP', 587: 'SMTP TLS', 465: 'SMTPS', 110: 'POP3', 143: 'IMAP', 22: 'SSH', 23: 'Telnet', 3389: 'RDP', 5900: 'VNC', 21: 'FTP', 20: 'FTP Data', 989: 'FTPS', 990: 'FTPS Control', 3306: 'MySQL', 5432: 'PostgreSQL', 6379: 'Redis', 27017: 'MongoDB', 1194: 'OpenVPN', 51820: 'WireGuard', 500: 'IPSec IKE', 4500: 'IPSec NAT-T', 67: 'DHCP', 68: 'DHCP Client', 123: 'NTP', 6881: 'BitTorrent', 6882: 'BitTorrent', 6883: 'BitTorrent', 9993: 'ZeroTier VPN',
|
||||
};
|
||||
|
||||
async function fetchAgents(siteUuid = null) {
|
||||
// Use /data/stats/top/agent/download — the only working agent-listing endpoint
|
||||
// in Netify Informatics API. /data/agents returns 404 (only on portal API).
|
||||
// filter_interval max = 525600 min (365 days) per Netify API validation.
|
||||
// Use max interval to catch ALL agents including low-traffic / inactive ones.
|
||||
// DB cross-check removed — was blocking new agents from appearing.
|
||||
const data = await netifyFetch('/data/stats/top/agent/download', {
|
||||
filter_interval: 525600,
|
||||
settings_limit: 1000000,
|
||||
}, null, siteUuid);
|
||||
if (!data || !Array.isArray(data)) return [];
|
||||
|
||||
const list = data.map(r => ({
|
||||
id: r.agent?.id,
|
||||
uuid: r.agent?.uuid || r.agent?.serial,
|
||||
serial: r.agent?.serial,
|
||||
label: r.agent?.label || r.agent?.serial,
|
||||
provisioned: true,
|
||||
activated: true,
|
||||
last_seen_at: r.agent?.last_seen_at ?? null,
|
||||
})).filter(a => a.uuid);
|
||||
|
||||
// Populate the agentMap (uuid → id) for downstream filter_agents usage
|
||||
for (const a of list) {
|
||||
if (a.uuid && a.id) agentMap[a.uuid] = a.id;
|
||||
}
|
||||
|
||||
return list;
|
||||
}
|
||||
|
||||
async function fetchBandwidthSummary(interval = 1440, agentUuid = null, siteUuid = null) {
|
||||
const [dlData, ulData, flowsData] = await Promise.all([
|
||||
netifyFetch('/data/stats/top/local_ip/download', { filter_interval: interval, settings_limit: 500 }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/local_ip/upload', { filter_interval: interval, settings_limit: 500 }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/local_ip/flow_count', { filter_interval: interval, settings_limit: 500 }, agentUuid, siteUuid),
|
||||
]);
|
||||
const bandwidth_down = dlData?.reduce((s, r) => s + (r.download ?? 0), 0) ?? 0;
|
||||
const bandwidth_up = ulData?.reduce((s, r) => s + (r.upload ?? 0), 0) ?? 0;
|
||||
const total_devices = dlData?.length ?? 0;
|
||||
const active_flows = flowsData?.reduce((s, r) => s + (r.flows ?? r.flow_count ?? 0), 0) ?? 0;
|
||||
return { bandwidth_down, bandwidth_up, total_devices, active_flows, total_threats: 0 };
|
||||
}
|
||||
|
||||
async function fetchTopApps(interval = 1440, limit = 200, agentUuid = null, siteUuid = null) {
|
||||
const [dlData, ulData] = await Promise.all([
|
||||
netifyFetch('/data/stats/top/application/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/application/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
|
||||
]);
|
||||
if (!dlData) return null;
|
||||
const ulMap = {};
|
||||
if (ulData) {
|
||||
for (const r of ulData) {
|
||||
const id = r.application?.id;
|
||||
if (id) ulMap[id] = r.upload ?? 0;
|
||||
}
|
||||
}
|
||||
return dlData.map(r => ({
|
||||
application: { id: r.application?.id ?? null, label: r.application?.label ?? 'Unknown', tag: r.application?.tag ?? null },
|
||||
download: r.download ?? 0, upload: ulMap[r.application?.id] ?? 0, flows: r.flows ?? 0,
|
||||
}));
|
||||
}
|
||||
|
||||
async function fetchDiscoveredDevices(interval = 1440, limit = 500, agentUuid = null, siteUuid = null) {
|
||||
const [dlData, ulData] = await Promise.all([
|
||||
netifyFetch('/data/stats/top/local_ip/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/local_ip/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
|
||||
]);
|
||||
if (!dlData) return null;
|
||||
const ulMap = {};
|
||||
if (ulData) {
|
||||
for (const r of ulData) {
|
||||
const ip = r.local_ip?.address ?? String(r.local_ip);
|
||||
ulMap[ip] = r.upload ?? 0;
|
||||
}
|
||||
}
|
||||
return dlData.map(r => {
|
||||
const ip = r.local_ip?.address ?? String(r.local_ip ?? '');
|
||||
return {
|
||||
ip_address: ip, mac_address: r.local_mac ?? null, device_label: r.device_label ?? ip, device_type: r.device_type ?? null,
|
||||
os_label: r.os_label ?? null, manufacturer: r.manufacturer ?? null, download: r.download ?? 0, upload: ulMap[ip] ?? 0,
|
||||
flows: r.flows ?? 0, last_seen: r.last_seen_at?.date ?? null,
|
||||
};
|
||||
}).filter(d => d.ip_address);
|
||||
}
|
||||
|
||||
async function fetchDeviceApps(ipAddress, interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
|
||||
const ipFilter = JSON.stringify([ipAddress]);
|
||||
const [dlData, ulData] = await Promise.all([
|
||||
netifyFetch('/data/stats/top/application/download', { filter_interval: interval, settings_limit: limit, filter_local_ips: ipFilter }, agentUuid, siteUuid),
|
||||
netifyFetch('/data/stats/top/application/upload', { filter_interval: interval, settings_limit: limit, filter_local_ips: ipFilter }, agentUuid, siteUuid),
|
||||
]);
|
||||
if (!dlData || !Array.isArray(dlData)) return [];
|
||||
const ulMap = {};
|
||||
if (ulData && Array.isArray(ulData)) {
|
||||
for (const r of ulData) {
|
||||
const id = r.application?.id;
|
||||
if (id) ulMap[id] = r.upload ?? 0;
|
||||
}
|
||||
}
|
||||
return dlData.map(r => ({
|
||||
app_label: r.application?.label ?? 'Unknown',
|
||||
app_id: r.application?.id ?? null,
|
||||
download: r.download ?? 0,
|
||||
upload: ulMap[r.application?.id] ?? 0,
|
||||
flows: r.flows ?? 0,
|
||||
})).filter(a => a.download > 0 || a.upload > 0);
|
||||
}
|
||||
|
||||
async function fetchCyberThreats(agentUuid = null, siteUuid = null) {
|
||||
const ipRepData = await netifyFetch('/data/stats/top/remote_ip/download', { filter_interval: 1440, settings_limit: 50 }, agentUuid, siteUuid);
|
||||
if (!ipRepData || !Array.isArray(ipRepData)) return [];
|
||||
const SUSPICIOUS_PORTS = new Set([23, 4444, 1337, 6667, 31337, 12345, 54321, 4899, 5554, 9999]);
|
||||
const threats = [];
|
||||
for (const r of ipRepData) {
|
||||
const ip = r.remote_ip?.address ?? null;
|
||||
const port = r.remote_port ?? 0;
|
||||
if (ip && SUSPICIOUS_PORTS.has(port)) {
|
||||
threats.push({
|
||||
threat_type: `Suspicious Port ${port}`,
|
||||
severity: 'High',
|
||||
src_ip: null,
|
||||
dst_ip: ip,
|
||||
dst_port: port,
|
||||
protocol: r.ip_protocol?.label ?? null,
|
||||
description: `Suspicious outbound connection to ${ip}:${port}`,
|
||||
event_at: new Date().toISOString(),
|
||||
});
|
||||
}
|
||||
}
|
||||
return threats;
|
||||
}
|
||||
|
||||
async function fetchEvents(limit = 100, agentUuid = null, siteUuid = null) {
|
||||
const raw = await netifyFetch('/event/events', { settings_limit: limit }, agentUuid, siteUuid);
|
||||
if (!raw || !Array.isArray(raw)) return [];
|
||||
return raw.map(r => {
|
||||
let msg = r.label || '';
|
||||
if (r.description) {
|
||||
try {
|
||||
const descObj = JSON.parse(r.description);
|
||||
msg = descObj.default || r.label || '';
|
||||
if (descObj.tags) {
|
||||
for (const k in descObj.tags) {
|
||||
const tagVal = descObj.tags[k];
|
||||
const val = Array.isArray(tagVal) ? (tagVal[0] === 'Unknown' && tagVal[1] ? tagVal[1] : tagVal[0]) : tagVal;
|
||||
msg = msg.replace(`{{ ${k} }}`, val).replace(`{{${k}}}`, val);
|
||||
}
|
||||
}
|
||||
} catch (e) {
|
||||
msg = r.description;
|
||||
}
|
||||
}
|
||||
let sevLabel = 'Info';
|
||||
if (r.severity >= 30) sevLabel = 'Critical';
|
||||
else if (r.severity >= 20) sevLabel = 'High';
|
||||
else if (r.severity >= 10) sevLabel = 'Warning';
|
||||
let srcIp = null;
|
||||
if (r.description) {
|
||||
try {
|
||||
const descObj = JSON.parse(r.description);
|
||||
srcIp = descObj.tags?.device_ip || descObj.tags?.ip || null;
|
||||
} catch {}
|
||||
}
|
||||
return {
|
||||
event_id: r.id || null,
|
||||
event_type: r.basename || 'unknown',
|
||||
severity: sevLabel,
|
||||
description: msg,
|
||||
category_label: r.category?.label || 'Intelligence',
|
||||
ip_address: srcIp,
|
||||
mac_address: r.additional?.device?.mac?.address || null,
|
||||
event_at: r.created_at?.date ? new Date(r.created_at.date) : new Date()
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
async function syncApplicationDictionary() {
|
||||
const mongoose = require('mongoose');
|
||||
const { LookupApp } = require('./models/Schemas');
|
||||
|
||||
console.log('[Netify] Fetching application catalog...');
|
||||
const allApps = await netifyFetch('/lookup/applications', { settings_limit: 5000 });
|
||||
if (!allApps || !Array.isArray(allApps)) {
|
||||
console.error('[Netify] Failed to fetch application dictionary.');
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(`[Netify] Application catalog fetched successfully. Got ${allApps.length} apps.`);
|
||||
if (allApps.length === 0) return;
|
||||
|
||||
console.log(`[Netify] Syncing ${allApps.length} application definitions to MongoDB...`);
|
||||
await LookupApp.deleteMany({});
|
||||
|
||||
const batchSize = 100;
|
||||
for (let i = 0; i < allApps.length; i += batchSize) {
|
||||
const batch = allApps.slice(i, i + batchSize);
|
||||
await LookupApp.insertMany(batch.map(app => ({
|
||||
id: app.id,
|
||||
name: app.name,
|
||||
label: app.label,
|
||||
tag: app.tag,
|
||||
description: app.description,
|
||||
full_name: app.full_name || app.application?.full_label || null,
|
||||
favicon: app.favicon || app.application?.favicon || null,
|
||||
icon: app.icon || app.application?.icon || null,
|
||||
logo: app.logo || app.application?.logo || null,
|
||||
application_category: {
|
||||
id: app.application_category?.id,
|
||||
name: app.application_category?.name,
|
||||
label: app.application_category?.label,
|
||||
tag: app.application_category?.tag
|
||||
}
|
||||
})));
|
||||
}
|
||||
console.log('[Netify] ✓ Application dictionary sync completed.');
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
fetchAgents,
|
||||
fetchBandwidthSummary,
|
||||
fetchTopApps,
|
||||
fetchDiscoveredDevices,
|
||||
fetchDeviceApps,
|
||||
fetchCyberThreats,
|
||||
fetchEvents,
|
||||
syncApplicationDictionary,
|
||||
PORT_SERVICE_MAP
|
||||
};
|
||||
+6
-1
@@ -4,6 +4,7 @@
|
||||
|
||||
const cron = require('node-cron');
|
||||
const { collectAllAgents, collectSpecificAgent, collectSpecificAgents } = require('./collector');
|
||||
const { syncApplicationDictionary } = require('./netifyClient');
|
||||
|
||||
// Mode: 'all' = collect all agents, 'agent' = collect one specific agent, 'agents' = collect multiple agents
|
||||
const COLLECT_MODE = process.env.PROXY_COLLECT_MODE || 'all';
|
||||
@@ -97,7 +98,11 @@ function startScheduler() {
|
||||
// Initial run immediately on startup (async, do not block server start)
|
||||
setTimeout(async () => {
|
||||
// 1. Sync dictionary first
|
||||
// await netifyClient.syncApplicationDictionary();
|
||||
try {
|
||||
await syncApplicationDictionary();
|
||||
} catch (err) {
|
||||
console.error('[Scheduler] Error syncing application dictionary:', err.message);
|
||||
}
|
||||
|
||||
// 2. Start normal telemetry collection
|
||||
runCollection().catch(err => console.error('[Scheduler] Initial run error:', err.message));
|
||||
|
||||
Reference in new issue
Block a user