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