← 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.`);