feat(source2): push all latest files - device labeling, help system, proxy docs, database isolation fix

- Added DOKUMENTASI-FILTER-PER-SITE.md (site isolation docs)
- Fixed start-with-env.js to force-load .env.production
- Fixed MONGODB_URI hostname from mongodb-netify to mongodb.prod.proit.id
- Updated .gitignore to exclude sensitive scripts and credential files
- Minor UI and labeling improvements
This commit is contained in:
rafif committed 2026-07-29 14:14:29 +07:00
1 parent a403f752f3
commit dd4c8f6876
385 files changed
+31016 -8268

No files matched your search

+20 -20
View File
@@ -1,20 +1,20 @@
FROM oven/bun:1-alpine
WORKDIR /app
# Install dependencies first (layer caching)
# bun install is compatible with npm package.json / package-lock.json
COPY package*.json ./
RUN bun install --production
# Copy application source
COPY . .
# Expose proxy REST API port
EXPOSE 4000
# Health check
HEALTHCHECK --interval=30s --timeout=10s --start-period=15s --retries=3 \
CMD bun -e "require('http').get('http://localhost:4000/health', r => r.statusCode === 200 ? process.exit(0) : process.exit(1)).on('error', () => process.exit(1))"
CMD ["bun", "run", "index.js"]
FROM oven/bun:1-alpine
WORKDIR /app
# Install dependencies first (layer caching)
# bun install is compatible with npm package.json / package-lock.json
COPY package*.json ./
RUN bun install --production
# Copy application source
COPY . .
# Expose proxy REST API port
EXPOSE 4000
# Health check
HEALTHCHECK --interval=30s --timeout=10s --start-period=15s --retries=3 \
CMD bun -e "require('http').get('http://localhost:4000/health', r => r.statusCode === 200 ? process.exit(0) : process.exit(1)).on('error', () => process.exit(1))"
CMD ["bun", "run", "index.js"]
+108 -108
View File
@@ -1,108 +1,108 @@
# BackOne DPI Proxy — Deployment Reference
Standalone Docker image for collecting Netify DPI data and writing to MongoDB.
---
## Environment Variables
### Required
| Variable | Description |
|---|---|
| `NETIFY_SITE_UUIDS` | Comma-separated Netify site UUIDs (or single `NETIFY_SITE_UUID`) |
| `NETIFY_TOKEN` | Netify JWT token (or `NETIFY_JWT_TOKEN` / `NETIFY_API_KEY`) |
### MongoDB
| Variable | Default | Description |
|---|---|---|
| `MONGODB_URI` | `mongodb://127.0.0.1:27017/backone_dpi` | MongoDB connection string |
### Collection Mode
| Variable | Default | Description |
|---|---|---|
| `PROXY_COLLECT_MODE` | `all` | `all` = all agents, `agent` = single agent, `agents` = list of agents |
| `PROXY_AGENT_UUID` | _(none)_ | Single agent UUID (required if `mode=agent`) |
| `PROXY_AGENT_UUIDS` | _(none)_ | Comma-separated agent UUIDs (required if `mode=agents`) |
| `PROXY_AGENT_DELAY_MS` | `5000` | Delay (ms) between each agent collection to avoid rate-limiting |
### Scheduling
| Variable | Default | Description |
|---|---|---|
| `PROXY_CRON_SCHEDULE` | `*/5 * * * *` | Cron expression for collection interval |
| `PROXY_CAPACITY_LOG_INTERVAL_MS` | `86400000` | How often to log DB capacity usage (default: 24h) |
### Limits
| Variable | Default | Description |
|---|---|---|
| `PROXY_FLOW_LIMIT` | `1000000` | Max flows to fetch per agent per cycle |
| `PROXY_PORT` | `4000` | REST API listen port |
### Netify API
| Variable | Default | Description |
|---|---|---|
| `NETIFY_INFORMATICS_BASE_URL` | `https://informatics.netify.ai/api/v1` | Netify API base URL |
---
## Docker Run Example
```bash
docker run -d \
--name backone_proxy \
-p 4000:4000 \
-e MONGODB_URI=mongodb://host.docker.internal:27017/backone_dpi \
-e NETIFY_SITE_UUIDS=site-uuid-1,site-uuid-2 \
-e NETIFY_TOKEN=your-jwt-token \
-e PROXY_COLLECT_MODE=agents \
-e PROXY_AGENT_UUIDS=agent-uuid-1,agent-uuid-2,agent-uuid-3 \
backone-proxy
```
## Docker Compose Example
```yaml
services:
proxy:
build: ./proxy
container_name: backone_proxy
restart: always
ports:
- "4000:4000"
environment:
- MONGODB_URI=mongodb://mongodb:27017/backone_dpi
- PROXY_PORT=4000
- PROXY_COLLECT_MODE=all
- PROXY_COLLECT_MODE=${PROXY_COLLECT_MODE:-all}
- PROXY_AGENT_UUID=${PROXY_AGENT_UUID:-}
- PROXY_AGENT_UUIDS=${PROXY_AGENT_UUIDS:-}
- PROXY_AGENT_DELAY_MS=${PROXY_AGENT_DELAY_MS:-5000}
- PROXY_CRON_SCHEDULE=${PROXY_CRON_SCHEDULE:-*/5 * * * *}
- NETIFY_SITE_UUIDS=${NETIFY_SITE_UUIDS}
- NETIFY_TOKEN=${NETIFY_TOKEN}
- NETIFY_INFORMATICS_BASE_URL=${NETIFY_INFORMATICS_BASE_URL:-https://informatics.netify.ai/api/v1}
```
## REST API Endpoints
| Method | Endpoint | Description |
|---|---|---|
| `GET` | `/health` | Liveness check (MongoDB status) |
| `GET` | `/status` | Scheduler status, mode, last run |
| `GET` | `/agents` | List agent UUIDs in MongoDB |
| `POST` | `/collect/all` | Manual trigger — all agents |
| `POST` | `/collect/:uuid` | Manual trigger — single agent |
| `POST` | `/collect/agents` | Manual trigger — multiple agents `{"uuids":[...], "delay_ms":5000}` |
| `GET` | `/latest` | Latest data from all collections (debug) |
| `GET` | `/domain-details?domain=...` | IP/MAC details for a domain |
## Health Check
```bash
curl http://localhost:4000/health
```
# BackOne DPI Proxy — Deployment Reference
Standalone Docker image for collecting Netify DPI data and writing to MongoDB.
---
## Environment Variables
### Required
| Variable | Description |
|---|---|
| `NETIFY_SITE_UUIDS` | Comma-separated Netify site UUIDs (or single `NETIFY_SITE_UUID`) |
| `NETIFY_TOKEN` | Netify JWT token (or `NETIFY_JWT_TOKEN` / `NETIFY_API_KEY`) |
### MongoDB
| Variable | Default | Description |
|---|---|---|
| `MONGODB_URI` | `mongodb://127.0.0.1:27017/backone_dpi` | MongoDB connection string |
### Collection Mode
| Variable | Default | Description |
|---|---|---|
| `PROXY_COLLECT_MODE` | `all` | `all` = all agents, `agent` = single agent, `agents` = list of agents |
| `PROXY_AGENT_UUID` | _(none)_ | Single agent UUID (required if `mode=agent`) |
| `PROXY_AGENT_UUIDS` | _(none)_ | Comma-separated agent UUIDs (required if `mode=agents`) |
| `PROXY_AGENT_DELAY_MS` | `5000` | Delay (ms) between each agent collection to avoid rate-limiting |
### Scheduling
| Variable | Default | Description |
|---|---|---|
| `PROXY_CRON_SCHEDULE` | `*/5 * * * *` | Cron expression for collection interval |
| `PROXY_CAPACITY_LOG_INTERVAL_MS` | `86400000` | How often to log DB capacity usage (default: 24h) |
### Limits
| Variable | Default | Description |
|---|---|---|
| `PROXY_FLOW_LIMIT` | `1000000` | Max flows to fetch per agent per cycle |
| `PROXY_PORT` | `4000` | REST API listen port |
### Netify API
| Variable | Default | Description |
|---|---|---|
| `NETIFY_INFORMATICS_BASE_URL` | `https://informatics.netify.ai/api/v1` | Netify API base URL |
---
## Docker Run Example
```bash
docker run -d \
--name backone_proxy \
-p 4000:4000 \
-e MONGODB_URI=mongodb://host.docker.internal:27017/backone_dpi \
-e NETIFY_SITE_UUIDS=site-uuid-1,site-uuid-2 \
-e NETIFY_TOKEN=your-jwt-token \
-e PROXY_COLLECT_MODE=agents \
-e PROXY_AGENT_UUIDS=agent-uuid-1,agent-uuid-2,agent-uuid-3 \
backone-proxy
```
## Docker Compose Example
```yaml
services:
proxy:
build: ./proxy
container_name: backone_proxy
restart: always
ports:
- "4000:4000"
environment:
- MONGODB_URI=mongodb://mongodb:27017/backone_dpi
- PROXY_PORT=4000
- PROXY_COLLECT_MODE=all
- PROXY_COLLECT_MODE=${PROXY_COLLECT_MODE:-all}
- PROXY_AGENT_UUID=${PROXY_AGENT_UUID:-}
- PROXY_AGENT_UUIDS=${PROXY_AGENT_UUIDS:-}
- PROXY_AGENT_DELAY_MS=${PROXY_AGENT_DELAY_MS:-5000}
- PROXY_CRON_SCHEDULE=${PROXY_CRON_SCHEDULE:-*/5 * * * *}
- NETIFY_SITE_UUIDS=${NETIFY_SITE_UUIDS}
- NETIFY_TOKEN=${NETIFY_TOKEN}
- NETIFY_INFORMATICS_BASE_URL=${NETIFY_INFORMATICS_BASE_URL:-https://informatics.netify.ai/api/v1}
```
## REST API Endpoints
| Method | Endpoint | Description |
|---|---|---|
| `GET` | `/health` | Liveness check (MongoDB status) |
| `GET` | `/status` | Scheduler status, mode, last run |
| `GET` | `/agents` | List agent UUIDs in MongoDB |
| `POST` | `/collect/all` | Manual trigger — all agents |
| `POST` | `/collect/:uuid` | Manual trigger — single agent |
| `POST` | `/collect/agents` | Manual trigger — multiple agents `{"uuids":[...], "delay_ms":5000}` |
| `GET` | `/latest` | Latest data from all collections (debug) |
| `GET` | `/domain-details?domain=...` | IP/MAC details for a domain |
## Health Check
```bash
curl http://localhost:4000/health
```
+34
View File
@@ -0,0 +1,34 @@
// check_device.js - Detailed check of 10.6.10.44 records
const path = require('path');
require('dotenv').config({ path: path.join(__dirname, '..', '.env.local') });
const mongoose = require('mongoose');
async function run() {
await mongoose.connect(process.env.MONGODB_URI || 'mongodb://localhost:27017/backone');
const db = mongoose.connection.db;
// Get all records for 10.6.10.44
const docs = await db.collection('devicestats')
.find({ ip_address: '10.6.10.44' })
.sort({ timestamp: 1 })
.toArray();
console.log(`Total docs for 10.6.10.44: ${docs.length}`);
docs.forEach((d, i) => {
console.log(`\n--- Doc ${i + 1} ---`);
console.log(' _id: ', d._id);
console.log(' agent_uuid: ', d.agent_uuid);
console.log(' timestamp: ', d.timestamp);
console.log(' created_at: ', d.created_at);
console.log(' updated_at: ', d.updated_at);
console.log(' download: ', d.download);
console.log(' device_label:', d.device_label);
});
// Check if there are different agent_uuids
const agents = [...new Set(docs.map(d => d.agent_uuid))];
console.log('\nDistinct agent_uuids for this IP:', agents);
await mongoose.disconnect();
}
run().catch(err => { console.error(err.message); process.exit(1); });
+63
View File
@@ -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); });
+73
View File
@@ -0,0 +1,73 @@
// proxy/clean_devicestat_duplicates.js
// ─────────────────────────────────────────────────────────────────────────────
// One-time cleanup script to deduplicate historical DeviceStat records.
// Keeps only the LATEST document per (agent_uuid, ip_address) pair,
// removing all older duplicates accumulated before the upsert fix.
//
// Usage: node proxy/clean_devicestat_duplicates.js
// ─────────────────────────────────────────────────────────────────────────────
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://localhost:27017/backone';
const DeviceStatSchema = new mongoose.Schema({
timestamp: { type: Date },
agent_uuid: { type: String },
site_uuid: { type: String },
ip_address: { type: String },
mac_address: { type: String },
device_label: String,
device_type: String,
os_label: String,
manufacturer: String,
download: Number,
upload: Number,
flows: Number,
last_seen: String,
}, { timestamps: { createdAt: 'created_at', updatedAt: 'updated_at' } });
const DeviceStat = mongoose.model('DeviceStat', DeviceStatSchema);
async function run() {
console.log('[Cleanup] Connecting to MongoDB...');
await mongoose.connect(MONGODB_URI);
console.log('[Cleanup] Connected.');
// Find all unique (agent_uuid, ip_address) combinations
const groups = await DeviceStat.aggregate([
{ $group: {
_id: { agent_uuid: '$agent_uuid', ip_address: '$ip_address' },
ids: { $push: '$_id' },
timestamps: { $push: '$timestamp' },
count: { $sum: 1 },
}},
{ $match: { count: { $gt: 1 } } },
]);
console.log(`[Cleanup] Found ${groups.length} (agent_uuid, ip_address) pairs with duplicates.`);
let totalDeleted = 0;
for (const group of groups) {
// Sort the ids by matching timestamps - keep the latest
const paired = group.ids.map((id, i) => ({ id, ts: group.timestamps[i] }));
paired.sort((a, b) => new Date(b.ts) - new Date(a.ts));
// Keep the first (newest), delete the rest
const toDelete = paired.slice(1).map(p => p.id);
const result = await DeviceStat.deleteMany({ _id: { $in: toDelete } });
totalDeleted += result.deletedCount;
}
const remaining = await DeviceStat.countDocuments();
console.log(`[Cleanup] Done. Deleted ${totalDeleted} duplicate DeviceStat records.`);
console.log(`[Cleanup] Remaining DeviceStat documents: ${remaining}`);
await mongoose.disconnect();
}
run().catch(err => {
console.error('[Cleanup] Fatal error:', err.message);
process.exit(1);
});
+5 -5
View File
@@ -8,8 +8,8 @@ const mongoose = require('mongoose');
const path = require('path');
require('dotenv').config({ path: path.join(__dirname, '../../../..', '.env.local') });
const SIAB_UUID = '6681452d_9cae_4ff4_8ae8_0d504774265e';
const NEXUS_UUID = 'd7902405_0dc2_458b_8584_ed4d24b64f24';
const SIAB_UUID = '6681452d_9cae_4ff4_8ae8_0d504774265e';
const OFFICE_UUID = '1959bb55_045b_47c7_bbdd_f33b7db197b9';
// Definitive SIAB agent list (from most recent collector run)
const SIAB_AGENTS = ['F6-2V-DT-8A', 'YW-6I-61-LL', '2F-TF-1D-GK', '1R-79-J9-YE', '8A-V3-PB-85'];
@@ -27,13 +27,13 @@ async function cleanup() {
for (const colName of collections) {
const col = db.collection(colName);
// 1. Delete SIAB agents that are stored under NEXUS site_uuid
// 1. Delete SIAB agents that are stored under OFFICE site_uuid
const r1 = await col.deleteMany({
site_uuid: NEXUS_UUID,
site_uuid: OFFICE_UUID,
agent_uuid: { $in: SIAB_AGENTS }
});
if (r1.deletedCount > 0) {
console.log(`[${colName}] Removed ${r1.deletedCount} docs (SIAB agents from NEXUS)`);
console.log(`[${colName}] Removed ${r1.deletedCount} docs (SIAB agents from OFFICE)`);
totalDeleted += r1.deletedCount;
}
+263 -262
View File
@@ -1,262 +1,263 @@
// 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 { 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 };
}
module.exports = {
collectAllAgents,
collectSpecificAgent,
collectSpecificAgents
};
// 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 } = require('./collectorHelperDpi2');
const { collectThreats, collectEvents } = require('./collectorHelperDpi3');
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;
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(5, 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();
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);
}
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(5, 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);
// ── 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(', ')}`);
// ── 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 {
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();
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);
}
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 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) {
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 };
}
module.exports = {
collectAllAgents,
collectSpecificAgent,
collectSpecificAgents
};
+167 -167
View File
@@ -1,167 +1,167 @@
// proxy/collectorHelper.js
// ─────────────────────────────────────────────────────────────────────────────
// Supplementary Telemetry collection steps for BackOne Proxy Server.
// Split from collector.js to satisfy the 256-line file size limit.
// ─────────────────────────────────────────────────────────────────────────────
const {
AppCategoryStat,
TlsVersionStat,
TlsCipherStat,
TlsSecurityStat,
CountryStat,
ProtocolStat,
SslSubjectAltNameStat,
SniHostnameStat,
SslServerCnStat,
QuicHostnameStat,
} = require('./models/Schemas');
async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify, label) {
try {
// 2b. App Categories
const categories = await netify.fetchTopAppCategories(1440, 50, agentUuid, SITE_UUID);
if (categories && categories.length > 0) {
const catDocs = categories.map(c => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
category_label: c.category_label,
download: c.download || 0,
upload: c.upload || 0,
flows: c.flows || 0,
}));
await AppCategoryStat.insertMany(catDocs);
console.log(`[Collector] ✓ ${catDocs.length} categories saved for ${label}`);
}
// 2c. TLS Versions
const tlsVersions = await netify.fetchTlsVersions(1440, 50, agentUuid, SITE_UUID);
if (tlsVersions && tlsVersions.length > 0) {
const tvDocs = tlsVersions.map(v => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
tls_version: v.tls_version,
download: v.download || 0,
upload: v.upload || 0,
flows: v.flows || 0,
}));
await TlsVersionStat.insertMany(tvDocs);
console.log(`[Collector] ✓ ${tvDocs.length} TLS versions saved for ${label}`);
}
// 2d. TLS Ciphers
const tlsCiphers = await netify.fetchTlsCiphers(1440, 50, agentUuid, SITE_UUID);
if (tlsCiphers && tlsCiphers.length > 0) {
const tcDocs = tlsCiphers.map(c => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
tls_cipher: c.tls_cipher,
download: c.download || 0,
upload: c.upload || 0,
flows: c.flows || 0,
}));
await TlsCipherStat.insertMany(tcDocs);
console.log(`[Collector] ✓ ${tcDocs.length} TLS ciphers saved for ${label}`);
}
// 2e. TLS Security
const tlsSecurity = await netify.fetchTlsSecurity(1440, 50, agentUuid, SITE_UUID);
if (tlsSecurity && tlsSecurity.length > 0) {
const tsDocs = tlsSecurity.map(s => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
tls_security: s.tls_security,
download: s.download || 0,
upload: s.upload || 0,
flows: s.flows || 0,
}));
await TlsSecurityStat.insertMany(tsDocs);
console.log(`[Collector] ✓ ${tsDocs.length} TLS security stats saved for ${label}`);
}
// 2f. Top Countries
const countries = await netify.fetchTopCountries(1440, 100, agentUuid, SITE_UUID);
if (countries && countries.length > 0) {
const coDocs = countries.map(c => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
country_code: c.country_code,
country_name: c.country_name || '',
download: c.download || 0,
upload: c.upload || 0,
flows: c.flows || 0,
}));
await CountryStat.insertMany(coDocs);
console.log(`[Collector] ✓ ${coDocs.length} countries saved for ${label}`);
}
// 2g. Top Protocols
const protocols = await netify.fetchTopProtocols(1440, 50, agentUuid, SITE_UUID);
if (protocols && protocols.length > 0) {
const protoDocs = protocols.map(p => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
protocol_label: p.protocol_label,
download: p.download || 0,
upload: p.upload || 0,
flows: p.flows || 0,
}));
await ProtocolStat.insertMany(protoDocs);
console.log(`[Collector] ✓ ${protoDocs.length} protocols saved for ${label}`);
}
// 2h. SNI Hostnames
const snis = await netify.fetchSniHostnames(1440, 10000, agentUuid, SITE_UUID);
if (snis && snis.length > 0) {
const sniDocs = snis.map(s => ({
timestamp, agent_uuid: agentUuid, site_uuid: SITE_UUID,
sni_hostname: (s.sni_hostname && String(s.sni_hostname).trim() !== '') ? s.sni_hostname : 'Unknown',
download: s.download || 0, upload: s.upload || 0, flows: s.flows || 0,
}));
await SniHostnameStat.insertMany(sniDocs);
console.log(`[Collector] ✓ ${sniDocs.length} SNI hostnames saved for ${label}`);
}
/* -- COMMENTED OUT DUE TO NETIFY API HTTP 422 (UNSUPPORTED TIER) --
// 2i. SSL Server Common Names
const cns = await netify.fetchSslServerCn(1440, 50, agentUuid);
if (cns && cns.length > 0) {
const cnDocs = cns.map(c => ({
timestamp, agent_uuid: agentUuid, site_uuid: SITE_UUID,
ssl_server_cn: c.ssl_server_cn, download: c.download || 0, upload: c.upload || 0, flows: c.flows || 0,
}));
await SslServerCnStat.insertMany(cnDocs);
console.log(`[Collector] ✓ ${cnDocs.length} SSL Server CNs saved for ${label}`);
}
// 2j. QUIC Hostnames
const quics = await netify.fetchQuicHostnames(1440, 50, agentUuid);
if (quics && quics.length > 0) {
const quicDocs = quics.map(q => ({
timestamp, agent_uuid: agentUuid, site_uuid: SITE_UUID,
quic_hostname: q.quic_hostname, download: q.download || 0, upload: q.upload || 0, flows: q.flows || 0,
}));
await QuicHostnameStat.insertMany(quicDocs);
console.log(`[Collector] ✓ ${quicDocs.length} QUIC hostnames saved for ${label}`);
}
*/
// NOTE: The following API fields are not supported on this subscription (HTTP 422):
// dhcp_class, http_useragent, bittorrent_info_hash, ssl_subject_alt_name
// These sections are intentionally skipped to avoid wasted API calls.
// Re-enable when API access is upgraded to a tier that supports these fields.
} catch (err) {
console.error(`[CollectorHelper] Error saving secondary telemetry:`, err.message);
}
}
module.exports = {
collectSecondaryTelemetry,
};
// proxy/collectorHelper.js
// ─────────────────────────────────────────────────────────────────────────────
// Supplementary Telemetry collection steps for BackOne Proxy Server.
// Split from collector.js to satisfy the 256-line file size limit.
// ─────────────────────────────────────────────────────────────────────────────
const {
AppCategoryStat,
TlsVersionStat,
TlsCipherStat,
TlsSecurityStat,
CountryStat,
ProtocolStat,
SslSubjectAltNameStat,
SniHostnameStat,
SslServerCnStat,
QuicHostnameStat,
} = require('./models/Schemas');
async function collectSecondaryTelemetry(agentUuid, timestamp, SITE_UUID, netify, label) {
try {
// 2b. App Categories
const categories = await netify.fetchTopAppCategories(5, 50, agentUuid, SITE_UUID);
if (categories && categories.length > 0) {
const catDocs = categories.map(c => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
category_label: c.category_label,
download: c.download || 0,
upload: c.upload || 0,
flows: c.flows || 0,
}));
await AppCategoryStat.insertMany(catDocs);
console.log(`[Collector] ✓ ${catDocs.length} categories saved for ${label}`);
}
// 2c. TLS Versions
const tlsVersions = await netify.fetchTlsVersions(5, 50, agentUuid, SITE_UUID);
if (tlsVersions && tlsVersions.length > 0) {
const tvDocs = tlsVersions.map(v => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
tls_version: v.tls_version,
download: v.download || 0,
upload: v.upload || 0,
flows: v.flows || 0,
}));
await TlsVersionStat.insertMany(tvDocs);
console.log(`[Collector] ✓ ${tvDocs.length} TLS versions saved for ${label}`);
}
// 2d. TLS Ciphers
const tlsCiphers = await netify.fetchTlsCiphers(5, 50, agentUuid, SITE_UUID);
if (tlsCiphers && tlsCiphers.length > 0) {
const tcDocs = tlsCiphers.map(c => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
tls_cipher: c.tls_cipher,
download: c.download || 0,
upload: c.upload || 0,
flows: c.flows || 0,
}));
await TlsCipherStat.insertMany(tcDocs);
console.log(`[Collector] ✓ ${tcDocs.length} TLS ciphers saved for ${label}`);
}
// 2e. TLS Security
const tlsSecurity = await netify.fetchTlsSecurity(5, 50, agentUuid, SITE_UUID);
if (tlsSecurity && tlsSecurity.length > 0) {
const tsDocs = tlsSecurity.map(s => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
tls_security: s.tls_security,
download: s.download || 0,
upload: s.upload || 0,
flows: s.flows || 0,
}));
await TlsSecurityStat.insertMany(tsDocs);
console.log(`[Collector] ✓ ${tsDocs.length} TLS security stats saved for ${label}`);
}
// 2f. Top Countries
const countries = await netify.fetchTopCountries(5, 100, agentUuid, SITE_UUID);
if (countries && countries.length > 0) {
const coDocs = countries.map(c => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
country_code: c.country_code,
country_name: c.country_name || '',
download: c.download || 0,
upload: c.upload || 0,
flows: c.flows || 0,
}));
await CountryStat.insertMany(coDocs);
console.log(`[Collector] ✓ ${coDocs.length} countries saved for ${label}`);
}
// 2g. Top Protocols
const protocols = await netify.fetchTopProtocols(5, 50, agentUuid, SITE_UUID);
if (protocols && protocols.length > 0) {
const protoDocs = protocols.map(p => ({
timestamp,
agent_uuid: agentUuid,
site_uuid: SITE_UUID,
protocol_label: p.protocol_label,
download: p.download || 0,
upload: p.upload || 0,
flows: p.flows || 0,
}));
await ProtocolStat.insertMany(protoDocs);
console.log(`[Collector] ✓ ${protoDocs.length} protocols saved for ${label}`);
}
// 2h. SNI Hostnames
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,
sni_hostname: (s.sni_hostname && String(s.sni_hostname).trim() !== '') ? s.sni_hostname : 'Unknown',
download: s.download || 0, upload: s.upload || 0, flows: s.flows || 0,
}));
await SniHostnameStat.insertMany(sniDocs);
console.log(`[Collector] ✓ ${sniDocs.length} SNI hostnames saved for ${label}`);
}
/* -- COMMENTED OUT DUE TO NETIFY API HTTP 422 (UNSUPPORTED TIER) --
// 2i. SSL Server Common Names
const cns = await netify.fetchSslServerCn(1440, 50, agentUuid);
if (cns && cns.length > 0) {
const cnDocs = cns.map(c => ({
timestamp, agent_uuid: agentUuid, site_uuid: SITE_UUID,
ssl_server_cn: c.ssl_server_cn, download: c.download || 0, upload: c.upload || 0, flows: c.flows || 0,
}));
await SslServerCnStat.insertMany(cnDocs);
console.log(`[Collector] ✓ ${cnDocs.length} SSL Server CNs saved for ${label}`);
}
// 2j. QUIC Hostnames
const quics = await netify.fetchQuicHostnames(1440, 50, agentUuid);
if (quics && quics.length > 0) {
const quicDocs = quics.map(q => ({
timestamp, agent_uuid: agentUuid, site_uuid: SITE_UUID,
quic_hostname: q.quic_hostname, download: q.download || 0, upload: q.upload || 0, flows: q.flows || 0,
}));
await QuicHostnameStat.insertMany(quicDocs);
console.log(`[Collector] ✓ ${quicDocs.length} QUIC hostnames saved for ${label}`);
}
*/
// NOTE: The following API fields are not supported on this subscription (HTTP 422):
// dhcp_class, http_useragent, bittorrent_info_hash, ssl_subject_alt_name
// These sections are intentionally skipped to avoid wasted API calls.
// Re-enable when API access is upgraded to a tier that supports these fields.
} catch (err) {
console.error(`[CollectorHelper] Error saving secondary telemetry:`, err.message);
}
}
module.exports = {
collectSecondaryTelemetry,
};
+46 -61
View File
@@ -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) {
@@ -115,16 +115,34 @@ async function collectFlows(agentUuid, timestamp, SITE_UUID, netify, label, ipTo
const blacklistedCategories = new Set(blacklistRules.filter(r => r.type === 'category').map(r => r.value.toLowerCase()));
const blacklistedDomains = new Set(blacklistRules.filter(r => r.type === 'domain').map(r => r.value.toLowerCase()));
const flowIdsInBatch = flowDocs.map(f => f.flow_id).filter(Boolean);
const existingFlowThreats = new Set(
await Threat.find({ flow_id: { $in: flowIdsInBatch } }).distinct('flow_id')
);
const threatDocs = [];
const eventDocs = [];
for (const f of flowDocs) {
let isViolation = false;
let categoryLabel = "";
// Check if domain is blacklisted
if (f.domain && blacklistedDomains.has(f.domain.toLowerCase())) {
isViolation = true;
} else if (f.app_label && blacklistedDomains.has(f.app_label.toLowerCase())) {
isViolation = true;
for (const r of blacklistRules) {
if (r.type === 'domain') {
const val = r.value.toLowerCase();
// Direct domain match
if (f.domain && f.domain.toLowerCase().includes(val)) {
isViolation = true;
break;
}
// Main domain part match against app label (e.g. "google" from "google.com")
const mainDomainPart = val.split('.')[0];
if (mainDomainPart && f.app_label && f.app_label.toLowerCase().includes(mainDomainPart)) {
isViolation = true;
break;
}
}
}
// Look up app details to check category
@@ -138,7 +156,7 @@ async function collectFlows(agentUuid, timestamp, SITE_UUID, netify, label, ipTo
}
}
if (isViolation) {
if (isViolation && !existingFlowThreats.has(f.flow_id)) {
threatDocs.push({
timestamp,
agent_uuid: f.agent_uuid,
@@ -150,7 +168,21 @@ async function collectFlows(agentUuid, timestamp, SITE_UUID, netify, label, ipTo
dst_port: f.dst_port,
protocol: f.protocol,
description: `Access to blacklisted app/domain: ${f.app_label} (${f.domain || 'N/A'})${categoryLabel ? ' - Category: ' + categoryLabel : ''}`,
event_at: new Date().toISOString()
event_at: new Date().toISOString(),
flow_id: f.flow_id
});
eventDocs.push({
timestamp,
agent_uuid: f.agent_uuid,
site_uuid: f.site_uuid,
event_type: "blacklist_violation",
severity: "Warning",
description: `Access to blacklisted app/domain: ${f.app_label} (${f.domain || 'N/A'})${categoryLabel ? ' - Category: ' + categoryLabel : ''}`,
ip_address: f.src_ip,
mac_address: f.src_mac,
event_at: new Date(),
flow_id: f.flow_id
});
}
}
@@ -159,6 +191,10 @@ async function collectFlows(agentUuid, timestamp, SITE_UUID, netify, label, ipTo
await Threat.insertMany(threatDocs);
console.log(`[Collector] ✓ ${threatDocs.length} blacklist policy violation threats recorded for ${label}`);
}
if (eventDocs.length > 0) {
await Event.insertMany(eventDocs);
console.log(`[Collector] ✓ ${eventDocs.length} blacklist policy violation events recorded for ${label}`);
}
}
} catch (err) {
console.error('[Collector] Blacklist detection failed:', err.message);
@@ -169,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
};
+61
View File
@@ -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
};
+4 -4
View File
@@ -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}`);
}
+44
View File
@@ -0,0 +1,44 @@
// diagnostic.js - Run: node proxy/diagnostic.js
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://localhost:27017/backone';
async function run() {
await mongoose.connect(MONGODB_URI);
const db = mongoose.connection.db;
// 1. Total count
const total = await db.collection('devicestats').countDocuments();
console.log('=== DeviceStat Total:', total);
// 2. Count for 10.6.10.44
const specific = await db.collection('devicestats').countDocuments({ ip_address: '10.6.10.44' });
console.log('=== Count for 10.6.10.44:', specific);
// 3. Sample doc for 10.6.10.44
const sample = await db.collection('devicestats').findOne({ ip_address: '10.6.10.44' });
console.log('=== Sample doc for 10.6.10.44:', JSON.stringify(sample, null, 2));
// 4. Duplicate groups (top 10)
const dups = await db.collection('devicestats').aggregate([
{ $group: { _id: { agent_uuid: '$agent_uuid', ip_address: '$ip_address' }, count: { $sum: 1 } } },
{ $match: { count: { $gt: 1 } } },
{ $sort: { count: -1 } },
{ $limit: 10 }
]).toArray();
console.log('=== Top duplicate groups:', JSON.stringify(dups, null, 2));
// 5. Check what collection name is actually used
const collections = await db.listCollections().toArray();
console.log('=== Collections:', collections.map(c => c.name));
// 6. Check indexes on devicestats
const indexes = await db.collection('devicestats').indexes();
console.log('=== Indexes on devicestats:', JSON.stringify(indexes, null, 2));
await mongoose.disconnect();
}
run().catch(err => { console.error('ERROR:', err.message); process.exit(1); });
+79
View File
@@ -0,0 +1,79 @@
// fix_devicestat_index.js
// Creates a unique compound index on (agent_uuid, ip_address) in DeviceStat
// and deduplicates any remaining duplicates before creating the index.
// Run: node proxy/fix_devicestat_index.js
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://localhost:27017/backone';
async function run() {
console.log('[Fix] Connecting to MongoDB...');
await mongoose.connect(MONGODB_URI);
const db = mongoose.connection.db;
const col = db.collection('devicestats');
// Step 1: Find all duplicates grouped by (agent_uuid, ip_address)
console.log('[Fix] Scanning for duplicates...');
const groups = await col.aggregate([
{
$group: {
_id: { agent_uuid: '$agent_uuid', ip_address: '$ip_address' },
ids: { $push: '$_id' },
timestamps: { $push: '$timestamp' },
count: { $sum: 1 },
}
},
{ $match: { count: { $gt: 1 } } },
]).toArray();
console.log(`[Fix] Found ${groups.length} duplicate groups.`);
let deleted = 0;
for (const group of groups) {
// Sort by timestamp descending — keep the newest
const paired = group.ids.map((id, i) => ({ id, ts: new Date(group.timestamps[i] || 0) }));
paired.sort((a, b) => b.ts - a.ts);
const toDelete = paired.slice(1).map(p => p.id);
const result = await col.deleteMany({ _id: { $in: toDelete } });
deleted += result.deletedCount;
}
console.log(`[Fix] Deleted ${deleted} duplicate documents.`);
const remaining = await col.countDocuments();
console.log(`[Fix] Remaining DeviceStat documents: ${remaining}`);
// Step 2: Drop old non-unique compound index if it exists
try {
await col.dropIndex('agent_uuid_1_timestamp_-1_ip_address_1');
console.log('[Fix] Dropped old compound index.');
} catch (e) {
console.log('[Fix] Old index not found or already dropped:', e.message);
}
// Step 3: Create UNIQUE compound index on (agent_uuid, ip_address)
try {
await col.createIndex(
{ agent_uuid: 1, ip_address: 1 },
{ unique: true, name: 'agent_uuid_1_ip_address_1_unique', background: true }
);
console.log('[Fix] Created unique index on (agent_uuid, ip_address).');
} catch (e) {
console.error('[Fix] Failed to create unique index:', e.message);
}
// Step 4: Verify indexes
const indexes = await col.indexes();
console.log('[Fix] Current indexes:');
indexes.forEach(idx => console.log(` - ${idx.name}: ${JSON.stringify(idx.key)} ${idx.unique ? '[UNIQUE]' : ''}`));
// Step 5: Verify 10.6.10.44
const cnt = await col.countDocuments({ ip_address: '10.6.10.44' });
console.log(`\n[Fix] Count for 10.6.10.44: ${cnt} (should be 1)`);
await mongoose.disconnect();
console.log('[Fix] Done.');
}
run().catch(err => { console.error('[Fix] Fatal:', err.message); process.exit(1); });
+13 -5
View File
@@ -1,10 +1,16 @@
// proxy/index.js
// ─────────────────────────────────────────────────────────────────────────────
// Polyfill global crypto for Node 18 compatibility (required by mongodb driver)
if (typeof globalThis.crypto === 'undefined') {
globalThis.crypto = require('crypto');
}
// BackOne Proxy Server - Entry Point
// ─────────────────────────────────────────────────────────────────────────────
const path = require('path');
require('dotenv').config({ path: path.join(__dirname, '..', '.env.local') });
const envFile = process.env.NODE_ENV === 'production' ? '.env.production' : '.env.local';
require('dotenv').config({ path: path.join(__dirname, '..', envFile) });
const express = require('express');
const cors = require('cors');
@@ -57,11 +63,13 @@ async function main() {
console.log('╚════════════════════════════════════════════════╝\n');
console.log(`[Proxy] Mode: ${process.env.PROXY_COLLECT_MODE || 'all'}`);
// Start REST API server first so liveness probes remain active
app.listen(PORT, '0.0.0.0', () => {
console.log(`\n🚀 Proxy REST API running at http://0.0.0.0:${PORT}`);
// Bind to 127.0.0.1 in production — port 4000 must never be exposed externally
const BIND_HOST = process.env.NODE_ENV === 'production' ? '127.0.0.1' : '0.0.0.0';
app.listen(PORT, BIND_HOST, () => {
console.log(`\n🚀 Proxy REST API running at http://${BIND_HOST}:${PORT}`);
console.log(` GET /health → liveness check`);
console.log(` GET /status → scheduler + DB status\n`);
console.log(` GET /status → scheduler + DB status`);
console.log(`🔒 Security : Bound to ${BIND_HOST} (internal only in production)\n`);
});
const connected = await connectDB();
+199 -182
View File
@@ -1,182 +1,199 @@
// proxy/models/Schemas.js
// MongoDB schemas shared between the proxy server (write) and backend (read).
// Each document is tagged with agent_uuid + site_uuid for tenant isolation.
//
// IMPORTANT: Indexes are set for common query patterns:
// - timestamp (for time-range queries)
// - agent_uuid (for per-tenant filtering)
// - site_uuid (for site-level aggregation)
const mongoose = require('mongoose');
const baseOptions = {
timestamps: { createdAt: 'created_at', updatedAt: 'updated_at' }
};
// ─── Bandwidth Summary (per agent, per collection cycle) ───────────────────────
const SummarySchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true }, // null = global/all agents
site_uuid: { type: String, index: true },
bandwidth_down: Number,
bandwidth_up: Number,
active_flows: Number,
download_speed: Number,
upload_speed: Number,
total_devices: Number,
total_threats: Number,
packet_drops: Number,
peak_flow_rate: Number,
cpu_usage: Number,
memory_usage: Number,
queue_depth: Number,
}, baseOptions);
// ─── Top Applications (per agent) ─────────────────────────────────────────────
const AppStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
app_label: { type: String, required: true },
download: Number,
upload: Number,
flows: Number,
}, baseOptions);
// ─── Protocol Statistics (per agent) ──────────────────────────────────────────
const ProtocolStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
protocol_label: { type: String, required: true },
download: Number,
upload: Number,
flows: Number,
}, baseOptions);
// ─── Discovered Devices (per agent, includes IP + MAC + device info) ───────────
const DeviceStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
ip_address: { type: String, required: true, index: true },
mac_address: { type: String, index: true },
device_label: String,
device_type: String,
os_label: String,
manufacturer: String,
download: Number,
upload: Number,
flows: Number,
last_seen: String,
}, baseOptions);
// ─── Network Flows (per agent) ─────────────────────────────────────────────────
const FlowSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
flow_id: String,
src_ip: { type: String, index: true },
src_mac: { type: String, index: true }, // indexed for MAC-to-IP resolution in events
dst_ip: { type: String, index: true },
dst_port: Number,
protocol: String,
app_label: String,
domain: String,
download: Number,
upload: Number,
first_seen: String,
last_seen: String,
}, baseOptions);
// ─── Cyber Threats (per agent) ─────────────────────────────────────────────────
const ThreatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
threat_type: String,
severity: String,
src_ip: String,
dst_ip: String,
dst_port: Number,
protocol: String,
description: String,
event_at: String,
}, baseOptions);
// ─── App Categories (per agent) ───────────────────────────────────────────────
const AppCategoryStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
category_label: { type: String, required: true },
download: Number,
upload: Number,
flows: Number,
}, baseOptions);
// ─── System Events (per agent) ─────────────────────────────────────────────────
const EventSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
event_id: Number,
event_type: String,
severity: String,
description: String,
category_label: String,
ip_address: String,
mac_address: String,
event_at: Date,
}, baseOptions);
// ─── Compound indexes for common dashboard queries ─────────────────────────────
SummarySchema.index({ agent_uuid: 1, timestamp: -1 });
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, app_label: 1, timestamp: -1 });
ThreatSchema.index({ agent_uuid: 1, timestamp: -1 });
AppCategoryStatSchema.index({ agent_uuid: 1, timestamp: -1 });
EventSchema.index({ agent_uuid: 1, timestamp: -1 });
FlowSchema.index({ site_uuid: 1, src_mac: 1, timestamp: -1 }); // for MAC-to-IP resolution
// ── Per-Device Per-Application Stats ────────────────────────────────────────
// Collected from DPI API: /data/stats/top/application/download with filter_local_ips
// Allows showing "YouTube 134GB" in Device Detail modal per specific IP
const DeviceAppStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
ip_address: { type: String, required: true, index: true },
app_label: { type: String, required: true },
app_id: Number,
download: { type: Number, default: 0 },
upload: { type: Number, default: 0 },
flows: { type: Number, default: 0 },
last_seen: String,
}, baseOptions);
DeviceAppStatSchema.index({ agent_uuid: 1, ip_address: 1, timestamp: -1 });
DeviceAppStatSchema.index({ ip_address: 1, app_label: 1, timestamp: -1 });
DeviceAppStatSchema.index({ site_uuid: 1, app_label: 1, timestamp: -1 });
const telemetrySchemas = require('./SchemasTelemetry');
const auxSchemas = require('./SchemasAux');
module.exports = {
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),
AppCategoryStat: mongoose.model('AppCategoryStat', AppCategoryStatSchema),
Event: mongoose.model('Event', EventSchema),
...auxSchemas,
...telemetrySchemas,
};
// proxy/models/Schemas.js
// MongoDB schemas shared between the proxy server (write) and backend (read).
// Each document is tagged with agent_uuid + site_uuid for tenant isolation.
//
// IMPORTANT: Indexes are set for common query patterns:
// - timestamp (for time-range queries)
// - agent_uuid (for per-tenant filtering)
// - site_uuid (for site-level aggregation)
const mongoose = require('mongoose');
const baseOptions = {
timestamps: { createdAt: 'created_at', updatedAt: 'updated_at' }
};
// ─── Bandwidth Summary (per agent, per collection cycle) ───────────────────────
const SummarySchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true }, // null = global/all agents
site_uuid: { type: String, index: true },
bandwidth_down: Number,
bandwidth_up: Number,
active_flows: Number,
download_speed: Number,
upload_speed: Number,
total_devices: Number,
total_threats: Number,
packet_drops: Number,
peak_flow_rate: Number,
cpu_usage: Number,
memory_usage: Number,
queue_depth: Number,
}, baseOptions);
// ─── Top Applications (per agent) ─────────────────────────────────────────────
const AppStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
app_label: { type: String, required: true },
download: Number,
upload: Number,
flows: Number,
}, baseOptions);
// ─── Protocol Statistics (per agent) ──────────────────────────────────────────
const ProtocolStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
protocol_label: { type: String, required: true },
download: Number,
upload: Number,
flows: Number,
}, baseOptions);
// ─── Discovered Devices (per agent, includes IP + MAC + device info) ───────────
const DeviceStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
ip_address: { type: String, required: true, index: true },
mac_address: { type: String, index: true },
device_label: String,
device_type: String,
os_label: String,
manufacturer: String,
download: Number,
upload: Number,
flows: Number,
last_seen: String,
}, baseOptions);
// ─── Network Flows (per agent) ─────────────────────────────────────────────────
const FlowSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
flow_id: String,
src_ip: { type: String, index: true },
src_mac: { type: String, index: true }, // indexed for MAC-to-IP resolution in events
dst_ip: { type: String, index: true },
dst_port: Number,
protocol: String,
app_label: String,
domain: String,
download: Number,
upload: Number,
first_seen: String,
last_seen: String,
}, baseOptions);
// ─── Cyber Threats (per agent) ─────────────────────────────────────────────────
const ThreatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
threat_type: String,
severity: String,
src_ip: String,
dst_ip: String,
dst_port: Number,
protocol: String,
description: String,
event_at: String,
flow_id: { type: String, index: true },
}, baseOptions);
// ─── App Categories (per agent) ───────────────────────────────────────────────
const AppCategoryStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
category_label: { type: String, required: true },
download: Number,
upload: Number,
flows: Number,
}, baseOptions);
// ─── System Events (per agent) ─────────────────────────────────────────────────
const EventSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
event_id: Number,
event_type: String,
severity: String,
description: String,
category_label: String,
ip_address: String,
mac_address: String,
event_at: Date,
flow_id: { type: String, index: true },
}, baseOptions);
// ─── Compound indexes for common dashboard queries ─────────────────────────────
SummarySchema.index({ agent_uuid: 1, timestamp: -1 });
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 });
EventSchema.index({ agent_uuid: 1, timestamp: -1 });
FlowSchema.index({ site_uuid: 1, src_mac: 1, timestamp: -1 }); // for MAC-to-IP resolution
// ── Per-Device Per-Application Stats ────────────────────────────────────────
// Collected from DPI API: /data/stats/top/application/download with filter_local_ips
// Allows showing "YouTube 134GB" in Device Detail modal per specific IP
const DeviceAppStatSchema = new mongoose.Schema({
timestamp: { type: Date, required: true, index: true, expires: '7d' },
agent_uuid: { type: String, index: true },
site_uuid: { type: String, index: true },
ip_address: { type: String, required: true, index: true },
app_label: { type: String, required: true },
app_id: Number,
download: { type: Number, default: 0 },
upload: { type: Number, default: 0 },
flows: { type: Number, default: 0 },
last_seen: String,
}, baseOptions);
DeviceAppStatSchema.index({ agent_uuid: 1, ip_address: 1, timestamp: -1 });
DeviceAppStatSchema.index({ ip_address: 1, app_label: 1, timestamp: -1 });
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),
DeviceAppStat: mongoose.model('DeviceAppStat', DeviceAppStatSchema),
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,
};
+89
View File
@@ -0,0 +1,89 @@
// proxy/netifyAgentFetcher.js
// ─────────────────────────────────────────────────────────────────────────────
// Fetches active agents from Netify Informatics API.
// Uses a two-strategy approach to handle endpoints that may timeout (Source 2).
// ─────────────────────────────────────────────────────────────────────────────
const { netifyFetch, agentMap } = require('./netifyClientCore');
// Strategy 1 timeout: 15s (agent/download can be slow on small deployments)
const PRIMARY_TIMEOUT_MS = 15000;
async function fetchAgents(siteUuid = null) {
// Strategy 1: Use /data/stats/top/agent/download (standard Netify endpoint)
// filter_interval reduced to 31 days to lessen query load vs. old 365-day value.
const primaryPromise = (async () => {
try {
const data = await netifyFetch('/data/stats/top/agent/download', {
filter_interval: 44640, // 31 days
settings_limit: 1000000,
}, null, siteUuid);
if (data && Array.isArray(data) && data.length > 0) return data;
return null;
} catch (e) {
return null;
}
})();
const timeoutPromise = new Promise(resolve =>
setTimeout(() => resolve(null), PRIMARY_TIMEOUT_MS)
);
const primaryData = await Promise.race([primaryPromise, timeoutPromise]);
if (primaryData && Array.isArray(primaryData) && primaryData.length > 0) {
const list = primaryData.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 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;
}
// Strategy 2: Fallback — discover agents from /data/flows
// Useful for Source 2 where /data/stats/top/agent/download consistently times out.
console.log('[fetchAgents] Primary endpoint timeout/empty. Using flows-based agent discovery...');
try {
const flowData = await netifyFetch('/data/flows', {
settings_limit: 100,
}, null, siteUuid);
if (!flowData || !Array.isArray(flowData)) return [];
// Extract unique agent_uuids from flow records
const seen = new Set();
const agentList = [];
for (const flow of flowData) {
const uuid = flow.agent_uuid;
if (uuid && !seen.has(uuid)) {
seen.add(uuid);
agentList.push({
id: null,
uuid: uuid,
serial: uuid,
label: uuid,
provisioned: true,
activated: true,
last_seen_at: flow.last_seen_at ?? null,
});
}
}
console.log(`[fetchAgents] Fallback discovered ${agentList.length} agent(s) from flows.`);
return agentList;
} catch (e) {
console.error('[fetchAgents] Fallback also failed:', e.message);
return [];
}
}
module.exports = { fetchAgents };
+16 -227
View File
@@ -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,
};
+219
View File
@@ -0,0 +1,219 @@
// proxy/netifyClientStats.js
// ─────────────────────────────────────────────────────────────────────────────
// Supplementary fetchers split from netifyClient.js to satisfy the 256-line limit.
// ─────────────────────────────────────────────────────────────────────────────
const { netifyFetch } = require('./netifyClientCore');
const { fetchAgents } = require('./netifyAgentFetcher');
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 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
};
+234 -234
View File
@@ -1,234 +1,234 @@
// proxy/netifyTelemetry.js
// ─────────────────────────────────────────────────────────────────────────────
// Supplementary Telemetry endpoints wrapper for BackOne Proxy Server
// Split from netifyClient.js to strictly respect the 256-line file size limit.
// ─────────────────────────────────────────────────────────────────────────────
const { netifyFetch } = require('./netifyClientCore');
async function fetchTopAppCategories(interval = 1440, limit = 15, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/application_category/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/application_category/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const key = r.application_category?.label ?? r.application_category;
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const label = r.application_category?.label ?? String(r.application_category ?? 'Unknown');
return {
category_label : label,
download : r.download ?? 0,
upload : ulMap[label] ?? 0,
flows : r.flows ?? 0,
};
});
}
async function fetchTlsVersions(interval = 1440, limit = 15, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/tls_version/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/tls_version/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const key = r.tls_version?.label ?? r.tls_version?.code ?? r.tls_version;
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const label = r.tls_version?.label ?? r.tls_version?.code ?? String(r.tls_version ?? 'Unknown');
return {
tls_version : label,
download : r.download ?? 0,
upload : ulMap[label] ?? 0,
flows : r.flows ?? 0,
};
});
}
async function fetchTlsCiphers(interval = 1440, limit = 15, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/tls_cipher/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/tls_cipher/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const key = r.tls_cipher?.label ?? r.tls_cipher?.code ?? r.tls_cipher;
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const label = r.tls_cipher?.label ?? r.tls_cipher?.code ?? String(r.tls_cipher ?? 'Unknown');
return {
tls_cipher : label,
download : r.download ?? 0,
upload : ulMap[label] ?? 0,
flows : r.flows ?? 0,
};
});
}
async function fetchTlsSecurity(interval = 1440, limit = 15, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/tls_security/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/tls_security/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const key = r.tls_security?.label ?? r.tls_security?.code ?? r.tls_security;
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const label = r.tls_security?.label ?? r.tls_security?.code ?? String(r.tls_security ?? 'Unknown');
return {
tls_security : label,
download : r.download ?? 0,
upload : ulMap[label] ?? 0,
flows : r.flows ?? 0,
};
});
}
async function fetchTopCountries(interval = 1440, limit = 100, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/country/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/country/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const cc = r.country?.code;
if (cc) ulMap[cc] = r.upload ?? 0;
}
}
return dlData.map(r => {
return {
country_code: r.country?.code ?? 'Unknown',
country_name: r.country?.label ?? '',
download: r.download ?? 0,
upload: ulMap[r.country?.code] ?? 0,
flows: r.flows ?? 0,
};
}).filter(r => r.country_code);
}
async function fetchTopProperty(fieldName, interval, limit, agentUuid, siteUuid) {
const [dlData, ulData] = await Promise.all([
netifyFetch(`/data/stats/top/${fieldName}/download`, { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch(`/data/stats/top/${fieldName}/upload`, { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const item = r[fieldName];
const key = item?.hash ?? item?.name ?? item?.label ?? String(item ?? '');
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const item = r[fieldName];
const key = item?.hash ?? item?.name ?? item?.label ?? String(item || 'Unknown');
const label = item?.label ?? key;
return {
key,
label,
download: r.download ?? 0,
upload: ulMap[key] ?? ulMap[label] ?? 0,
flows: r.flows ?? 0,
};
});
}
function mapProp(data, keyName) {
return data.map(d => ({
[keyName]: d.key,
download: d.download,
upload: d.upload,
flows: d.flows,
}));
}
async function fetchDhcpFingerprints(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('dhcp_class', interval, limit, agentUuid, siteUuid), 'fingerprint');
}
async function fetchHttpUserAgents(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('http_useragent', interval, limit, agentUuid, siteUuid), 'user_agent');
}
async function fetchBittorrentHashes(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
const data = await fetchTopProperty('bittorrent_info_hash', interval, limit, agentUuid, siteUuid);
return data.map(d => ({
info_hash: d.key,
label: d.label,
download: d.download,
upload: d.upload,
flows: d.flows,
}));
}
async function fetchSniHostnames(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('tls_sni', interval, limit, agentUuid, siteUuid), 'sni_hostname');
}
async function fetchSslServerCn(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ssl_server_cn', interval, limit, agentUuid, siteUuid), 'ssl_server_cn');
}
async function fetchQuicHostnames(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('quic_hostname', interval, limit, agentUuid, siteUuid), 'quic_hostname');
}
async function fetchSshClients(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ssh_client', interval, limit, agentUuid, siteUuid), 'ssh_client');
}
async function fetchSshServers(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ssh_server', interval, limit, agentUuid, siteUuid), 'ssh_server');
}
async function fetchMdnsHostnames(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('mdns_hostname', interval, limit, agentUuid, siteUuid), 'mdns_hostname');
}
async function fetchTopProtocols(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ip_protocol', interval, limit, agentUuid, siteUuid), 'protocol_label');
}
async function fetchSslSubjectAltNames(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ssl_subject_alt_name', interval, limit, agentUuid, siteUuid), 'alt_name');
}
module.exports = {
fetchTopAppCategories,
fetchTlsVersions,
fetchTlsCiphers,
fetchTlsSecurity,
fetchTopCountries,
fetchDhcpFingerprints,
fetchHttpUserAgents,
fetchBittorrentHashes,
fetchSniHostnames,
fetchSslServerCn,
fetchQuicHostnames,
fetchSshClients,
fetchSshServers,
fetchMdnsHostnames,
fetchTopProtocols,
fetchSslSubjectAltNames,
};
// proxy/netifyTelemetry.js
// ─────────────────────────────────────────────────────────────────────────────
// Supplementary Telemetry endpoints wrapper for BackOne Proxy Server
// Split from netifyClient.js to strictly respect the 256-line file size limit.
// ─────────────────────────────────────────────────────────────────────────────
const { netifyFetch } = require('./netifyClientCore');
async function fetchTopAppCategories(interval = 1440, limit = 15, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/application_category/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/application_category/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const key = r.application_category?.label ?? r.application_category;
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const label = r.application_category?.label ?? String(r.application_category ?? 'Unknown');
return {
category_label : label,
download : r.download ?? 0,
upload : ulMap[label] ?? 0,
flows : r.flows ?? 0,
};
});
}
async function fetchTlsVersions(interval = 1440, limit = 15, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/tls_version/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/tls_version/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const key = r.tls_version?.label ?? r.tls_version?.code ?? r.tls_version;
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const label = r.tls_version?.label ?? r.tls_version?.code ?? String(r.tls_version ?? 'Unknown');
return {
tls_version : label,
download : r.download ?? 0,
upload : ulMap[label] ?? 0,
flows : r.flows ?? 0,
};
});
}
async function fetchTlsCiphers(interval = 1440, limit = 15, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/tls_cipher/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/tls_cipher/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const key = r.tls_cipher?.label ?? r.tls_cipher?.code ?? r.tls_cipher;
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const label = r.tls_cipher?.label ?? r.tls_cipher?.code ?? String(r.tls_cipher ?? 'Unknown');
return {
tls_cipher : label,
download : r.download ?? 0,
upload : ulMap[label] ?? 0,
flows : r.flows ?? 0,
};
});
}
async function fetchTlsSecurity(interval = 1440, limit = 15, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/tls_security/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/tls_security/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const key = r.tls_security?.label ?? r.tls_security?.code ?? r.tls_security;
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const label = r.tls_security?.label ?? r.tls_security?.code ?? String(r.tls_security ?? 'Unknown');
return {
tls_security : label,
download : r.download ?? 0,
upload : ulMap[label] ?? 0,
flows : r.flows ?? 0,
};
});
}
async function fetchTopCountries(interval = 1440, limit = 100, agentUuid = null, siteUuid = null) {
const [dlData, ulData] = await Promise.all([
netifyFetch('/data/stats/top/country/download', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch('/data/stats/top/country/upload', { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const cc = r.country?.code;
if (cc) ulMap[cc] = r.upload ?? 0;
}
}
return dlData.map(r => {
return {
country_code: r.country?.code ?? 'Unknown',
country_name: r.country?.label ?? '',
download: r.download ?? 0,
upload: ulMap[r.country?.code] ?? 0,
flows: r.flows ?? 0,
};
}).filter(r => r.country_code);
}
async function fetchTopProperty(fieldName, interval, limit, agentUuid, siteUuid) {
const [dlData, ulData] = await Promise.all([
netifyFetch(`/data/stats/top/${fieldName}/download`, { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
netifyFetch(`/data/stats/top/${fieldName}/upload`, { filter_interval: interval, settings_limit: limit }, agentUuid, siteUuid),
]);
if (!dlData) return [];
const ulMap = {};
if (ulData) {
for (const r of ulData) {
const item = r[fieldName];
const key = item?.hash ?? item?.name ?? item?.label ?? String(item ?? '');
if (key) ulMap[key] = r.upload ?? 0;
}
}
return dlData.map(r => {
const item = r[fieldName];
const key = item?.hash ?? item?.name ?? item?.label ?? String(item || 'Unknown');
const label = item?.label ?? key;
return {
key,
label,
download: r.download ?? 0,
upload: ulMap[key] ?? ulMap[label] ?? 0,
flows: r.flows ?? 0,
};
});
}
function mapProp(data, keyName) {
return data.map(d => ({
[keyName]: d.key,
download: d.download,
upload: d.upload,
flows: d.flows,
}));
}
async function fetchDhcpFingerprints(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('dhcp_class', interval, limit, agentUuid, siteUuid), 'fingerprint');
}
async function fetchHttpUserAgents(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('http_useragent', interval, limit, agentUuid, siteUuid), 'user_agent');
}
async function fetchBittorrentHashes(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
const data = await fetchTopProperty('bittorrent_info_hash', interval, limit, agentUuid, siteUuid);
return data.map(d => ({
info_hash: d.key,
label: d.label,
download: d.download,
upload: d.upload,
flows: d.flows,
}));
}
async function fetchSniHostnames(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('tls_sni', interval, limit, agentUuid, siteUuid), 'sni_hostname');
}
async function fetchSslServerCn(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ssl_server_cn', interval, limit, agentUuid, siteUuid), 'ssl_server_cn');
}
async function fetchQuicHostnames(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('quic_hostname', interval, limit, agentUuid, siteUuid), 'quic_hostname');
}
async function fetchSshClients(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ssh_client', interval, limit, agentUuid, siteUuid), 'ssh_client');
}
async function fetchSshServers(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ssh_server', interval, limit, agentUuid, siteUuid), 'ssh_server');
}
async function fetchMdnsHostnames(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('mdns_hostname', interval, limit, agentUuid, siteUuid), 'mdns_hostname');
}
async function fetchTopProtocols(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ip_protocol', interval, limit, agentUuid, siteUuid), 'protocol_label');
}
async function fetchSslSubjectAltNames(interval = 1440, limit = 50, agentUuid = null, siteUuid = null) {
return mapProp(await fetchTopProperty('ssl_subject_alt_name', interval, limit, agentUuid, siteUuid), 'alt_name');
}
module.exports = {
fetchTopAppCategories,
fetchTlsVersions,
fetchTlsCiphers,
fetchTlsSecurity,
fetchTopCountries,
fetchDhcpFingerprints,
fetchHttpUserAgents,
fetchBittorrentHashes,
fetchSniHostnames,
fetchSslServerCn,
fetchQuicHostnames,
fetchSshClients,
fetchSshServers,
fetchMdnsHostnames,
fetchTopProtocols,
fetchSslSubjectAltNames,
};
+1312 -1312
View File
File diff suppressed because it is too large. Load diff
+19 -19
View File
@@ -1,19 +1,19 @@
{
"name": "backone-proxy",
"version": "1.0.0",
"description": "BackOne DPI Proxy Server - Fetches from DPI API, filters per agent_uuid, stores to MongoDB",
"main": "index.js",
"scripts": {
"start": "node index.js",
"start:bun": "bun run index.js",
"dev": "nodemon index.js"
},
"dependencies": {
"axios": "^1.6.2",
"cors": "^2.8.5",
"dotenv": "^16.3.1",
"express": "^4.18.2",
"mongoose": "^8.0.3",
"node-cron": "^3.0.3"
}
}
{
"name": "backone-proxy",
"version": "1.0.0",
"description": "BackOne DPI Proxy Server - Fetches from DPI API, filters per agent_uuid, stores to MongoDB",
"main": "index.js",
"scripts": {
"start": "node index.js",
"start:bun": "bun run index.js",
"dev": "nodemon index.js"
},
"dependencies": {
"axios": "^1.6.2",
"cors": "^2.8.5",
"dotenv": "^16.3.1",
"express": "^4.18.2",
"mongoose": "^8.0.3",
"node-cron": "^3.0.3"
}
}
+129 -124
View File
@@ -1,124 +1,129 @@
// proxy/scheduler.js
// Cron scheduler for automatic data collection from DPI API
// Runs every 5 minutes, collecting data for all agents or a specific agent.
const cron = require('node-cron');
const { collectAllAgents, collectSpecificAgent, collectSpecificAgents } = require('./collector');
// Mode: 'all' = collect all agents, 'agent' = collect one specific agent, 'agents' = collect multiple agents
const COLLECT_MODE = process.env.PROXY_COLLECT_MODE || 'all';
const SPECIFIC_AGENT = process.env.PROXY_AGENT_UUID || null;
const SPECIFIC_AGENTS = (process.env.PROXY_AGENT_UUIDS || '').split(',').map(s => s.trim()).filter(Boolean);
const AGENT_DELAY_MS = parseInt(process.env.PROXY_AGENT_DELAY_MS || '5000');
const CRON_SCHEDULE = process.env.PROXY_CRON_SCHEDULE || '*/5 * * * *';
// Capacity logging is an expensive full-scan aggregation. Run it at most once per
// interval (default 24h) instead of every collection cycle to reduce CPU/DB load.
const CAPACITY_LOG_INTERVAL_MS = parseInt(process.env.PROXY_CAPACITY_LOG_INTERVAL_MS || String(24 * 60 * 60 * 1000));
let isRunning = false;
let lastRunAt = null;
let lastRunResult = null;
let runCount = 0;
let lastCapacityLogAt = 0;
/**
* Execute one collection cycle (called by cron and manual trigger).
* Prevents concurrent runs with isRunning guard.
*/
async function runCollection() {
if (isRunning) {
console.log('[Scheduler] Skipping - previous run still in progress');
return { skipped: true, reason: 'already_running' };
}
isRunning = true;
lastRunAt = new Date();
runCount++;
try {
let result;
if (COLLECT_MODE === 'agent' && SPECIFIC_AGENT) {
console.log(`[Scheduler] Run #${runCount} - Mode: SPECIFIC AGENT (${SPECIFIC_AGENT})`);
result = await collectSpecificAgent(SPECIFIC_AGENT);
} else if (COLLECT_MODE === 'agents' && SPECIFIC_AGENTS.length > 0) {
console.log(`[Scheduler] Run #${runCount} - Mode: SPECIFIC AGENTS (${SPECIFIC_AGENTS.join(', ')})`);
result = await collectSpecificAgents(SPECIFIC_AGENTS, AGENT_DELAY_MS);
} else {
console.log(`[Scheduler] Run #${runCount} - Mode: ALL AGENTS`);
result = await collectAllAgents();
}
lastRunResult = { ...result, run_count: runCount };
// Log MongoDB database capacity usage (expensive full-scan aggregation).
// Only run periodically (default: every 24h) to avoid high CPU/DB load each cycle.
const now = Date.now();
if (now - lastCapacityLogAt >= CAPACITY_LOG_INTERVAL_MS) {
lastCapacityLogAt = now;
const { logCapacityStats } = require('./db/capacityTracker');
await logCapacityStats(`[PROXY] [MongoDB] Capacity Used after Run #${runCount}:`);
}
return lastRunResult;
} catch (err) {
console.error('[Scheduler] Unhandled error during collection:', err.message);
lastRunResult = { success: false, error: err.message, run_count: runCount };
return lastRunResult;
} finally {
isRunning = false;
}
}
/**
* Start the scheduler (cron job + immediate first run).
*/
function startScheduler() {
console.log(`[Scheduler] Starting proxy data collector`);
const modeLabel = COLLECT_MODE === 'agent'
? `SPECIFIC AGENT (${SPECIFIC_AGENT})`
: COLLECT_MODE === 'agents'
? `SPECIFIC AGENTS (${SPECIFIC_AGENTS.join(', ')})`
: 'ALL AGENTS';
console.log(`[Scheduler] Mode : ${modeLabel}`);
console.log(`[Scheduler] Schedule : ${CRON_SCHEDULE} (every 5 minutes by default)`);
// Validate cron expression
if (!cron.validate(CRON_SCHEDULE)) {
console.error(`[Scheduler] Invalid cron expression: "${CRON_SCHEDULE}". Using default.`);
}
// Start recurring cron job
cron.schedule(CRON_SCHEDULE, () => {
runCollection().catch(err => console.error('[Scheduler] Cron error:', err.message));
});
console.log('[Scheduler] Cron job registered. Starting initial collection...');
// Initial run immediately on startup (async, do not block server start)
setTimeout(async () => {
// 1. Sync dictionary first
// await netifyClient.syncApplicationDictionary();
// 2. Start normal telemetry collection
runCollection().catch(err => console.error('[Scheduler] Initial run error:', err.message));
}, 2000);
}
/**
* Get current scheduler status (for REST API endpoint).
*/
function getStatus() {
return {
is_running: isRunning,
run_count: runCount,
last_run_at: lastRunAt?.toISOString() ?? null,
collect_mode: COLLECT_MODE,
agent_uuid: SPECIFIC_AGENT,
agent_uuids: COLLECT_MODE === 'agents' ? SPECIFIC_AGENTS : [],
agent_delay_ms: AGENT_DELAY_MS,
cron_schedule: CRON_SCHEDULE,
last_result: lastRunResult,
};
}
module.exports = { startScheduler, runCollection, getStatus };
// proxy/scheduler.js
// Cron scheduler for automatic data collection from DPI API
// Runs every 5 minutes, collecting data for all agents or a specific agent.
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';
const SPECIFIC_AGENT = process.env.PROXY_AGENT_UUID || null;
const SPECIFIC_AGENTS = (process.env.PROXY_AGENT_UUIDS || '').split(',').map(s => s.trim()).filter(Boolean);
const AGENT_DELAY_MS = parseInt(process.env.PROXY_AGENT_DELAY_MS || '5000');
const CRON_SCHEDULE = process.env.PROXY_CRON_SCHEDULE || '*/5 * * * *';
// Capacity logging is an expensive full-scan aggregation. Run it at most once per
// interval (default 24h) instead of every collection cycle to reduce CPU/DB load.
const CAPACITY_LOG_INTERVAL_MS = parseInt(process.env.PROXY_CAPACITY_LOG_INTERVAL_MS || String(24 * 60 * 60 * 1000));
let isRunning = false;
let lastRunAt = null;
let lastRunResult = null;
let runCount = 0;
let lastCapacityLogAt = 0;
/**
* Execute one collection cycle (called by cron and manual trigger).
* Prevents concurrent runs with isRunning guard.
*/
async function runCollection() {
if (isRunning) {
console.log('[Scheduler] Skipping - previous run still in progress');
return { skipped: true, reason: 'already_running' };
}
isRunning = true;
lastRunAt = new Date();
runCount++;
try {
let result;
if (COLLECT_MODE === 'agent' && SPECIFIC_AGENT) {
console.log(`[Scheduler] Run #${runCount} - Mode: SPECIFIC AGENT (${SPECIFIC_AGENT})`);
result = await collectSpecificAgent(SPECIFIC_AGENT);
} else if (COLLECT_MODE === 'agents' && SPECIFIC_AGENTS.length > 0) {
console.log(`[Scheduler] Run #${runCount} - Mode: SPECIFIC AGENTS (${SPECIFIC_AGENTS.join(', ')})`);
result = await collectSpecificAgents(SPECIFIC_AGENTS, AGENT_DELAY_MS);
} else {
console.log(`[Scheduler] Run #${runCount} - Mode: ALL AGENTS`);
result = await collectAllAgents();
}
lastRunResult = { ...result, run_count: runCount };
// Log MongoDB database capacity usage (expensive full-scan aggregation).
// Only run periodically (default: every 24h) to avoid high CPU/DB load each cycle.
const now = Date.now();
if (now - lastCapacityLogAt >= CAPACITY_LOG_INTERVAL_MS) {
lastCapacityLogAt = now;
const { logCapacityStats } = require('./db/capacityTracker');
await logCapacityStats(`[PROXY] [MongoDB] Capacity Used after Run #${runCount}:`);
}
return lastRunResult;
} catch (err) {
console.error('[Scheduler] Unhandled error during collection:', err.message);
lastRunResult = { success: false, error: err.message, run_count: runCount };
return lastRunResult;
} finally {
isRunning = false;
}
}
/**
* Start the scheduler (cron job + immediate first run).
*/
function startScheduler() {
console.log(`[Scheduler] Starting proxy data collector`);
const modeLabel = COLLECT_MODE === 'agent'
? `SPECIFIC AGENT (${SPECIFIC_AGENT})`
: COLLECT_MODE === 'agents'
? `SPECIFIC AGENTS (${SPECIFIC_AGENTS.join(', ')})`
: 'ALL AGENTS';
console.log(`[Scheduler] Mode : ${modeLabel}`);
console.log(`[Scheduler] Schedule : ${CRON_SCHEDULE} (every 5 minutes by default)`);
// Validate cron expression
if (!cron.validate(CRON_SCHEDULE)) {
console.error(`[Scheduler] Invalid cron expression: "${CRON_SCHEDULE}". Using default.`);
}
// Start recurring cron job
cron.schedule(CRON_SCHEDULE, () => {
runCollection().catch(err => console.error('[Scheduler] Cron error:', err.message));
});
console.log('[Scheduler] Cron job registered. Starting initial collection...');
// Initial run immediately on startup (async, do not block server start)
setTimeout(async () => {
// 1. Sync dictionary first
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));
}, 2000);
}
/**
* Get current scheduler status (for REST API endpoint).
*/
function getStatus() {
return {
is_running: isRunning,
run_count: runCount,
last_run_at: lastRunAt?.toISOString() ?? null,
collect_mode: COLLECT_MODE,
agent_uuid: SPECIFIC_AGENT,
agent_uuids: COLLECT_MODE === 'agents' ? SPECIFIC_AGENTS : [],
agent_delay_ms: AGENT_DELAY_MS,
cron_schedule: CRON_SCHEDULE,
last_result: lastRunResult,
};
}
module.exports = { startScheduler, runCollection, getStatus };
+71
View File
@@ -0,0 +1,71 @@
// test_api.js - Tests the actual metadata-detail API endpoint
// Run: node proxy/test_api.js
const http = require('http');
function request(path) {
return new Promise((resolve, reject) => {
const options = {
hostname: 'localhost',
port: 3001,
path,
method: 'GET',
};
const req = http.request(options, res => {
let body = '';
res.on('data', chunk => body += chunk);
res.on('end', () => {
try { resolve({ status: res.statusCode, data: JSON.parse(body) }); }
catch (e) { resolve({ status: res.statusCode, raw: body }); }
});
});
req.on('error', reject);
req.end();
});
}
async function main() {
// Test 1: netbios_hostname for 10.6.10.44
console.log('\n=== TEST 1: netbios_hostname=10.6.10.44 ===');
try {
const r1 = await request('/api/dashboard/metadata-detail?type=netbios_hostname&value=10.6.10.44');
console.log('Status:', r1.status);
if (r1.data) {
console.log('Count:', r1.data.count);
console.log('First 3 rows:', JSON.stringify(r1.data.data?.slice(0, 3), null, 2));
} else {
console.log('Raw:', r1.raw?.slice(0, 500));
}
} catch (e) {
console.log('ERROR (maybe backend is on different port):', e.message);
}
// Test 2: Try port 3000 (Next.js API routes)
console.log('\n=== TEST 2: via Next.js port 3000 ===');
try {
const r2 = await request('/api/dashboard/metadata-detail?type=netbios_hostname&value=10.6.10.44');
const options2 = { hostname: 'localhost', port: 3000, path: '/api/dashboard/metadata-detail?type=netbios_hostname&value=10.6.10.44', method: 'GET' };
const r3 = await new Promise((resolve, reject) => {
const req = http.request(options2, res => {
let body = '';
res.on('data', chunk => body += chunk);
res.on('end', () => {
try { resolve({ status: res.statusCode, data: JSON.parse(body) }); }
catch (e) { resolve({ status: res.statusCode, raw: body?.slice(0, 500) }); }
});
});
req.on('error', reject);
req.end();
});
console.log('Port 3000 - Status:', r3.status);
if (r3.data) {
console.log('Count:', r3.data.count);
console.log('First 3 rows:', JSON.stringify(r3.data.data?.slice(0, 3), null, 2));
} else {
console.log('Raw:', r3.raw);
}
} catch (e) {
console.log('Port 3000 ERROR:', e.message);
}
}
main().catch(console.error);
+1
View File
@@ -0,0 +1 @@
const mongoose = require('mongoose'); require('dotenv').config({ path: '../.env.local' }); const { LookupApp } = require('./models/Schemas'); async function test() { await mongoose.connect(process.env.MONGODB_URI || 'mongodb://127.0.0.1:27017/backone_dpi'); const sample = await LookupApp.findOne({ tag: /youtube/i }).lean(); console.log(JSON.stringify(sample, null, 2)); await mongoose.disconnect(); } test().catch(console.error);
+89
View File
@@ -0,0 +1,89 @@
// test_metadata_detail.js - Test after restart
// Run: node proxy/test_metadata_detail.js
const http = require('http');
function post(path, body) {
return new Promise((resolve, reject) => {
const data = JSON.stringify(body);
const options = {
hostname: 'localhost', port: 3001, path, method: 'POST',
headers: { 'Content-Type': 'application/json', 'Content-Length': Buffer.byteLength(data) },
};
const req = http.request(options, res => {
let buf = '';
res.on('data', c => buf += c);
res.on('end', () => {
try { resolve({ status: res.statusCode, headers: res.headers, data: JSON.parse(buf) }); }
catch { resolve({ status: res.statusCode, raw: buf }); }
});
});
req.on('error', reject);
req.write(data);
req.end();
});
}
function get(path, cookie) {
return new Promise((resolve, reject) => {
const options = {
hostname: 'localhost', port: 3001, path, method: 'GET',
headers: cookie ? { Cookie: cookie } : {},
};
const req = http.request(options, res => {
let buf = '';
res.on('data', c => buf += c);
res.on('end', () => {
try { resolve({ status: res.statusCode, data: JSON.parse(buf) }); }
catch { resolve({ status: res.statusCode, raw: buf?.slice(0, 300) }); }
});
});
req.on('error', reject);
req.end();
});
}
async function main() {
// Step 1: Login to get cookie
console.log('=== Step 1: Login ===');
const login = await post('/api/auth/login', { username: 'admin', password: 'admin123' });
console.log('Login status:', login.status);
const setCookie = login.headers?.['set-cookie'];
let cookie = '';
if (setCookie) {
cookie = setCookie.map(c => c.split(';')[0]).join('; ');
console.log('Cookie obtained:', cookie.slice(0, 60) + '...');
} else {
console.log('No cookie received. Auth response:', JSON.stringify(login.data));
// Try with a known admin credential
}
// Step 2: Test health
console.log('\n=== Step 2: Health Check ===');
const health = await get('/api/health', cookie);
console.log('Health:', health.status, JSON.stringify(health.data));
// Step 3: Test metadata-detail netbios_hostname
console.log('\n=== Step 3: metadata-detail netbios_hostname=10.6.10.44 ===');
const r = await get('/api/dashboard/metadata-detail?type=netbios_hostname&value=10.6.10.44', cookie);
console.log('Status:', r.status);
if (r.data) {
console.log('Count (should be 1):', r.data.count);
console.log('Data:', JSON.stringify(r.data.data, null, 2));
} else {
console.log('Raw:', r.raw);
}
// Step 4: Verify DB has 1 doc for 10.6.10.44
console.log('\n=== Step 4: Direct DB verification ===');
const mongoose = require('mongoose');
const path = require('path');
require('dotenv').config({ path: path.join(__dirname, '..', '.env.local') });
await mongoose.connect(process.env.MONGODB_URI || 'mongodb://localhost:27017/backone');
const db = mongoose.connection.db;
const cnt = await db.collection('devicestats').countDocuments({ ip_address: '10.6.10.44' });
console.log('DB count for 10.6.10.44:', cnt, '(expected: 1)');
await mongoose.disconnect();
}
main().catch(console.error);
+1
View File
@@ -0,0 +1 @@
const mongoose = require('mongoose'); require('dotenv').config({ path: '../.env.local' }); const { syncApplicationDictionary } = require('./netifyClient'); const { LookupApp } = require('./models/Schemas'); async function test() { await mongoose.connect(process.env.MONGODB_URI || 'mongodb://127.0.0.1:27017/backone_dpi'); console.log('Connected to DB'); await syncApplicationDictionary(); const count = await LookupApp.countDocuments(); console.log('Total LookupApps in DB:', count); await mongoose.disconnect(); } test().catch(console.error);
+56
View File
@@ -0,0 +1,56 @@
// verify_fix.js - Directly verify the MongoDB aggregation returns correct results
// This simulates what the backend metadata-detail endpoint does after the fix.
// Run: node proxy/verify_fix.js
const path = require('path');
require('dotenv').config({ path: path.join(__dirname, '..', '.env.local') });
const mongoose = require('mongoose');
async function run() {
await mongoose.connect(process.env.MONGODB_URI || 'mongodb://localhost:27017/backone');
const db = mongoose.connection.db;
const col = db.collection('devicestats');
const testValues = ['10.6.10.44'];
// Also find other common device_labels to test
const sample = await col.find({}).limit(20).toArray();
const labels = [...new Set(sample.map(d => d.device_label).filter(Boolean))];
console.log('Sample device_labels to test:', labels.slice(0, 5));
for (const value of [...testValues, ...labels.slice(0, 3)]) {
// Count raw docs matching
const rawCount = await col.countDocuments({ device_label: value });
// Simulate the new aggregation pipeline
const aggResult = await col.aggregate([
{ $match: { device_label: value } },
{ $sort: { timestamp: -1 } },
{ $group: {
_id: '$ip_address',
mac_address: { $first: '$mac_address' },
device_label: { $first: '$device_label' },
download: { $max: '$download' },
upload: { $max: '$upload' },
}},
{ $sort: { download: -1 } },
]).toArray();
const status = rawCount > aggResult.length ? '✅ FIXED (was duplicated)' : '✓ OK';
console.log(`\ndevice_label="${value}": raw=${rawCount} rows → aggregated=${aggResult.length} unique devices ${status}`);
if (aggResult.length > 0) {
console.log(' First result:', JSON.stringify(aggResult[0], null, 2));
}
}
// Verify unique index exists
const indexes = await col.indexes();
const uniqueIdx = indexes.find(i => i.unique && i.key.agent_uuid && i.key.ip_address);
console.log('\n=== Unique Index on (agent_uuid, ip_address):', uniqueIdx ? `✅ EXISTS (${uniqueIdx.name})` : '❌ MISSING');
// Final count
const total = await col.countDocuments();
console.log('=== Total DeviceStat docs:', total);
await mongoose.disconnect();
}
run().catch(err => { console.error(err.message); process.exit(1); });