spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * News events → news_items (deduplicated by URL) with entity linking (operators, countries, facilities, MW).3 * Items carrying an eventType with significance ≥ 40 also produce an `events` row (entityType news_event).4 */5import { sql, textArray } from "@dci/db";6import { cleanText, classifyAiEvidence, normalizeName, parseAllMw, stableId, type EventType, type NormalizedNewsEvent } from "@dci/core";7import { countryFromTextSafe } from "../connectors/news/locations-lexicon.js";8import { addRef, bump, isoOrNull, knownCountries, type IngestContext, type Tx } from "./common.js";9import { CANONICAL_OPERATORS, findCanonicalOperator } from "./canonical-operators.js";10import { findEventCluster, recordEvent } from "./events.js";11import { geocodeCity } from "./geocode.js";12import { assignMetro } from "./metros.js";13import { lookupOperatorId, resolveOperator } from "./operators.js";1415const EVENT_SIGNIFICANCE: Partial<Record<EventType, number>> = {16 acquisition: 75,17 closure: 70,18 construction_started: 65,19 facility_opened: 65,20 expansion_announced: 60,21 project_announced: 60,22 planning_approved: 55,23 planning_filed: 45,24 power_agreement: 55,25 investment_announced: 60,26 cloud_region_announced: 60,27 cloud_region_launched: 55,28 incident: 70,29 phase_announced: 50,30 land_acquired: 60,31 grid_connection: 65,32 grid_constraint: 70,33 utility_event: 50,34 operator_expansion: 65,35 customer_agreement: 45,36 partnership: 35,37 executive_change: 20,38 project_delayed: 70,39 project_cancelled: 80,40 news: 30,41 page_changed: 10,42};4344function escapeRe(s: string): string {45 return s.replace(/[.*+?^${}()|[\]\\]/g, "\\$&");46}4748/** Copyright guard for third-party publishers: the stored summary never exceeds this many characters. */49export const THIRD_PARTY_SUMMARY_MAX = 400;5051/** Trim on a word boundary, appending an ellipsis (same behaviour as the news parser's `shortSummary`). */52export function shortSummary(s: string | null | undefined, max = THIRD_PARTY_SUMMARY_MAX): string | null {53 const t = cleanText(s);54 if (!t) return null;55 if (t.length <= max) return t;56 const cut = t.slice(0, max - 1);57 return `${cut.slice(0, Math.max(cut.lastIndexOf(" "), max - 40)).trim()}…`;58}5960/** Operator lexicon: curated canonical names/aliases + operators already in the database (loaded once per batch). */61async function operatorLexicon(tx: Tx, ctx: IngestContext): Promise<Array<{ id: string | null; name: string; patterns: RegExp[] }>> {62 if (ctx.caches.operatorLexicon) return ctx.caches.operatorLexicon;63 const entries: Array<{ id: string | null; name: string; patterns: RegExp[] }> = [];64 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"));65 for (const c of CANONICAL_OPERATORS) entries.push({ id: null, name: c.name, patterns: mk([c.name, ...c.aliases]) });66 const rows = await tx.execute(sql`select id, name, aliases from operators order by created_at asc limit 5000`);67 const byNorm = new Map(entries.map((e) => [normalizeName(e.name), e]));68 for (const r of rows) {69 const name = String(r.name);70 const aliases = Array.isArray(r.aliases) ? (r.aliases as string[]) : [];71 const e = byNorm.get(normalizeName(name));72 if (e) e.id = String(r.id);73 else entries.push({ id: String(r.id), name, patterns: mk([name, ...aliases]) });74 }75 ctx.caches.operatorLexicon = entries;76 return entries;77}7879export interface NewsLinks {80 operatorIds: string[];81 operatorNames: string[];82 countryIso2s: string[];83 facilityIds: string[];84 mw: number | null;85}8687export async function linkNews(tx: Tx, ctx: IngestContext, n: NormalizedNewsEvent): Promise<NewsLinks> {88 const text = `${n.title}\n${n.summary ?? ""}`;89 const lex = await operatorLexicon(tx, ctx);90 const operatorNames = new Set<string>();91 for (const e of lex) if (e.patterns.some((re) => re.test(text))) operatorNames.add(e.name);92 const operatorIds: string[] = [];93 // lexicon hits are curated or already-known operators → safe to resolve (creates the curated record if missing)94 for (const name of [...operatorNames].slice(0, 12)) {95 try {96 const op = await resolveOperator(tx, ctx, { name });97 if (!operatorIds.includes(op.id)) operatorIds.push(op.id);98 } catch {99 /* skip unresolvable mention */100 }101 }102 // connector-provided mentions are free text: link only to operators we already know, never create from news103 for (const m of n.mentions?.operators ?? []) {104 const name = m?.trim();105 if (!name) continue;106 const id = findCanonicalOperator(name) ? (await resolveOperator(tx, ctx, { name })).id : await lookupOperatorId(tx, name);107 if (id) {108 if (!operatorIds.includes(id)) operatorIds.push(id);109 operatorNames.add(findCanonicalOperator(name)?.name ?? name);110 }111 }112 const known = await knownCountries(tx, ctx);113 const countries = new Set<string>();114 // the hardened matcher: "North America" is not the US, "Georgia" is a US state, "… Jordan" is a person115 const fromText = countryFromTextSafe(text);116 if (fromText && known.has(fromText)) countries.add(fromText);117 for (const c of n.mentions?.countriesIso2 ?? []) if (c && known.has(c.toUpperCase())) countries.add(c.toUpperCase());118 const facilityIds: string[] = [];119 for (const fname of (n.mentions?.facilities ?? []).slice(0, 10)) {120 const norm = normalizeName(fname);121 if (!norm) continue;122 const rows = operatorIds.length123 ? await tx.execute(sql`select id from facilities where merged_into is null and normalized_name = ${norm} and operator_id in ${operatorIds} limit 2`)124 : await tx.execute(sql`select id from facilities where merged_into is null and normalized_name = ${norm} limit 2`);125 if (rows.length === 1) facilityIds.push(String(rows[0]!.id));126 }127 // `mentions.mw` is already portfolio-filtered by the news parser; the fallback only reads the headline + summary128 const mwList = n.mentions?.mw?.filter((v) => typeof v === "number" && v > 0 && v <= 20_000) ?? [];129 const mw = mwList.length ? Math.max(...mwList) : (parseAllMw(text).filter((v) => v <= 20_000)[0] ?? null);130 return { operatorIds, operatorNames: [...operatorNames], countryIso2s: [...countries], facilityIds, mw };131}132133export async function ingestNews(tx: Tx, ctx: IngestContext, n: NormalizedNewsEvent): Promise<void> {134 const title = cleanText(n.title);135 const url = n.url?.trim();136 if (!title || !url) throw new Error(`news_event ${n.key}: title and url are required`);137 const links = await linkNews(tx, ctx, n);138 const eventType = n.eventType ?? null;139 const significance = eventType ? (EVENT_SIGNIFICANCE[eventType] ?? 30) : links.mw != null && links.mw >= 100 ? 40 : 20;140 const id = stableId("news", url);141 const publishedAt = isoOrNull(n.publishedAt);142 const mentions = { ...(n.mentions ?? {}), linkedOperators: links.operatorNames };143 // third-party publishers (kind = news): title, date, link and a ≤ 400-character summary only — never more text144 const summary = ctx.run.sourceKind === "news" ? shortSummary(n.summary, THIRD_PARTY_SUMMARY_MAX) : cleanText(n.summary);145 const ai = n.isAi ?? ["confirmed", "likely"].includes(classifyAiEvidence(`${title} ${n.summary ?? ""}`).level);146 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)147 values (${id}, ${ctx.run.sourceId}, ${ctx.run.connectorId}, ${url}, ${title.slice(0, 500)}, ${publishedAt}, ${summary}, ${n.pageType ?? "unknown"}, ${eventType}, ${JSON.stringify(mentions)}::jsonb,148 ${textArray(links.operatorIds)}::text[], ${textArray(links.countryIso2s)}::text[], ${textArray(links.facilityIds)}::text[], ${links.mw}, ${significance}, ${n.projectClass ?? null})149 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),150 event_type = coalesce(excluded.event_type, news_items.event_type), mentions = excluded.mentions, operator_ids = excluded.operator_ids, country_iso2s = excluded.country_iso2s,151 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)152 returning (xmax = 0) as inserted`);153 const isNew = Boolean(inserted[0]?.inserted);154 if (isNew) ctx.stats.created++;155 else ctx.stats.unchanged++;156 bump(ctx, "news_event");157 addRef(ctx, "news_item", id);158159 // market linkage: a named city → metro (city-level geocode, never a facility location)160 const city = n.mentions?.cities?.[0] ?? null;161 let metroId: string | null = null;162 if (city) {163 const g = await geocodeCity(tx, { city, countryIso2: links.countryIso2s[0] ?? null });164 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;165 if (metroId && !ctx.run.dryRun) await tx.execute(sql`update news_items set metro_id = ${metroId} where id = ${id}`);166 }167 // grid / utility constraint reports (moratoria, delays, load caps, new transmission, regulation) → power layer168 const gridKind = classifyGridConstraint(`${title}. ${n.summary ?? ""}`);169 let effectiveEventType: EventType | null = eventType;170 if (gridKind && (eventType == null || eventType === "news" || eventType === "power_agreement" || eventType === "utility_event")) effectiveEventType = ctx.run.sourceKind === "utility" && gridKind !== "moratorium" ? "utility_event" : "grid_constraint";171 const effectiveSignificance = effectiveEventType && effectiveEventType !== eventType ? (EVENT_SIGNIFICANCE[effectiveEventType] ?? significance) : significance;172173 if (effectiveEventType && effectiveSignificance >= 40) {174 const day = publishedAt ? publishedAt.slice(0, 10) : ctx.day;175 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 });176 // one event per article: fingerprint on the URL, not on the day, so a re-crawl never duplicates it177 const eventId = await recordEvent(tx, ctx, {178 entityType: "news_event",179 entityId: id,180 eventType: effectiveEventType,181 title,182 summary,183 newValue: { url, mw: links.mw, operators: links.operatorNames, countries: links.countryIso2s, gridKind, projectClass: n.projectClass ?? null },184 significance: cluster.joined && ctx.run.sourceKind === "news" ? Math.max(20, effectiveSignificance - 15) : effectiveSignificance,185 confidence: n.provenance.confidence,186 effectiveDate: publishedAt ? publishedAt.slice(0, 10) : null,187 url,188 countryIso2: links.countryIso2s[0] ?? null,189 operatorId: links.operatorIds[0] ?? null,190 metroId,191 fingerprint: `news:${id}:${effectiveEventType}`,192 isAi: ai,193 clusterId: cluster.clusterId,194 });195 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`);196 if (gridKind && effectiveEventType && ["grid_constraint", "utility_event"].includes(effectiveEventType) && !ctx.run.dryRun) {197 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)198 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})199 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)`);200 }201 }202}203204export const NEWS_EVENT_SIGNIFICANCE = EVENT_SIGNIFICANCE;205206/** Grid / utility constraint vocabulary → grid_constraints.kind (docs/CLAIMS.md, power layer). */207export const GRID_CONSTRAINT_RULES: Array<[RegExp, string]> = [208 [/\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"],209 [/\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"],210 [/\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"],211 [/\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"],212 [/\b(large[- ]load (?:tariff|rule|rules|queue|request|customers?|study|interconnection)|interconnection queue|load (?:interconnection )?queue)\b/i, "large_load_queue"],213 [/\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"],214 [/\b(new substation|substation (?:approved|planned|to be built|construction|expansion|upgrade)|builds? (?:a |two |three )?(?:new )?substations?)\b/i, "new_substation"],215 [/\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"],216];217export function classifyGridConstraint(text: string): string | null {218 for (const [re, kind] of GRID_CONSTRAINT_RULES) if (re.test(text)) return kind;219 return null;220}221