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