← back to Re Flyer Aggregator
scripts/ingest-assets.mjs
85 lines
#!/usr/bin/env node
// TK-10708 Ingest classified marketing-asset findings into the LOCAL reflyers staging
// DB (external_marketing_asset + asset_subject_link). Idempotent per (source_landing_url,
// doc_number). Also registers the Tier-2 self-generated recaps as assets (--recaps).
//
// Usage:
// node scripts/ingest-assets.mjs data/findings-2026-08-19.json
// node scripts/ingest-assets.mjs --recaps # register generated out/spotlight-*.html
// node scripts/ingest-assets.mjs --report # print aggregated counts
import { execFileSync } from 'node:child_process';
import { readFileSync, readdirSync } from 'node:fs';
import { fileURLToPath } from 'node:url';
import { dirname, join } from 'node:path';
const ROOT = join(dirname(fileURLToPath(import.meta.url)), '..');
const DB = process.env.REFLYERS_DB || 'reflyers'; // promote to usre via REFLYERS_DB=usre
const psql = sql => execFileSync('psql', [DB, '-t', '-A', '-c', sql], { encoding: 'utf8' }).trim();
const lit = v => (v === null || v === undefined || v === '') ? 'NULL' : `'${String(v).replace(/'/g, "''")}'`;
const litB = v => (v === null || v === undefined) ? 'NULL' : (v ? 'true' : 'false');
const args = process.argv.slice(2);
const STATS = { inserted: 0, skipped: 0 };
function insertAsset(a) {
// idempotency: skip if same landing_url already present
if (a.source_landing_url) {
const dup = psql(`SELECT id FROM external_marketing_asset WHERE source_landing_url=${lit(a.source_landing_url)} LIMIT 1`);
if (dup) { STATS.skipped++; return +String(dup).split('\n')[0].trim(); }
}
STATS.inserted++;
const id = psql(`INSERT INTO external_marketing_asset
(asset_type,tier,title,description,source_name,source_landing_url,document_url,discovery_method,rights_basis,access_status,robots_ok,local_path,file_format,last_verified_at)
VALUES (${lit(a.asset_type)},${a.tier|0},${lit(a.title)},${lit(a.description)},${lit(a.source_name)},${lit(a.source_landing_url)},${lit(a.document_url)},${lit(a.discovery_method)},${lit(a.rights_basis)},${lit(a.access_status)},${litB(a.robots_ok)},${lit(a.local_path)},${lit(a.file_format)},now())
RETURNING id`);
// -t -A still prints the "INSERT 0 1" status after the RETURNING value -> take line 1.
return +String(id).split('\n')[0].trim();
}
function linkSubject(assetId, s) {
const dup = psql(`SELECT id FROM asset_subject_link WHERE asset_id=${assetId} AND county_fips=${lit(s.county_fips)} AND coalesce(doc_number,'')=${lit(s.doc_number||'')} LIMIT 1`);
if (dup) return;
psql(`INSERT INTO asset_subject_link (asset_id,subject_type,county_fips,ain,doc_number,match_confidence)
VALUES (${assetId},'deal',${lit(s.county_fips)},NULL,${lit(s.doc_number)},${s.match_confidence ?? 'NULL'})`);
}
if (args.includes('--report')) {
console.log('=== assets by tier x rights_basis ===');
console.log(psql(`SELECT tier, rights_basis, count(*) FROM external_marketing_asset GROUP BY 1,2 ORDER BY 1,2`));
console.log('\n=== assets per deal (top) ===');
console.log(psql(`SELECT l.doc_number, count(*) AS assets, round(avg(l.match_confidence),2) AS avg_conf
FROM asset_subject_link l GROUP BY 1 ORDER BY 2 DESC`));
console.log('\n=== total ===', psql(`SELECT count(*) FROM external_marketing_asset`), 'assets,',
psql(`SELECT count(*) FROM asset_subject_link`), 'links');
process.exit(0);
}
if (args.includes('--recaps')) {
// Register generated recaps. Each out/spotlight-<doc>.html is a Tier-2 asset for that deal.
const outDir = join(ROOT, 'out');
const files = readdirSync(outDir).filter(f => /^spotlight-.*\.html$/.test(f));
let n = 0;
for (const f of files) {
const m = f.match(/^spotlight-([0-9-]+)\.html$/); // per-doc files only (skip top-N aggregates)
if (!m) continue;
const doc = m[1];
const id = insertAsset({ asset_type:'deal_recap', tier:2, title:`Self-generated closed-sale recap (${doc})`,
source_name:'self-generated', discovery_method:'generated', rights_basis:'self_generated',
access_status:'internal', local_path:`out/${f}`, file_format:'html' });
linkSubject(id, { county_fips:'12086', doc_number:doc, match_confidence:1.0 });
n++;
}
console.log(`Registered ${n} per-deal recap asset(s).`);
process.exit(0);
}
const file = args[0];
if (!file) { console.error('Usage: ingest-assets.mjs <findings.json> | --recaps | --report'); process.exit(1); }
const findings = JSON.parse(readFileSync(join(ROOT, file), 'utf8'));
let assets = 0, links = 0;
for (const rec of findings) {
const id = insertAsset(rec.asset); assets++;
for (const s of (rec.subjects || [])) { linkSubject(id, s); links++; }
}
console.log(`Ingested from ${file}: ${STATS.inserted} new asset(s), ${STATS.skipped} already-present (idempotent), ${links} subject-link upserts.`);