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