import type { EventDTO, SourceKind } from "@dci/core"; import { pg, andAll, eventCols, eventJoins, page, likePattern, type Fragment } from "../lib/sql.js"; import { int, reqStr, type Row } from "../lib/rows.js"; import { asSourceKind, eventDto } from "../lib/dto.js"; import { entityKey, resolveEntityRefs } from "../lib/resolve.js"; export interface EventFilters { type?: string[]; country?: string; operator?: string; // slug metro?: string; // slug project?: string; // slug or id entityType?: string; entityId?: string; minSignificance?: number; significance?: "major" | "medium" | "minor"; confidence?: string[]; sourceKind?: string[]; ai?: boolean; since?: string; until?: string; q?: string; /** collapse events sharing a cluster id into one row (default true) */ dedupe?: boolean; reviewStatus?: string; page?: number; perPage?: number; } /** Primary-source order used to pick the representative member of a cluster. */ const SOURCE_RANK = ["operator", "government", "filing", "utility", "cloud_provider", "registry", "secondary", "news", "dataset", "community"]; export const EVENTS_DEDUPE_METHODOLOGY = "dedupe=true (default) collapses events that share a cluster_id (documents describing the same announcement) into one row: the member from the most authoritative source kind (operator > government > filing > utility > cloud provider > registry > secondary > news > dataset > community), then the highest significance, then the earliest detection. evidenceCount = number of documents in the cluster; otherSources lists the other members (or same-day events of the same operator and type when no cluster id exists). Significance bands: major ≥ 75, medium 45–74, minor < 45."; function sourceRankExpr(sql: ReturnType): Fragment { return sql`(case coalesce(e.source_kind, s.kind) ${sql.unsafe(SOURCE_RANK.map((k, i) => `when '${k}' then ${i}`).join(" "))} else 99 end)`; } /** Map rows to EventDTOs with {slug,name} resolved per entity type. */ export async function toEventDtos(rows: Row[]): Promise { const refs = await resolveEntityRefs(rows.map((r) => ({ type: String(r.entity_type ?? ""), id: r.entity_id == null ? null : String(r.entity_id) }))); return rows.map((r) => eventDto(r, refs.get(entityKey(r.entity_type, r.entity_id)) ?? null)); } /** * Fill evidenceCount / otherSources for a page of events in one query: other members of the same cluster, or — when * the event has no cluster id — same-day events with the same operator and event type. */ export async function attachOtherSources(dtos: EventDTO[]): Promise { if (!dtos.length) return dtos; const sql = pg(); const ids = dtos.map((d) => d.id); const rows = await sql` select x.id as for_id, e.id, s.name as source_name, coalesce(e.source_kind, s.kind) as source_kind, e.url from events x join events e on e.id <> x.id and e.review_status <> 'rejected' and ( (x.cluster_id is not null and e.cluster_id = x.cluster_id) or (x.cluster_id is null and x.operator_id is not null and e.operator_id = x.operator_id and e.event_type = x.event_type and e.detected_at::date = x.detected_at::date) ) left join sources s on s.id = e.source_id where x.id = any(${ids}) order by x.id, e.detected_at asc`; const by = new Map>(); const seen = new Set(); for (const r of rows) { const forId = reqStr(r.for_id); const url = reqStr(r.url); const k = `${forId}|${url}`; if (seen.has(k)) continue; seen.add(k); const list = by.get(forId) ?? []; if (list.length < 10) list.push({ sourceName: reqStr(r.source_name, "unknown source"), sourceKind: asSourceKind(r.source_kind), url }); by.set(forId, list); } for (const d of dtos) { const others = by.get(d.id) ?? []; d.otherSources = others; d.evidenceCount = Math.max(d.evidenceCount ?? 1, others.length + 1); } return dtos; } function conds(f: EventFilters): Fragment[] { const sql = pg(); const c: Fragment[] = []; if (f.type?.length) c.push(sql`e.event_type = any(${f.type})`); if (f.country) c.push(sql`e.country_iso2 = ${f.country.toUpperCase()}`); if (f.operator) c.push(sql`(o.slug = ${f.operator} or o.id = ${f.operator})`); if (f.metro) c.push(sql`e.metro_id in (select id from metros where slug = ${f.metro} or id = ${f.metro})`); if (f.project) c.push(sql`(e.project_id in (select id from projects where slug = ${f.project} or id = ${f.project}) or (e.entity_type = 'project' and e.entity_id in (select id from projects where slug = ${f.project} or id = ${f.project})))`); if (f.entityType) c.push(sql`e.entity_type = ${f.entityType}`); if (f.entityId) c.push(sql`e.entity_id = ${f.entityId}`); if (f.minSignificance != null) c.push(sql`e.significance >= ${f.minSignificance}`); if (f.significance === "major") c.push(sql`e.significance >= 75`); if (f.significance === "medium") c.push(sql`e.significance >= 45 and e.significance < 75`); if (f.significance === "minor") c.push(sql`e.significance < 45`); if (f.confidence?.length) c.push(sql`e.confidence = any(${f.confidence})`); if (f.sourceKind?.length) c.push(sql`coalesce(e.source_kind, s.kind) = any(${f.sourceKind})`); if (f.ai === true) c.push(sql`e.is_ai`); if (f.ai === false) c.push(sql`not e.is_ai`); if (f.since) c.push(sql`e.detected_at >= ${f.since}::timestamptz`); if (f.until) c.push(sql`e.detected_at <= ${f.until}::timestamptz`); if (f.q) { const t = f.q.trim(); if (t) c.push(sql`(e.title ilike ${likePattern(t)} or e.summary ilike ${likePattern(t)})`); } if (f.reviewStatus) c.push(sql`e.review_status = ${f.reviewStatus}`); else c.push(sql`e.review_status <> 'rejected'`); return c; } export async function listEvents(f: EventFilters): Promise<{ items: EventDTO[]; total: number; page: number; perPage: number }> { const sql = pg(); const pg_ = page(f.page, f.perPage, 100, 50); const dedupe = f.dedupe !== false; const where = andAll(sql, conds(f)); const rows = dedupe ? await sql` with ranked as ( select ${eventCols(sql)}, row_number() over (partition by coalesce(e.cluster_id, e.id) order by ${sourceRankExpr(sql)} asc, e.significance desc, e.detected_at asc, e.id asc) as rn from events e ${eventJoins(sql)} where ${where} ) select *, count(*) over() as total from ranked where rn = 1 order by detected_at desc, id desc limit ${pg_.perPage} offset ${pg_.offset}` : await sql` select ${eventCols(sql)}, count(*) over() as total from events e ${eventJoins(sql)} where ${where} order by e.detected_at desc, e.id desc limit ${pg_.perPage} offset ${pg_.offset}`; const total = rows.length ? int(rows[0]!.total) : 0; const items = await attachOtherSources(await toEventDtos(rows)); return { items, total, page: pg_.page, perPage: pg_.perPage }; } export async function getEvent(id: string): Promise { const sql = pg(); const rows = await sql`select ${eventCols(sql)} from events e ${eventJoins(sql)} where e.id = ${id} limit 1`; if (!rows.length) return null; const dtos = await attachOtherSources(await toEventDtos(rows)); return dtos[0] ?? null; } /** Arbitrary condition over `events e` (joined with sources s, operators o, metros em, projects ep), newest first. */ export async function eventsWhere(cond: Fragment, limit = 20): Promise { const sql = pg(); const rows = await sql` select ${eventCols(sql)} from events e ${eventJoins(sql)} where (${cond}) and e.review_status <> 'rejected' order by e.detected_at desc limit ${limit}`; return toEventDtos(rows); } /** Events attached to an entity (entity_type + entity_id), newest first. */ export async function eventsForEntity(entityType: string, entityId: string, limit = 50): Promise { const sql = pg(); return eventsWhere(sql`e.entity_type = ${entityType} and e.entity_id = ${entityId}`, limit); } /** Recent events for an operator (by operator_id or entity = operator). */ export async function eventsForOperator(operatorId: string, limit = 20): Promise { const sql = pg(); return eventsWhere(sql`e.operator_id = ${operatorId} or (e.entity_type = 'operator' and e.entity_id = ${operatorId})`, limit); } export async function eventsForCountry(iso2: string, limit = 20): Promise { const sql = pg(); return eventsWhere(sql`e.country_iso2 = ${iso2}`, limit); } export async function eventsForMetro(metroId: string, limit = 20): Promise { const sql = pg(); return eventsWhere(sql`e.metro_id = ${metroId} or (e.entity_type = 'facility' and e.entity_id in (select id from facilities where metro_id = ${metroId}))`, limit); } export async function eventsForProject(projectId: string, limit = 50): Promise { const sql = pg(); return eventsWhere(sql`e.project_id = ${projectId} or (e.entity_type = 'project' and e.entity_id = ${projectId})`, limit); } export async function latestEvents(limit = 20, minSignificance = 0): Promise { const sql = pg(); return eventsWhere(sql`e.significance >= ${minSignificance}`, limit); }