pipeline/aggregate-core.js
// Aggregate + GIS cores — pure rollups over the day's data, shared runtimes.
import { GEO_COUNTRIES } from "./geo.js";
import { lookupPlace } from "./gazetteer.js";
export function buildTopicsTree({ essences, events = null, topics }) {
// essences: all lines incl. failure records; tree counts reports, then events
const tName = Object.fromEntries(topics.map((t) => [t.id, t.name]));
const sName = {};
for (const t of topics) for (const s of t.subtopics) sName[s] = s.replaceAll("-", " ");
const tree = topics.filter((t) => t.id !== "other").map((t) => ({
topic: t.id, name: t.name, count: 0, events: 0, score: 0,
subtopics: t.subtopics.map((s) => ({ subtopic: s, name: sName[s], count: 0, events: 0, score: 0, essences: [], event_refs: [] })),
}));
const other = { topic: "other", name: tName.other || "Other", count: 0, events: 0, score: 0, subtopics: [{ subtopic: "other", name: "Other", count: 0, events: 0, score: 0, essences: [], event_refs: [] }] };
tree.push(other);
const tIx = Object.fromEntries(tree.map((t, i) => [t.topic, i]));
const ref = (e) => ({ id: e.id, essence: e.essence, source_id: e.article.source_id, source_name: e.article.source_name, url: e.article.url, title: e.article.title, event_type: e.event_type, uncertainty: e.uncertainty });
for (const e of essences) {
if (e.failure) continue;
const primary = e.topics[0] || { topic: "other", subtopic: "other" };
const t = tree[tIx[primary.topic] ?? tIx.other];
t.count++;
let sub = t.subtopics.find((s) => s.subtopic === primary.subtopic);
if (!sub) {
sub = { subtopic: primary.subtopic, name: sName[primary.subtopic] || primary.subtopic.replaceAll("-", " "), count: 0, events: 0, score: 0, essences: [], event_refs: [] };
t.subtopics.push(sub);
}
sub.count++;
if (sub.essences.length < 3) sub.essences.push(ref(e));
}
if (events) {
const byId = Object.fromEntries(essences.filter((e) => !e.failure).map((e) => [e.id, e]));
for (const ev of events) {
const t = tree[tIx[ev.topic] ?? tIx.other];
t.events++;
t.score += ev.weight;
let sub = t.subtopics.find((s) => s.subtopic === ev.subtopic);
if (!sub) {
sub = { subtopic: ev.subtopic, name: sName[ev.subtopic] || ev.subtopic.replaceAll("-", " "), count: 0, events: 0, score: 0, essences: [], event_refs: [] };
t.subtopics.push(sub);
}
sub.events++;
sub.score += ev.weight;
const members = ev.members.map((id) => byId[id]).filter(Boolean);
sub.event_refs.push({
id: ev.id, weight: ev.weight, reports: ev.reports,
thread_id: ev.thread_id, thread_days: ev.thread_days, continues: ev.continues,
first_seen_at: ev.first_seen_at || null, last_update_at: ev.last_update_at || null,
sources: ev.sources, facets: ev.facets.map((f) => ({ label: f.label, members: f.members })),
essences: members.slice(0, 8).map(ref),
});
}
}
const sortKey = (t) => t.score || t.count;
return { tree: tree.filter((t) => t.count > 0).sort((a, b) => sortKey(b) - sortKey(a)), ref };
}
export function buildGis({ essences, events = null }) {
const byCountry = {};
for (const e of essences) {
if (e.failure) continue;
for (const loc of e.locations || []) {
const c = GEO_COUNTRIES[loc.country_code];
if (!c) continue;
if (!byCountry[loc.country_code]) byCountry[loc.country_code] = { country_code: loc.country_code, name: c.name, lat: c.lat, lon: c.lon, count: 0, events: null, essences: [] };
if (byCountry[loc.country_code].essences.length < 5) {
byCountry[loc.country_code].essences.push({ id: e.id, essence: e.essence, source_id: e.article.source_id, source_name: e.article.source_name, url: e.article.url, title: e.article.title });
}
byCountry[loc.country_code].count++;
}
}
let eventCount = 0;
if (events) {
for (const cc of Object.keys(byCountry)) { byCountry[cc].events = 0; byCountry[cc].last_update_at = null; }
for (const ev of events) {
for (const loc of ev.locations || []) {
const c = GEO_COUNTRIES[loc.country_code];
if (!c || !byCountry[loc.country_code]) continue;
const bc = byCountry[loc.country_code];
bc.events++;
bc.topic_weights = bc.topic_weights || {};
bc.topic_weights[ev.topic] = (bc.topic_weights[ev.topic] || 0) + ev.weight;
if (ev.last_update_at && (!bc.last_update_at || ev.last_update_at > bc.last_update_at)) bc.last_update_at = ev.last_update_at;
if (ev.first_seen_at && (!bc.first_seen_at || ev.first_seen_at < bc.first_seen_at)) bc.first_seen_at = ev.first_seen_at;
eventCount++;
}
}
}
// City-level placement: events fuse reports, so places fuse events.
// Unresolved place names stay on their country placement (the stated limit).
const places = {};
if (events) {
const byId = Object.fromEntries(essences.filter((e) => !e.failure).map((e) => [e.id, e]));
for (const ev of events) {
for (const loc of ev.locations || []) {
const hit = lookupPlace(loc.name, loc.country_code);
if (!hit) continue;
const key = `${loc.country_code}|${String(hit[3]).toLowerCase()}`;
if (!places[key]) {
places[key] = { name: hit[3], country_code: loc.country_code, lat: hit[0], lon: hit[1], precision: "place", events: 0, reports: 0, weight: 0, essences: [] };
}
const p = places[key];
p.events++; p.reports += ev.reports; p.weight += ev.weight;
p.topic_weights = p.topic_weights || {};
p.topic_weights[ev.topic] = (p.topic_weights[ev.topic] || 0) + ev.weight;
if (ev.last_update_at && (!p.last_update_at || ev.last_update_at > p.last_update_at)) p.last_update_at = ev.last_update_at;
if (ev.first_seen_at && (!p.first_seen_at || ev.first_seen_at < p.first_seen_at)) p.first_seen_at = ev.first_seen_at;
if (p.essences.length < 2) {
const e = byId[ev.seed];
if (e) p.essences.push({ id: e.id, essence: e.essence, source_id: e.article.source_id, source_name: e.article.source_name, url: e.article.url, title: e.article.title });
}
}
}
}
const dominant = (w) => {
const entries = Object.entries(w || {}).sort((a, b) => b[1] - a[1]);
return entries.length ? { topic: entries[0][0], topics: entries.slice(0, 3).map(([topic, weight]) => ({ topic, weight: Math.round(weight * 10) / 10 })) } : { topic: null, topics: [] };
};
for (const bc of Object.values(byCountry)) {
const d = dominant(bc.topic_weights);
bc.dominant_topic = d.topic;
bc.topics = d.topics;
delete bc.topic_weights;
}
for (const p of Object.values(places)) {
const d = dominant(p.topic_weights);
p.dominant_topic = d.topic;
p.topics = d.topics;
delete p.topic_weights;
}
return {
resolution: "country",
countries: Object.values(byCountry).sort((a, b) => (b.events ?? b.count) - (a.events ?? a.count)),
places: Object.values(places).sort((a, b) => b.weight - a.weight),
eventPlacements: eventCount,
};
}