← back to AbramsOS
scripts/ingest-mode-claims-corpus.js
77 lines
#!/usr/bin/env node
// Corpus (George-path) ingester: reads the local Mode newsletter corpus files
// (data/mode-emails-raw.js + data/mode-emails-new.js — bodies pulled via George), parses them
// with lib/mode-claims, and upserts distinct settlements into settlement_claim.
//
// This is the durable fallback for when the Gmail OAuth *connector* isn't set up yet
// (connector_account has 0 rows), which is why scripts/ingest-mode-claims.js can't run.
// Same schema + same idempotent UNIQUE(user_id,slug) upsert as that script. No submits, no browser. $0.
//
// Usage: node scripts/ingest-mode-claims-corpus.js
const path = require('path');
const db = require('../lib/db');
const { parseModeEmail } = require('../lib/mode-claims');
const { id } = require('../lib/ids');
const USER = process.env.ABRAMSOS_USER_ID || 'user_steve';
function loadCorpus() {
const out = [];
// Load EVERY data/mode-emails-*.js corpus shard (raw + new + any george-refreshed shard),
// so newly fetched mailers are picked up automatically without editing this file.
const fs = require('fs');
const dataDir = path.join(__dirname, '..', 'data');
let files = [];
try {
files = fs.readdirSync(dataDir)
.filter(f => /^mode-emails-.*\.js$/.test(f))
.sort() // stable order across shards
.map(f => path.join(dataDir, f));
} catch (e) { console.warn(`[corpus] readdir ${dataDir}: ${e.message}`); }
for (const f of files) {
try { const arr = require(f); if (Array.isArray(arr)) out.push(...arr); }
catch (e) { console.warn(`[corpus] skip ${f}: ${e.message}`); }
}
// de-dupe emails by id, newest wins
const byId = new Map();
for (const e of out) byId.set(e.id, e);
return [...byId.values()];
}
async function main() {
const emails = loadCorpus();
console.log(`[corpus] ${emails.length} newsletter bodies loaded`);
let seen = 0, inserted = 0, updated = 0;
for (const e of emails) {
for (const c of parseModeEmail(e.body, { emailId: e.id, emailDate: e.date })) {
seen++;
const r = await db.query(
`INSERT INTO settlement_claim
(id,user_id,slug,name,mode_url,payout_text,payout_max_cents,deadline,proof_required,
no_claim_required,eligibility_question,category,source_email_id,source_email_date,raw_block)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15)
ON CONFLICT (user_id,slug) DO UPDATE SET
mode_url=COALESCE(EXCLUDED.mode_url,settlement_claim.mode_url),
name=COALESCE(EXCLUDED.name,settlement_claim.name),
deadline=COALESCE(EXCLUDED.deadline,settlement_claim.deadline),
payout_text=COALESCE(EXCLUDED.payout_text,settlement_claim.payout_text),
payout_max_cents=GREATEST(COALESCE(EXCLUDED.payout_max_cents,0),COALESCE(settlement_claim.payout_max_cents,0)),
proof_required=COALESCE(settlement_claim.proof_required,EXCLUDED.proof_required),
no_claim_required=(settlement_claim.no_claim_required OR EXCLUDED.no_claim_required),
eligibility_question=COALESCE(settlement_claim.eligibility_question,EXCLUDED.eligibility_question),
updated_at=now()
RETURNING (xmax=0) AS inserted`,
[id('claim'), USER, c.slug, c.name, c.mode_url, c.payout_text, c.payout_max_cents, c.deadline,
c.proof_required, c.no_claim_required, c.eligibility_question, c.category, c.source_email_id, c.source_email_date, c.raw_block]);
if (r.rows[0]?.inserted) inserted++; else updated++;
}
}
// auto-expire past-deadline rows that were never acted on
const exp = await db.query(
`UPDATE settlement_claim SET fill_state='expired',updated_at=now()
WHERE user_id=$1 AND deadline < current_date AND fill_state IN ('none','queued')`, [USER]);
console.log(`[corpus] parsed ${seen} blocks → ${inserted} new, ${updated} updated; ${exp.rowCount} expired past-deadline`);
process.exit(0);
}
main().catch(e => { console.error('FATAL', e); process.exit(1); });