source.fact.ngo a coherence.ngo project

pipeline/cluster-core.js

raw ↗ · AGPL-3.0

// 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 }, }; }