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