SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
2 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%

Ingest: deterministic event clustering, grid-constraint detection + grid_constraints rows, operator_expansion events, news metro linkage; repair script quality-apply.ts; docs/CLAIMS.md; API gets worker URL for the extraction debugger proxy

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Simon-Pierre Boucher committed 20 days ago (Sep 12, 2026) parent c75dcca

16 changed files +1,076 −53

added apps/api/src/lib/csv.ts +10 −0
@@ -0,0 +1,10 @@
1 +/** Minimal RFC 4180 CSV escaper (no dependency). */
2 +export function csvCell(v: unknown): string {
3 + if (v == null) return "";
4 + const s = typeof v === "object" ? JSON.stringify(v) : String(v);
5 + return /[",\r\n]/.test(s) ? `"${s.replace(/"/g, '""')}"` : s;
6 +}
7 +
8 +export function csvLine(values: unknown[]): string {
9 + return `${values.map(csvCell).join(",")}\r\n`;
10 +}
modified apps/api/src/lib/dto.ts +188 −8
@@ -1,7 +1,7 @@
1 1 /** Row → DTO mappers (api-types contract). Inputs are raw rows produced by the fragments in sql.ts. */
2 −import type { CloudRegionSummary, ConfidenceLevel, EventDTO, EventType, EntityType, FacilityStatus, FacilitySummary, FacilityType, GeoPrecision, IxpSummary, ProjectSummary, ProvenanceDTO, ReviewStatus, SourceKind, SourceRef } from "@dci/core";
3 −import { CONFIDENCE_LEVELS, FACILITY_STATUSES, FACILITY_TYPES, GEO_PRECISIONS, SOURCE_KINDS } from "@dci/core";
4 −import { bool, int, iso, num, ref, reqIso, reqStr, str, type Row } from "./rows.js";
2 +import type { AiEvidence, AuthorityTier, ClaimDTO, ClaimScope, ClaimStatus, CloudRegionSummary, ConfidenceLevel, EventDTO, EventType, EntityType, FacilityStatus, FacilitySummary, FacilityType, GeoPrecision, GridConstraintDTO, IxpSummary, ProjectClass, ProjectSummary, ProvenanceDTO, QualityFlagDTO, ReviewStatus, SourceKind, SourceRef } from "@dci/core";
3 +import { AI_EVIDENCE_LEVELS, CLAIM_SCOPES, CLAIM_STATUSES, CONFIDENCE_LEVELS, FACILITY_STATUSES, FACILITY_TYPES, GEO_PRECISIONS, PROJECT_CLASSES, SOURCE_KINDS, CAPACITY_PREDICATE_LABEL, INVESTMENT_PREDICATE_LABEL, significanceBand } from "@dci/core";
4 +import { bool, int, iso, json, num, ref, reqIso, reqStr, str, type Row } from "./rows.js";
5 5
6 6 export function asStatus(v: unknown): FacilityStatus {
7 7 const s = str(v);
@@ -23,9 +23,42 @@ export function asSourceKind(v: unknown): SourceKind {
23 23 const s = str(v);
24 24 return s && (SOURCE_KINDS as readonly string[]).includes(s) ? (s as SourceKind) : "secondary";
25 25 }
26 +export function asAiEvidence(v: unknown): AiEvidence {
27 + const s = str(v);
28 + return s && (AI_EVIDENCE_LEVELS as readonly string[]).includes(s) ? (s as AiEvidence) : "unknown";
29 +}
30 +export function asScope(v: unknown): ClaimScope | null {
31 + const s = str(v);
32 + return s && (CLAIM_SCOPES as readonly string[]).includes(s) ? (s as ClaimScope) : null;
33 +}
34 +export function asClaimStatus(v: unknown): ClaimStatus {
35 + const s = str(v);
36 + return s && (CLAIM_STATUSES as readonly string[]).includes(s) ? (s as ClaimStatus) : "current";
37 +}
38 +export function asProjectClass(v: unknown): ProjectClass | null {
39 + const s = str(v);
40 + return s && (PROJECT_CLASSES as readonly string[]).includes(s) ? (s as ProjectClass) : null;
41 +}
42 +export function asRecordScope(v: unknown): "building" | "facility" | "campus" {
43 + const s = str(v);
44 + return s === "building" || s === "campus" ? s : "facility";
45 +}
46 +export function asRedistribution(v: unknown): SourceRef["redistribution"] {
47 + const s = str(v);
48 + return s === "allowed" || s === "attribution" || s === "restricted" || s === "unknown" ? s : null;
49 +}
50 +export function asSeverity(v: unknown): "info" | "warn" | "critical" {
51 + const s = str(v);
52 + return s === "info" || s === "critical" ? s : "warn";
53 +}
54 +
55 +/** Human label of a claim predicate (capacity / investment vocabularies, else the predicate itself). */
56 +export function predicateLabel(p: string): string {
57 + return (CAPACITY_PREDICATE_LABEL as Record<string, string>)[p] ?? (INVESTMENT_PREDICATE_LABEL as Record<string, string>)[p] ?? p.replace(/_/g, " ");
58 +}
26 59
27 60 export function facilitySummary(r: Row): FacilitySummary {
28 − return {
61 + const out: FacilitySummary = {
29 62 id: reqStr(r.id),
30 63 slug: reqStr(r.slug),
31 64 name: reqStr(r.name),
@@ -43,6 +76,7 @@ export function facilitySummary(r: Row): FacilitySummary {
43 76 itCapacityMw: num(r.it_capacity_mw),
44 77 totalPowerMw: num(r.total_power_mw),
45 78 plannedPowerMw: num(r.planned_power_mw),
79 + mwIsEstimate: bool(r.mw_is_estimate),
46 80 isAi: bool(r.is_ai),
47 81 isHyperscale: bool(r.is_hyperscale),
48 82 confidence: asConfidence(r.confidence),
@@ -52,10 +86,21 @@ export function facilitySummary(r: Row): FacilitySummary {
52 86 carriersCount: num(r.carriers_count),
53 87 ixpCount: num(r.ixp_count),
54 88 networksCount: num(r.networks_count),
89 + recordScope: asRecordScope(r.record_scope),
90 + parentFacility: ref(r, "parent"),
91 + aiEvidence: asAiEvidence(r.ai_evidence),
92 + capacityScope: asScope(r.capacity_scope),
93 + capacitySemantics: str(r.capacity_semantics),
94 + utilityCapacityMw: num(r.utility_capacity_mw),
95 + gridConnectionMw: num(r.grid_connection_mw),
96 + ultimateCampusMw: num(r.ultimate_campus_mw),
97 + sourceCount: int(r.source_count),
55 98 };
99 + return out;
56 100 }
57 101
58 102 export function projectSummary(r: Row): ProjectSummary {
103 + const ai = asAiEvidence(r.ai_evidence);
59 104 return {
60 105 id: reqStr(r.id),
61 106 slug: reqStr(r.slug),
@@ -77,6 +122,23 @@ export function projectSummary(r: Row): ProjectSummary {
77 122 phaseCount: num(r.phase_count),
78 123 confidence: asConfidence(r.confidence),
79 124 lastUpdate: reqIso(r.last_update ?? r.updated_at),
125 + geoPrecision: asPrecision(r.geo_precision),
126 + projectClass: asProjectClass(r.project_class),
127 + evidenceLevel: (() => { const s = str(r.evidence_level); return s === "strong" || s === "weak" || s === "none" ? s : null; })(),
128 + isAi: bool(r.is_ai) || ai === "confirmed" || ai === "likely",
129 + aiEvidence: ai,
130 + capacityScope: asScope(r.capacity_scope),
131 + capacitySemantics: str(r.capacity_semantics),
132 + investmentScope: asScope(r.investment_scope),
133 + investmentSemantics: str(r.investment_semantics),
134 + investmentCurrency: str(r.investment_currency),
135 + investmentOriginal: num(r.investment_original),
136 + developer: ref(r, "dev"),
137 + tenant: ref(r, "ten"),
138 + constructionStartedOn: str(r.construction_started_on),
139 + approvedOn: str(r.approved_on),
140 + permitFiledOn: str(r.permit_filed_on),
141 + openedOn: str(r.opened_on),
80 142 };
81 143 }
82 144
@@ -102,7 +164,7 @@ export function cloudRegionSummary(r: Row): CloudRegionSummary {
102 164 }
103 165
104 166 export function ixpSummary(r: Row): IxpSummary {
105 − return {
167 + const out: IxpSummary = {
106 168 id: reqStr(r.id),
107 169 slug: reqStr(r.slug),
108 170 name: reqStr(r.name),
@@ -113,10 +175,15 @@ export function ixpSummary(r: Row): IxpSummary {
113 175 networkCount: num(r.network_count),
114 176 facilityCount: int(r.facility_count),
115 177 };
178 + if (r.met_id !== undefined) out.metro = ref(r, "met");
179 + if (r.lat !== undefined) out.lat = num(r.lat);
180 + if (r.lng !== undefined) out.lng = num(r.lng);
181 + return out;
116 182 }
117 183
118 184 export function eventDto(r: Row, entity?: { slug: string; name: string } | null): EventDTO {
119 − return {
185 + const sig = int(r.significance, 30);
186 + const out: EventDTO = {
120 187 id: reqStr(r.id),
121 188 entityType: reqStr(r.entity_type, "facility") as EntityType,
122 189 entityId: str(r.entity_id),
@@ -132,12 +199,20 @@ export function eventDto(r: Row, entity?: { slug: string; name: string } | null)
132 199 url: reqStr(r.url),
133 200 title: reqStr(r.title),
134 201 summary: str(r.summary),
135 − significance: int(r.significance, 30),
202 + significance: sig,
136 203 confidence: asConfidence(r.confidence),
137 204 reviewStatus: (str(r.review_status) ?? "auto") as ReviewStatus,
138 205 countryIso2: str(r.country_iso2),
139 206 operator: ref(r, "op"),
207 + metro: ref(r, "met"),
208 + project: ref(r, "prj"),
209 + isAi: bool(r.is_ai),
210 + significanceBand: significanceBand(sig),
211 + evidenceCount: Math.max(1, int(r.evidence_count, 1)),
212 + clusterId: str(r.cluster_id),
213 + documentId: str(r.document_id),
140 214 };
215 + return out;
141 216 }
142 217
143 218 export function provenanceDto(r: Row): ProvenanceDTO {
@@ -154,11 +229,105 @@ export function provenanceDto(r: Row): ProvenanceDTO {
154 229 confidence: asConfidence(r.confidence),
155 230 isEstimate: bool(r.is_estimate),
156 231 method: str(r.method),
232 + isWinner: bool(r.is_winner),
233 + scope: str(r.scope),
234 + runId: str(r.run_id),
235 + documentId: str(r.document_id),
236 + };
237 +}
238 +
239 +export function claimDto(r: Row, isWinner = false): ClaimDTO {
240 + const tier = str(r.authority_tier);
241 + return {
242 + id: reqStr(r.id),
243 + predicate: reqStr(r.predicate),
244 + label: predicateLabel(reqStr(r.predicate)),
245 + value: num(r.value),
246 + valueText: str(r.value_text),
247 + unit: str(r.unit),
248 + scope: asScope(r.scope) ?? "unknown",
249 + scopeReason: str(r.scope_reason),
250 + sourceId: reqStr(r.source_id),
251 + sourceName: reqStr(r.source_name, reqStr(r.source_id)),
252 + sourceKind: asSourceKind(r.source_kind),
253 + url: reqStr(r.url),
254 + documentId: str(r.document_id),
255 + publishedAt: str(r.published_at),
256 + retrievedAt: reqIso(r.retrieved_at),
257 + confidence: asConfidence(r.confidence),
258 + isEstimate: bool(r.is_estimate),
259 + authorityTier: (tier && /^[A-E]$/.test(tier) ? tier : "D") as AuthorityTier,
260 + evidenceText: str(r.evidence_text),
261 + parserName: str(r.parser_name),
262 + parserVersion: str(r.parser_version),
263 + status: asClaimStatus(r.status),
264 + rejectionReason: str(r.rejection_reason),
265 + isWinner,
266 + firstObserved: reqIso(r.first_observed),
267 + lastObserved: reqIso(r.last_observed),
268 + };
269 +}
270 +
271 +export function gridConstraintDto(r: Row): GridConstraintDTO {
272 + return {
273 + id: reqStr(r.id),
274 + kind: reqStr(r.kind, "regulation"),
275 + title: reqStr(r.title),
276 + summary: str(r.summary),
277 + effectiveDate: str(r.effective_date),
278 + url: reqStr(r.url),
279 + sourceName: str(r.source_name),
280 + sourceKind: r.source_kind == null ? null : asSourceKind(r.source_kind),
281 + metro: ref(r, "met"),
282 + countryIso2: str(r.country_iso2),
283 + confidence: asConfidence(r.confidence),
284 + eventId: str(r.event_id),
285 + };
286 +}
287 +
288 +const GRID_KIND_RULES: Array<[RegExp, string]> = [
289 + [/\b(moratorium|pause on|halt(ed|s)?|ban on)\b/i, "moratorium"],
290 + [/\b(load cap|cap on|capped)\b/i, "load_cap"],
291 + [/\b(large[- ]load (queue|tariff|request)|interconnection queue|queue)\b/i, "large_load_queue"],
292 + [/\b(delay(ed|s)?|wait(ing)? (time|list)|backlog|until 20\d\d)\b/i, "grid_delay"],
293 + [/\b(new substation|substation)\b/i, "new_substation"],
294 + [/\b(transmission (line|project|upgrade)|new transmission)\b/i, "new_transmission"],
295 + [/\b(restrict(ion|ed)|curtail|limit(ed|s)?)\b/i, "capacity_restriction"],
296 +];
297 +
298 +/** A grid / utility / power event rendered as a GridConstraintDTO (kind inferred from the title, never invented beyond it). */
299 +export function gridConstraintFromEvent(e: EventDTO): GridConstraintDTO {
300 + const text = `${e.title} ${e.summary ?? ""}`;
301 + const kind = e.eventType === "power_agreement" ? "power_agreement" : GRID_KIND_RULES.find(([re]) => re.test(text))?.[1] ?? (e.eventType === "grid_connection" ? "new_substation" : "regulation");
302 + return { id: e.id, kind, title: e.title, summary: e.summary, effectiveDate: e.effectiveDate, url: e.url, sourceName: e.sourceName, sourceKind: e.sourceKind, metro: e.metro ?? null, countryIso2: e.countryIso2, confidence: e.confidence, eventId: e.id };
303 +}
304 +
305 +export function qualityFlagDto(r: Row, entity: { slug: string; name: string } | null = null): QualityFlagDTO {
306 + const st = str(r.status);
307 + return {
308 + id: reqStr(r.id),
309 + entityType: reqStr(r.entity_type),
310 + entityId: reqStr(r.entity_id),
311 + entity,
312 + claimId: str(r.claim_id),
313 + code: reqStr(r.code),
314 + severity: asSeverity(r.severity),
315 + field: str(r.field),
316 + message: reqStr(r.message),
317 + details: json<Record<string, unknown> | null>(r.details, null),
318 + priority: int(r.priority),
319 + status: st === "resolved" || st === "dismissed" ? st : "open",
320 + resolution: str(r.resolution),
321 + resolvedBy: str(r.resolved_by),
322 + resolvedAt: iso(r.resolved_at),
323 + runId: str(r.run_id),
324 + createdAt: reqIso(r.created_at),
325 + updatedAt: reqIso(r.updated_at),
157 326 };
158 327 }
159 328
160 329 export function sourceRef(r: Row): SourceRef {
161 − return { id: reqStr(r.id), name: reqStr(r.name), domain: reqStr(r.domain), kind: asSourceKind(r.kind), url: str(r.url), license: str(r.license), attribution: str(r.attribution) };
330 + return { id: reqStr(r.id), name: reqStr(r.name), domain: reqStr(r.domain), kind: asSourceKind(r.kind), url: str(r.url), license: str(r.license), attribution: str(r.attribution), redistribution: asRedistribution(r.redistribution), attributionRequired: r.attribution_required == null ? null : bool(r.attribution_required) };
162 331 }
163 332
164 333 export function statusBreakdown(rows: Row[], statusKey = "status", countKey = "n"): Partial<Record<FacilityStatus, number>> {
@@ -189,3 +358,14 @@ export function growthSeries(rows: Row[]): Array<{ year: number; facilities: num
189 358 }
190 359 return out;
191 360 }
361 +
362 +export function share(part: number, total: number, digits = 3): number {
363 + if (!total) return 0;
364 + const f = 10 ** digits;
365 + return Math.round((part / total) * f) / f;
366 +}
367 +
368 +export function round2(v: number | null | undefined): number | null {
369 + if (v == null || !Number.isFinite(v)) return null;
370 + return Math.round(v * 100) / 100;
371 +}
added apps/api/src/lib/nearby.ts +99 −0
@@ -0,0 +1,99 @@
1 +/**
2 + * Shared "what is around this point" helper (facilities, projects, IXPs, cloud regions, metros) — haversine on a bbox
3 + * pre-filter, no PostGIS. Layers without a connector (cable landing stations, substations, power plants) are returned
4 + * EMPTY with a note: nothing is ever faked.
5 + */
6 +import type { NearbyInfrastructure } from "@dci/core";
7 +import { pg, facilityJoins, facilitySummaryCols, projectJoins, projectSummaryCols, cloudRegionCols, haversineFor, withinBbox, projectLive, type Fragment } from "./sql.js";
8 +import { num, str, reqStr, round, type Row } from "./rows.js";
9 +import { cloudRegionSummary, facilitySummary, ixpSummary, projectSummary } from "./dto.js";
10 +
11 +export const NEARBY_TYPES = ["facilities", "projects", "ixps", "cloud_regions", "metros"] as const;
12 +export type NearbyType = (typeof NEARBY_TYPES)[number];
13 +export const NEARBY_MAX_RADIUS_KM = 200;
14 +export const NEARBY_NOTE = "Distances are great-circle (haversine) from the centre; city-level coordinates are approximate (see geoPrecision). IXPs without published coordinates are placed at their metro reference point (lat/lng from the metro, flagged in `note`). Cable landing stations, substations and power plants are empty until a licensed connector for those layers exists — they are never inferred.";
15 +
16 +export interface NearbyOptions {
17 + lat: number;
18 + lng: number;
19 + radiusKm?: number;
20 + types?: NearbyType[];
21 + /** exclude the entity the query is made for */
22 + excludeFacilityId?: string | null;
23 + excludeProjectId?: string | null;
24 + limitPerType?: number;
25 +}
26 +
27 +export async function nearby(o: NearbyOptions): Promise<NearbyInfrastructure> {
28 + const sql = pg();
29 + const radiusKm = Math.min(NEARBY_MAX_RADIUS_KM, Math.max(0.1, o.radiusKm ?? 25));
30 + const types = new Set<NearbyType>(o.types?.length ? o.types : NEARBY_TYPES);
31 + const limit = Math.min(500, Math.max(1, o.limitPerType ?? 50));
32 + const { lat, lng } = o;
33 + const fDist = haversineFor(sql, lat, lng, sql`f.lat`, sql`f.lng`);
34 + const pDist = haversineFor(sql, lat, lng, sql`p.lat`, sql`p.lng`);
35 + const rDist = haversineFor(sql, lat, lng, sql`r.lat`, sql`r.lng`);
36 + const mDist = haversineFor(sql, lat, lng, sql`m.lat`, sql`m.lng`);
37 + const xLat: Fragment = sql`coalesce(xm.lat, xm2.lat)`;
38 + const xLng: Fragment = sql`coalesce(xm.lng, xm2.lng)`;
39 + const xDist = haversineFor(sql, lat, lng, xLat, xLng);
40 + const none = Promise.resolve([] as Row[]);
41 + const [facs, prjs, ixps, crs, mets] = await Promise.all([
42 + types.has("facilities")
43 + ? sql<Row[]>`select ${facilitySummaryCols(sql)}, ${fDist} as distance_km from facilities f ${facilityJoins(sql)}
44 + where f.merged_into is null and f.lat is not null and f.lng is not null ${o.excludeFacilityId ? sql`and f.id <> ${o.excludeFacilityId}` : sql``}
45 + and ${withinBbox(sql, lat, lng, radiusKm, sql`f.lat`, sql`f.lng`)} and ${fDist} <= ${radiusKm}
46 + order by ${fDist} asc limit ${limit}`
47 + : none,
48 + types.has("projects")
49 + ? sql<Row[]>`select ${projectSummaryCols(sql)}, ${pDist} as distance_km from projects p ${projectJoins(sql)}
50 + where ${projectLive(sql)} and p.lat is not null and p.lng is not null ${o.excludeProjectId ? sql`and p.id <> ${o.excludeProjectId}` : sql``}
51 + and ${withinBbox(sql, lat, lng, radiusKm, sql`p.lat`, sql`p.lng`)} and ${pDist} <= ${radiusKm}
52 + order by ${pDist} asc limit ${limit}`
53 + : none,
54 + types.has("ixps")
55 + ? sql<Row[]>`select x.id, x.slug, x.name, x.name_long, x.city, x.country_iso2, x.website, x.network_count,
56 + (select count(*)::int from facility_ixps fx where fx.ixp_id = x.id) as facility_count,
57 + ${xLat} as lat, ${xLng} as lng, ${xDist} as distance_km, (xm.id is not null) as via_metro,
58 + coalesce(xm.id, xm2.id) as met_id, coalesce(xm.slug, xm2.slug) as met_slug, coalesce(xm.name, xm2.name) as met_name
59 + from ixps x
60 + left join metros xm on xm.id = x.metro_id
61 + left join lateral (select m2.id, m2.slug, m2.name, m2.lat, m2.lng from metros m2 where x.metro_id is null and x.country_iso2 = m2.country_iso2 and x.city is not null and (lower(m2.name) = lower(x.city) or exists (select 1 from unnest(m2.aliases) a where lower(a) = lower(x.city))) limit 1) xm2 on true
62 + where ${xLat} is not null and ${withinBbox(sql, lat, lng, radiusKm, xLat, xLng)} and ${xDist} <= ${radiusKm}
63 + order by ${xDist} asc limit ${limit}`
64 + : none,
65 + types.has("cloud_regions")
66 + ? sql<Row[]>`select ${cloudRegionCols(sql)}, ${rDist} as distance_km from cloud_regions r join operators pr on pr.id = r.provider_id
67 + where r.lat is not null and r.lng is not null and ${withinBbox(sql, lat, lng, radiusKm, sql`r.lat`, sql`r.lng`)} and ${rDist} <= ${radiusKm}
68 + order by ${rDist} asc limit ${limit}`
69 + : none,
70 + types.has("metros")
71 + ? sql<Row[]>`select m.id, m.slug, m.name, m.country_iso2, ${mDist} as distance_km from metros m
72 + where ${withinBbox(sql, lat, lng, radiusKm, sql`m.lat`, sql`m.lng`)} and ${mDist} <= ${radiusKm}
73 + order by ${mDist} asc limit ${limit}`
74 + : none,
75 + ]);
76 + const d = (r: Row) => round(num(r.distance_km), 2) ?? 0;
77 + const ixpViaMetro = ixps.filter((r) => r.via_metro === true).length;
78 + return {
79 + center: { lat, lng },
80 + radiusKm,
81 + facilities: facs.map((r) => ({ ...facilitySummary(r), distanceKm: d(r) })),
82 + projects: prjs.map((r) => ({ ...projectSummary(r), distanceKm: d(r) })),
83 + ixps: ixps.map((r) => ({ ...ixpSummary(r), lat: num(r.lat), lng: num(r.lng), distanceKm: d(r) })),
84 + cloudRegions: crs.map((r) => ({ ...cloudRegionSummary(r), distanceKm: d(r) })),
85 + metros: mets.map((r) => ({ id: reqStr(r.id), slug: reqStr(r.slug), name: reqStr(r.name), countryIso2: reqStr(r.country_iso2), distanceKm: d(r) })),
86 + landingStations: [],
87 + substations: [],
88 + powerPlants: [],
89 + note: `${NEARBY_NOTE}${ixpViaMetro ? ` ${ixpViaMetro} IXP${ixpViaMetro > 1 ? "s" : ""} placed at metro coordinates (no published IXP location).` : ""}`,
90 + };
91 +}
92 +
93 +export function parseNearbyTypes(v: string | undefined): NearbyType[] | undefined {
94 + if (!v) return undefined;
95 + const out = v.split(",").map((s) => s.trim().toLowerCase().replace(/-/g, "_")).filter((s): s is NearbyType => (NEARBY_TYPES as readonly string[]).includes(s));
96 + return out.length ? out : undefined;
97 +}
98 +
99 +export { str as _str };
added apps/api/src/lib/quality.ts +148 −0
@@ -0,0 +1,148 @@
1 +/**
2 + * Claim-first read helpers shared by detail payloads and the /history, /claims, /provenance endpoints:
3 + * claims with their winner flag, dated history points (observed provenance + claims + change events) and the
4 + * DataQualitySummary (sources, completeness, claims by status, open quality flags, pending duplicate matches).
5 + */
6 +import type { ClaimDTO, DataQualitySummary, EntityHistory, EventDTO, HistoryPoint, ProvenanceDTO } from "@dci/core";
7 +import { CAPACITY_COLUMN } from "@dci/core";
8 +import { pg, claimCols, eventCols, eventJoins } from "./sql.js";
9 +import { bool, int, iso, num, reqStr, str, type Row } from "./rows.js";
10 +import { asSeverity, claimDto, eventDto, provenanceDto } from "./dto.js";
11 +
12 +const MW_FIELDS = new Set(["itCapacityMw", "totalPowerMw", "plannedPowerMw", "utilityCapacityMw", "gridConnectionMw", "ultimateCampusMw", "plannedMw"]);
13 +const MW_EVENTS = new Set(["capacity_changed", "planned_capacity_changed"]);
14 +
15 +/** Claims about a subject, winner flag derived from the displayed column value (same predicate column, same value). */
16 +export async function claimsFor(subjectType: string, subjectId: string, opts: { status?: string[]; predicate?: string; limit?: number } = {}): Promise<ClaimDTO[]> {
17 + const sql = pg();
18 + const rows = await sql<Row[]>`
19 + select ${claimCols(sql)} from claims k left join sources s on s.id = k.source_id
20 + where k.subject_type = ${subjectType} and k.subject_id = ${subjectId}
21 + ${opts.status?.length ? sql`and k.status = any(${opts.status})` : sql``}
22 + ${opts.predicate ? sql`and k.predicate = ${opts.predicate}` : sql``}
23 + order by k.status = 'current' desc, k.authority_tier asc, k.last_observed desc limit ${opts.limit ?? 500}`;
24 + const winners = await winnerValues(subjectType, subjectId);
25 + return rows.map((r) => {
26 + const pred = reqStr(r.predicate);
27 + const col = subjectType === "project" ? (pred === "planned_power_mw" || pred === "phase_mw" || pred === "it_capacity_mw" ? "plannedMw" : /usd$/.test(pred) ? "investmentUsd" : null) : (CAPACITY_COLUMN as Record<string, string | null>)[pred] ?? null;
28 + const v = num(r.value);
29 + const isWinner = str(r.status) === "current" && col != null && v != null && winners.get(col) != null && Math.abs(winners.get(col)! - v) < 1e-9;
30 + return claimDto(r, isWinner);
31 + });
32 +}
33 +
34 +async function winnerValues(subjectType: string, subjectId: string): Promise<Map<string, number | null>> {
35 + const sql = pg();
36 + const out = new Map<string, number | null>();
37 + if (subjectType === "facility" || subjectType === "campus") {
38 + const r = (await sql<Row[]>`select it_capacity_mw, total_power_mw, planned_power_mw, utility_capacity_mw, grid_connection_mw, ultimate_campus_mw from facilities where id = ${subjectId}`)[0];
39 + if (r) { out.set("itCapacityMw", num(r.it_capacity_mw)); out.set("totalPowerMw", num(r.total_power_mw)); out.set("plannedPowerMw", num(r.planned_power_mw)); out.set("utilityCapacityMw", num(r.utility_capacity_mw)); out.set("gridConnectionMw", num(r.grid_connection_mw)); out.set("ultimateCampusMw", num(r.ultimate_campus_mw)); }
40 + } else if (subjectType === "project") {
41 + const r = (await sql<Row[]>`select planned_mw, investment_usd from projects where id = ${subjectId}`)[0];
42 + if (r) { out.set("plannedMw", num(r.planned_mw)); out.set("investmentUsd", num(r.investment_usd)); }
43 + }
44 + return out;
45 +}
46 +
47 +/** All provenance rows (current and superseded) for an entity — the /provenance endpoint. */
48 +export async function provenanceAll(entityType: string, entityId: string, limit = 1000): Promise<ProvenanceDTO[]> {
49 + const sql = pg();
50 + const rows = await sql<Row[]>`
51 + select p.field, p.value, p.source_id, s.name as source_name, s.kind as source_kind, p.url, p.first_observed, p.last_observed, p.retrieved_at, p.confidence, p.is_estimate, p.method, p.is_winner, p.scope, p.run_id, p.document_id, p.is_current
52 + from provenance p left join sources s on s.id = p.source_id
53 + where p.entity_type = ${entityType} and p.entity_id = ${entityId}
54 + order by p.is_current desc, p.field, p.last_observed desc limit ${limit}`;
55 + return rows.map(provenanceDto);
56 +}
57 +
58 +function pointFromProvenance(p: ProvenanceDTO): HistoryPoint {
59 + return { date: p.firstObserved, field: p.field, predicate: null, value: p.value, sourceId: p.sourceId, sourceName: p.sourceName, sourceKind: p.sourceKind, url: p.url, kind: "observed" };
60 +}
61 +function pointFromClaim(c: ClaimDTO, field: string): HistoryPoint {
62 + return { date: c.publishedAt && /^\d{4}/.test(c.publishedAt) ? c.publishedAt : c.firstObserved, field, predicate: c.predicate, value: c.value ?? c.valueText, sourceId: c.sourceId, sourceName: c.sourceName, sourceKind: c.sourceKind, url: c.url, claimId: c.id, kind: "claim" };
63 +}
64 +function pointFromEvent(e: EventDTO, field: string): HistoryPoint {
65 + return { date: e.effectiveDate && /^\d{4}-\d{2}-\d{2}/.test(e.effectiveDate) ? e.effectiveDate : e.detectedAt, field, predicate: null, value: e.newValue, oldValue: e.oldValue, sourceId: e.sourceId, sourceName: e.sourceName, sourceKind: e.sourceKind, url: e.url, eventId: e.id, kind: "changed" };
66 +}
67 +
68 +function fieldOfEvent(e: EventDTO): string {
69 + switch (e.eventType) {
70 + case "planned_capacity_changed": return "plannedPowerMw";
71 + case "capacity_changed": return "itCapacityMw";
72 + case "status_changed": case "project_status_changed": return "status";
73 + case "opening_date_changed": return "openedOn";
74 + case "operator_changed": return "operatorId";
75 + case "owner_changed": return "ownerId";
76 + default: return e.eventType;
77 + }
78 +}
79 +
80 +const byDate = (a: HistoryPoint, b: HistoryPoint) => (a.date < b.date ? -1 : a.date > b.date ? 1 : 0);
81 +
82 +/** MW history for a facility / project: winner + non-winner observations, claims and capacity_changed events, dated. */
83 +export function capacityHistory(provenance: ProvenanceDTO[], claims: ClaimDTO[], events: EventDTO[]): HistoryPoint[] {
84 + const out: HistoryPoint[] = [];
85 + for (const p of provenance) if (MW_FIELDS.has(p.field)) out.push(pointFromProvenance(p));
86 + for (const c of claims) if (/mw$/.test(c.predicate) && c.status !== "rejected") out.push(pointFromClaim(c, (CAPACITY_COLUMN as Record<string, string | null>)[c.predicate] ?? c.predicate));
87 + for (const e of events) if (MW_EVENTS.has(e.eventType)) out.push(pointFromEvent(e, fieldOfEvent(e)));
88 + return out.sort(byDate);
89 +}
90 +
91 +/** Full field history (EntityHistory): every field's dated observations + claims, and the change list. */
92 +export function entityHistory(entityType: EntityHistory["entityType"], entityId: string, provenance: ProvenanceDTO[], claims: ClaimDTO[], events: EventDTO[]): EntityHistory {
93 + const fields: Record<string, HistoryPoint[]> = {};
94 + const push = (f: string, p: HistoryPoint) => { (fields[f] ??= []).push(p); };
95 + for (const p of provenance) push(p.field, pointFromProvenance(p));
96 + for (const c of claims) if (c.status !== "rejected") push((CAPACITY_COLUMN as Record<string, string | null>)[c.predicate] ?? (/usd$/.test(c.predicate) ? "investmentUsd" : c.predicate), pointFromClaim(c, (CAPACITY_COLUMN as Record<string, string | null>)[c.predicate] ?? c.predicate));
97 + const changes: HistoryPoint[] = [];
98 + for (const e of events) {
99 + if (e.oldValue == null && e.newValue == null) continue;
100 + if (!/changed|opened|started|approved|filed|cancelled|delayed|acquisition|closure/.test(e.eventType)) continue;
101 + const f = fieldOfEvent(e);
102 + const hp = pointFromEvent(e, f);
103 + changes.push(hp);
104 + push(f, hp);
105 + }
106 + for (const k of Object.keys(fields)) fields[k]!.sort(byDate);
107 + return { entityType, entityId, fields, changes: changes.sort(byDate) };
108 +}
109 +
110 +/** Events for a subject (entity_type + id, plus project_id for projects), newest first. */
111 +export async function eventsForSubject(subjectType: string, subjectId: string, limit = 200): Promise<EventDTO[]> {
112 + const sql = pg();
113 + const rows = await sql<Row[]>`select ${eventCols(sql)} from events e ${eventJoins(sql)}
114 + where ((e.entity_type = ${subjectType} and e.entity_id = ${subjectId}) ${subjectType === "project" ? sql`or e.project_id = ${subjectId}` : sql``}) and e.review_status <> 'rejected'
115 + order by e.detected_at desc limit ${limit}`;
116 + return rows.map((r) => eventDto(r, null));
117 +}
118 +
119 +/** DataQualitySummary from provenance / claims / quality_flags / entity_matches. */
120 +export async function dataQualityFor(entityType: string, entityId: string, completeness: number, lastVerified: string | null): Promise<DataQualitySummary> {
121 + const sql = pg();
122 + const [prov, cl, flags, dup] = await Promise.all([
123 + sql<Row[]>`select count(distinct p.source_id)::int as sources, count(distinct p.source_id) filter (where s.kind in ('operator','government','filing','utility','cloud_provider','registry'))::int as primary_sources, count(distinct p.field)::int as fields
124 + from provenance p left join sources s on s.id = p.source_id where p.entity_type = ${entityType} and p.entity_id = ${entityId} and p.is_current`,
125 + sql<Row[]>`select count(*)::int as total, count(*) filter (where status = 'current')::int as current, count(*) filter (where status = 'unscoped')::int as unscoped, count(*) filter (where status = 'review')::int as review from claims where subject_type = ${entityType} and subject_id = ${entityId}`,
126 + sql<Row[]>`select code, severity, message, field from quality_flags where entity_type = ${entityType} and entity_id = ${entityId} and status = 'open' order by priority desc, created_at desc limit 50`,
127 + entityType === "facility"
128 + ? sql<Row[]>`select 1 from entity_matches where status = 'pending' and (matched_facility_id = ${entityId} or candidate->>'createdFacilityId' = ${entityId}) limit 1`
129 + : Promise.resolve([] as Row[]),
130 + ]);
131 + const p = prov[0] ?? {};
132 + const c = cl[0] ?? {};
133 + return {
134 + sourceCount: int(p.sources),
135 + primarySourceCount: int(p.primary_sources),
136 + lastVerified: lastVerified ? iso(lastVerified) : null,
137 + completeness,
138 + fieldsWithProvenance: int(p.fields),
139 + claimsTotal: int(c.total),
140 + claimsCurrent: int(c.current),
141 + claimsUnscoped: int(c.unscoped),
142 + claimsInReview: int(c.review),
143 + openFlags: flags.map((f) => ({ code: reqStr(f.code), severity: asSeverity(f.severity), message: reqStr(f.message), field: str(f.field) })),
144 + pendingDuplicate: dup.length > 0,
145 + };
146 +}
147 +
148 +export { bool as _bool };
modified apps/api/src/lib/source-history.ts +23 −1
@@ -66,7 +66,7 @@ export async function sourceRefsFor(ids: Iterable<string>): Promise<SourceRef[]>
66 66 const list = [...new Set([...ids].filter(Boolean))];
67 67 if (!list.length) return [];
68 68 const sql = pg();
69 − const rows = await sql<Row[]>`select id, name, domain, kind, url, license, attribution from sources where id = any(${list}) order by name`;
69 + const rows = await sql<Row[]>`select id, name, domain, kind, url, license, attribution, redistribution, attribution_required from sources where id = any(${list}) order by name`;
70 70 return rows.map(sourceRef);
71 71 }
72 72
@@ -75,3 +75,25 @@ export function sourceIdsOf(...lists: Array<Array<{ sourceId: string }>>): Set<s
75 75 for (const l of lists) for (const x of l) if (x.sourceId) s.add(x.sourceId);
76 76 return s;
77 77 }
78 +
79 +/** Distinct sources behind a set of entities (current provenance rows), capped — for LIST envelopes. */
80 +export async function sourcesForEntities(entityType: string, ids: string[], cap = 30): Promise<SourceRef[]> {
81 + const list = [...new Set(ids.filter(Boolean))];
82 + if (!list.length) return [];
83 + const sql = pg();
84 + const rows = await sql<Row[]>`
85 + select s.id, s.name, s.domain, s.kind, s.url, s.license, s.attribution, s.redistribution, s.attribution_required, count(*) as n
86 + from provenance p join sources s on s.id = p.source_id
87 + where p.entity_type = ${entityType} and p.entity_id = any(${list}) and p.is_current
88 + group by s.id order by n desc, s.name limit ${cap}`;
89 + return rows.map(sourceRef);
90 +}
91 +
92 +/** Distinct sources of a set of source ids (events / map points), capped. */
93 +export async function sourcesForIds(ids: Iterable<string>, cap = 30): Promise<SourceRef[]> {
94 + const list = [...new Set([...ids].filter(Boolean))].slice(0, 500);
95 + if (!list.length) return [];
96 + const sql = pg();
97 + const rows = await sql<Row[]>`select id, name, domain, kind, url, license, attribution, redistribution, attribution_required from sources where id = any(${list}) order by name limit ${cap}`;
98 + return rows.map(sourceRef);
99 +}
modified apps/api/src/lib/sql.ts +95 −7
@@ -44,13 +44,57 @@ export function orAll(sql: Sql, conds: Fragment[]): Fragment {
44 44 return conds.slice(1).reduce<Fragment>((acc, c) => sql`${acc} or ${c}`, conds[0]!);
45 45 }
46 46
47 +/** Live project rows: false positives hidden by review and merged duplicates are never listed, counted or summed. */
48 +export function projectLive(sql: Sql, alias = "p"): Fragment {
49 + const a = sql(alias);
50 + return sql`(${a}.hidden = false and ${a}.merged_into is null)`;
51 +}
52 +
53 +/** AI evidence levels that qualify a facility / project as AI infrastructure (keyword mentions alone never do). */
54 +export const AI_LEVELS = ["confirmed", "likely"];
55 +
56 +/**
57 + * Containment-aware facility view (mirrors apps/worker/src/rankings.ts FACILITY_VIEW): a campus row whose buildings
58 + * publish their own MW contributes nothing, a building without a figure under a campus with one is "covered", and a
59 + * campus with building rows is not counted as a facility (its buildings are).
60 + */
61 +export function facilityView(sql: Sql): Fragment {
62 + return sql`(
63 + select f.*,
64 + exists (select 1 from facilities ch where ch.parent_facility_id = f.id and ch.merged_into is null) as has_children,
65 + exists (select 1 from facilities ch where ch.parent_facility_id = f.id and ch.merged_into is null and coalesce(ch.it_capacity_mw, ch.total_power_mw, ch.planned_power_mw) is not null) as children_have_mw,
66 + (f.parent_facility_id is not null and exists (select 1 from facilities pp where pp.id = f.parent_facility_id and pp.merged_into is null and coalesce(pp.it_capacity_mw, pp.total_power_mw, pp.planned_power_mw) is not null)
67 + and coalesce(f.it_capacity_mw, f.total_power_mw, f.planned_power_mw) is null) as covered_by_parent
68 + from facilities f where f.merged_into is null
69 + )`;
70 +}
71 +/** Known (operational) MW of a facilityView row `f`, null when its buildings carry the figure. */
72 +export function knownMwAgg(sql: Sql): Fragment {
73 + return sql`(case when f.children_have_mw then null else coalesce(f.it_capacity_mw, f.total_power_mw) end)`;
74 +}
75 +/** Pipeline MW of a facilityView row `f` (planned figure first). */
76 +export function pipelineMwAgg(sql: Sql): Fragment {
77 + return sql`(case when f.children_have_mw then null else coalesce(f.planned_power_mw, f.it_capacity_mw, f.total_power_mw) end)`;
78 +}
79 +/** Rows counted as one facility (campus rows with building rows are not). */
80 +export function countedAgg(sql: Sql): Fragment {
81 + return sql`(not f.has_children)`;
82 +}
83 +/** A facilityView row has a known MW figure (own figure or covered by its campus). */
84 +export function hasMwAgg(sql: Sql): Fragment {
85 + return sql`(coalesce(f.it_capacity_mw, f.total_power_mw, f.planned_power_mw) is not null or f.covered_by_parent)`;
86 +}
87 +
88 +export const METHODOLOGY_CONTAINMENT = "Aggregates are containment-aware: a campus and its buildings are never both counted or summed (buildings with their own figures win; otherwise the campus figure stands for them). Published figures only — no extrapolation. Coverage = share of counted facilities with any MW figure.";
89 +
47 90 /** Columns for a FacilitySummary from `facilities f` joined with operators o, metros m, countries c. */
48 91 export function facilitySummaryCols(sql: Sql): Fragment {
49 92 return sql`
50 93 f.id, f.slug, f.name, f.city, f.region_name, f.country_iso2, c.name as country_name,
51 94 f.lat, f.lng, f.geo_precision, f.status, f.facility_type, f.it_capacity_mw, f.total_power_mw, f.planned_power_mw,
52 95 f.is_ai, f.is_hyperscale, f.confidence, f.completeness, f.opened_on, f.last_verified, f.carriers_count, f.ixp_count, f.networks_count,
53 − f.mw_is_estimate,
96 + f.mw_is_estimate, f.record_scope, f.ai_evidence, f.capacity_scope, f.capacity_semantics, f.utility_capacity_mw, f.grid_connection_mw, f.ultimate_campus_mw, f.source_count,
97 + f.parent_facility_id as parent_id, (select x.slug from facilities x where x.id = f.parent_facility_id) as parent_slug, (select x.name from facilities x where x.id = f.parent_facility_id) as parent_name,
54 98 o.id as op_id, o.slug as op_slug, o.name as op_name,
55 99 m.id as met_id, m.slug as met_slug, m.name as met_name`;
56 100 }
@@ -65,8 +109,12 @@ export function facilityJoins(sql: Sql): Fragment {
65 109 /** Columns for a ProjectSummary from `projects p` joined with operators o, facilities pf, metros m. */
66 110 export function projectSummaryCols(sql: Sql): Fragment {
67 111 return sql`
68 − p.id, p.slug, p.name, p.city, p.region_name, p.country_iso2, p.lat, p.lng, p.status, p.announced_on, p.expected_opening,
112 + p.id, p.slug, p.name, p.city, p.region_name, p.country_iso2, p.lat, p.lng, p.geo_precision, p.status, p.announced_on, p.expected_opening,
69 113 p.planned_mw, p.investment_usd, p.acreage, p.phase_count, p.confidence, p.last_update, p.is_ai, p.description, p.source_url, p.updated_at,
114 + p.project_class, p.evidence_level, p.ai_evidence, p.capacity_scope, p.capacity_semantics, p.investment_scope, p.investment_semantics, p.investment_currency, p.investment_original,
115 + p.construction_started_on, p.approved_on, p.permit_filed_on, p.opened_on, p.hidden, p.merged_into, p.campus_id,
116 + p.developer_id as dev_id, (select x.slug from operators x where x.id = p.developer_id) as dev_slug, (select x.name from operators x where x.id = p.developer_id) as dev_name,
117 + p.tenant_id as ten_id, (select x.slug from operators x where x.id = p.tenant_id) as ten_slug, (select x.name from operators x where x.id = p.tenant_id) as ten_name,
70 118 o.id as op_id, o.slug as op_slug, o.name as op_name,
71 119 pf.id as fac_id, pf.slug as fac_slug, pf.name as fac_name,
72 120 m.id as met_id, m.slug as met_slug, m.name as met_name`;
@@ -90,15 +138,44 @@ export function cloudRegionCols(sql: Sql): Fragment {
90 138 export function eventCols(sql: Sql): Fragment {
91 139 return sql`
92 140 e.id, e.entity_type, e.entity_id, e.event_type, e.detected_at, e.effective_date, e.old_value, e.new_value, e.source_id, e.document_id, e.url, e.title, e.summary,
93 − e.significance, e.confidence, e.review_status, e.country_iso2, e.metro_id, e.project_id,
94 − s.name as source_name, s.kind as source_kind,
95 − o.id as op_id, o.slug as op_slug, o.name as op_name`;
141 + e.significance, e.confidence, e.review_status, e.country_iso2, e.metro_id, e.project_id, e.cluster_id, e.evidence_count, e.is_ai, e.run_id,
142 + s.name as source_name, coalesce(e.source_kind, s.kind) as source_kind,
143 + o.id as op_id, o.slug as op_slug, o.name as op_name,
144 + em.id as met_id, em.slug as met_slug, em.name as met_name,
145 + ep.id as prj_id, ep.slug as prj_slug, ep.name as prj_name`;
96 146 }
97 147
98 148 export function eventJoins(sql: Sql): Fragment {
99 149 return sql`
100 150 left join sources s on s.id = e.source_id
101 − left join operators o on o.id = e.operator_id`;
151 + left join operators o on o.id = e.operator_id
152 + left join metros em on em.id = e.metro_id
153 + left join projects ep on ep.id = e.project_id`;
154 +}
155 +
156 +/** Event types describing corporate life rather than physical infrastructure. */
157 +export const CORPORATE_EVENT_TYPES = ["acquisition", "investment_announced", "partnership", "executive_change", "customer_agreement", "power_agreement"];
158 +/** Event types about electricity supply and grid constraints. */
159 +export const POWER_EVENT_TYPES = ["power_agreement", "grid_connection", "grid_constraint", "utility_event"];
160 +export const GRID_EVENT_TYPES = ["grid_constraint", "utility_event", "power_agreement"];
161 +
162 +/** Columns for a ClaimDTO from `claims k` joined with sources s. */
163 +export function claimCols(sql: Sql): Fragment {
164 + return sql`
165 + k.id, k.subject_type, k.subject_id, k.predicate, k.value, k.value_text, k.unit, k.scope, k.scope_reason, k.source_id, k.connector_id, k.document_id, k.url, k.published_at, k.retrieved_at,
166 + k.confidence, k.is_estimate, k.authority_tier, k.evidence_text, k.parser_name, k.parser_version, k.run_id, k.status, k.rejection_reason, k.first_observed, k.last_observed,
167 + s.name as source_name, s.kind as source_kind`;
168 +}
169 +
170 +/** Columns for a GridConstraintDTO from `grid_constraints g` joined with sources s and metros gm. */
171 +export function gridConstraintCols(sql: Sql): Fragment {
172 + return sql`
173 + g.id, g.kind, g.title, g.summary, g.effective_date, g.url, g.country_iso2, g.confidence, g.event_id, g.source_id,
174 + s.name as source_name, s.kind as source_kind,
175 + gm.id as met_id, gm.slug as met_slug, gm.name as met_name`;
176 +}
177 +export function gridConstraintJoins(sql: Sql): Fragment {
178 + return sql`left join sources s on s.id = g.source_id left join metros gm on gm.id = g.metro_id`;
102 179 }
103 180
104 181 export interface Page { page: number; perPage: number; offset: number }
@@ -117,7 +194,18 @@ export function bboxAround(lat: number, lng: number, km: number): { minLat: numb
117 194
118 195 /** Haversine distance in km as SQL (facility alias `f`). */
119 196 export function haversineExpr(sql: Sql, lat: number, lng: number): Fragment {
120 − return sql`(2 * 6371.0088 * asin(sqrt(power(sin(radians(f.lat - ${lat}) / 2), 2) + cos(radians(${lat})) * cos(radians(f.lat)) * power(sin(radians(f.lng - ${lng}) / 2), 2))))`;
197 + return haversineFor(sql, lat, lng, sql`f.lat`, sql`f.lng`);
198 +}
199 +
200 +/** Haversine distance in km as SQL for arbitrary lat / lng columns. */
201 +export function haversineFor(sql: Sql, lat: number, lng: number, latCol: Fragment, lngCol: Fragment): Fragment {
202 + return sql`(2 * 6371.0088 * asin(sqrt(least(1.0, power(sin(radians(${latCol} - ${lat}) / 2), 2) + cos(radians(${lat})) * cos(radians(${latCol})) * power(sin(radians(${lngCol} - ${lng}) / 2), 2)))))`;
203 +}
204 +
205 +/** `lat between … and lng between …` for a bbox around a point (pre-filter before haversine). */
206 +export function withinBbox(sql: Sql, lat: number, lng: number, km: number, latCol: Fragment, lngCol: Fragment): Fragment {
207 + const b = bboxAround(lat, lng, km);
208 + return sql`(${latCol} between ${b.minLat} and ${b.maxLat} and ${lngCol} between ${b.minLng} and ${b.maxLng})`;
121 209 }
122 210
123 211 /** `left(col, 4)::int` when the partial date starts with a year, else null. */
added apps/worker/probe-fixture.ts +28 −0
@@ -0,0 +1,28 @@
1 +// TEMPORARY probe (deleted after fixture authoring): run archived bodies through extract→normalize→validate.
2 +// usage: node ../../node_modules/tsx/dist/cli.mjs probe-fixture.ts <connector> <file> <url> [contentType]
3 +import { readFileSync } from "node:fs";
4 +import { GenericConnector, createTestContext, getImplementation, parseConnectorConfig, validateEntities, type RawDocument } from "@dci/connectors";
5 +import { register as reg1 } from "./src/connectors/operators1/index.js";
6 +import { register as reg2 } from "./src/connectors/operators2/index.js";
7 +import { register as regNews } from "./src/connectors/news/index.js";
8 +import { articleContent } from "./src/connectors/news/article-parser.js";
9 +import { extractAnnouncement } from "./src/connectors/news/extract-project.js";
10 +
11 +const [connectorId, file, url, contentType = "text/html; charset=utf-8"] = process.argv.slice(2) as [string, string, string, string?];
12 +reg1(); reg2(); regNews();
13 +const cfg = parseConnectorConfig(readFileSync(`../../config/connectors/${connectorId}.yaml`, "utf8"));
14 +const body = readFileSync(file);
15 +const doc: RawDocument = { url, finalUrl: url, fetchedAt: new Date().toISOString(), status: 200, contentType, body, text: body.toString("utf8"), headers: {}, etag: null, lastModified: null, notModified: false, fetcher: "cache", level: 1, durationMs: 0, credits: 0 };
16 +const impl = cfg.implementation ? getImplementation(cfg.implementation) : undefined;
17 +const connector = impl ? impl(cfg) : new GenericConnector(cfg);
18 +const ctx = createTestContext(cfg);
19 +const recs = await connector.extract(ctx, doc);
20 +const ents = await connector.normalize(ctx, recs);
21 +const report = validateEntities(ents);
22 +const strip = (e: Record<string, unknown>) => { const { provenance: _p, description: _d, ...rest } = e; return rest; };
23 +console.log(JSON.stringify({ records: recs.length, entities: ents.map((e) => strip(e as unknown as Record<string, unknown>)), report }, null, 1));
24 +if (cfg.kind === "news" || cfg.kind === "government") {
25 + const c = await articleContent(doc);
26 + const a = extractAnnouncement(c.title, c.text, { publishedAt: c.published });
27 + console.log("ANNOUNCEMENT", JSON.stringify({ title: c.title, published: c.published, class: a.classification, headlineMw: a.headlineMw, status: a.status, operator: a.operator?.name, location: a.location, money: a.money, mwAll: a.mwAll }, null, 1));
28 +}
modified apps/worker/src/ingest/events.ts +53 −0
@@ -199,3 +199,56 @@ export async function emitDiffEvents(tx: Tx, ctx: IngestContext, i: DiffEventsIn
199 199 export function discoverySignificance(mw: number | null, pipeline: boolean): number {
200 200 return (mw != null && mw >= 50) || pipeline ? 65 : 40;
201 201 }
202 +
203 +/**
204 + * Event clustering (docs/CLAIMS.md): one underlying announcement covered by an operator PR and several outlets must
205 + * be ONE event with several evidence documents. Deterministic signals only: same event family, same operator (or
206 + * same country + same city when no operator), publication within ±3 days, MW overlap (±10 %) when both sides have one,
207 + * or ≥ 0.6 headline token overlap. Returns the cluster id to store on the new event (existing cluster → the primary
208 + * event's evidence_count is bumped and, when the incoming source is primary and the existing one is not, the cluster
209 + * primary is re-pointed).
210 + */
211 +const EVENT_FAMILY: Record<string, string> = {
212 + project_announced: "announcement", expansion_announced: "announcement", investment_announced: "announcement", phase_announced: "announcement",
213 + construction_started: "construction", planning_filed: "planning", planning_approved: "planning", land_acquired: "planning",
214 + facility_opened: "opening", cloud_region_announced: "cloud", cloud_region_launched: "cloud",
215 + acquisition: "corporate", partnership: "corporate", customer_agreement: "corporate", executive_change: "corporate", operator_expansion: "corporate",
216 + power_agreement: "power", grid_connection: "power", grid_constraint: "power", utility_event: "power",
217 + project_delayed: "setback", project_cancelled: "setback", closure: "setback", incident: "setback",
218 +};
219 +const STOP = new Set(["the", "a", "an", "of", "in", "at", "to", "for", "and", "on", "with", "its", "new", "data", "center", "centre", "centers", "centres", "campus", "mw", "gw", "announces", "announced", "plans", "build", "project", "facility", "datacenter", "datacentre"]);
220 +export function headlineTokens(s: string): Set<string> { return new Set(s.toLowerCase().normalize("NFKD").replace(/[^a-z0-9 ]/g, " ").split(/\s+/).filter((t) => t.length >= 3 && !STOP.has(t) && !/^\d+$/.test(t))); }
221 +export function headlineOverlap(a: string, b: string): number { const x = headlineTokens(a), y = headlineTokens(b); if (!x.size || !y.size) return 0; let n = 0; for (const t of x) if (y.has(t)) n++; return n / Math.min(x.size, y.size); }
222 +
223 +export interface ClusterProbe { eventType: EventType; title: string; operatorId?: string | null; countryIso2?: string | null; metroId?: string | null; city?: string | null; mw?: number | null; day: string; sourceKind: string }
224 +
225 +export async function findEventCluster(tx: Tx, ctx: IngestContext, p: ClusterProbe): Promise<{ clusterId: string; joined: boolean }> {
226 + const family = EVENT_FAMILY[p.eventType] ?? p.eventType;
227 + const types = Object.entries(EVENT_FAMILY).filter(([, f]) => f === family).map(([t]) => t);
228 + if (!types.includes(p.eventType)) types.push(p.eventType);
229 + if (!p.operatorId && !p.metroId && !p.countryIso2) return { clusterId: newId("cluster"), joined: false };
230 + const rows = await tx.execute(sql`
231 + select id, cluster_id, title, new_value, source_kind, evidence_count, operator_id, metro_id, country_iso2 from events
232 + where event_type in ${types} and detected_at > now() - interval '45 days'
233 + and coalesce(effective_date, detected_at::date::text) between (${p.day}::date - 3)::text and (${p.day}::date + 3)::text
234 + and (${p.operatorId ?? null}::text is not null and operator_id = ${p.operatorId ?? null}
235 + or (${p.operatorId ?? null}::text is null and operator_id is null and ${p.countryIso2 ?? null}::text is not null and country_iso2 = ${p.countryIso2 ?? null} and (${p.metroId ?? null}::text is null or metro_id is null or metro_id = ${p.metroId ?? null})))
236 + order by detected_at asc limit 50`);
237 + for (const r of rows) {
238 + const nv = (r.new_value ?? {}) as Record<string, unknown>;
239 + const mwOther = typeof nv.mw === "number" ? nv.mw : typeof nv.plannedMw === "number" ? nv.plannedMw : Array.isArray(nv.mw) ? Number((nv.mw as number[])[0]) : null;
240 + const mwOk = p.mw != null && mwOther != null ? Math.abs(p.mw - mwOther) <= 0.1 * Math.max(p.mw, mwOther) : null;
241 + if (mwOk === false) continue;
242 + const overlap = headlineOverlap(p.title, String(r.title));
243 + const sameCity = !!p.city && new RegExp(`\\b${p.city.replace(/[.*+?^${}()|[\]\\]/g, "\\$&")}\\b`, "i").test(String(r.title));
244 + if (mwOk === true || overlap >= 0.6 || (sameCity && overlap >= 0.3)) {
245 + const clusterId = r.cluster_id ? String(r.cluster_id) : newId("cluster");
246 + if (!ctx.run.dryRun) {
247 + if (!r.cluster_id) await tx.execute(sql`update events set cluster_id = ${clusterId} where id = ${r.id}`);
248 + await tx.execute(sql`update events set evidence_count = evidence_count + 1 where cluster_id = ${clusterId}`);
249 + }
250 + return { clusterId, joined: true };
251 + }
252 + }
253 + return { clusterId: newId("cluster"), joined: false };
254 +}
modified apps/worker/src/ingest/facilities.ts +31 −1
@@ -8,7 +8,7 @@
8 8 * Field merge follows the authority policy in match.ts; coordinates never lose precision.
9 9 */
10 10 import { sql, textArray } from "@dci/db";
11 −import { cleanText, geohash, newId, normalizeName, validLatLng, type ConfidenceLevel, type NormalizedFacility, type Provenance } from "@dci/core";
11 +import { cleanText, geohash, newId, normalizeName, sha256, validLatLng, type ConfidenceLevel, type NormalizedFacility, type Provenance } from "@dci/core";
12 12 import { addRef, bump, isPrimaryKind, provenanceFor, safeCountry, uniqueSlug, uniqStrings, normalizeWebsite, type IngestContext, type Tx } from "./common.js";
13 13 import { isHyperscalerName } from "./canonical-operators.js";
14 14 import { inferOperatorFromName } from "./operator-inference.js";
@@ -501,6 +501,7 @@ export async function ingestFacility(tx: Tx, ctx: IngestContext, nf: NormalizedF
501 501 const conf = (await tx.execute(sql`select confidence from facilities where id = ${id}`))[0]?.confidence as ConfidenceLevel | undefined;
502 502 if (isNew) {
503 503 ctx.stats.created++;
504 + await maybeOperatorExpansion(tx, ctx, { entityType: "facility", entityId: id, entityName: String(next.name), operatorId: (next.operatorId as string | null) ?? null, operatorName: operator?.name ?? null, countryIso2: (next.countryIso2 as string | null) ?? null, metroId: (next.metroId as string | null) ?? null, url, confidence: conf ?? "moderate", isAi: !!next.isAi });
504 505 const mw = bestMw({ itCapacityMw: next.itCapacityMw as number | null, totalPowerMw: next.totalPowerMw as number | null, plannedPowerMw: next.plannedPowerMw as number | null });
505 506 const pipeline = isPipelineStatus(next.status as string | null);
506 507 const opName = operator?.name ?? null;
@@ -630,3 +631,32 @@ export async function previewFacilityResolution(tx: Tx, ctx: IngestContext, nf:
630 631 const how = byKey ? "key" : byExt ? "external_id" : !best ? "create" : best.reasons.includes("rule:campus-vs-building") ? "campus_link" : decide(best.score);
631 632 return { how, matchedId: byKey ?? byExt ?? (best && (how === "merge") ? best.id : null), candidates: scored };
632 633 }
634 +
635 +/**
636 + * First record of an operator in a country (or a metro) → `operator_expansion` event ("Which companies are expanding
637 + * into new countries?"). Deterministic: counts other live facilities / projects of the operator in that country.
638 + */
639 +export async function maybeOperatorExpansion(tx: Tx, ctx: IngestContext, i: { entityType: "facility" | "project"; entityId: string; entityName: string; operatorId: string | null; operatorName: string | null; countryIso2: string | null; metroId: string | null; url: string; confidence: ConfidenceLevel; isAi: boolean }): Promise<void> {
640 + if (!i.operatorId || !i.countryIso2) return;
641 + const others = (await tx.execute(sql`select (select count(*) from facilities f where f.operator_id = ${i.operatorId} and f.country_iso2 = ${i.countryIso2} and f.merged_into is null and f.id <> ${i.entityId})::int
642 + + (select count(*) from projects p where p.operator_id = ${i.operatorId} and p.country_iso2 = ${i.countryIso2} and p.merged_into is null and not p.hidden and p.id <> ${i.entityId})::int as n`))[0];
643 + if (Number(others?.n ?? 0) > 0) return;
644 + const country = (await tx.execute(sql`select name from countries where iso2 = ${i.countryIso2}`))[0]?.name;
645 + await recordEvent(tx, ctx, {
646 + entityType: i.entityType,
647 + entityId: i.entityId,
648 + eventType: "operator_expansion",
649 + title: `${i.operatorName ?? "Operator"} enters ${String(country ?? i.countryIso2)}: ${i.entityName}`.slice(0, 300),
650 + summary: `First ${i.entityType} of ${i.operatorName ?? "this operator"} indexed in ${String(country ?? i.countryIso2)}.`,
651 + newValue: { countryIso2: i.countryIso2, operatorId: i.operatorId, [i.entityType === "facility" ? "facilityId" : "projectId"]: i.entityId },
652 + significance: 70,
653 + confidence: i.confidence,
654 + url: i.url,
655 + countryIso2: i.countryIso2,
656 + operatorId: i.operatorId,
657 + metroId: i.metroId,
658 + projectId: i.entityType === "project" ? i.entityId : null,
659 + isAi: i.isAi,
660 + fingerprint: sha256(`operator_expansion|${i.operatorId}|${i.countryIso2}`),
661 + });
662 +}
modified apps/worker/src/ingest/news.ts +49 −7
@@ -7,7 +7,9 @@ import { cleanText, classifyAiEvidence, normalizeName, parseAllMw, stableId, typ
7 7 import { countryFromTextSafe } from "../connectors/news/locations-lexicon.js";
8 8 import { addRef, bump, isoOrNull, knownCountries, type IngestContext, type Tx } from "./common.js";
9 9 import { CANONICAL_OPERATORS, findCanonicalOperator } from "./canonical-operators.js";
10 −import { recordEvent } from "./events.js";
10 +import { findEventCluster, recordEvent } from "./events.js";
11 +import { geocodeCity } from "./geocode.js";
12 +import { assignMetro } from "./metros.js";
11 13 import { lookupOperatorId, resolveOperator } from "./operators.js";
12 14
13 15 const EVENT_SIGNIFICANCE: Partial<Record<EventType, number>> = {
@@ -154,25 +156,65 @@ export async function ingestNews(tx: Tx, ctx: IngestContext, n: NormalizedNewsEv
154 156 bump(ctx, "news_event");
155 157 addRef(ctx, "news_item", id);
156 158
157 − if (eventType && significance >= 40) {
159 + // 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 layer
168 + 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;
172 +
173 + 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 });
158 176 // one event per article: fingerprint on the URL, not on the day, so a re-crawl never duplicates it
159 − await recordEvent(tx, ctx, {
177 + const eventId = await recordEvent(tx, ctx, {
160 178 entityType: "news_event",
161 179 entityId: id,
162 − eventType,
180 + eventType: effectiveEventType,
163 181 title,
164 182 summary,
165 − newValue: { url, mw: links.mw, operators: links.operatorNames, countries: links.countryIso2s },
166 − significance,
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,
167 185 confidence: n.provenance.confidence,
168 186 effectiveDate: publishedAt ? publishedAt.slice(0, 10) : null,
169 187 url,
170 188 countryIso2: links.countryIso2s[0] ?? null,
171 189 operatorId: links.operatorIds[0] ?? null,
172 − fingerprint: `news:${id}:${eventType}`,
190 + metroId,
191 + fingerprint: `news:${id}:${effectiveEventType}`,
173 192 isAi: ai,
193 + clusterId: cluster.clusterId,
174 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 + }
175 201 }
176 202 }
177 203
178 204 export const NEWS_EVENT_SIGNIFICANCE = EVENT_SIGNIFICANCE;
205 +
206 +/** Grid / utility constraint vocabulary → grid_constraints.kind (docs/CLAIMS.md, power layer). */
207 +export 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 +];
217 +export 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 +}
modified apps/worker/src/ingest/projects.ts +2 −0
@@ -42,6 +42,7 @@ import { operatorNames, resolveOperator } from "./operators.js";
42 42 import { backingObservation, loadCurrentProvenance, writeProvenance, type CurrentProvenance, type ObservedField } from "./provenance.js";
43 43 import { markWinners, recordCapacityClaim, recordInvestmentClaim, reviewPriority, writeQualityFlags, writeTextClaim } from "./claims.js";
44 44 import { resolveCampus } from "./campuses.js";
45 +import { maybeOperatorExpansion } from "./facilities.js";
45 46
46 47 type Row = Record<string, unknown>;
47 48
@@ -429,6 +430,7 @@ export async function ingestProject(tx: Tx, ctx: IngestContext, p: NormalizedPro
429 430 isAi: !!next.isAi,
430 431 reviewStatus: p.evidenceLevel === "weak" ? "pending" : "auto",
431 432 });
433 + await maybeOperatorExpansion(tx, ctx, { entityType: "project", entityId: id, entityName: String(next.name), operatorId: (next.operatorId as string | null) ?? null, operatorName: operator?.name ?? null, countryIso2: (next.countryIso2 as string | null) ?? null, metroId: (next.metroId as string | null) ?? null, url, confidence, isAi: !!next.isAi });
432 434 } else {
433 435 const names = await operatorNames(tx, ctx, [before.operatorId as string | null, next.operatorId as string | null]);
434 436 await emitDiffEvents(tx, ctx, { entityType: "project", entityId: id, entityName: String(next.name), before, after: next, specs: TRACKED_PROJECT_FIELDS, url, confidence, countryIso2: next.countryIso2 as string | null, operatorId: next.operatorId as string | null, metroId: next.metroId as string | null, projectId: id, operatorNames: names });
modified deploy/compose.data.yml +2 −0
@@ -114,6 +114,8 @@ services:
114 114 NEXT_PUBLIC_API_URL: ${NEXT_PUBLIC_API_URL:-https://www.datacenterindex.io}
115 115 DCI_API_CACHE_ENTRIES: ${DCI_API_CACHE_ENTRIES:-5000}
116 116 DCI_TRUST_PROXY: "1"
117 + # internal worker HTTP (extraction debugger /trace, /data-gaps) — proxied behind the admin token
118 + DCI_WORKER_URL: http://worker-maint:8320
117 119 depends_on:
118 120 postgres: { condition: service_healthy }
119 121 redis: { condition: service_healthy }
added docs/CLAIMS.md +100 −0
@@ -0,0 +1,100 @@
1 +# Claim-first data architecture
2 +
3 +Since 2026-09-12 every numerical figure DataCenterIndex extracts is stored as a **claim** before it can become a displayed
4 +value. A claim is one assertion by one document about one subject, with its **scope**, its **semantics**, the **sentence**
5 +that supports it, the **parser** that produced it and the **field-level authority** of its source. Reconciliation then
6 +decides which claim (if any) populates a column. Nothing is deleted when sources disagree.
7 +
8 +## Claim
9 +
10 +| column | meaning |
11 +|---|---|
12 +| `subject_type`, `subject_id` | facility / campus / project / operator / market / country |
13 +| `predicate` | what is asserted — `it_capacity_mw`, `critical_power_mw`, `utility_capacity_mw`, `grid_connection_mw`, `current_power_mw`, `planned_power_mw`, `ultimate_campus_mw`, `phase_mw`, `project_investment_usd`, `campus_investment_usd`, `company_investment_usd`, `country_program_usd`, `multi_year_capex_usd`, `deal_value_usd`, `status`, … |
14 +| `value`, `value_text`, `unit` | the figure (MW, USD) or text |
15 +| `scope`, `scope_reason` | `building` · `facility` · `campus` · `metro` · `country` · `portfolio` · `company` · `unknown`, and the rule that decided |
16 +| `evidence_text`, `evidence_start`, `evidence_end` | the supporting sentence and its offsets in the parsed text — **no sentence, no numerical claim** (structured spec tables record the extraction method instead) |
17 +| `source_id`, `url`, `document_id`, `published_at`, `retrieved_at` | where and when |
18 +| `authority_tier` | A–E for THIS field (see below) |
19 +| `parser_name`, `parser_version`, `run_id` | who extracted it (reprocessing + rollback) |
20 +| `status` | `current` · `superseded` (same source, same page, new value) · `rejected` (review / rollback) · `review` (critical sanity flag) · `unscoped` (figure describes a portfolio, a company, a country, a metro, or its scope is unknown) |
21 +
22 +## Scope rules
23 +
24 +- Only **site scopes** (`building`, `facility`, `campus`) may populate a facility's or project's capacity / investment
25 + columns. Portfolio, company, country, metro and unknown scopes are stored and shown in the evidence drawer, never
26 + summed.
27 +- Scope is classified deterministically from the sentence (`packages/core/src/claims.ts` `classifyScope`): company →
28 + portfolio → country → metro → campus → building → facility → the record's own scope when the sentence has no signal
29 + (structured parsers) → unknown.
30 +- A campus figure is never written to a building record; a building's figure never stands for its campus.
31 +
32 +## Capacity semantics
33 +
34 +IT load, critical power, utility capacity, grid connection, current power, planned capacity, ultimate build-out and
35 +phase capacity are different predicates and different columns (`it_capacity_mw`, `total_power_mw`, `planned_power_mw`,
36 +`utility_capacity_mw`, `grid_connection_mw`, `ultimate_campus_mw`). The displayed figure carries `capacity_scope` and
37 +`capacity_semantics` so the UI can say what "300 MW" means. When a sentence gives no cue the predicate defaults from the
38 +field/record status and the claim is flagged `mw_semantics_default` (info).
39 +
40 +## Field-level authority
41 +
42 +One universal source rank is wrong. `FIELD_AUTHORITY` in `packages/core/src/claims.ts` gives a tier per field:
43 +
44 +| field | A | B | C | D | E |
45 +|---|---|---|---|---|---|
46 +| capacity | utility, government, filing | operator, cloud provider | registry, secondary | news | dataset, community |
47 +| geometry | government, community (OSM) | registry, operator, dataset | cloud provider, utility, filing | secondary | news |
48 +| interconnection | registry (PeeringDB) | operator | community, dataset | secondary, news, government, filing | utility, cloud |
49 +| status | operator, government, filing, cloud provider | utility, registry | secondary, news | dataset, community | — |
50 +| investment | filing, government | operator, cloud provider | utility, secondary | news | registry, dataset, community |
51 +
52 +Estimates and LLM-only extractions drop one tier; a human review raises one.
53 +
54 +## Sanity engines
55 +
56 +`capacitySanity` / `investmentSanity` run on every claim. Blocking flags (`blocks: true`) keep the figure out of the
57 +columns: non-site scope, campus figure on a building, ≥ 20 GW (market statistic), invalid value, deal value / capex /
58 +country programme as a project investment, > $500 B. Non-blocking critical flags (single site > 1 000 MW without a campus
59 +designation, > $50 B on one site, money and MW sharing a number in one sentence, 5× change) assign the figure but open a
60 +`quality_flags` row prioritised by impact (`reviewPriority`).
61 +
62 +## Announcement classification (projects)
63 +
64 +Before a news article can create a project, `classifyProjectEvent` labels it: `NEW_BUILD`, `EXPANSION`,
65 +`CONSTRUCTION_START`, `PERMIT`, `LAND_ACQUISITION`, `GRID_CONNECTION` (physical — may create / modify a project);
66 +`POWER_AGREEMENT`, `FINANCING`, `ACQUISITION`, `PARTNERSHIP`, `CUSTOMER_AGREEMENT` (associated — attach an event to an
67 +identifiable project, never create one); `EXECUTIVE_APPOINTMENT`, `SUSTAINABILITY`, `PRODUCT_NEWS`,
68 +`GENERAL_COMPANY_NEWS`, `UNKNOWN` (never a project). Creation additionally needs the evidence threshold: (explicit
69 +name or operator) + location + a development verb. Company domiciles ("Denver-based", "headquartered in") are stripped
70 +before locating the project.
71 +
72 +## Lifecycle state machine
73 +
74 +`projectTransition(from, to)`: forward moves along rumored → proposed → announced → permitting → approved →
75 +under_construction → partially_operational → operational are accepted; `delayed` / `cancelled` are side branches;
76 +leaving `delayed` resumes; a backward move needs a source that outranks the stored one, otherwise it is flagged
77 +(`status_backward`). Stage dates (`permit_filed_on`, `approved_on`, `construction_started_on`, `opened_on`) are set from
78 +the announcement that moved the stage.
79 +
80 +## Winners, history, rollback
81 +
82 +`provenance.is_winner` marks the observation backing each displayed value; `provenance.run_id`, `claims.run_id`,
83 +`events.run_id` and `document_versions.run_id` tie every change to a connector run (`POST /api/admin/runs/:id/rollback`).
84 +Capacity history = winner changes + claims + `capacity_changed` events. Daily `entity_snapshots` keep global totals,
85 +per-operator and per-country totals, rankings and stage counts for "as of" views and the automated regression checks
86 +(facilities −5 %, known MW +20 %, one operator +10 GW/day, one connector > 500 records/day, > 200 location changes/day).
87 +
88 +## Containment (campus vs building)
89 +
90 +`facilities.parent_facility_id` + `record_scope` (building / facility / campus). The matcher's `rule:campus-vs-building`
91 +now links instead of flagging a duplicate. Every aggregate uses `FACILITY_VIEW` (apps/worker/src/rankings.ts): a campus
92 +whose buildings publish figures contributes no MW itself; a building without a figure under a campus with one is
93 +"covered"; a campus with buildings is not counted as an extra facility.
94 +
95 +## Debugging
96 +
97 +`pnpm dci trace <doc-id|url>` (or `GET /api/admin/documents/:id/trace`) shows every stage for one document: raw fetch →
98 +parsed text → structured records → normalized entities → claims with scope/semantics/evidence → match candidates →
99 +dry-run reconciliation → resulting changes. `pnpm dci quality` runs the sweep, `pnpm dci snapshot` the snapshots +
100 +regression checks, `pnpm dci gaps` the data-gap counts, `pnpm dci quarantine <id> on` puts a connector in preview mode.
modified packages/connectors/src/robots.ts +124 −27
@@ -1,13 +1,55 @@
1 −import { DirectFetcher } from "./fetchers.js";
1 +import { DEFAULT_UA, DirectFetcher } from "./fetchers.js";
2 +
3 +/**
4 + * robots.txt evaluator (RFC 9309 subset): User-agent groups, Allow/Disallow longest-match with `*` wildcards and the
5 + * `$` end anchor, percent-decoded paths, Crawl-delay, Sitemap. Cached per origin (bounded, LRU-ish).
6 + *
7 + * Which group applies — the crawler identifies itself with the bot token (`DCI_USER_AGENT`, default
8 + * "datacenterindexbot") at L1. We evaluate the group addressed to that token AND the `*` group: a URL is fetched only
9 + * when both allow it (conservative: a site that disallows everybody but grants a named exception is honoured, a site
10 + * that disallows our token is honoured whatever `*` says). The L2 "browser identity" is the same direct transport with
11 + * a browser User-Agent; it is only ever used AFTER the bot token was allowed by robots — the browser UA is never the
12 + * identity robots.txt is evaluated for, and never a way around a disallow.
13 + *
14 + * Fetch outcomes: 200 → rules; 404/410 → allow everything (standard); any other status or a transport error →
15 + * keep the previously cached rules when there are some (even expired), otherwise the origin is "unknown" for one hour
16 + * and callers must restrict themselves to L1/L2 fetches (see `isAllowedByRobots().unknown`).
17 + */
18 +export interface RobotsGroup { allow: string[]; disallow: string[]; crawlDelay: number | null }
19 +export interface Rules extends RobotsGroup {
20 + /** group addressed to our bot token (null when the file has none) */
21 + agent: RobotsGroup | null;
22 + /** the `*` group (null when the file has none) */
23 + star: RobotsGroup | null;
24 + sitemaps: string[];
25 + fetchedAt: number;
26 + /** robots.txt could not be retrieved (429/5xx/timeout) and nothing was cached: only L1/L2 fetches are allowed */
27 + unknown?: boolean;
28 + /** cache lifetime for this entry (ms) */
29 + ttlMs?: number;
30 +}
2 31
3 −/** Minimal robots.txt evaluator (User-agent groups, Allow/Disallow longest-match, Crawl-delay, Sitemap). Cached per host for 24 h. */
4 −interface Rules { allow: string[]; disallow: string[]; crawlDelay: number | null; sitemaps: string[]; fetchedAt: number }
5 −const cache = new Map<string, Rules>();
6 32 const TTL = 24 * 3_600_000;
33 +const UNKNOWN_TTL = 3_600_000;
34 +export const ROBOTS_CACHE_MAX = 5_000;
35 +const cache = new Map<string, Rules>();
36 +
37 +/** Product token of the crawler's User-Agent ("DataCenterIndexBot/0.1 (+…)" → "datacenterindexbot"). */
38 +export function botToken(ua: string = DEFAULT_UA): string {
39 + const m = ua.trim().match(/^([A-Za-z0-9_.-]+)/);
40 + return (m?.[1] ?? "datacenterindexbot").toLowerCase();
41 +}
42 +
43 +function emptyGroup(): RobotsGroup { return { allow: [], disallow: [], crawlDelay: null }; }
7 44
8 −export function parseRobots(text: string, ua = "datacenterindexbot"): Rules {
45 +function safeDecode(s: string): string {
46 + try { return decodeURIComponent(s); } catch { return s; }
47 +}
48 +
49 +export function parseRobots(text: string, ua: string = botToken()): Rules {
50 + const token = ua.toLowerCase();
9 51 const lines = text.split(/\r?\n/).map((l) => l.replace(/#.*/, "").trim()).filter(Boolean);
10 − const groups: Array<{ agents: string[]; allow: string[]; disallow: string[]; crawlDelay: number | null }> = [];
52 + const groups: Array<{ agents: string[]; rules: RobotsGroup }> = [];
11 53 const sitemaps: string[] = [];
12 54 let cur: (typeof groups)[number] | null = null;
13 55 let lastWasAgent = false;
@@ -17,7 +59,7 @@ export function parseRobots(text: string, ua = "datacenterindexbot"): Rules {
17 59 const k = line.slice(0, i).trim().toLowerCase();
18 60 const v = line.slice(i + 1).trim();
19 61 if (k === "user-agent") {
20 − if (!cur || !lastWasAgent) { cur = { agents: [], allow: [], disallow: [], crawlDelay: null }; groups.push(cur); }
62 + if (!cur || !lastWasAgent) { cur = { agents: [], rules: emptyGroup() }; groups.push(cur); }
21 63 cur.agents.push(v.toLowerCase());
22 64 lastWasAgent = true;
23 65 continue;
@@ -25,41 +67,96 @@ export function parseRobots(text: string, ua = "datacenterindexbot"): Rules {
25 67 lastWasAgent = false;
26 68 if (k === "sitemap") { sitemaps.push(v); continue; }
27 69 if (!cur) continue;
28 − if (k === "allow") cur.allow.push(v);
29 − else if (k === "disallow") cur.disallow.push(v);
30 − else if (k === "crawl-delay") { const d = Number(v); if (Number.isFinite(d)) cur.crawlDelay = d; }
70 + if (k === "allow") cur.rules.allow.push(v);
71 + else if (k === "disallow") cur.rules.disallow.push(v);
72 + else if (k === "crawl-delay") { const d = Number(v); if (Number.isFinite(d) && d >= 0) cur.rules.crawlDelay = d; }
31 73 }
32 − const mine = groups.find((g) => g.agents.some((a) => a !== "*" && ua.includes(a))) ?? groups.find((g) => g.agents.includes("*"));
33 − return { allow: mine?.allow ?? [], disallow: mine?.disallow ?? [], crawlDelay: mine?.crawlDelay ?? null, sitemaps, fetchedAt: Date.now() };
74 + // most specific match for our token: exact token, then a token that is a prefix of ours ("datacenterindex") or vice versa
75 + const merge = (gs: Array<{ rules: RobotsGroup }>): RobotsGroup | null => {
76 + if (!gs.length) return null;
77 + const out = emptyGroup();
78 + for (const g of gs) { out.allow.push(...g.rules.allow); out.disallow.push(...g.rules.disallow); if (g.rules.crawlDelay != null) out.crawlDelay = out.crawlDelay == null ? g.rules.crawlDelay : Math.max(out.crawlDelay, g.rules.crawlDelay); }
79 + return out;
80 + };
81 + const exact = groups.filter((g) => g.agents.includes(token));
82 + const partial = exact.length ? [] : groups.filter((g) => g.agents.some((a) => a !== "*" && a.length >= 3 && (token.includes(a) || a.includes(token))));
83 + const agent = merge(exact.length ? exact : partial);
84 + const star = merge(groups.filter((g) => g.agents.includes("*")));
85 + const effective = agent ?? star ?? emptyGroup();
86 + return { ...effective, agent, star, sitemaps, fetchedAt: Date.now(), ttlMs: TTL };
34 87 }
35 88
36 −function matchLen(pattern: string, path: string): number {
89 +/**
90 + * Length of the pattern when it matches the path (longest match wins), -1 otherwise. `*` matches any run of
91 + * characters, a trailing `$` anchors the end of the path; both sides are percent-decoded before comparison.
92 + */
93 +export function matchLen(pattern: string, path: string): number {
37 94 if (!pattern) return -1;
38 − const re = new RegExp("^" + pattern.split("*").map((s) => s.replace(/[.*+?^${}()|[\]\\]/g, "\\$&")).join(".*") + (pattern.endsWith("$") ? "" : ""));
39 − return re.test(path) ? pattern.length : -1;
95 + const anchored = pattern.endsWith("$");
96 + const body = anchored ? pattern.slice(0, -1) : pattern;
97 + const re = new RegExp("^" + body.split("*").map((s) => safeDecode(s).replace(/[.*+?^${}()|[\]\\]/g, "\\$&")).join(".*") + (anchored ? "$" : ""));
98 + return re.test(safeDecode(path)) ? pattern.length : -1;
40 99 }
41 100
42 −export function robotsAllows(rules: Rules, url: string): boolean {
43 − const u = new URL(url);
44 − const path = u.pathname + u.search;
101 +export function groupAllows(g: RobotsGroup | null | undefined, path: string): boolean {
102 + if (!g) return true;
45 103 let bestAllow = -1, bestDis = -1;
46 − for (const p of rules.allow) bestAllow = Math.max(bestAllow, matchLen(p, path));
47 − for (const p of rules.disallow) bestDis = Math.max(bestDis, matchLen(p, path));
104 + for (const p of g.allow) bestAllow = Math.max(bestAllow, matchLen(p, path));
105 + for (const p of g.disallow) bestDis = Math.max(bestDis, matchLen(p, path));
48 106 if (bestDis < 0) return true;
49 107 return bestAllow >= bestDis;
50 108 }
51 109
52 −export async function getRobots(url: string): Promise<Rules> {
53 − const origin = new URL(url).origin;
110 +/** Both the group addressed to our token (when present) and the `*` group must allow the URL. */
111 +export function robotsAllows(rules: Rules, url: string): boolean {
112 + const u = new URL(url);
113 + const path = u.pathname + u.search;
114 + if (rules.agent || rules.star) return groupAllows(rules.agent, path) && groupAllows(rules.star, path);
115 + // legacy / hand-built rules without group detail
116 + return groupAllows({ allow: rules.allow, disallow: rules.disallow, crawlDelay: rules.crawlDelay }, path);
117 +}
118 +
119 +function cacheGet(origin: string): Rules | undefined {
54 120 const hit = cache.get(origin);
55 − if (hit && Date.now() - hit.fetchedAt < TTL) return hit;
56 − const doc = await new DirectFetcher(1).fetch(`${origin}/robots.txt`, { timeoutMs: 15_000, maxBytes: 512 * 1024, accept: "text/plain,*/*" });
57 − const rules = doc.status === 200 && doc.text ? parseRobots(doc.text) : { allow: [], disallow: [], crawlDelay: null, sitemaps: [], fetchedAt: Date.now() };
121 + if (hit) { cache.delete(origin); cache.set(origin, hit); } // refresh recency
122 + return hit;
123 +}
124 +function cacheSet(origin: string, rules: Rules): void {
125 + if (cache.has(origin)) cache.delete(origin);
58 126 cache.set(origin, rules);
127 + while (cache.size > ROBOTS_CACHE_MAX) { const oldest = cache.keys().next().value; if (oldest === undefined) break; cache.delete(oldest); }
128 +}
129 +/** Test hook: drop every cached robots.txt (also lets tests seed an origin). */
130 +export function resetRobotsCache(seed?: Array<[origin: string, rules: Rules]>): void {
131 + cache.clear();
132 + for (const [o, r] of seed ?? []) cacheSet(o, r);
133 +}
134 +export function robotsCacheSize(): number { return cache.size; }
135 +
136 +export const ALLOW_ALL: Omit<Rules, "fetchedAt"> = { allow: [], disallow: [], crawlDelay: null, agent: null, star: null, sitemaps: [] };
137 +
138 +export async function getRobots(url: string, fetchImpl: (robotsUrl: string) => Promise<{ status: number; text: string; error?: { code: string } | null }> = defaultFetch): Promise<Rules> {
139 + const origin = new URL(url).origin;
140 + const hit = cacheGet(origin);
141 + if (hit && Date.now() - hit.fetchedAt < (hit.ttlMs ?? TTL)) return hit;
142 + const doc = await fetchImpl(`${origin}/robots.txt`).catch(() => ({ status: 0, text: "", error: { code: "fetch_threw" } }));
143 + let rules: Rules;
144 + if (!doc.error && doc.status === 200) rules = parseRobots(doc.text);
145 + else if (!doc.error && (doc.status === 404 || doc.status === 410)) rules = { ...ALLOW_ALL, fetchedAt: Date.now(), ttlMs: TTL };
146 + else if (hit) rules = { ...hit, fetchedAt: Date.now(), ttlMs: UNKNOWN_TTL }; // keep the last known rules a while longer
147 + else rules = { ...ALLOW_ALL, fetchedAt: Date.now(), unknown: true, ttlMs: UNKNOWN_TTL };
148 + cacheSet(origin, rules);
59 149 return rules;
60 150 }
61 151
62 −export async function isAllowedByRobots(url: string): Promise<{ allowed: boolean; crawlDelay: number | null; sitemaps: string[] }> {
152 +async function defaultFetch(robotsUrl: string): Promise<{ status: number; text: string; error?: { code: string } | null }> {
153 + const d = await new DirectFetcher(1).fetch(robotsUrl, { timeoutMs: 15_000, maxBytes: 512 * 1024, accept: "text/plain,*/*" });
154 + return { status: d.status, text: d.text, error: d.error ?? null };
155 +}
156 +
157 +export interface RobotsDecision { allowed: boolean; crawlDelay: number | null; sitemaps: string[]; /** robots.txt unavailable and nothing cached: allow direct (L1/L2) fetches only */ unknown: boolean }
158 +
159 +export async function isAllowedByRobots(url: string): Promise<RobotsDecision> {
63 160 const rules = await getRobots(url);
64 − return { allowed: robotsAllows(rules, url), crawlDelay: rules.crawlDelay, sitemaps: rules.sitemaps };
161 + return { allowed: robotsAllows(rules, url), crawlDelay: rules.crawlDelay, sitemaps: rules.sitemaps, unknown: Boolean(rules.unknown) };
65 162 }
modified packages/core/src/claims.ts +2 −2
@@ -318,12 +318,12 @@ export const PHYSICAL_CLASSES: ReadonlySet<ProjectClass> = new Set<ProjectClass>
318 318 export const ASSOCIATED_CLASSES: ReadonlySet<ProjectClass> = new Set<ProjectClass>(["POWER_AGREEMENT", "FINANCING", "ACQUISITION", "PARTNERSHIP", "CUSTOMER_AGREEMENT"]);
319 319
320 320 const CLASS_RULES: Array<[ProjectClass, RegExp]> = [
321 − ["EXECUTIVE_APPOINTMENT", /\b(appoint(?:s|ed|ment)?|names? [A-Z][\w'-]+ [A-Z][\w'-]+ (?:as|to)|joins? (?:as|the board|its board)|(?:new|hires?|promotes?|welcomes?|taps?) (?:[A-Z][\w'-]+ [A-Z][\w'-]+ as )?(?:CEO|CFO|COO|CTO|CRO|CIO|CMO|chief|president|head of|director|vice president|VP|managing director|general manager|board member|chairman|chair)|leadership (?:change|team|appointments?)|steps? down|retire(?:s|ment)|passing of|obituary|in memoriam|executive (?:team|hire|appointment))\b/i],
321 + ["EXECUTIVE_APPOINTMENT", /\b(appoint(?:s|ed|ment)?|names? [A-Z][\w'-]+ [A-Z][\w'-]+ (?:as|to)|joins? (?:as|the board|its board)|joins? [A-Z][\w&.'-]+(?: [A-Z][\w&.'-]+){0,3} as (?:its )?(?:new )?(?:CEO|CFO|COO|CTO|CRO|CIO|CMO|chief|president|head|director|vice president|VP|SVP|EVP|managing director|general manager|chairman|chair|partner)|(?:new|hires?|promotes?|welcomes?|taps?) (?:[A-Z][\w'-]+ [A-Z][\w'-]+ as )?(?:CEO|CFO|COO|CTO|CRO|CIO|CMO|chief|president|head of|director|vice president|VP|managing director|general manager|board member|chairman|chair)|leadership (?:change|team|appointments?)|steps? down|retire(?:s|ment)|passing of|obituary|in memoriam|executive (?:team|hire|appointment))\b/i],
322 322 ["GENERAL_COMPANY_NEWS", /\b(market (?:to|will|set to|expected to|projected to) (?:surpass|reach|hit|grow|exceed|top)|market size|cagr|forecast(?:s|ed)? (?:to|that)|report(?: finds| shows| reveals|:)|survey|study (?:finds|shows|reveals)|(?:index|trends?|predictions?|outlook|white ?paper|e-?book|webinar|podcast|interview|q&a|faq|explained|101\b|guide to|tips|best practices|case study|customer spotlight|award|shortlist|finalist|recogni[sz]ed|named (?:a |one of |to )?(?:top|best|leader))\b|quarterly results|half[- ]year results|annual results|earnings|revenue|ebitda|guidance|investor (?:day|presentation)|lawsuit|sues?\b|court|settlement|opinion|op-ed|why |how |what )/i],
323 323 ["POWER_AGREEMENT", /\b(ppa\b|power purchase agreement|(?:virtual|corporate|long-term|\d+-year) (?:power|energy) (?:purchase )?agreement|(?:energy|power|electricity) (?:supply|offtake) (?:agreement|deal|contract)|offtake|(?:signs?|signed|inks?|inked|secures?|secured) (?:a )?(?:\d+[ -]?(?:mw|gw) )?(?:of )?(?:solar|wind|nuclear|geothermal|hydro|gas|renewable|clean|carbon-free) (?:power|energy|capacity|supply)|(?:solar|wind|nuclear|geothermal|hydro|gas|battery|fuel cell|smr|reactor)s? (?:to power|will power|powering|deal|plant to supply|farm to supply)|behind-the-meter|on-site (?:generation|power plant|gas plant|solar)|microgrid|small modular reactor)\b/i],
324 324 ["SUSTAINABILITY", /\b(sustainability (?:report|goals?|targets?|commitment|strategy)|net[- ]zero|carbon[- ]neutral|carbon[- ]free|esg\b|renewable (?:energy )?(?:certificates?|credits?|target|goal|commitment)|100% renewable|water (?:positive|stewardship|usage|conservation)|emissions? (?:reduction|target)|leed (?:gold|platinum|certif)|green (?:building|certification)|energy efficiency (?:program|initiative)|climate (?:pledge|commitment)|biodiversity|tree planting|heat (?:reuse|recovery) (?:program|scheme|initiative))\b/i],
325 325 ["GRID_CONNECTION", /\b(grid connection|grid-?connected|interconnection (?:agreement|request|capacity|queue|study)|connection agreement|(?:new |expanded )?substation|transmission (?:line|upgrade|project|capacity)|energi[sz]ation|energi[sz]ed|large[- ]load (?:tariff|request|agreement|customer)|load (?:request|study|letter)|(?:secured|reserved|allocated|approved) (?:\d+[ -]?(?:mw|gw) of )?(?:grid |utility )?(?:power|capacity) (?:from|with) (?:the utility|[A-Z][\w&]+ (?:energy|power|electric|utility)))\b/i],
326 − ["FINANCING", /\b(financ(?:es|ing|ed)|refinanc\w+|raises? (?:\$|€|£|us\$)|raised (?:\$|€|£|us\$)|series [a-e]\b|funding (?:round|arranged|secured)|loan|credit facility|green bond|bonds?\b|notes? offering|debt (?:facility|package|raise)|equity (?:raise|investment|stake)|ipo\b|capital raise|construction (?:loan|financing)|securitization|abs\b|term loan|(?:secures?|secured|closes?|closed) (?:\$|€|£|us\$)[\d.,]+ ?(?:m|bn|b|million|billion)? (?:in )?(?:financing|loan|facility|debt|funding))\b/i],
326 + ["FINANCING", /\b(financ(?:es|ing|ed)|refinanc\w+|(?:debt|notes?|bond|equity|securit\w+) (?:issuance|offering|placement|raise)|raises? (?:\$|€|£|us\$)|raised (?:\$|€|£|us\$)|series [a-e]\b|funding (?:round|arranged|secured)|loan|credit facility|green bond|bonds?\b|notes? offering|debt (?:facility|package|raise)|equity (?:raise|investment|stake)|ipo\b|capital raise|construction (?:loan|financing)|securitization|abs\b|term loan|(?:secures?|secured|closes?|closed) (?:\$|€|£|us\$)[\d.,]+ ?(?:m|bn|b|million|billion)? (?:in )?(?:financing|loan|facility|debt|funding))\b/i],
327 327 ["ACQUISITION", /\b(acqui(?:res?|red|sition) (?:of )?(?:[A-Z][\w&'.-]+ )+(?:group|holdings|inc|ltd|llc|limited|corp|sa|ag|plc|pte|data ?cent(?:er|re)s?|portfolio|platform|company|business)|(?:acquires?|acquired|acquiring|buys?|bought|buying|purchases?|purchased|purchasing|takes? over|took over) (?:a |an |the |its |two |three |four |five |six |\d+ )?(?:\d+[ -]?(?:mw|gw) )?(?:(?:hyperscale|colocation|carrier-neutral|edge|operational|existing) )?(?:data ?cent(?:er|re)s?|facility|facilities|portfolio|platform|operator|business|company|stake|campus from|site from)|acquisition of|takeover|buyout|merger|merges? with|to be acquired|stake in|sells?|sold|divest\w+|sale of)\b/i],
328 328 ["CUSTOMER_AGREEMENT", /\b((?:signs?|signed|secures?|secured|inks?|inked|lands?|landed|wins?|won|announces?) (?:a |an |its |the )?(?:multi-year |long-term |major |anchor |hyperscale |\d+[ -]?(?:mw|gw) )?(?:colocation |hosting |capacity |wholesale |pre-?)?(?:lease|leases|tenant|customer|contract|agreement|deal) (?:with|for|from)|(?:lease|leases|leased) (?:\d+[ -]?(?:mw|gw)|[\d,]+ (?:sq|square))|(?:selects?|selected|chooses?|chose|picks?|picked|taps?|tapped) [A-Z][\w&'.-]+ (?:for|as|to)|(?:hosting|colocation|lease|leasing|services?|master|supply|framework|capacity) agreements?|pre-?leas\w+|fully leased|anchor tenant|moves? into|deploys? (?:in|at|with)|expands? (?:its )?(?:footprint|presence) (?:in|at|with) [A-Z][\w&'.-]+(?:'s)? (?:data ?cent(?:er|re)|facility|campus))\b/i],
329 329 ["PARTNERSHIP", /\b(partnership|partners? with|joint venture|\bjv\b|join(?:s|ed)? forces|team(?:s|ed)? up|collaborat(?:es?|ion|ing)|alliance|memorandum of understanding|\bmou\b|strategic (?:agreement|relationship|cooperation)|framework agreement)\b/i],
added scripts/quality-apply.ts +122 −0
@@ -0,0 +1,122 @@
1 +/**
2 + * One-shot data repair after the claim-first upgrade (safe to re-run; every step is idempotent and reversible —
3 + * projects are hidden, never deleted; figures are nulled only when no site-scoped current claim backs them and the
4 + * old value stays in provenance).
5 + *
6 + * set -a; source .env; set +a
7 + * node node_modules/tsx/dist/cli.mjs scripts/quality-apply.ts # report only
8 + * node node_modules/tsx/dist/cli.mjs scripts/quality-apply.ts --hide-false-positives # hide vetoed projects
9 + * node node_modules/tsx/dist/cli.mjs scripts/quality-apply.ts --null-unbacked-mw # after `dci reprocess … --stale`
10 + * node node_modules/tsx/dist/cli.mjs scripts/quality-apply.ts --fix-slugs # non-Latin operator slugs ("item")
11 + * node node_modules/tsx/dist/cli.mjs scripts/quality-apply.ts --link-campuses # OSM buildings inside a campus footprint
12 + * node node_modules/tsx/dist/cli.mjs scripts/quality-apply.ts --all
13 + */
14 +import { closeDb, getDb, sql } from "@dci/db";
15 +import { classifyProjectEvent, PHYSICAL_CLASSES, slugify, sha256 } from "@dci/core";
16 +import { isNonProjectTitle, OPENING_TITLE_RE } from "../apps/worker/src/connectors/news/extract-project.js";
17 +import { hideProject } from "../apps/worker/src/ingest/projects.js";
18 +
19 +const args = process.argv.slice(2);
20 +const ALL = args.includes("--all");
21 +const has = (f: string) => ALL || args.includes(f);
22 +
23 +async function hideFalsePositives(): Promise<void> {
24 + const db = getDb();
25 + const rows = await db.execute<{ id: string; name: string; slug: string; planned_mw: number | null; project_class: string | null; status: string; title: string | null; summary: string | null; operator_id: string | null; city: string | null; country_iso2: string | null }>(sql`
26 + select p.id, p.name, p.slug, p.planned_mw, p.project_class, p.status, n.title, n.summary, p.operator_id, p.city, p.country_iso2
27 + from projects p left join news_items n on n.url = p.source_url
28 + where p.merged_into is null and not p.hidden`);
29 + let hidden = 0, kept = 0;
30 + const reasons: Record<string, number> = {};
31 + for (const r of rows) {
32 + const title = r.title ?? r.name;
33 + const cls = classifyProjectEvent({ title, lead: r.summary ?? "", hasOperator: !!r.operator_id, hasLocation: !!(r.city || r.country_iso2), hasExplicitName: !/ \| |—| - /.test(r.name) && r.name.length < 90, status: r.status });
34 + let reason: string | null = null;
35 + if (!PHYSICAL_CLASSES.has(cls.class)) reason = `class:${cls.class}`;
36 + else if (isNonProjectTitle(title)) reason = "veto:title";
37 + else if (OPENING_TITLE_RE.test(title) && !/\b(to|will|set to|plans? to|due to|expected to|slated to) (open|be online|go live)\b/i.test(title)) reason = "existing-facility";
38 + else if (cls.evidence.strength === "none") reason = "evidence:none";
39 + if (!reason) { kept++; continue; }
40 + reasons[reason] = (reasons[reason] ?? 0) + 1;
41 + console.log(` hide ${r.slug.padEnd(70)} ${String(r.planned_mw ?? "").padStart(7)} MW ${reason} «${title.slice(0, 80)}»`);
42 + if (has("--hide-false-positives")) await db.transaction((tx) => hideProject(tx, r.id, `${reason}: ${title.slice(0, 200)}`, "quality-apply"));
43 + hidden++;
44 + }
45 + console.log(`\nprojects: ${kept} kept · ${hidden} ${has("--hide-false-positives") ? "hidden" : "would be hidden"} · ${JSON.stringify(reasons)}`);
46 +}
47 +
48 +async function nullUnbackedMw(): Promise<void> {
49 + const db = getDb();
50 + // projects whose planned_mw / investment is not backed by any current site-scoped claim (after reprocessing) → null
51 + const mw = await db.execute<{ id: string; slug: string; planned_mw: number }>(sql`
52 + select p.id, p.slug, p.planned_mw from projects p where p.merged_into is null and p.planned_mw is not null
53 + and exists (select 1 from claims c where c.subject_type = 'project' and c.subject_id = p.id and c.predicate like '%_mw')
54 + and not exists (select 1 from claims c where c.subject_type = 'project' and c.subject_id = p.id and c.predicate like '%_mw' and c.status = 'current' and c.scope in ('building','facility','campus') and abs(c.value - p.planned_mw) < 0.01)`);
55 + for (const r of mw) console.log(` null planned_mw ${r.slug} (${r.planned_mw} MW: only unscoped / rejected claims back it)`);
56 + const inv = await db.execute<{ id: string; slug: string; investment_usd: number }>(sql`
57 + select p.id, p.slug, p.investment_usd from projects p where p.merged_into is null and p.investment_usd is not null
58 + and exists (select 1 from claims c where c.subject_type = 'project' and c.subject_id = p.id and c.predicate like '%_usd')
59 + and not exists (select 1 from claims c where c.subject_type = 'project' and c.subject_id = p.id and c.predicate = 'project_investment_usd' and c.status = 'current' and c.scope in ('building','facility','campus') and abs(c.value - p.investment_usd) < 1)`);
60 + for (const r of inv) console.log(` null investment_usd ${r.slug} ($${(r.investment_usd / 1e9).toFixed(2)}B: only unscoped / deal-value claims back it)`);
61 + if (has("--null-unbacked-mw")) {
62 + for (const r of mw) await db.execute(sql`update projects set planned_mw = null, capacity_scope = null, capacity_semantics = null, updated_at = now() where id = ${r.id}`);
63 + for (const r of inv) await db.execute(sql`update projects set investment_usd = null, investment_scope = null, investment_semantics = null, updated_at = now() where id = ${r.id}`);
64 + }
65 + console.log(`\nunbacked figures: ${mw.length} planned_mw · ${inv.length} investment_usd ${has("--null-unbacked-mw") ? "nulled" : "would be nulled"}`);
66 +}
67 +
68 +async function fixSlugs(): Promise<void> {
69 + const db = getDb();
70 + const rows = await db.execute<{ id: string; name: string; slug: string; website: string | null; aliases: string[] }>(sql`select id, name, slug, website, aliases from operators where slug ~ '^(item|op|operator)(-\\d+)?$' or slug = ''`);
71 + for (const r of rows) {
72 + const latin = (r.aliases ?? []).find((a) => /^[A-Za-z0-9 .&'-]+$/.test(a) && slugify(a)) ?? (r.website ? r.website.replace(/^https?:\/\/(www\.)?/, "").split("/")[0]?.split(".")[0] : null);
73 + const base = latin ? slugify(latin) : `operator-${sha256(r.name).slice(0, 8)}`;
74 + let slug = base, i = 2;
75 + while ((await db.execute(sql`select 1 from operators where slug = ${slug} and id <> ${r.id}`)).length) slug = `${base}-${i++}`;
76 + console.log(` slug ${r.slug} → ${slug} (${r.name})`);
77 + if (has("--fix-slugs")) await db.execute(sql`update operators set slug = ${slug}, updated_at = now() where id = ${r.id}`);
78 + }
79 + console.log(`\nslugs: ${rows.length} ${has("--fix-slugs") ? "fixed" : "would be fixed"}`);
80 +}
81 +
82 +async function linkCampuses(): Promise<void> {
83 + const db = getDb();
84 + // Containment, not proximity: a building joins a campus only when it is within 400 m AND (same non-null operator, or a
85 + // nameless / operator-less record whose name shares a distinctive token with the campus name). Nearest campus wins.
86 + const rows = await db.execute<{ campus_id: string; campus: string; building_id: string; building: string; m: number; same_op: boolean; campus_operator: string | null }>(sql`
87 + with pairs as (
88 + select c.id as campus_id, c.name as campus, b.id as building_id, b.name as building,
89 + 6371000 * acos(least(1, cos(radians(c.lat)) * cos(radians(b.lat)) * cos(radians(b.lng - c.lng)) + sin(radians(c.lat)) * sin(radians(b.lat)))) as m,
90 + (b.operator_id is not null and b.operator_id = c.operator_id) as same_op,
91 + (select o.name from operators o where o.id = c.operator_id) as campus_operator,
92 + row_number() over (partition by b.id order by 6371000 * acos(least(1, cos(radians(c.lat)) * cos(radians(b.lat)) * cos(radians(b.lng - c.lng)) + sin(radians(c.lat)) * sin(radians(b.lat))))) as rn
93 + from facilities c join facilities b on b.id <> c.id and b.merged_into is null and b.parent_facility_id is null and b.lat is not null
94 + and abs(b.lat - c.lat) < 0.005 and abs(b.lng - c.lng) < 0.007
95 + and b.name !~* '(campus|park|complex|hub|cluster|mega ?site)'
96 + where c.merged_into is null and c.lat is not null and c.name ~* '(campus|park|complex|hub|cluster|mega ?site)')
97 + select campus_id, campus, building_id, building, round(m) as m, same_op, campus_operator from pairs where rn = 1 and m <= 400 order by campus, m`);
98 + const GENERIC = new Set(["data", "center", "centre", "campus", "park", "complex", "hub", "cluster", "site", "mega", "the", "of", "at", "dc", "datacenter", "datacentre", "building", "and", "&", "inc", "ltd", "gmbh", "sa", "llc", "co", "corp", "global", "digital", "cloud", "services", "technologies", "technology", "group", "holdings"]);
99 + const tokens = (s: string) => new Set(s.toLowerCase().normalize("NFKD").replace(/[^a-z0-9 ]/g, " ").split(/\s+/).filter((t) => t.length >= 3 && !GENERIC.has(t)));
100 + // an operator-less record (OSM building tagged only with a name) joins when its name carries the campus operator's name
101 + const keep = rows.filter((r) => { if (r.same_op) return true; if (!r.campus_operator) return false; const a = tokens(r.building), b = tokens(r.campus_operator); return b.size > 0 && [...b].every((t) => a.has(t)); });
102 + for (const r of keep) console.log(` ${r.building.padEnd(50)} → building of ${r.campus} (${r.m} m${r.same_op ? ", same operator" : ", operator name in building name"})`);
103 + if (has("--link-campuses")) {
104 + for (const r of keep) await db.execute(sql`update facilities set parent_facility_id = ${r.campus_id}, record_scope = 'building', updated_at = now() where id = ${r.building_id} and parent_facility_id is null`);
105 + const campusIds = [...new Set(keep.map((r) => r.campus_id))];
106 + if (campusIds.length) await db.execute(sql`update facilities set record_scope = 'campus', updated_at = now() where id in ${campusIds}`);
107 + }
108 + console.log(`\ncampus links: ${keep.length} ${has("--link-campuses") ? "applied" : "would be applied"} (${rows.length - keep.length} nearby pairs rejected: different or unknown operator)`);
109 +}
110 +
111 +async function main(): Promise<void> {
112 + console.log("── projects: announcement class veto");
113 + await hideFalsePositives();
114 + console.log("\n── unbacked figures");
115 + await nullUnbackedMw();
116 + console.log("\n── operator slugs");
117 + await fixSlugs();
118 + console.log("\n── campus containment");
119 + await linkCampuses();
120 + await closeDb();
121 +}
122 +main().catch(async (e) => { console.error(e); await closeDb(); process.exit(1); });
123