/** * News events → news_items (deduplicated by URL) with entity linking (operators, countries, facilities, MW). * Items carrying an eventType with significance ≥ 40 also produce an `events` row (entityType news_event). */ import { sql, textArray } from "@dci/db"; import { cleanText, classifyAiEvidence, normalizeName, parseAllMw, stableId, type EventType, type NormalizedNewsEvent } from "@dci/core"; import { countryFromTextSafe } from "../connectors/news/locations-lexicon.js"; import { addRef, bump, isoOrNull, knownCountries, type IngestContext, type Tx } from "./common.js"; import { CANONICAL_OPERATORS, findCanonicalOperator } from "./canonical-operators.js"; import { findEventCluster, recordEvent } from "./events.js"; import { geocodeCity } from "./geocode.js"; import { assignMetro } from "./metros.js"; import { lookupOperatorId, resolveOperator } from "./operators.js"; const EVENT_SIGNIFICANCE: Partial> = { acquisition: 75, closure: 70, construction_started: 65, facility_opened: 65, expansion_announced: 60, project_announced: 60, planning_approved: 55, planning_filed: 45, power_agreement: 55, investment_announced: 60, cloud_region_announced: 60, cloud_region_launched: 55, incident: 70, phase_announced: 50, land_acquired: 60, grid_connection: 65, grid_constraint: 70, utility_event: 50, operator_expansion: 65, customer_agreement: 45, partnership: 35, executive_change: 20, project_delayed: 70, project_cancelled: 80, news: 30, page_changed: 10, }; function escapeRe(s: string): string { return s.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"); } /** Copyright guard for third-party publishers: the stored summary never exceeds this many characters. */ export const THIRD_PARTY_SUMMARY_MAX = 400; /** Trim on a word boundary, appending an ellipsis (same behaviour as the news parser's `shortSummary`). */ export function shortSummary(s: string | null | undefined, max = THIRD_PARTY_SUMMARY_MAX): string | null { const t = cleanText(s); if (!t) return null; if (t.length <= max) return t; const cut = t.slice(0, max - 1); return `${cut.slice(0, Math.max(cut.lastIndexOf(" "), max - 40)).trim()}…`; } /** Operator lexicon: curated canonical names/aliases + operators already in the database (loaded once per batch). */ async function operatorLexicon(tx: Tx, ctx: IngestContext): Promise> { if (ctx.caches.operatorLexicon) return ctx.caches.operatorLexicon; const entries: Array<{ id: string | null; name: string; patterns: RegExp[] }> = []; const mk = (terms: string[]) => terms.filter((t) => t && t.length >= 3 && !/^(io|dc|it|the|group|data|center|centre|cloud)$/i.test(t)).map((t) => new RegExp(`(^|[^A-Za-z0-9])${escapeRe(t)}(?![A-Za-z0-9])`, "i")); for (const c of CANONICAL_OPERATORS) entries.push({ id: null, name: c.name, patterns: mk([c.name, ...c.aliases]) }); const rows = await tx.execute(sql`select id, name, aliases from operators order by created_at asc limit 5000`); const byNorm = new Map(entries.map((e) => [normalizeName(e.name), e])); for (const r of rows) { const name = String(r.name); const aliases = Array.isArray(r.aliases) ? (r.aliases as string[]) : []; const e = byNorm.get(normalizeName(name)); if (e) e.id = String(r.id); else entries.push({ id: String(r.id), name, patterns: mk([name, ...aliases]) }); } ctx.caches.operatorLexicon = entries; return entries; } export interface NewsLinks { operatorIds: string[]; operatorNames: string[]; countryIso2s: string[]; facilityIds: string[]; mw: number | null; } export async function linkNews(tx: Tx, ctx: IngestContext, n: NormalizedNewsEvent): Promise { const text = `${n.title}\n${n.summary ?? ""}`; const lex = await operatorLexicon(tx, ctx); const operatorNames = new Set(); for (const e of lex) if (e.patterns.some((re) => re.test(text))) operatorNames.add(e.name); const operatorIds: string[] = []; // lexicon hits are curated or already-known operators → safe to resolve (creates the curated record if missing) for (const name of [...operatorNames].slice(0, 12)) { try { const op = await resolveOperator(tx, ctx, { name }); if (!operatorIds.includes(op.id)) operatorIds.push(op.id); } catch { /* skip unresolvable mention */ } } // connector-provided mentions are free text: link only to operators we already know, never create from news for (const m of n.mentions?.operators ?? []) { const name = m?.trim(); if (!name) continue; const id = findCanonicalOperator(name) ? (await resolveOperator(tx, ctx, { name })).id : await lookupOperatorId(tx, name); if (id) { if (!operatorIds.includes(id)) operatorIds.push(id); operatorNames.add(findCanonicalOperator(name)?.name ?? name); } } const known = await knownCountries(tx, ctx); const countries = new Set(); // the hardened matcher: "North America" is not the US, "Georgia" is a US state, "… Jordan" is a person const fromText = countryFromTextSafe(text); if (fromText && known.has(fromText)) countries.add(fromText); for (const c of n.mentions?.countriesIso2 ?? []) if (c && known.has(c.toUpperCase())) countries.add(c.toUpperCase()); const facilityIds: string[] = []; for (const fname of (n.mentions?.facilities ?? []).slice(0, 10)) { const norm = normalizeName(fname); if (!norm) continue; const rows = operatorIds.length ? await tx.execute(sql`select id from facilities where merged_into is null and normalized_name = ${norm} and operator_id in ${operatorIds} limit 2`) : await tx.execute(sql`select id from facilities where merged_into is null and normalized_name = ${norm} limit 2`); if (rows.length === 1) facilityIds.push(String(rows[0]!.id)); } // `mentions.mw` is already portfolio-filtered by the news parser; the fallback only reads the headline + summary const mwList = n.mentions?.mw?.filter((v) => typeof v === "number" && v > 0 && v <= 20_000) ?? []; const mw = mwList.length ? Math.max(...mwList) : (parseAllMw(text).filter((v) => v <= 20_000)[0] ?? null); return { operatorIds, operatorNames: [...operatorNames], countryIso2s: [...countries], facilityIds, mw }; } export async function ingestNews(tx: Tx, ctx: IngestContext, n: NormalizedNewsEvent): Promise { const title = cleanText(n.title); const url = n.url?.trim(); if (!title || !url) throw new Error(`news_event ${n.key}: title and url are required`); const links = await linkNews(tx, ctx, n); const eventType = n.eventType ?? null; const significance = eventType ? (EVENT_SIGNIFICANCE[eventType] ?? 30) : links.mw != null && links.mw >= 100 ? 40 : 20; const id = stableId("news", url); const publishedAt = isoOrNull(n.publishedAt); const mentions = { ...(n.mentions ?? {}), linkedOperators: links.operatorNames }; // third-party publishers (kind = news): title, date, link and a ≤ 400-character summary only — never more text const summary = ctx.run.sourceKind === "news" ? shortSummary(n.summary, THIRD_PARTY_SUMMARY_MAX) : cleanText(n.summary); const ai = n.isAi ?? ["confirmed", "likely"].includes(classifyAiEvidence(`${title} ${n.summary ?? ""}`).level); const inserted = await tx.execute(sql`insert into news_items (id, source_id, connector_id, url, title, published_at, summary, page_type, event_type, mentions, operator_ids, country_iso2s, facility_ids, mw, significance, project_class) values (${id}, ${ctx.run.sourceId}, ${ctx.run.connectorId}, ${url}, ${title.slice(0, 500)}, ${publishedAt}, ${summary}, ${n.pageType ?? "unknown"}, ${eventType}, ${JSON.stringify(mentions)}::jsonb, ${textArray(links.operatorIds)}::text[], ${textArray(links.countryIso2s)}::text[], ${textArray(links.facilityIds)}::text[], ${links.mw}, ${significance}, ${n.projectClass ?? null}) on conflict (url) do update set title = excluded.title, summary = coalesce(excluded.summary, news_items.summary), published_at = coalesce(excluded.published_at, news_items.published_at), event_type = coalesce(excluded.event_type, news_items.event_type), mentions = excluded.mentions, operator_ids = excluded.operator_ids, country_iso2s = excluded.country_iso2s, facility_ids = excluded.facility_ids, mw = coalesce(excluded.mw, news_items.mw), significance = greatest(excluded.significance, news_items.significance), project_class = coalesce(excluded.project_class, news_items.project_class) returning (xmax = 0) as inserted`); const isNew = Boolean(inserted[0]?.inserted); if (isNew) ctx.stats.created++; else ctx.stats.unchanged++; bump(ctx, "news_event"); addRef(ctx, "news_item", id); // market linkage: a named city → metro (city-level geocode, never a facility location) const city = n.mentions?.cities?.[0] ?? null; let metroId: string | null = null; if (city) { const g = await geocodeCity(tx, { city, countryIso2: links.countryIso2s[0] ?? null }); metroId = g?.metroId ?? (g ? (await assignMetro(tx, { lat: g.lat, lng: g.lng, city, countryIso2: links.countryIso2s[0] ?? null })).metroId : null) ?? (await assignMetro(tx, { city, countryIso2: links.countryIso2s[0] ?? null })).metroId; if (metroId && !ctx.run.dryRun) await tx.execute(sql`update news_items set metro_id = ${metroId} where id = ${id}`); } // grid / utility constraint reports (moratoria, delays, load caps, new transmission, regulation) → power layer const gridKind = classifyGridConstraint(`${title}. ${n.summary ?? ""}`); let effectiveEventType: EventType | null = eventType; if (gridKind && (eventType == null || eventType === "news" || eventType === "power_agreement" || eventType === "utility_event")) effectiveEventType = ctx.run.sourceKind === "utility" && gridKind !== "moratorium" ? "utility_event" : "grid_constraint"; const effectiveSignificance = effectiveEventType && effectiveEventType !== eventType ? (EVENT_SIGNIFICANCE[effectiveEventType] ?? significance) : significance; if (effectiveEventType && effectiveSignificance >= 40) { const day = publishedAt ? publishedAt.slice(0, 10) : ctx.day; const cluster = await findEventCluster(tx, ctx, { eventType: effectiveEventType, title, operatorId: links.operatorIds[0] ?? null, countryIso2: links.countryIso2s[0] ?? null, metroId, city, mw: links.mw, day, sourceKind: ctx.run.sourceKind }); // one event per article: fingerprint on the URL, not on the day, so a re-crawl never duplicates it const eventId = await recordEvent(tx, ctx, { entityType: "news_event", entityId: id, eventType: effectiveEventType, title, summary, newValue: { url, mw: links.mw, operators: links.operatorNames, countries: links.countryIso2s, gridKind, projectClass: n.projectClass ?? null }, significance: cluster.joined && ctx.run.sourceKind === "news" ? Math.max(20, effectiveSignificance - 15) : effectiveSignificance, confidence: n.provenance.confidence, effectiveDate: publishedAt ? publishedAt.slice(0, 10) : null, url, countryIso2: links.countryIso2s[0] ?? null, operatorId: links.operatorIds[0] ?? null, metroId, fingerprint: `news:${id}:${effectiveEventType}`, isAi: ai, clusterId: cluster.clusterId, }); if (eventId && !ctx.run.dryRun && !cluster.joined) await tx.execute(sql`update events set cluster_id = ${cluster.clusterId} where id = ${eventId} and cluster_id is null`); if (gridKind && effectiveEventType && ["grid_constraint", "utility_event"].includes(effectiveEventType) && !ctx.run.dryRun) { await tx.execute(sql`insert into grid_constraints (id, metro_id, country_iso2, kind, title, summary, effective_date, source_id, document_id, url, event_id, confidence) values (${stableId("constraint", url)}, ${metroId}, ${links.countryIso2s[0] ?? null}, ${gridKind}, ${title.slice(0, 300)}, ${summary}, ${publishedAt ? publishedAt.slice(0, 10) : null}, ${ctx.run.sourceId}, ${ctx.doc?.documentId ?? null}, ${url}, ${eventId}, ${n.provenance.confidence}) on conflict (id) do update set metro_id = coalesce(excluded.metro_id, grid_constraints.metro_id), country_iso2 = coalesce(excluded.country_iso2, grid_constraints.country_iso2), kind = excluded.kind, title = excluded.title, summary = excluded.summary, event_id = coalesce(excluded.event_id, grid_constraints.event_id)`); } } } export const NEWS_EVENT_SIGNIFICANCE = EVENT_SIGNIFICANCE; /** Grid / utility constraint vocabulary → grid_constraints.kind (docs/CLAIMS.md, power layer). */ export const GRID_CONSTRAINT_RULES: Array<[RegExp, string]> = [ [/\b(moratorium|moratoria|pause on (?:new )?(?:data cent(?:er|re)|connections?|hookups?)|halt(?:s|ed)? (?:new )?(?:connections?|hookups?|data cent(?:er|re) (?:approvals|permits))|freeze on (?:new )?(?:connections?|data cent(?:er|re))|ban on (?:new )?data cent(?:er|re))\b/i, "moratorium"], [/\b(load cap|capacity cap|cap on (?:new )?(?:load|connections?)|(?:limit|restrict)(?:s|ed|ing)? (?:new )?(?:load|connections?|power (?:supply|allocation)))\b/i, "load_cap"], [/\b(grid (?:constraint|constraints|congestion|bottleneck|shortfall|shortage|capacity (?:limit|shortage|constraint))|(?:power|electricity|capacity) (?:shortage|shortfall|constraints?|crunch)|no (?:grid |spare )?capacity (?:until|before|available)|capacity restrictions?)\b/i, "capacity_restriction"], [/\b((?:grid|interconnection|connection|hookup|energi[sz]ation) (?:delays?|wait(?:ing)? (?:times?|lists?)|queue|backlog)|wait (?:up to |of )?\d+ years? for (?:power|grid|a connection)|(?:delays?|pushed back) (?:grid|power) connections?)\b/i, "grid_delay"], [/\b(large[- ]load (?:tariff|rule|rules|queue|request|customers?|study|interconnection)|interconnection queue|load (?:interconnection )?queue)\b/i, "large_load_queue"], [/\b(new (?:transmission|high-voltage|\d+ ?kv) (?:line|lines|project|corridor)|transmission (?:line|upgrade|expansion|project|build-?out)|\d+ ?kv (?:line|transmission))\b/i, "new_transmission"], [/\b(new substation|substation (?:approved|planned|to be built|construction|expansion|upgrade)|builds? (?:a |two |three )?(?:new )?substations?)\b/i, "new_substation"], [/\b((?:energy|electricity|grid|utility|power) (?:regulation|regulations|rule|rules|tariff|tariffs|law|legislation|bill|policy)|regulator[s']? (?:approv|reject|propos|order)\w*|public (?:utility|service) commission (?:approv|reject|order|rul)\w*|(?:ofgem|ferc|puc|psc|acer|cru|eirgrid|nesO|national grid) (?:approv|reject|propos|order|rul|announc)\w*)\b/i, "regulation"], ]; export function classifyGridConstraint(text: string): string | null { for (const [re, kind] of GRID_CONSTRAINT_RULES) if (re.test(text)) return kind; return null; }