scripts/aggregate.js
#!/usr/bin/env node
// aggregate.js — thin Node wrapper: shared aggregate core + local storage.
import fs from "node:fs";
import path from "node:path";
import { DATA_ROOT, loadTopics, nowIso, parseArgs, readJsonl } from "./lib.js";
import { buildTopicsTree } from "../pipeline/aggregate-core.js";
const { args } = parseArgs(process.argv.slice(2), { date: { takes: "value", default: null } });
const date = args.date || new Date().toISOString().slice(0, 10);
const essFile = path.join(DATA_ROOT, "distilled", date, "essence.jsonl");
const rawFile = path.join(DATA_ROOT, "raw", date, "articles.jsonl");
const srcFile = path.join(DATA_ROOT, "raw", date, "sources.json");
const eventsFile = path.join(DATA_ROOT, "aggregates", date, "events.json");
const outDir = path.join(DATA_ROOT, "aggregates", date);
const essences = readJsonl(essFile);
if (!essences.length) {
console.log(`aggregate ${date}: no essence records yet, skipping`);
process.exit(0);
}
const articles = readJsonl(rawFile);
const sourceStatus = fs.existsSync(srcFile) ? JSON.parse(fs.readFileSync(srcFile, "utf8")) : {};
const events = fs.existsSync(eventsFile) ? JSON.parse(fs.readFileSync(eventsFile, "utf8")).events : null;
const topics = loadTopics();
const { tree } = buildTopicsTree({ essences, events, topics });
fs.mkdirSync(outDir, { recursive: true });
fs.writeFileSync(path.join(outDir, "topics.json"), JSON.stringify({ date, generated_at: nowIso(), tree }, null, 1) + "\n");
// Day stats: unique counts, source health, spend from the ledger (day's lines)
const ledgerPath = path.join(DATA_ROOT, "ledger.jsonl");
let tokensIn = 0, tokensOut = 0, usd = 0;
for (const l of readJsonl(ledgerPath)) if (l.date === date) { tokensIn += l.tokens_in || 0; tokensOut += l.tokens_out || 0; usd += l.cost_usd || 0; }
const seenOk = new Set(), failedIds = new Set();
for (const e of essences) {
if (e.failure) failedIds.add(e.article.raw_id);
else if (!seenOk.has(e.article.raw_id)) seenOk.add(e.article.raw_id);
}
fs.writeFileSync(path.join(outDir, "day.json"), JSON.stringify({
date,
generated_at: nowIso(),
counts: {
articles: articles.length,
essences: seenOk.size,
failures: [...failedIds].filter((id) => !seenOk.has(id)).length,
events: events ? events.length : null,
},
sources: sourceStatus,
spend: { tokens_in: tokensIn, tokens_out: tokensOut, usd },
}, null, 1) + "\n");
console.log(`aggregate ${date}: ${essences.length} essence lines across ${tree.length} topics${events ? `, ${events.length} events` : ""} -> aggregates/${date}/`);
for (const t of tree) console.log(` ${t.name}: ${t.count} reports${events ? ` / ${t.events} events (score ${t.score.toFixed(1)})` : ""}`);