SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
6 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
9.1 KB · 184 lines typescript
Raw Blame History
1import type { EventDTO, SourceKind } from "@dci/core";2import { pg, andAll, eventCols, eventJoins, page, likePattern, type Fragment } from "../lib/sql.js";3import { int, reqStr, type Row } from "../lib/rows.js";4import { asSourceKind, eventDto } from "../lib/dto.js";5import { entityKey, resolveEntityRefs } from "../lib/resolve.js";67export interface EventFilters {8  type?: string[];9  country?: string;10  operator?: string; // slug11  metro?: string; // slug12  project?: string; // slug or id13  entityType?: string;14  entityId?: string;15  minSignificance?: number;16  significance?: "major" | "medium" | "minor";17  confidence?: string[];18  sourceKind?: string[];19  ai?: boolean;20  since?: string;21  until?: string;22  q?: string;23  /** collapse events sharing a cluster id into one row (default true) */24  dedupe?: boolean;25  reviewStatus?: string;26  page?: number;27  perPage?: number;28}2930/** Primary-source order used to pick the representative member of a cluster. */31const SOURCE_RANK = ["operator", "government", "filing", "utility", "cloud_provider", "registry", "secondary", "news", "dataset", "community"];3233export 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.";3435function sourceRankExpr(sql: ReturnType<typeof pg>): Fragment {36  return sql`(case coalesce(e.source_kind, s.kind) ${sql.unsafe(SOURCE_RANK.map((k, i) => `when '${k}' then ${i}`).join(" "))} else 99 end)`;37}3839/** Map rows to EventDTOs with {slug,name} resolved per entity type. */40export async function toEventDtos(rows: Row[]): Promise<EventDTO[]> {41  const refs = await resolveEntityRefs(rows.map((r) => ({ type: String(r.entity_type ?? ""), id: r.entity_id == null ? null : String(r.entity_id) })));42  return rows.map((r) => eventDto(r, refs.get(entityKey(r.entity_type, r.entity_id)) ?? null));43}4445/**46 * Fill evidenceCount / otherSources for a page of events in one query: other members of the same cluster, or — when47 * the event has no cluster id — same-day events with the same operator and event type.48 */49export async function attachOtherSources(dtos: EventDTO[]): Promise<EventDTO[]> {50  if (!dtos.length) return dtos;51  const sql = pg();52  const ids = dtos.map((d) => d.id);53  const rows = await sql<Row[]>`54    select x.id as for_id, e.id, s.name as source_name, coalesce(e.source_kind, s.kind) as source_kind, e.url55    from events x56    join events e on e.id <> x.id and e.review_status <> 'rejected' and (57      (x.cluster_id is not null and e.cluster_id = x.cluster_id)58      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)59    )60    left join sources s on s.id = e.source_id61    where x.id = any(${ids})62    order by x.id, e.detected_at asc`;63  const by = new Map<string, Array<{ sourceName: string; sourceKind: SourceKind; url: string }>>();64  const seen = new Set<string>();65  for (const r of rows) {66    const forId = reqStr(r.for_id);67    const url = reqStr(r.url);68    const k = `${forId}|${url}`;69    if (seen.has(k)) continue;70    seen.add(k);71    const list = by.get(forId) ?? [];72    if (list.length < 10) list.push({ sourceName: reqStr(r.source_name, "unknown source"), sourceKind: asSourceKind(r.source_kind), url });73    by.set(forId, list);74  }75  for (const d of dtos) {76    const others = by.get(d.id) ?? [];77    d.otherSources = others;78    d.evidenceCount = Math.max(d.evidenceCount ?? 1, others.length + 1);79  }80  return dtos;81}8283function conds(f: EventFilters): Fragment[] {84  const sql = pg();85  const c: Fragment[] = [];86  if (f.type?.length) c.push(sql`e.event_type = any(${f.type})`);87  if (f.country) c.push(sql`e.country_iso2 = ${f.country.toUpperCase()}`);88  if (f.operator) c.push(sql`(o.slug = ${f.operator} or o.id = ${f.operator})`);89  if (f.metro) c.push(sql`e.metro_id in (select id from metros where slug = ${f.metro} or id = ${f.metro})`);90  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})))`);91  if (f.entityType) c.push(sql`e.entity_type = ${f.entityType}`);92  if (f.entityId) c.push(sql`e.entity_id = ${f.entityId}`);93  if (f.minSignificance != null) c.push(sql`e.significance >= ${f.minSignificance}`);94  if (f.significance === "major") c.push(sql`e.significance >= 75`);95  if (f.significance === "medium") c.push(sql`e.significance >= 45 and e.significance < 75`);96  if (f.significance === "minor") c.push(sql`e.significance < 45`);97  if (f.confidence?.length) c.push(sql`e.confidence = any(${f.confidence})`);98  if (f.sourceKind?.length) c.push(sql`coalesce(e.source_kind, s.kind) = any(${f.sourceKind})`);99  if (f.ai === true) c.push(sql`e.is_ai`);100  if (f.ai === false) c.push(sql`not e.is_ai`);101  if (f.since) c.push(sql`e.detected_at >= ${f.since}::timestamptz`);102  if (f.until) c.push(sql`e.detected_at <= ${f.until}::timestamptz`);103  if (f.q) { const t = f.q.trim(); if (t) c.push(sql`(e.title ilike ${likePattern(t)} or e.summary ilike ${likePattern(t)})`); }104  if (f.reviewStatus) c.push(sql`e.review_status = ${f.reviewStatus}`);105  else c.push(sql`e.review_status <> 'rejected'`);106  return c;107}108109export async function listEvents(f: EventFilters): Promise<{ items: EventDTO[]; total: number; page: number; perPage: number }> {110  const sql = pg();111  const pg_ = page(f.page, f.perPage, 100, 50);112  const dedupe = f.dedupe !== false;113  const where = andAll(sql, conds(f));114  const rows = dedupe115    ? await sql<Row[]>`116        with ranked as (117          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 rn118          from events e ${eventJoins(sql)}119          where ${where}120        )121        select *, count(*) over() as total from ranked where rn = 1122        order by detected_at desc, id desc123        limit ${pg_.perPage} offset ${pg_.offset}`124    : await sql<Row[]>`125        select ${eventCols(sql)}, count(*) over() as total126        from events e ${eventJoins(sql)}127        where ${where}128        order by e.detected_at desc, e.id desc129        limit ${pg_.perPage} offset ${pg_.offset}`;130  const total = rows.length ? int(rows[0]!.total) : 0;131  const items = await attachOtherSources(await toEventDtos(rows));132  return { items, total, page: pg_.page, perPage: pg_.perPage };133}134135export async function getEvent(id: string): Promise<EventDTO | null> {136  const sql = pg();137  const rows = await sql<Row[]>`select ${eventCols(sql)} from events e ${eventJoins(sql)} where e.id = ${id} limit 1`;138  if (!rows.length) return null;139  const dtos = await attachOtherSources(await toEventDtos(rows));140  return dtos[0] ?? null;141}142143/** Arbitrary condition over `events e` (joined with sources s, operators o, metros em, projects ep), newest first. */144export async function eventsWhere(cond: Fragment, limit = 20): Promise<EventDTO[]> {145  const sql = pg();146  const rows = await sql<Row[]>`147    select ${eventCols(sql)} from events e ${eventJoins(sql)}148    where (${cond}) and e.review_status <> 'rejected'149    order by e.detected_at desc limit ${limit}`;150  return toEventDtos(rows);151}152153/** Events attached to an entity (entity_type + entity_id), newest first. */154export async function eventsForEntity(entityType: string, entityId: string, limit = 50): Promise<EventDTO[]> {155  const sql = pg();156  return eventsWhere(sql`e.entity_type = ${entityType} and e.entity_id = ${entityId}`, limit);157}158159/** Recent events for an operator (by operator_id or entity = operator). */160export async function eventsForOperator(operatorId: string, limit = 20): Promise<EventDTO[]> {161  const sql = pg();162  return eventsWhere(sql`e.operator_id = ${operatorId} or (e.entity_type = 'operator' and e.entity_id = ${operatorId})`, limit);163}164165export async function eventsForCountry(iso2: string, limit = 20): Promise<EventDTO[]> {166  const sql = pg();167  return eventsWhere(sql`e.country_iso2 = ${iso2}`, limit);168}169170export async function eventsForMetro(metroId: string, limit = 20): Promise<EventDTO[]> {171  const sql = pg();172  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);173}174175export async function eventsForProject(projectId: string, limit = 50): Promise<EventDTO[]> {176  const sql = pg();177  return eventsWhere(sql`e.project_id = ${projectId} or (e.entity_type = 'project' and e.entity_id = ${projectId})`, limit);178}179180export async function latestEvents(limit = 20, minSignificance = 0): Promise<EventDTO[]> {181  const sql = pg();182  return eventsWhere(sql`e.significance >= ${minSignificance}`, limit);183}184