source.fact.ngo a coherence.ngo project

scripts/cluster.js

raw ↗ · AGPL-3.0

#!/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 */ }