pipeline/cluster-core.js
// Event layer core — shared by both runtimes. Pure data-in/data-out except the
// injected `ask(messages)` AI callback and `onSpend` ledger hook.
//
// Phase A: deterministic candidates (word-set similarity, union-find).
// Phase B: one AI call per candidate -> events + facets; failures degrade to
// the deterministic grouping, never to data loss.
// Phase C: cross-day thread matching vs the previous day's events (one batched
// AI call); conservative high-bar linking without a verdict.
import { hashId, jaccard, words, extractJson, eventWeight } from "./util.js";
import { CLUSTER_PROMPT, THREADS_PROMPT } from "./prompts.js";
const MAX_CANDIDATE_MEMBERS = 12; // oversized candidates: extras fall back deterministically
const THREAD_SHORTLIST = 24;
const eventWords = (ev, byId) => {
const set = new Set();
for (const m of ev.members.slice(0, 6)) {
const e = byId.get(m);
if (e) for (const w of words(`${e.article.title} ${e.essence}`)) set.add(w);
}
return set;
};
const shortText = (ev, byId, side) =>
`[${side}] topic=${ev.topic} · ` + ev.members.slice(0, 2).map((m) => {
const e = byId.get(m);
return e ? `${e.article.source_name} · ${e.article.title}: ${e.essence}` : "";
}).filter(Boolean).join(" / ");
function mkEvent(members, facets) {
const memberIds = new Set(members.map((e) => e.id));
const sources = [...new Set(members.map((e) => e.article.source_id))];
const topicVotes = {};
for (const e of members) {
const t = e.topics[0]?.topic || "other";
topicVotes[t] = (topicVotes[t] || 0) + 1;
}
const topic = Object.entries(topicVotes).sort((a, b) => b[1] - a[1])[0][0];
const seed = members[0];
const locations = {};
for (const e of members) for (const l of e.locations || []) locations[l.country_code] = l.name;
const cleanFacets = facets
.map((f) => ({ label: f.label, members: [...new Set(f.members)].filter((m) => memberIds.has(m)) }))
.filter((f) => f.members.length)
.map((f) => ({ id: hashId("fct", seed.id + f.label), label: f.label, members: f.members }));
// Freshness: first-seen and last-update from the member records' own timestamps.
const times = members.map((e) => e.distillation?.distilled_at).filter(Boolean).sort();
return {
id: hashId("evt", seed.id),
seed: seed.id,
members: members.map((e) => e.id),
sources,
topic,
subtopic: seed.topics[0]?.subtopic || "other",
facets: cleanFacets,
locations: Object.entries(locations).map(([country_code, name]) => ({ country_code, name })),
reports: members.length,
first_seen_at: times[0] || null,
last_update_at: times[times.length - 1] || null,
weight: eventWeight(sources.length, cleanFacets.length, 1),
};
}
export async function runCluster({ essences, prevEvents, prevEssences, ask, onSpend, budget }) {
// essences: ok records only (no failures), in append order
const essById = new Map(essences.map((e) => [e.id, e]));
/* ---------- phase A ---------- */
const sig = essences.map((e) => ({
e,
w: words(`${e.article.title} ${e.essence}`),
sub: e.topics[0]?.subtopic,
ccs: new Set((e.locations || []).map((l) => l.country_code)),
}));
const parent = sig.map((_, i) => i);
const find = (i) => { while (parent[i] !== i) { parent[i] = parent[parent[i]]; i = parent[i]; } return i; };
const union = (a, b) => { const ra = find(a), rb = find(b); if (ra !== rb) parent[rb] = ra; };
let pairs = 0;
for (let i = 0; i < sig.length; i++) {
for (let j = i + 1; j < sig.length; j++) {
const a = sig[i], b = sig[j];
const jac = jaccard(a.w, b.w);
const sharedCountry = [...a.ccs].some((c) => b.ccs.has(c));
if (jac >= 0.3 || (a.sub === b.sub && sharedCountry && jac >= 0.18)) { union(i, j); pairs++; }
}
}
const clusters = new Map();
for (let i = 0; i < sig.length; i++) {
const r = find(i);
if (!clusters.has(r)) clusters.set(r, []);
clusters.get(r).push(i);
}
const candidates = [...clusters.values()].filter((g) => g.length > 1);
/* ---------- phase B ---------- */
const facetsOf = new Map();
let aiCalls = 0, fallbacks = 0;
const withinBudget = () => budget == null || !budget.exhausted();
if (candidates.length && withinBudget()) {
for (let ci = 0; ci < candidates.length; ci++) {
const members = candidates[ci].slice(0, MAX_CANDIDATE_MEMBERS).map((i) => essences[i]);
if (members.length < 2) continue;
if (!withinBudget()) { fallbacks += candidates.length - ci; break; }
const list = members.map((e) => `[${e.id}] ${e.article.source_name} · ${e.article.title}: ${e.essence}`).join("\n");
try {
const r = await ask([{ role: "system", content: CLUSTER_PROMPT }, { role: "user", content: list }], { temperature: 0.2, maxTokens: 3000 });
if (onSpend) onSpend({ script: "cluster", ...r });
aiCalls++;
const out = extractJson(r.text);
const ids = new Set(members.map((e) => e.id));
const ok = Array.isArray(out.groups) && out.groups.every((g) =>
Array.isArray(g.members) && g.members.every((m) => ids.has(m)) &&
(g.facets == null || Array.isArray(g.facets)));
if (!ok) throw new Error("shape validation failed");
const groups = [];
for (const g of out.groups) {
const facets = (g.facets && g.facets.length ? g.facets : [{ label: "corroboration", members: g.members }])
.filter((f) => Array.isArray(f.members) && f.members.length && f.label)
.map((f) => ({ label: String(f.label).slice(0, 80), members: f.members.filter((m) => ids.has(m)) }));
groups.push({ members: g.members, facets });
}
facetsOf.set(ci, groups);
} catch (err) {
fallbacks++;
facetsOf.set(ci, [{
members: members.map((e) => e.id),
facets: [{ label: "corroboration", members: members.map((e) => e.id) }],
}]);
}
}
}
/* ---------- assemble ---------- */
const events = [];
const assigned = new Set();
for (const [ci, groups] of facetsOf) {
const clusterMap = new Map(candidates[ci].map((i) => [essences[i].id, essences[i]]));
const claimed = new Set();
for (const g of groups) {
const members = g.members.filter((m) => clusterMap.has(m) && !claimed.has(m));
for (const id of members) { assigned.add(id); claimed.add(id); }
if (members.length) events.push(mkEvent(members.map((id) => clusterMap.get(id)), g.facets));
}
for (const id of candidates[ci].map((i) => essences[i].id)) {
if (!claimed.has(id)) { assigned.add(id); events.push(mkEvent([clusterMap.get(id)], [])); }
}
}
for (const e of essences) {
if (!assigned.has(e.id)) { assigned.add(e.id); events.push(mkEvent([e], [])); }
}
// sanity: every essence in exactly one event
const seen = new Set();
for (const ev of events) for (const m of ev.members) {
if (seen.has(m)) throw new Error(`essence ${m} assigned to two events`);
seen.add(m);
}
/* ---------- phase C: threads ---------- */
let continuations = 0, threadAiCalls = 0;
if (prevEvents && prevEvents.length) {
const prevById = prevEssences;
const pairsList = [];
for (const ev of events) {
const wt = eventWords(ev, essById);
if (!wt.size) continue;
for (const pev of prevEvents) {
const wy = eventWords(pev, prevById);
if (!wy.size) continue;
const jac = jaccard(wt, wy);
const sharedCountry = (ev.locations || []).some((l) => (pev.locations || []).some((p) => p.country_code === l.country_code));
if ((jac >= 0.18 && ev.topic === pev.topic) || jac >= 0.28 || (sharedCountry && ev.subtopic === pev.subtopic && jac >= 0.12)) {
pairsList.push({ today: ev, prev: pev, jac });
}
}
}
pairsList.sort((a, b) => b.jac - a.jac);
const shortlist = pairsList.slice(0, THREAD_SHORTLIST);
let verdicts = null;
if (shortlist.length && withinBudget()) {
try {
const user = shortlist.map((p, i) =>
`PAIR ${i + 1}:\nTODAY: ${shortText(p.today, essById, "today")}\nYESTERDAY: ${shortText(p.prev, prevById, "yesterday")}`
).join("\n\n");
const r = await ask([{ role: "system", content: THREADS_PROMPT }, { role: "user", content: user }], { temperature: 0.2, maxTokens: 3000 });
if (onSpend) onSpend({ script: "threads", ...r });
threadAiCalls++;
const out = extractJson(r.text);
verdicts = new Map();
if (Array.isArray(out.verdicts)) {
for (const v of out.verdicts) if (Number.isInteger(v.pair) && v.pair >= 1 && v.pair <= shortlist.length) verdicts.set(v.pair, !!v.continues);
}
} catch { verdicts = null; }
}
const usedPrev = new Set(), usedToday = new Set();
shortlist.forEach((p, i) => {
const confirmed = verdicts ? verdicts.get(i + 1) === true : p.jac >= 0.30;
if (!confirmed || usedToday.has(p.today.id) || usedPrev.has(p.prev.id)) return;
p.today.continues = p.prev.id;
p.today.thread_id = p.prev.thread_id || p.prev.id;
p.today.thread_days = (p.prev.thread_days || 1) + 1;
p.today.first_seen_at = p.prev.first_seen_at || p.today.first_seen_at; // the thread's origin, not today's
usedPrev.add(p.prev.id);
usedToday.add(p.today.id);
continuations++;
});
}
// thread defaults + sustention into weight
for (const ev of events) {
if (!ev.thread_id) { ev.thread_id = ev.id; ev.thread_days = 1; }
if (ev.continues === undefined) ev.continues = null;
ev.weight = eventWeight(ev.sources.length, ev.facets.length, ev.thread_days);
}
return {
events,
clustering: { candidates: candidates.length, ai_calls: aiCalls, fallbacks, thread_ai_calls: threadAiCalls, continuations },
};
}