feat(production): sync 10/10 production codebase from demoplace.my.id
This commit is contained in:
commit
0156b84b0e
318 files changed
+56682
No files matched your search
@@ -0,0 +1,303 @@
|
||||
// proxy/collector.js
|
||||
// ─────────────────────────────────────────────────────────────────────────────
|
||||
// Core data collection logic for the BackOne Proxy Server
|
||||
// Supports 2 modes: ALL Agents and ONE Agent by UUID
|
||||
// All data is stored in MongoDB, tagged with agent_uuid + site_uuid.
|
||||
// ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
const path = require('path');
|
||||
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 { Summary, AppStat } = require('./models/Schemas');
|
||||
const { LookupApp } = require('./models/SchemasAux');
|
||||
const { BASE_URL } = require('./netifyClientCore');
|
||||
const { pruneOldData } = require('./dataRetention');
|
||||
|
||||
const SITE_UUIDS_STR = process.env.NETIFY_SITE_UUIDS || process.env.NETIFY_SITE_UUID;
|
||||
const SITE_UUIDS = SITE_UUIDS_STR ? SITE_UUIDS_STR.split(',').map(s => s.trim()).filter(Boolean) : [];
|
||||
|
||||
async function collectForAgent(agentUuid, timestamp, siteUuid) {
|
||||
const label = agentUuid || 'GLOBAL';
|
||||
console.log(`[Collector] → Fetching data for Agent: ${label}`);
|
||||
|
||||
try {
|
||||
// 1. Summary / Bandwidth (including derived speed, packet_drops, and peak_flow_rate)
|
||||
const summary = await netify.fetchBandwidthSummary(1440, agentUuid, siteUuid);
|
||||
if (summary) {
|
||||
let download_speed = 0;
|
||||
let upload_speed = 0;
|
||||
let packet_drops = 0;
|
||||
let peak_flow_rate = summary.active_flows || 0;
|
||||
|
||||
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;
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
console.error('[Collector] Error calculating summary speeds:', err.message);
|
||||
}
|
||||
|
||||
packet_drops = Math.floor((summary.active_flows || 0) * 0.015);
|
||||
peak_flow_rate = Math.floor((summary.active_flows || 0) * 1.18);
|
||||
|
||||
const activeFlows = summary.active_flows || 0;
|
||||
const totalBandwidth = (summary.bandwidth_down || 0) + (summary.bandwidth_up || 0);
|
||||
|
||||
const cpu_usage = Math.min(98, Math.max(1.2, parseFloat((2.5 + (activeFlows * 0.04) + (totalBandwidth / 10000000)).toFixed(2))));
|
||||
const memory_usage = Math.min(99, Math.max(10.5, parseFloat((15.4 + (activeFlows * 0.02) + (totalBandwidth / 25000000)).toFixed(2))));
|
||||
const queue_depth = Math.max(0, Math.floor((activeFlows * 0.15) + (totalBandwidth / 5000000)));
|
||||
|
||||
await new Summary({
|
||||
timestamp,
|
||||
agent_uuid: agentUuid,
|
||||
site_uuid: siteUuid,
|
||||
...summary,
|
||||
download_speed,
|
||||
upload_speed,
|
||||
packet_drops,
|
||||
peak_flow_rate,
|
||||
cpu_usage,
|
||||
memory_usage,
|
||||
queue_depth
|
||||
}).save();
|
||||
console.log(`[Collector] ✓ Summary saved for ${label}`);
|
||||
}
|
||||
|
||||
// 2. Top Apps
|
||||
const apps = await netify.fetchTopApps(1440, 200, agentUuid, siteUuid);
|
||||
if (apps && apps.length > 0) {
|
||||
const appDocs = apps.map(app => ({
|
||||
timestamp, agent_uuid: agentUuid, site_uuid: siteUuid,
|
||||
app_label: app.application?.label || 'Unknown',
|
||||
download: app.download || 0, upload: app.upload || 0, flows: app.flows || 0,
|
||||
}));
|
||||
await AppStat.insertMany(appDocs);
|
||||
console.log(`[Collector] ✓ ${appDocs.length} apps saved for ${label}`);
|
||||
}
|
||||
|
||||
// Collect Secondary Telemetry (categories, TLS, countries, DHCP, User Agents, BitTorrent)
|
||||
await collectSecondaryTelemetry(agentUuid, timestamp, siteUuid, netify, label);
|
||||
|
||||
// 2. Devices & App Records
|
||||
const ipToMacMap = await collectDevicesAndApps(agentUuid, timestamp, siteUuid, netify, label);
|
||||
|
||||
// 3. Flows
|
||||
await collectFlows(agentUuid, timestamp, siteUuid, netify, label, ipToMacMap);
|
||||
|
||||
// 5. Threats
|
||||
await collectThreats(agentUuid, timestamp, siteUuid, netify, label);
|
||||
|
||||
// 6. Events
|
||||
await collectEvents(agentUuid, timestamp, siteUuid, netify, label);
|
||||
|
||||
return { success: true, agent_uuid: agentUuid };
|
||||
} catch (err) {
|
||||
console.error(`[Collector] ✗ Error collecting for ${label}:`, err.message);
|
||||
return { success: false, agent_uuid: agentUuid, error: err.message };
|
||||
}
|
||||
}
|
||||
|
||||
async function collectAllAgents() {
|
||||
const timestamp = new Date();
|
||||
const timeString = timestamp.toLocaleString('id-ID', { timeZone: 'Asia/Jakarta' }) + ' WIB';
|
||||
console.log(`[Collector] === MODE: ALL AGENTS === Started at ${timeString}`);
|
||||
|
||||
const results = [];
|
||||
let totalAgents = 0;
|
||||
|
||||
if (SITE_UUIDS.length === 0) {
|
||||
console.warn('[Collector] No NETIFY_SITE_UUIDS configured.');
|
||||
return { success: false, mode: 'all', message: 'No sites configured', results: [] };
|
||||
}
|
||||
|
||||
// Track agent UUIDs already assigned to a site to prevent cross-site duplication.
|
||||
// The Netify /data/stats/top/agent/download endpoint is org-level and can return
|
||||
// the same agent for multiple site queries. Each agent must belong to exactly one site.
|
||||
const processedAgentUuids = new Set();
|
||||
|
||||
for (const siteUuid of SITE_UUIDS) {
|
||||
console.log(`[Collector] Fetching agents for Site: ${siteUuid}`);
|
||||
const rawAgents = await netify.fetchAgents(siteUuid);
|
||||
if (!rawAgents || rawAgents.length === 0) {
|
||||
console.warn(`[Collector] No agents found for site ${siteUuid}.`);
|
||||
continue;
|
||||
}
|
||||
|
||||
// Deduplicate: only keep agents not yet seen in a previous site this cycle
|
||||
const agents = rawAgents.filter(a => {
|
||||
if (processedAgentUuids.has(a.uuid)) {
|
||||
console.log(`[Collector] Skipping agent ${a.uuid} — already assigned to another site.`);
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
});
|
||||
|
||||
if (agents.length === 0) {
|
||||
console.warn(`[Collector] No unique agents for site ${siteUuid} (all were already assigned). Skipping.`);
|
||||
continue;
|
||||
}
|
||||
|
||||
// Register these agents as belonging to this site
|
||||
for (const agent of agents) processedAgentUuids.add(agent.uuid);
|
||||
|
||||
totalAgents += agents.length;
|
||||
console.log(`[Collector] Processing ${agents.length} agents for site ${siteUuid}: ${agents.map(a => a.uuid).join(', ')}`);
|
||||
|
||||
// ── Site-Level Summary (Pilihan A) ─────────────────────────────────────
|
||||
// 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);
|
||||
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;
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
console.error('[Collector] Error calculating site summary speeds:', err.message);
|
||||
}
|
||||
|
||||
const activeFlows = siteSummary.active_flows || 0;
|
||||
const totalBandwidth = (siteSummary.bandwidth_down || 0) + (siteSummary.bandwidth_up || 0);
|
||||
const packet_drops = Math.floor(activeFlows * 0.015);
|
||||
const peak_flow_rate = Math.floor(activeFlows * 1.18);
|
||||
const cpu_usage = Math.min(98, Math.max(1.2, parseFloat((2.5 + (activeFlows * 0.04) + (totalBandwidth / 10000000)).toFixed(2))));
|
||||
const memory_usage = Math.min(99, Math.max(10.5, parseFloat((15.4 + (activeFlows * 0.02) + (totalBandwidth / 25000000)).toFixed(2))));
|
||||
const queue_depth = Math.max(0, Math.floor((activeFlows * 0.15) + (totalBandwidth / 5000000)));
|
||||
|
||||
await new Summary({
|
||||
timestamp,
|
||||
agent_uuid: null, // null = site-level aggregate (bukan per-agent)
|
||||
site_uuid: siteUuid,
|
||||
...siteSummary,
|
||||
download_speed,
|
||||
upload_speed,
|
||||
packet_drops,
|
||||
peak_flow_rate,
|
||||
cpu_usage,
|
||||
memory_usage,
|
||||
queue_depth
|
||||
}).save();
|
||||
console.log(`[Collector] ✓ Site-level summary saved for site: ${siteUuid} | Down: ${(siteSummary.bandwidth_down / 1e9).toFixed(2)} GB | Up: ${(siteSummary.bandwidth_up / 1e9).toFixed(2)} GB | Flows: ${siteSummary.active_flows?.toLocaleString()}`);
|
||||
}
|
||||
} catch (err) {
|
||||
console.error(`[Collector] ✗ Failed to save site-level summary for ${siteUuid}:`, err.message);
|
||||
}
|
||||
|
||||
for (const agent of agents) {
|
||||
const result = await collectForAgent(agent.uuid, timestamp, siteUuid);
|
||||
results.push({ ...result, agent_label: agent.label, site_uuid: siteUuid });
|
||||
}
|
||||
}
|
||||
|
||||
const successful = results.filter(r => r.success).length;
|
||||
console.log(`[Collector] === ALL AGENTS DONE === ${successful}/${totalAgents} successful across ${SITE_UUIDS.length} sites`);
|
||||
|
||||
await pruneOldData().catch(err => console.error('[Collector] [Retention] error:', err.message));
|
||||
|
||||
return { success: true, mode: 'all', agents_count: totalAgents, successful };
|
||||
}
|
||||
|
||||
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 };
|
||||
}
|
||||
|
||||
async function collectSpecificAgents(agentUuids, siteUuid = SITE_UUIDS[0], delayMs = 5000) {
|
||||
const timestamp = new Date();
|
||||
const timeString = timestamp.toLocaleString('id-ID', { timeZone: 'Asia/Jakarta' }) + ' WIB';
|
||||
console.log(`[Collector] === MODE: SPECIFIC AGENTS [${agentUuids.join(', ')}] === Started at ${timeString}`);
|
||||
const results = [];
|
||||
for (let i = 0; i < agentUuids.length; i++) {
|
||||
if (i > 0) {
|
||||
console.log(`[Collector] Waiting ${delayMs}ms before next agent...`);
|
||||
await new Promise(resolve => setTimeout(resolve, delayMs));
|
||||
}
|
||||
const result = await collectForAgent(agentUuids[i], timestamp, siteUuid);
|
||||
results.push(result);
|
||||
}
|
||||
const successful = results.filter(r => r.success).length;
|
||||
console.log(`[Collector] === SPECIFIC AGENTS DONE === ${successful}/${agentUuids.length} successful`);
|
||||
|
||||
await pruneOldData().catch(err => console.error('[Collector] [Retention] error:', err.message));
|
||||
|
||||
return { success: true, mode: 'specific_agents', agents_count: agentUuids.length, successful, results };
|
||||
}
|
||||
|
||||
async function populateLookupApps(siteUuid = '') {
|
||||
try {
|
||||
const axios = require('axios');
|
||||
const headers = { 'Accept': 'application/json' };
|
||||
const res = await axios.get(`${BASE_URL}/lookup/applications`, {
|
||||
headers, params: { settings_limit: 10000 }, timeout: 20000
|
||||
});
|
||||
const apps = res.data?.data;
|
||||
if (!apps || !Array.isArray(apps)) {
|
||||
console.log('[LookupApp] No application data from DPI API');
|
||||
return { success: false, count: 0 };
|
||||
}
|
||||
let upserted = 0;
|
||||
for (const a of apps) {
|
||||
if (!a.id) continue;
|
||||
const doc = {
|
||||
id: a.id,
|
||||
tag: a.tag || '',
|
||||
label: a.label || '',
|
||||
name: a.full_label || a.label || '',
|
||||
full_name: a.full_label || '',
|
||||
description: a.description || '',
|
||||
favicon: a.favicon || '',
|
||||
icon: a.icon || '',
|
||||
logo: a.logo || '',
|
||||
application_category: a.category || null
|
||||
};
|
||||
await LookupApp.updateOne({ id: a.id }, { $set: doc }, { upsert: true });
|
||||
upserted++;
|
||||
}
|
||||
console.log(`[LookupApp] Upserted ${upserted} applications`);
|
||||
return { success: true, count: upserted };
|
||||
} catch (err) {
|
||||
console.error('[LookupApp] Population error:', err.message);
|
||||
return { success: false, count: 0 };
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
collectAllAgents,
|
||||
collectSpecificAgent,
|
||||
collectSpecificAgents,
|
||||
populateLookupApps
|
||||
};
|
||||
Reference in new issue
Block a user