← back to Butlr

lib/scheduling.js

184 lines

// lib/scheduling.js — JSON-backed scheduling layer for Butlr.
//
// Two kinds of scheduled work:
//   (1) one-shot — schedule_at is a single ISO timestamp; row fires once
//   (2) recurring — cron_expr (cron-style 5-field) computes next_fire_at on each tick
//
// The worker tick runs every 60s. For each row whose next_fire_at <= now, it
// dispatches via the same /api/external/place-call path used for immediate
// calls (honors repeat-call gate, dry-run env, etc.).
//
// Hard rule — scheduled calls inherit ALL gating that immediate calls have.
// MAX_CALLS_PER_PHONE applies. DNC list applies. consent_recording must be
// true. No bypass.

const fs = require('fs');
const path = require('path');
const crypto = require('crypto');

const FILE = path.join(__dirname, '..', 'data', 'scheduled-calls.json');

function readAll() {
  try { return JSON.parse(fs.readFileSync(FILE, 'utf8')); }
  catch (e) { if (e.code === 'ENOENT') return []; throw e; }
}
function writeAll(rows) {
  fs.writeFileSync(FILE, JSON.stringify(rows, null, 2));
}
function genId() {
  return crypto.randomBytes(6).toString('base64url');
}

// Tiny cron evaluator — supports 5-field expressions with `*`, `*/N`, `N`,
// `N-M`, `A,B,C`. Fields: minute hour day-of-month month day-of-week.
// Returns the next Date >= `from` matching the expression.
function parseField(field, min, max) {
  if (field === '*') return null;  // match-all
  const out = new Set();
  for (const part of field.split(',')) {
    let m;
    if ((m = part.match(/^\*\/(\d+)$/))) {
      const step = parseInt(m[1], 10);
      for (let i = min; i <= max; i += step) out.add(i);
    } else if ((m = part.match(/^(\d+)-(\d+)$/))) {
      for (let i = parseInt(m[1], 10); i <= parseInt(m[2], 10); i++) out.add(i);
    } else if (/^\d+$/.test(part)) {
      out.add(parseInt(part, 10));
    } else {
      throw new Error('bad cron field: ' + field);
    }
  }
  return out;
}
function nextCronFire(expr, from = new Date()) {
  const parts = String(expr).trim().split(/\s+/);
  if (parts.length !== 5) throw new Error('cron must be 5 fields (m h dom mon dow)');
  const [minF, hourF, domF, monF, dowF] = parts;
  const sets = {
    minute: parseField(minF, 0, 59),
    hour:   parseField(hourF, 0, 23),
    dom:    parseField(domF, 1, 31),
    mon:    parseField(monF, 1, 12),
    dow:    parseField(dowF, 0, 6),
  };
  const d = new Date(from);
  d.setSeconds(0, 0);
  d.setMinutes(d.getMinutes() + 1);  // strictly forward
  // Brute-force scan up to 1 year ahead — fine for sane crons.
  for (let i = 0; i < 366 * 24 * 60; i++) {
    if ((!sets.minute || sets.minute.has(d.getMinutes())) &&
        (!sets.hour   || sets.hour.has(d.getHours())) &&
        (!sets.dom    || sets.dom.has(d.getDate())) &&
        (!sets.mon    || sets.mon.has(d.getMonth() + 1)) &&
        (!sets.dow    || sets.dow.has(d.getDay()))) {
      return new Date(d);
    }
    d.setMinutes(d.getMinutes() + 1);
  }
  return null;
}

// Public API
function listScheduled(userId, opts = {}) {
  const all = readAll();
  let rows = userId ? all.filter(r => r.user_id === userId) : all;
  if (opts.status) rows = rows.filter(r => r.status === opts.status);
  if (opts.from)   rows = rows.filter(r => r.next_fire_at && new Date(r.next_fire_at) >= new Date(opts.from));
  if (opts.to)     rows = rows.filter(r => r.next_fire_at && new Date(r.next_fire_at) <= new Date(opts.to));
  return rows.sort((a, b) => new Date(a.next_fire_at || 0) - new Date(b.next_fire_at || 0));
}

function addScheduled({ user_id, business_name, business_phone, goal, callback_phone, callback_name, schedule_at, cron_expr, category, max_hold_minutes, max_spend_cents, consent_recording }) {
  if (!business_phone) throw new Error('business_phone required');
  if (!schedule_at && !cron_expr) throw new Error('schedule_at or cron_expr required');

  let next;
  if (cron_expr) {
    next = nextCronFire(cron_expr);
    if (!next) throw new Error('cron expression never fires in next year');
  } else {
    next = new Date(schedule_at);
    if (isNaN(next.getTime())) throw new Error('schedule_at not parseable');
    if (next < new Date()) throw new Error('schedule_at is in the past');
  }

  const row = {
    id: 'S' + genId(),
    user_id: user_id || null,
    business_name: business_name || '',
    business_phone,
    goal: goal || '',
    callback_phone: callback_phone || '',
    callback_name: callback_name || '',
    category: category || 'other',
    max_hold_minutes: max_hold_minutes || 5,
    max_spend_cents: max_spend_cents || 50,
    consent_recording: !!consent_recording,
    schedule_at: cron_expr ? null : next.toISOString(),
    cron_expr: cron_expr || null,
    next_fire_at: next.toISOString(),
    status: 'scheduled',
    fire_count: 0,
    last_fire_at: null,
    last_call_id: null,
    created_at: new Date().toISOString(),
  };
  const all = readAll();
  all.push(row);
  writeAll(all);
  return row;
}

function cancelScheduled(id, userId) {
  const all = readAll();
  const r = all.find(x => x.id === id);
  if (!r) return false;
  if (userId && r.user_id !== userId) return false;
  r.status = 'canceled';
  r.next_fire_at = null;
  writeAll(all);
  return true;
}

// Worker tick — finds due rows and dispatches them. Returns array of fired IDs.
// `dispatchFn` is injected so this stays test-friendly + doesn't import
// circular deps into data.js.
async function tickWorker(dispatchFn) {
  const all = readAll();
  const now = new Date();
  const fired = [];
  for (const r of all) {
    if (r.status !== 'scheduled') continue;
    if (!r.next_fire_at) continue;
    if (new Date(r.next_fire_at) > now) continue;

    try {
      const result = await dispatchFn(r);
      r.fire_count = (r.fire_count || 0) + 1;
      r.last_fire_at = now.toISOString();
      r.last_call_id = (result && result.call_id) || null;

      if (r.cron_expr) {
        // Recurring — compute next fire, leave status='scheduled'
        const next = nextCronFire(r.cron_expr);
        r.next_fire_at = next ? next.toISOString() : null;
        if (!r.next_fire_at) r.status = 'done';
      } else {
        // One-shot — mark done
        r.status = 'done';
        r.next_fire_at = null;
      }
      fired.push({ id: r.id, ok: true, call_id: r.last_call_id });
    } catch (e) {
      r.status = 'failed';
      r.last_fire_at = now.toISOString();
      r.last_error = String(e.message || e).slice(0, 240);
      fired.push({ id: r.id, ok: false, error: r.last_error });
    }
  }
  if (fired.length) writeAll(all);
  return fired;
}

module.exports = { listScheduled, addScheduled, cancelScheduled, tickWorker, nextCronFire };