← back to Watches
scripts/price-scheduler.js
75 lines
/**
* Price Update Scheduler
* Runs price aggregation on a schedule (every 6 hours)
*/
import cron from 'node-cron';
import priceAggregator from '../services/price-aggregator.js';
import pool from '../database/connection.js';
const UPDATE_INTERVAL = '0 */6 * * *'; // Every 6 hours
let isRunning = false;
let lastRunTime = null;
let lastRunStats = null;
async function getWatches() {
const result = await pool.query('SELECT id, reference, model FROM watches WHERE is_active = true');
return result.rows;
}
async function runPriceUpdate() {
if (isRunning) {
console.log('Price update already running, skipping');
return;
}
isRunning = true;
const startTime = Date.now();
console.log(`[${new Date().toISOString()}] Starting price update...`);
try {
const watches = await getWatches();
console.log(`Found ${watches.length} watches to update`);
const results = await priceAggregator.aggregateAllWatches(watches);
const successful = results.filter(r => r.hasData).length;
const failed = results.filter(r => r.error).length;
lastRunStats = {
total: watches.length,
successful,
failed,
duration: Date.now() - startTime
};
console.log(`Price update complete: ${successful} successful, ${failed} failed, ${lastRunStats.duration}ms`);
} catch (error) {
console.error('Price update failed:', error);
lastRunStats = { error: error.message };
} finally {
isRunning = false;
lastRunTime = new Date().toISOString();
}
}
export function startScheduler() {
console.log(`Price scheduler started: ${UPDATE_INTERVAL}`);
cron.schedule(UPDATE_INTERVAL, runPriceUpdate);
// Run immediately on start
setTimeout(runPriceUpdate, 5000);
}
export function getStatus() {
return { isRunning, lastRunTime, lastRunStats, schedule: UPDATE_INTERVAL };
}
export function triggerManualUpdate() {
runPriceUpdate();
return { message: 'Update triggered' };
}
export default { startScheduler, getStatus, triggerManualUpdate };