← back to AbramsOS

lib/scheduler.js

109 lines

// In-process cron. Single-instance pm2 process so no leader election needed.
// All jobs are read-only or generate-only (idempotent via UNIQUE indexes).

const fs = require('fs');
const path = require('path');
const cron = require('node-cron');
const db = require('./db');
const reminders = require('./reminder-engine');
const cpsc = require('./cpsc-fetcher');
const matcher = require('./recall-matcher');
const digest = require('./digest');
const { id } = require('./ids');

const LOG_FILE = path.join(__dirname, '..', 'logs', 'scheduler.log');

function logLine(line) {
  const stamp = new Date().toISOString();
  const msg = `${stamp}  ${line}\n`;
  fs.appendFileSync(LOG_FILE, msg);
  console.log('[scheduler]', line);
}

async function regenerateRemindersForAll() {
  try {
    const users = await db.query(`SELECT id FROM user_account`);
    let total = 0;
    for (const u of users.rows) {
      total += await reminders.generateForUser(u.id);
    }
    logLine(`reminders: ${total} new across ${users.rows.length} user(s)`);
  } catch (err) {
    logLine(`reminders FAILED: ${err.message}`);
  }
}

async function refreshCpscRecalls(days = 7) {
  try {
    const since = new Date(Date.now() - days * 86400e3);
    const raw = await cpsc.fetchRecalls({ since });
    const list = Array.isArray(raw) ? raw : raw?.results || [];
    let inserted = 0, updated = 0, matches = 0;
    for (const r of list) {
      const norm = cpsc.normalize(r);
      if (!norm.external_id) continue;
      const result = await db.query(
        `INSERT INTO recall_event (id, authority, external_id, published_at, title, hazard, remedy, url, product_keys_jsonb, raw_jsonb)
         VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
         ON CONFLICT (authority, external_id) DO UPDATE SET
           published_at = EXCLUDED.published_at, title = EXCLUDED.title,
           hazard = EXCLUDED.hazard, remedy = EXCLUDED.remedy, url = EXCLUDED.url,
           product_keys_jsonb = EXCLUDED.product_keys_jsonb, raw_jsonb = EXCLUDED.raw_jsonb
         RETURNING id, (xmax = 0) AS was_inserted`,
        [id('document'), norm.authority, norm.external_id, norm.published_at, norm.title, norm.hazard, norm.remedy, norm.url, JSON.stringify(norm.product_keys), JSON.stringify(norm.raw)]
      );
      const row = result.rows[0];
      if (row.was_inserted) inserted += 1; else updated += 1;
      const fullRow = await db.query(`SELECT * FROM recall_event WHERE id = $1`, [row.id]);
      const users = await db.query(`SELECT id FROM user_account`);
      for (const u of users.rows) {
        try { matches += await matcher.matchRecallAgainstUser(fullRow.rows[0], u.id); } catch (_) {}
      }
    }
    await db.query(
      `INSERT INTO recall_pull_log (authority, fetched, inserted, updated) VALUES ('CPSC', $1, $2, $3)`,
      [list.length, inserted, updated]
    );
    logLine(`cpsc: fetched=${list.length} inserted=${inserted} updated=${updated} matches=${matches}`);
  } catch (err) {
    logLine(`cpsc FAILED: ${err.message}`);
    try {
      await db.query(`INSERT INTO recall_pull_log (authority, fetched, inserted, updated, error) VALUES ('CPSC', 0, 0, 0, $1)`, [err.message]);
    } catch (_) {}
  }
}

async function sendDailyDigest() {
  try {
    const users = await db.query(`SELECT id FROM user_account`);
    for (const u of users.rows) {
      const r = await digest.sendDigest(u.id);
      logLine(`digest[${u.id}]: ${r.sent ? `sent (${r.total} deadlines, ${r.overdue} overdue) msg=${r.messageId}` : `skipped — ${r.skipped}`}`);
    }
  } catch (err) {
    logLine(`digest FAILED: ${err.message}`);
  }
}

let started = false;
function start() {
  if (started) return;
  if (process.env.NODE_ENV === 'test') return;
  if (process.env.SCHEDULER_DISABLED === '1') return;

  // Every 4 hours, regenerate reminders
  cron.schedule('11 */4 * * *', regenerateRemindersForAll);
  // Every 24 hours at 03:17, refresh CPSC recalls (last 7 days)
  cron.schedule('17 3 * * *', () => refreshCpscRecalls(7));
  // Daily 08:00 — email Steve his deadline digest. OFF by default: opt in with
  // DIGEST_ENABLED=1 (and set GEORGE_BASIC_AUTH). Double-gated so the committed
  // code never auto-sends until Steve explicitly turns it on.
  const digestOn = process.env.DIGEST_ENABLED === '1';
  if (digestOn) cron.schedule('0 8 * * *', sendDailyDigest, { timezone: 'America/Los_Angeles' });

  started = true;
  logLine(`scheduler started · reminders every 4h · cpsc daily 03:17 · digest ${digestOn ? 'daily 08:00' : 'DISABLED (set DIGEST_ENABLED=1)'}`);
}

module.exports = { start, regenerateRemindersForAll, refreshCpscRecalls };