← 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 };