← back to AbramsOS

lib/scheduler.js

187 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 autopilot = require('./claims-autopilot');
const claimsAlert = require('./claims-alert');
const amazonOrders = require('./amazon-orders');
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}`);
  }
}

async function runClaimsAutopilot() {
  // Stage/fill pass — failure here is a real problem, log as FAILED.
  let r;
  try {
    r = await autopilot.autoProcess();
    logLine(`claims-autopilot: staged=${r.staged.length} automatic=${r.automatic.length} already=${r.already.length}` +
      (r.profile_missing.length ? ` PROFILE-MISSING=${r.profile_missing.join(',')}` : ''));
  } catch (err) {
    logLine(`claims-autopilot FAILED: ${err.message}`);
    return;
  }

  // Alert pass — George 401/config issues are expected until cred is set; log as WARN, not FAILED.
  // CNCP parking-lot fallback fires ONLY when George is unavailable.
  let georgeOk = false;
  try {
    const a = await claimsAlert.sendClaimsAlert();
    logLine(`claims-alert: ${a.sent ? `SENT total=${a.total} urgent=${a.urgent}` : 'skipped — ' + a.skipped}`);
    georgeOk = true;
  } catch (err) {
    logLine(`claims-alert WARN (George unavailable): ${err.message} — pushing urgent items to CNCP`);
  }

  // CNCP parking-lot fallback: only when George fails so the board doesn't flood on every 6h run.
  if (!georgeOk) {
    try {
      const db2 = require('./db');
      const urgent = await db2.query(
        `SELECT name, deadline, (deadline::date - current_date) days_left, fill_state
           FROM settlement_claim WHERE user_id=$1 AND deadline IS NOT NULL
             AND deadline >= current_date AND deadline <= current_date + 7
             AND eligibility_state='eligible'
           ORDER BY deadline ASC`, [autopilot.USER]);
      if (urgent.rows.length) {
        const lines = urgent.rows.map(c =>
          `${c.days_left === 0 ? '🔴 TODAY' : c.days_left === 1 ? '🟡 TOMORROW' : `🟠 ${c.days_left}d`} · ${c.name} (${c.fill_state})`
        ).join('\n');
        const http = require('http');
        const body = JSON.stringify({ title: `⚖️ ${urgent.rows.length} claim(s) due ≤7 days`, body: lines, source: 'claims-autopilot' });
        await new Promise(res => {
          const req = http.request({ hostname: '127.0.0.1', port: 3333, path: '/api/parking-lot', method: 'POST', headers: { 'Content-Type': 'application/json', 'Content-Length': Buffer.byteLength(body) } }, res);
          req.on('error', res); req.end(body);
        }).catch(() => {});
        logLine(`claims-alert CNCP: posted ${urgent.rows.length} urgent claims to parking lot`);
      }
    } catch (e) { /* CNCP fallback is best-effort */ }
  }
}

async function runAmazonOrders() {
  try {
    const r = await amazonOrders.pollAndImport();
    logLine(`amazon-orders: ${r.skipped ? 'skipped — ' + r.skipped : `imported=${r.imported} reorders=${r.reorders} scanned=${r.scanned}`}`);
  } catch (err) {
    logLine(`amazon-orders 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' });

  // Settlement-claim autopilot — auto-fill every open claim + deadline-alert so
  // none lapse. Runs every 6h and once at boot. Staging is local/reversible and
  // always runs; the alert self-throttles + no-ops without George. NEVER submits.
  const claimsOn = process.env.CLAIMS_AUTOPILOT_DISABLED !== '1';
  if (claimsOn) {
    cron.schedule('23 */6 * * *', runClaimsAutopilot, { timezone: 'America/Los_Angeles' });
    setTimeout(() => { runClaimsAutopilot(); }, 15000); // boot pass, after DB warms
  }

  // Amazon order poller — every 30 min, import new order confirmations into
  // purchase + reorder_item (dedup by Gmail id). No-ops without George config.
  const ordersOn = process.env.AMAZON_ORDERS_DISABLED !== '1';
  if (ordersOn) {
    cron.schedule('*/30 * * * *', runAmazonOrders);
    setTimeout(() => { runAmazonOrders(); }, 25000);
  }

  started = true;
  logLine(`scheduler started · reminders every 4h · cpsc daily 03:17 · digest ${digestOn ? 'daily 08:00' : 'DISABLED'} · claims-autopilot ${claimsOn ? 'every 6h' : 'DISABLED'}`);
}

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