← back to AbramsOS

scripts/ingest-recalls.js

79 lines

#!/usr/bin/env node
// Backfill CPSC recalls for the last N days (default 365). Idempotent: ON CONFLICT skips.
// Usage: node scripts/ingest-recalls.js [days]

require('dotenv').config();
const db = require('../lib/db');
const { id } = require('../lib/ids');
const cpsc = require('../lib/cpsc-fetcher');
const matcher = require('../lib/recall-matcher');

const DEV_USER_ID = 'user_steve';

async function main() {
  const days = parseInt(process.argv[2] || '365', 10);
  const since = new Date(Date.now() - days * 24 * 60 * 60 * 1000);
  console.log(`[ingest] CPSC recalls since ${since.toISOString().slice(0, 10)} (${days} days)`);

  let raw;
  try {
    raw = await cpsc.fetchRecalls({ since });
  } catch (err) {
    console.error('[ingest] CPSC fetch failed:', err.message);
    await db.query(`INSERT INTO recall_pull_log (authority, fetched, inserted, updated, error) VALUES ('CPSC', 0, 0, 0, $1)`, [err.message]);
    process.exit(1);
  }

  const list = Array.isArray(raw) ? raw : raw?.results || [];
  console.log(`[ingest] received ${list.length} recalls from CPSC`);

  let inserted = 0;
  let updated = 0;
  let totalMatches = 0;

  for (const r of list) {
    const norm = cpsc.normalize(r);
    if (!norm.external_id) continue;
    const eventId = id('document'); // reuse ulid; we're storing in recall_event

    try {
      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`,
        [eventId, norm.authority, norm.external_id, norm.published_at, norm.title, norm.hazard, norm.remedy, norm.url, norm.product_keys, norm.raw]
      );
      const row = result.rows[0];
      if (row.was_inserted) inserted += 1; else updated += 1;

      // Try to match this recall against the user's purchases
      try {
        const fullRow = await db.query(`SELECT * FROM recall_event WHERE id = $1`, [row.id]);
        const matches = await matcher.matchRecallAgainstUser(fullRow.rows[0], DEV_USER_ID);
        totalMatches += matches;
      } catch (e) {
        console.warn('[match]', norm.external_id, e.message);
      }
    } catch (err) {
      console.warn('[insert]', norm.external_id, err.message);
    }
  }

  await db.query(
    `INSERT INTO recall_pull_log (authority, fetched, inserted, updated) VALUES ('CPSC', $1, $2, $3)`,
    [list.length, inserted, updated]
  );
  console.log(`[ingest] done. inserted=${inserted} updated=${updated} matches=${totalMatches}`);
  await db.pool.end();
}

main().catch((e) => { console.error(e); process.exit(1); });