← back to Stars of Design

scrapers/gmail-recon-ingest.js

137 lines

#!/usr/bin/env node
'use strict';
// Bulk-upsert sod_email_candidates from Gmail recon batches.
//
// Reads NDJSON from stdin (one candidate per line) OR a JSON-array file
// via --file=<path>. Each record looks like:
//   { email, display_name, domain, first_seen, last_seen,
//     thread_count, steve_replied, signal_score, last_signature, sample_subject }
//
// Idempotent. Re-running on the same threads updates rather than dup-creates.
// Signal score is *taken from input* (calculator lives in the producer);
// this script only persists.

require('dotenv').config();
const { Pool } = require('pg');
const fs = require('fs');
const readline = require('readline');

const args = process.argv.slice(2).reduce((m, a) => {
  const [k, v] = a.replace(/^--/, '').split('=');
  m[k] = v ?? true; return m;
}, {});
const FILE = args.file || null;
const SOURCE = args.source || 'gmail-recon';
const DRY = !!args.dry;

const pool = new Pool({
  connectionString: process.env.DATABASE_URL || 'postgresql:///dw_unified?host=/tmp&user=stevestudio2',
  max: 4,
});

const PERSONAL_ESPS = new Set(['gmail.com','hotmail.com','yahoo.com','outlook.com','icloud.com','aol.com','me.com','msn.com','protonmail.com','live.com','comcast.net','att.net','verizon.net']);
const KNOWN_NOISE   = new Set(['noreply','no-reply','donotreply','do-not-reply','support','help','info','contact','hello','team','newsletter','press','marketing','sales','accounts','billing','notifications','alert','alerts','update','updates','mailer-daemon','postmaster','reply','automated']);

function domainBonus(domain, localpart) {
  if (!domain) return 0;
  if (KNOWN_NOISE.has(localpart)) return -0.25;
  if (PERSONAL_ESPS.has(domain)) return 0.0;
  // Strong designer-flavored TLDs/hosts get a bump
  if (/(design|studio|interiors|architect|ateliers?|stylist|hospitality|hospitalitygroup)/i.test(domain)) return 0.20;
  if (/(asid|iida|aia|nkba|kbis|nycxdesign|wallpaperz?|holland|sherwin|kravet|schumacher|cowtan|romo|maharam|knoll|herman)/i.test(domain)) return 0.05;
  return 0.0;
}

async function upsert(rec) {
  const email   = String(rec.email || '').toLowerCase().trim();
  if (!email || !email.includes('@')) return { skipped: true };
  const [local, dom] = email.split('@');
  const domain = (rec.domain || dom).toLowerCase();
  const localpart = local.toLowerCase();
  const displayName = rec.display_name || null;
  const firstSeen = rec.first_seen || null;
  const lastSeen  = rec.last_seen  || null;
  const tc = Math.max(1, parseInt(rec.thread_count || 1, 10));
  const sr = Math.max(0, parseInt(rec.steve_replied || 0, 10));
  const lastSig = rec.last_signature || null;
  const sampleSubject = rec.sample_subject || null;

  // signal = base reply-ratio (capped) + domain bonus
  const baseRatio = sr / Math.max(1, tc);
  const score = Math.max(0, Math.min(1, baseRatio * 0.85 + domainBonus(domain, localpart) + (sr > 0 ? 0.10 : 0)));

  if (DRY) {
    console.log(JSON.stringify({ email, displayName, domain, tc, sr, score: score.toFixed(3) }));
    return { dry: true };
  }

  // Merge counts on conflict; never decrease counts.
  await pool.query(`
    INSERT INTO sod_email_candidates
      (email, display_name, domain, first_seen, last_seen, thread_count, steve_replied,
       signal_score, last_signature, notes)
    VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10)
    ON CONFLICT (email) DO UPDATE SET
      display_name   = COALESCE(EXCLUDED.display_name, sod_email_candidates.display_name),
      first_seen     = LEAST(COALESCE(sod_email_candidates.first_seen, EXCLUDED.first_seen), EXCLUDED.first_seen),
      last_seen      = GREATEST(COALESCE(sod_email_candidates.last_seen, EXCLUDED.last_seen), EXCLUDED.last_seen),
      thread_count   = sod_email_candidates.thread_count + EXCLUDED.thread_count,
      steve_replied  = sod_email_candidates.steve_replied + EXCLUDED.steve_replied,
      signal_score   = GREATEST(COALESCE(sod_email_candidates.signal_score, 0), EXCLUDED.signal_score),
      last_signature = COALESCE(EXCLUDED.last_signature, sod_email_candidates.last_signature),
      notes          = COALESCE(NULLIF(EXCLUDED.notes,''), sod_email_candidates.notes),
      updated_at     = NOW()
  `, [email, displayName, domain, firstSeen, lastSeen, tc, sr, score, lastSig, sampleSubject]);
  return { ok: true, score };
}

async function processStream(stream) {
  let runId = null;
  if (!DRY) {
    const r = await pool.query(
      `INSERT INTO sod_ingest_runs (run_kind, source, status) VALUES ('gmail-candidates', $1, 'running') RETURNING id`,
      [SOURCE]
    );
    runId = r.rows[0].id;
  }
  const rl = readline.createInterface({ input: stream, crlfDelay: Infinity });
  let n = 0, ok = 0, skipped = 0, errs = 0;
  for await (const line of rl) {
    const s = line.trim();
    if (!s || s.startsWith('#')) continue;
    n++;
    try {
      const rec = JSON.parse(s);
      const r = await upsert(rec);
      if (r.ok || r.dry) ok++; else skipped++;
    } catch (e) {
      errs++;
      console.error(`line ${n}: ${e.message}`);
    }
    if (n % 100 === 0) console.log(`  …${n} processed (ok=${ok} skip=${skipped} err=${errs})`);
  }
  if (runId) {
    await pool.query(
      `UPDATE sod_ingest_runs SET status='ok', finished_at=NOW(), records_in=$1, records_out=$2, notes=$3 WHERE id=$4`,
      [n, ok, `skip=${skipped} err=${errs}`, runId]
    );
  }
  console.log(`[gmail-recon-ingest] done. processed=${n} ok=${ok} skip=${skipped} err=${errs}`);
  await pool.end();
}

(async function main() {
  if (FILE && fs.existsSync(FILE)) {
    const ext = FILE.endsWith('.json') ? 'array' : 'ndjson';
    if (ext === 'array') {
      const arr = JSON.parse(fs.readFileSync(FILE, 'utf8'));
      const tmp = require('stream').Readable.from(arr.map(r => JSON.stringify(r) + '\n').join(''));
      await processStream(tmp);
    } else {
      await processStream(fs.createReadStream(FILE));
    }
  } else {
    await processStream(process.stdin);
  }
})().catch(e => { console.error('fatal:', e); process.exit(1); });