← back to Norma

agents/shared/cron-scheduler.js

181 lines

/**
 * Norma Cron Scheduler
 *
 * Simple cron wrapper using node-cron for agent task scheduling.
 * Tracks job state, last run times, and supports start/stop per job.
 */

const cron = require('node-cron');

class CronScheduler {
  /**
   * @param {string} agentName - Agent name for logging
   */
  constructor(agentName) {
    this.agentName = agentName;
    this.jobs = new Map();
  }

  /**
   * Add a cron job.
   *
   * @param {string} name - Unique job name (e.g. 'monitor-subreddits')
   * @param {string} schedule - Cron expression (e.g. every 30 min: 0,30 * * * *)
   * @param {Function} fn - Async function to execute on schedule
   * @param {Object} [options]
   * @param {boolean} [options.runOnStart=false] - Run immediately when started
   */
  add(name, schedule, fn, options = {}) {
    if (this.jobs.has(name)) {
      console.warn(`[${new Date().toISOString()}] [cron] Job "${name}" already exists, replacing`);
      this.stop(name);
    }

    if (!cron.validate(schedule)) {
      throw new Error(`Invalid cron expression: "${schedule}" for job "${name}"`);
    }

    const wrappedFn = async () => {
      const jobMeta = this.jobs.get(name);
      if (!jobMeta) return;

      jobMeta.lastRun = new Date();
      jobMeta.runCount++;

      try {
        await fn();
        jobMeta.lastStatus = 'success';
        jobMeta.lastError = null;
      } catch (err) {
        jobMeta.lastStatus = 'error';
        jobMeta.lastError = err.message;
        console.error(
          `[${new Date().toISOString()}] [cron] Job "${name}" failed:`, err.message
        );
      }
    };

    const task = cron.schedule(schedule, wrappedFn, {
      scheduled: false, // don't auto-start
    });

    this.jobs.set(name, {
      name,
      schedule,
      task,
      fn: wrappedFn,
      running: false,
      lastRun: null,
      lastStatus: null,
      lastError: null,
      runCount: 0,
      runOnStart: options.runOnStart || false,
    });

    console.log(`[${new Date().toISOString()}] [cron] [${this.agentName}] Added job: "${name}" (${schedule})`);
  }

  /**
   * Start a specific cron job.
   *
   * @param {string} name - Job name
   */
  start(name) {
    const job = this.jobs.get(name);
    if (!job) {
      console.warn(`[${new Date().toISOString()}] [cron] Job "${name}" not found`);
      return;
    }

    if (job.running) {
      console.warn(`[${new Date().toISOString()}] [cron] Job "${name}" is already running`);
      return;
    }

    job.task.start();
    job.running = true;
    console.log(`[${new Date().toISOString()}] [cron] [${this.agentName}] Started job: "${name}"`);

    // Run immediately if configured
    if (job.runOnStart) {
      job.fn().catch((err) => {
        console.error(`[${new Date().toISOString()}] [cron] Job "${name}" runOnStart failed:`, err.message);
      });
    }
  }

  /**
   * Start all registered cron jobs.
   */
  startAll() {
    for (const [name] of this.jobs) {
      this.start(name);
    }
    console.log(`[${new Date().toISOString()}] [cron] [${this.agentName}] All ${this.jobs.size} jobs started`);
  }

  /**
   * Stop a specific cron job.
   *
   * @param {string} name - Job name
   */
  stop(name) {
    const job = this.jobs.get(name);
    if (!job) {
      console.warn(`[${new Date().toISOString()}] [cron] Job "${name}" not found`);
      return;
    }

    job.task.stop();
    job.running = false;
    console.log(`[${new Date().toISOString()}] [cron] [${this.agentName}] Stopped job: "${name}"`);
  }

  /**
   * Stop all registered cron jobs.
   */
  stopAll() {
    for (const [name] of this.jobs) {
      this.stop(name);
    }
    console.log(`[${new Date().toISOString()}] [cron] [${this.agentName}] All jobs stopped`);
  }

  /**
   * Get status of all registered jobs.
   *
   * @returns {{ jobs: Array<{ name: string, schedule: string, running: boolean, lastRun: string|null, lastStatus: string|null, lastError: string|null, runCount: number }> }}
   */
  status() {
    const jobs = [];
    for (const [, job] of this.jobs) {
      jobs.push({
        name: job.name,
        schedule: job.schedule,
        running: job.running,
        lastRun: job.lastRun ? job.lastRun.toISOString() : null,
        lastStatus: job.lastStatus,
        lastError: job.lastError,
        runCount: job.runCount,
      });
    }
    return { agentName: this.agentName, jobs };
  }

  /**
   * Remove a job entirely.
   *
   * @param {string} name - Job name
   */
  remove(name) {
    const job = this.jobs.get(name);
    if (!job) return;

    job.task.stop();
    this.jobs.delete(name);
    console.log(`[${new Date().toISOString()}] [cron] [${this.agentName}] Removed job: "${name}"`);
  }
}

module.exports = { CronScheduler };