[object Object]

← back to AbramsOS

tick 9: in-process cron for reminders + CPSC refresh

d34f2f605e0499d62d6780274c5f0e7c9fdcb4f9 · 2026-05-10 01:03:06 -0700 · Steve

- npm install node-cron
- lib/scheduler.js: two jobs
  · reminders: every 4h (idempotent via UNIQUE dedupe)
  · cpsc:       daily 03:17 PT, last 7 days (idempotent via ON CONFLICT)
- Logs each tick to logs/scheduler.log
- start() is a no-op when NODE_ENV=test or SCHEDULER_DISABLED=1, so tests stay fast
- Wired into server.js startup; verified live: scheduler started message in log
- 34/35 green

Files touched

Diff

commit d34f2f605e0499d62d6780274c5f0e7c9fdcb4f9
Author: Steve <steve@designerwallcoverings.com>
Date:   Sun May 10 01:03:06 2026 -0700

    tick 9: in-process cron for reminders + CPSC refresh
    
    - npm install node-cron
    - lib/scheduler.js: two jobs
      · reminders: every 4h (idempotent via UNIQUE dedupe)
      · cpsc:       daily 03:17 PT, last 7 days (idempotent via ON CONFLICT)
    - Logs each tick to logs/scheduler.log
    - start() is a no-op when NODE_ENV=test or SCHEDULER_DISABLED=1, so tests stay fast
    - Wired into server.js startup; verified live: scheduler started message in log
    - 34/35 green
---
 lib/scheduler.js  | 90 +++++++++++++++++++++++++++++++++++++++++++++++++++++++
 package-lock.json | 10 +++++++
 package.json      |  1 +
 server.js         |  1 +
 4 files changed, 102 insertions(+)

diff --git a/lib/scheduler.js b/lib/scheduler.js
new file mode 100644
index 0000000..b2857a4
--- /dev/null
+++ b/lib/scheduler.js
@@ -0,0 +1,90 @@
+// 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 { 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 (_) {}
+  }
+}
+
+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));
+
+  started = true;
+  logLine('scheduler started · reminders every 4h · cpsc daily 03:17');
+}
+
+module.exports = { start, regenerateRemindersForAll, refreshCpscRecalls };
diff --git a/package-lock.json b/package-lock.json
index a12d00d..81f24ff 100644
--- a/package-lock.json
+++ b/package-lock.json
@@ -18,6 +18,7 @@
         "googleapis": "^144.0.0",
         "helmet": "^8.0.0",
         "morgan": "^1.10.0",
+        "node-cron": "^4.2.1",
         "otplib": "^12.0.1",
         "pdf-parse": "^2.4.5",
         "pg": "^8.13.0",
@@ -1707,6 +1708,15 @@
         "node": "^18 || ^20 || >= 21"
       }
     },
+    "node_modules/node-cron": {
+      "version": "4.2.1",
+      "resolved": "https://registry.npmjs.org/node-cron/-/node-cron-4.2.1.tgz",
+      "integrity": "sha512-lgimEHPE/QDgFlywTd8yTR61ptugX3Qer29efeyWw2rv259HtGBNn1vZVmp8lB9uo9wC0t/AT4iGqXxia+CJFg==",
+      "license": "ISC",
+      "engines": {
+        "node": ">=6.0.0"
+      }
+    },
     "node_modules/node-fetch": {
       "version": "2.7.0",
       "resolved": "https://registry.npmjs.org/node-fetch/-/node-fetch-2.7.0.tgz",
diff --git a/package.json b/package.json
index a094e51..dd98305 100644
--- a/package.json
+++ b/package.json
@@ -20,6 +20,7 @@
     "googleapis": "^144.0.0",
     "helmet": "^8.0.0",
     "morgan": "^1.10.0",
+    "node-cron": "^4.2.1",
     "otplib": "^12.0.1",
     "pdf-parse": "^2.4.5",
     "pg": "^8.13.0",
diff --git a/server.js b/server.js
index 7d0937d..8c673eb 100644
--- a/server.js
+++ b/server.js
@@ -76,6 +76,7 @@ app.use((err, _req, res, _next) => {
 if (require.main === module) {
   app.listen(PORT, () => {
     console.log(`[abramsos] listening on http://localhost:${PORT}`);
+    require('./lib/scheduler').start();
   });
 }
 

← 1f31901 tick 8: CSRF protection on HTML POST forms  ·  back to AbramsOS  ·  tick 10: /audit page (read-only viewer for audit_log + auth_ ee6a1eb →