125 lines
4.7 KiB
JavaScript
125 lines
4.7 KiB
JavaScript
// 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 };
|