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