scripts/cluster.js
#!/usr/bin/env node
// cluster.js — thin Node wrapper: shared cluster core + REST AI client.
// Events + facets + cross-day threads; every ok essence lands in exactly one event.
import fs from "node:fs";
import path from "node:path";
import { budgetCheck, DATA_ROOT, loadConfig, logSpend, nowIso, parseArgs, PROMPT_VERSION, prevDate as prevOf, readAllJsonl, readJsonl, runModel, spendToday } from "./lib.js";
import { runCluster } from "../pipeline/cluster-core.js";
const { args } = parseArgs(process.argv.slice(2), { date: { takes: "value", default: null } });
const config = loadConfig();
const date = args.date || new Date().toISOString().slice(0, 10);
const essFile = path.join(DATA_ROOT, "distilled", date, "essence.jsonl");
const outFile = path.join(DATA_ROOT, "aggregates", date, "events.json");
const essences = readJsonl(essFile).filter((e) => !e.failure);
if (!essences.length) {
console.log(`cluster ${date}: no essence records, skipping`);
process.exit(0);
}
// Previous day for thread matching (fixed once committed)
const pd = prevOf(date);
const prevFile = path.join(DATA_ROOT, "aggregates", pd, "events.json");
const prevEvents = fs.existsSync(prevFile) ? JSON.parse(fs.readFileSync(prevFile, "utf8")).events : null;
const prevEssences = prevEvents ? new Map(readJsonl(path.join(DATA_ROOT, "distilled", pd, "essence.jsonl")).filter((e) => !e.failure).map((e) => [e.id, e])) : null;
let spent = 0;
const budget = {
exhausted: () => config.daily_budget_usd - spendToday() <= 0 || !budgetCheck(config, "cluster").ok,
};
const ask = async (messages, { temperature, maxTokens }) => {
const r = await runModel(config.distill_model, messages, { temperature, max_tokens: maxTokens });
if (r.cost != null) spent += r.cost;
return r;
};
const onSpend = (s) => logSpend({ ...s, model: config.distill_model });
console.log(`cluster ${date}: ${essences.length} essences${prevEvents ? `, ${prevEvents.length} previous-day events` : ""}`);
const result = await runCluster({ essences, prevEvents, prevEssences, ask, onSpend, budget });
const out = {
date,
generated_at: nowIso(),
model: config.distill_model,
prompt_version: PROMPT_VERSION,
clustering: { ...result.clustering, previous_day: prevEvents ? pd : null },
events: result.events,
};
fs.mkdirSync(path.dirname(outFile), { recursive: true });
fs.writeFileSync(outFile, JSON.stringify(out, null, 1) + "\n");
const multi = result.events.filter((e) => e.reports > 1).length;
const faceted = result.events.filter((e) => e.facets.some((f) => f.label !== "corroboration")).length;
const threads = result.events.filter((e) => e.thread_days > 1);
console.log(`cluster ${date}: ${result.events.length} events (${multi} corroborated, ${faceted} with distinct facets) -> aggregates/${date}/events.json [AI calls: ${result.clustering.ai_calls}, fallbacks: ${result.clustering.fallbacks}, ~$${spent.toFixed(4)}]`);
if (threads.length) console.log(` threads: ${threads.length} events continue from previous days (longest: day ${Math.max(...threads.map((e) => e.thread_days))})`);
// unused essences warning keeps parity with the raw JSONL check
if (readAllJsonl) { /* keep import used for future tooling */ }