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%

Classifier: status-derived project classes need a headline verb + site noun; utility pipeline / contracted load / generation capex scopes

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

10 changed files +422 −51

modified apps/api/src/lib/nearby.ts +1 −3
@@ -5,7 +5,7 @@
5 5 */
6 6 import type { NearbyInfrastructure } from "@dci/core";
7 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";
8 +import { num, reqStr, round, type Row } from "./rows.js";
9 9 import { cloudRegionSummary, facilitySummary, ixpSummary, projectSummary } from "./dto.js";
10 10
11 11 export const NEARBY_TYPES = ["facilities", "projects", "ixps", "cloud_regions", "metros"] as const;
@@ -95,5 +95,3 @@ export function parseNearbyTypes(v: string | undefined): NearbyType[] | undefine
95 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 96 return out.length ? out : undefined;
97 97 }
98 −
99 −export { str as _str };
added apps/api/src/lib/pipeline.ts +207 −0
@@ -0,0 +1,207 @@
1 +/**
2 + * Containment-aware aggregates shared by operator / metro / country payloads, /compare and the dashboard:
3 + * PipelineBreakdown (facility lifecycle buckets + project stages), ExpansionVelocity, market concentration (HHI),
4 + * momentum components and the CoverageRow. Every MW figure is a sum of published site-scoped figures — no estimate.
5 + */
6 +import type { CoverageRow, ExpansionVelocity, MarketConcentration, MarketMomentum, PipelineBreakdown } from "@dci/core";
7 +import { pg, facilityView, knownMwAgg, pipelineMwAgg, countedAgg, hasMwAgg, projectLive, OPERATIONAL_SET, CONSTRUCTION_SET, PLANNED_SET, GRID_EVENT_TYPES, type Fragment, type Sql } from "./sql.js";
8 +import { int, num, reqStr, str, type Row } from "./rows.js";
9 +import { asStatus, round2, share } from "./dto.js";
10 +
11 +/** Scope of an aggregate: a condition on the facility view row `f` and on the project row `p`. */
12 +export interface Scope { facility: Fragment; project: Fragment; event?: Fragment; cloud?: Fragment }
13 +
14 +export function scopeFor(sql: Sql, s: { operatorId?: string; metroId?: string; countryIso2?: string }): Scope {
15 + if (s.operatorId) return { facility: sql`(f.operator_id = ${s.operatorId} or f.owner_id = ${s.operatorId})`, project: sql`p.operator_id = ${s.operatorId}`, event: sql`e.operator_id = ${s.operatorId}`, cloud: sql`r.provider_id = ${s.operatorId}` };
16 + if (s.metroId) return { facility: sql`f.metro_id = ${s.metroId}`, project: sql`p.metro_id = ${s.metroId}`, event: sql`(e.metro_id = ${s.metroId} or (e.entity_type = 'facility' and e.entity_id in (select id from facilities where metro_id = ${s.metroId})))`, cloud: sql`r.metro_id = ${s.metroId}` };
17 + if (s.countryIso2) return { facility: sql`f.country_iso2 = ${s.countryIso2}`, project: sql`p.country_iso2 = ${s.countryIso2}`, event: sql`e.country_iso2 = ${s.countryIso2}`, cloud: sql`r.country_iso2 = ${s.countryIso2}` };
18 + return { facility: sql`true`, project: sql`true`, event: sql`true`, cloud: sql`true` };
19 +}
20 +
21 +export async function pipelineBreakdown(scope: Scope): Promise<PipelineBreakdown> {
22 + const sql = pg();
23 + const known = knownMwAgg(sql), pipe = pipelineMwAgg(sql), counted = countedAgg(sql), hasMw = hasMwAgg(sql);
24 + const [fr, pr] = await Promise.all([
25 + sql<Row[]>`select
26 + count(*) filter (where ${counted})::int as total,
27 + count(*) filter (where ${counted} and ${hasMw})::int as with_mw,
28 + count(*) filter (where ${counted} and f.status = any(${OPERATIONAL_SET}))::int as op_n, sum(${known}) filter (where f.status = any(${OPERATIONAL_SET}))::float as op_mw,
29 + count(*) filter (where ${counted} and f.status = any(${CONSTRUCTION_SET}))::int as con_n, sum(${pipe}) filter (where f.status = any(${CONSTRUCTION_SET}))::float as con_mw,
30 + count(*) filter (where ${counted} and f.status in ('approved', 'permitting'))::int as app_n, sum(${pipe}) filter (where f.status in ('approved', 'permitting'))::float as app_mw,
31 + count(*) filter (where ${counted} and f.status in ('rumored', 'proposed', 'announced', 'delayed'))::int as ann_n, sum(${pipe}) filter (where f.status in ('rumored', 'proposed', 'announced', 'delayed'))::float as ann_mw
32 + from ${facilityView(sql)} f where ${scope.facility}`,
33 + sql<Row[]>`select p.status, count(*)::int as n, sum(p.planned_mw)::float as mw from projects p where ${projectLive(sql)} and ${scope.project} group by p.status order by n desc`,
34 + ]);
35 + const r = fr[0] ?? {};
36 + const total = int(r.total);
37 + return {
38 + operational: { count: int(r.op_n), mw: round2(num(r.op_mw)) },
39 + construction: { count: int(r.con_n), mw: round2(num(r.con_mw)) },
40 + approved: { count: int(r.app_n), mw: round2(num(r.app_mw)) },
41 + announced: { count: int(r.ann_n), mw: round2(num(r.ann_mw)) },
42 + projects: pr.map((x) => ({ status: asStatus(x.status), count: int(x.n), mw: round2(num(x.mw)) })),
43 + mwCoverage: total ? share(int(r.with_mw), total) : 0,
44 + };
45 +}
46 +
47 +const WINDOWS: Array<{ label: "12m" | "3y" | "5y"; months: number }> = [{ label: "12m", months: 12 }, { label: "3y", months: 36 }, { label: "5y", months: 60 }];
48 +
49 +/** Expansion velocity: new facilities (opened_on else first_seen), projects (announced_on else created_at), new countries / metros per window. */
50 +export async function expansionVelocity(scope: Scope): Promise<ExpansionVelocity> {
51 + const sql = pg();
52 + const known = knownMwAgg(sql), counted = countedAgg(sql);
53 + const openedDate = sql`(case when f.opened_on ~ '^\\d{4}-\\d{2}-\\d{2}' then f.opened_on::date when f.opened_on ~ '^\\d{4}-\\d{2}$' then (f.opened_on || '-01')::date when f.opened_on ~ '^\\d{4}$' then (f.opened_on || '-01-01')::date else f.first_seen::date end)`;
54 + const annDate = sql`(case when p.announced_on ~ '^\\d{4}-\\d{2}-\\d{2}' then p.announced_on::date when p.announced_on ~ '^\\d{4}-\\d{2}$' then (p.announced_on || '-01')::date when p.announced_on ~ '^\\d{4}$' then (p.announced_on || '-01-01')::date else p.created_at::date end)`;
55 + const windows = await Promise.all(WINDOWS.map(async (w) => {
56 + const since = sql`(current_date - make_interval(months => ${w.months}))`;
57 + const [fr, pr, geo] = await Promise.all([
58 + sql<Row[]>`select count(*) filter (where ${counted})::int as n, sum(${known})::float as mw from ${facilityView(sql)} f where ${scope.facility} and ${openedDate} >= ${since}`,
59 + sql<Row[]>`select count(*)::int as n, sum(p.planned_mw)::float as mw from projects p where ${projectLive(sql)} and ${scope.project} and ${annDate} >= ${since}`,
60 + sql<Row[]>`with fx as (select f.country_iso2, f.metro_id, min(${openedDate}) as first_d from ${facilityView(sql)} f where ${scope.facility} group by 1, 2)
61 + select count(distinct country_iso2) filter (where country_iso2 is not null and first_d >= ${since} and not exists (select 1 from fx f2 where f2.country_iso2 = fx.country_iso2 and f2.first_d < ${since}))::int as countries,
62 + count(distinct metro_id) filter (where metro_id is not null and first_d >= ${since} and not exists (select 1 from fx f2 where f2.metro_id = fx.metro_id and f2.first_d < ${since}))::int as metros
63 + from fx`,
64 + ]);
65 + return { label: w.label, newFacilities: int(fr[0]?.n), newProjects: int(pr[0]?.n), newCountries: int(geo[0]?.countries), newMetros: int(geo[0]?.metros), openedMw: round2(num(fr[0]?.mw)), announcedMw: round2(num(pr[0]?.mw)) };
66 + }));
67 + const years = await sql<Row[]>`select extract(year from ${openedDate})::int as year, bool_and(f.opened_on ~ '^\\d{4}') as opened_basis,
68 + count(*) filter (where ${counted})::int as n, array_agg(distinct f.country_iso2) filter (where f.country_iso2 is not null) as countries, array_agg(distinct f.metro_id) filter (where f.metro_id is not null) as metros
69 + from ${facilityView(sql)} f where ${scope.facility} group by 1 order by 1`;
70 + const seenC = new Set<string>(), seenM = new Set<string>();
71 + let cum = 0;
72 + const countriesOverTime: ExpansionVelocity["countriesOverTime"] = [];
73 + for (const y of years) {
74 + const year = int(y.year);
75 + if (!year) continue;
76 + for (const c of (y.countries as string[] | null) ?? []) seenC.add(c);
77 + for (const m of (y.metros as string[] | null) ?? []) seenM.add(m);
78 + cum += int(y.n);
79 + countriesOverTime.push({ year, countries: seenC.size, metros: seenM.size, facilities: cum, basis: y.opened_basis === true ? "opened" : "first_seen" });
80 + }
81 + return { windows, countriesOverTime };
82 +}
83 +
84 +export const CONCENTRATION_NOTE = "HHI = Σ (share × 100)² over operators, 0–10 000; 10 000 = one operator. Computed on counted facilities (campus rows with buildings excluded) and, separately, on known operational MW only when published figures exist — `coverage` is the share of facilities with a figure; treat the MW view as partial when coverage is low.";
85 +
86 +export async function concentration(scope: Scope): Promise<MarketConcentration> {
87 + const sql = pg();
88 + const known = knownMwAgg(sql), counted = countedAgg(sql), hasMw = hasMwAgg(sql);
89 + const rows = await sql<Row[]>`select o.id, o.slug, o.name, count(*) filter (where ${counted})::int as n, sum(${known}) filter (where f.status = any(${OPERATIONAL_SET}))::float as mw, count(*) filter (where ${counted} and ${hasMw})::int as with_mw
90 + from ${facilityView(sql)} f join operators o on o.id = f.operator_id where ${scope.facility} group by o.id, o.slug, o.name order by n desc, mw desc nulls last`;
91 + const totalN = rows.reduce((a, r) => a + int(r.n), 0);
92 + const totalCounted = totalN;
93 + const withMw = rows.reduce((a, r) => a + int(r.with_mw), 0);
94 + const totalMw = rows.reduce((a, r) => a + (num(r.mw) ?? 0), 0);
95 + const hhi = (shares: number[]) => Math.round(shares.reduce((a, s) => a + (s * 100) ** 2, 0));
96 + const fShares = rows.map((r) => (totalN ? int(r.n) / totalN : 0));
97 + const facilities = {
98 + hhi: hhi(fShares),
99 + top3Share: share(rows.slice(0, 3).reduce((a, r) => a + int(r.n), 0), totalN),
100 + top5Share: share(rows.slice(0, 5).reduce((a, r) => a + int(r.n), 0), totalN),
101 + top: rows.slice(0, 10).map((r) => ({ id: reqStr(r.id), slug: reqStr(r.slug), name: reqStr(r.name), count: int(r.n), share: share(int(r.n), totalN) })),
102 + };
103 + let knownMw: MarketConcentration["knownMw"] = null;
104 + if (totalMw > 0) {
105 + const byMw = rows.filter((r) => (num(r.mw) ?? 0) > 0).sort((a, b) => (num(b.mw) ?? 0) - (num(a.mw) ?? 0));
106 + knownMw = {
107 + hhi: hhi(byMw.map((r) => (num(r.mw) ?? 0) / totalMw)),
108 + top3Share: share(byMw.slice(0, 3).reduce((a, r) => a + (num(r.mw) ?? 0), 0), totalMw),
109 + top5Share: share(byMw.slice(0, 5).reduce((a, r) => a + (num(r.mw) ?? 0), 0), totalMw),
110 + coverage: share(withMw, totalCounted),
111 + top: byMw.slice(0, 10).map((r) => ({ id: reqStr(r.id), slug: reqStr(r.slug), name: reqStr(r.name), mw: round2(num(r.mw)) ?? 0, share: share(num(r.mw) ?? 0, totalMw) })),
112 + };
113 + }
114 + return { operatorCount: rows.length, facilities, knownMw, note: CONCENTRATION_NOTE };
115 +}
116 +
117 +/** 12-month momentum components (never collapsed into a score). */
118 +export async function momentum(scope: Scope): Promise<MarketMomentum> {
119 + const sql = pg();
120 + const since = sql`(current_date - interval '12 months')`;
121 + const annDate = sql`(case when p.announced_on ~ '^\\d{4}-\\d{2}-\\d{2}' then p.announced_on::date when p.announced_on ~ '^\\d{4}-\\d{2}$' then (p.announced_on || '-01')::date when p.announced_on ~ '^\\d{4}$' then (p.announced_on || '-01-01')::date else p.created_at::date end)`;
122 + const conDate = sql`(case when p.construction_started_on ~ '^\\d{4}-\\d{2}-\\d{2}' then p.construction_started_on::date when p.construction_started_on ~ '^\\d{4}-\\d{2}$' then (p.construction_started_on || '-01')::date when p.construction_started_on ~ '^\\d{4}$' then (p.construction_started_on || '-01-01')::date else null end)`;
123 + const openedDate = sql`(case when f.opened_on ~ '^\\d{4}-\\d{2}-\\d{2}' then f.opened_on::date when f.opened_on ~ '^\\d{4}-\\d{2}$' then (f.opened_on || '-01')::date when f.opened_on ~ '^\\d{4}$' then (f.opened_on || '-01-01')::date else null end)`;
124 + const [pa, pc, fo, ne, cr, ev] = await Promise.all([
125 + sql<Row[]>`select count(*)::int as n, sum(p.planned_mw)::float as mw from projects p where ${projectLive(sql)} and ${scope.project} and ${annDate} >= ${since}`,
126 + sql<Row[]>`select count(*)::int as n, sum(p.planned_mw)::float as mw from projects p where ${projectLive(sql)} and ${scope.project} and (${conDate} >= ${since} or (p.status = 'under_construction' and ${conDate} is null and exists (select 1 from events e where e.project_id = p.id and e.event_type in ('construction_started', 'project_status_changed') and e.new_value::text ilike '%under_construction%' and e.detected_at >= ${since})))`,
127 + sql<Row[]>`select count(*) filter (where ${countedAgg(sql)})::int as n from ${facilityView(sql)} f where ${scope.facility} and f.status = any(${OPERATIONAL_SET}) and ${openedDate} >= ${since}`,
128 + sql<Row[]>`with fx as (select f.operator_id, min(coalesce(${openedDate}, f.first_seen::date)) as first_d from ${facilityView(sql)} f where ${scope.facility} and f.operator_id is not null group by 1)
129 + select o.id, o.slug, o.name from fx join operators o on o.id = fx.operator_id where fx.first_d >= ${since} order by fx.first_d desc limit 20`,
130 + sql<Row[]>`select count(*)::int as n from cloud_regions r where ${scope.cloud ?? sql`true`} and ((r.launched_on ~ '^\\d{4}' and (case when r.launched_on ~ '^\\d{4}-\\d{2}-\\d{2}' then r.launched_on::date when r.launched_on ~ '^\\d{4}-\\d{2}$' then (r.launched_on || '-01')::date else (left(r.launched_on, 4) || '-01-01')::date end) >= ${since}) or r.created_at::date >= ${since})`,
131 + sql<Row[]>`select count(*)::int as total, count(*) filter (where e.event_type = any(${GRID_EVENT_TYPES}))::int as grid from events e where ${scope.event ?? sql`true`} and e.review_status <> 'rejected' and e.detected_at >= ${since}`,
132 + ]);
133 + return {
134 + window: "12m",
135 + projectsAnnounced: int(pa[0]?.n),
136 + projectsEnteredConstruction: int(pc[0]?.n),
137 + constructionMw: round2(num(pc[0]?.mw)),
138 + announcedMw: round2(num(pa[0]?.mw)),
139 + newEntrants: ne.map((r) => ({ id: reqStr(r.id), slug: reqStr(r.slug), name: reqStr(r.name) })),
140 + facilitiesOpened: int(fo[0]?.n),
141 + cloudRegionsAdded: int(cr[0]?.n),
142 + gridEvents: int(ev[0]?.grid),
143 + eventsTotal: int(ev[0]?.total),
144 + };
145 +}
146 +
147 +/** CoverageRow for a scope (containment-aware counts, share of facilities with each field known). */
148 +export async function coverageRow(key: string, name: string, slug: string, scope: Scope): Promise<CoverageRow> {
149 + const sql = pg();
150 + const counted = countedAgg(sql), hasMw = hasMwAgg(sql);
151 + const [fr, pr] = await Promise.all([
152 + sql<Row[]>`select count(*) filter (where ${counted})::int as n,
153 + count(*) filter (where ${counted} and ${hasMw})::int as with_mw,
154 + count(*) filter (where ${counted} and f.operator_id is not null)::int as with_op,
155 + count(*) filter (where ${counted} and f.lat is not null and f.geo_precision in ('exact', 'parcel', 'street'))::int as precise,
156 + count(*) filter (where ${counted} and f.lat is not null)::int as any_loc,
157 + count(*) filter (where ${counted} and f.status <> 'unknown')::int as with_status,
158 + count(*) filter (where ${counted} and f.opened_on ~ '^\\d{4}')::int as with_open,
159 + count(*) filter (where ${counted} and f.source_count >= 2)::int as multi,
160 + count(*) filter (where ${counted} and (coalesce(f.carriers_count, 0) > 0 or coalesce(f.ixp_count, 0) > 0 or exists (select 1 from facility_tenants t where t.facility_id = f.id) or exists (select 1 from facility_ixps x where x.facility_id = f.id)))::int as conn,
161 + count(*) filter (where ${counted} and exists (select 1 from provenance p join sources s on s.id = p.source_id where p.entity_type = 'facility' and p.entity_id = f.id and p.is_current and s.kind in ('operator','government','filing','utility','cloud_provider','registry')))::int as primary_src
162 + from ${facilityView(sql)} f where ${scope.facility}`,
163 + sql<Row[]>`select count(*)::int as n, count(*) filter (where p.lat is not null)::int as with_loc from projects p where ${projectLive(sql)} and ${scope.project}`,
164 + ]);
165 + const r = fr[0] ?? {};
166 + const n = int(r.n);
167 + return {
168 + key, name, slug,
169 + facilities: n,
170 + capacityCoverage: share(int(r.with_mw), n),
171 + operatorCoverage: share(int(r.with_op), n),
172 + preciseLocationCoverage: share(int(r.precise), n),
173 + anyLocationCoverage: share(int(r.any_loc), n),
174 + statusCoverage: share(int(r.with_status), n),
175 + openingDateCoverage: share(int(r.with_open), n),
176 + multiSourceCoverage: share(int(r.multi), n),
177 + projects: int(pr[0]?.n),
178 + projectsWithLocation: int(pr[0]?.with_loc),
179 + connectivityCoverage: share(int(r.conn), n),
180 + primarySourceShare: share(int(r.primary_src), n),
181 + };
182 +}
183 +
184 +/** Aggregates over the containment-aware view grouped by a dimension (for lists / dashboards). */
185 +export interface DimAgg { id: string; slug: string; name: string; countryIso2: string | null; facilities: number; operational: number; construction: number; planned: number; knownMw: number | null; constructionMw: number | null; plannedMw: number | null; withMw: number; operators: number; ai: number; hyperscale: number; coverage: number }
186 +
187 +export function dimAggFromRow(r: Row): DimAgg {
188 + const n = int(r.facilities);
189 + return { id: reqStr(r.id), slug: reqStr(r.slug), name: reqStr(r.name), countryIso2: str(r.country_iso2), facilities: n, operational: int(r.operational), construction: int(r.construction), planned: int(r.planned), knownMw: round2(num(r.known_mw)), constructionMw: round2(num(r.construction_mw)), plannedMw: round2(num(r.planned_mw)), withMw: int(r.with_mw), operators: int(r.operators), ai: int(r.ai), hyperscale: int(r.hyperscale), coverage: share(int(r.with_mw), n) };
190 +}
191 +
192 +/** Column list for a containment-aware aggregate over facilityView `f` (pair with a group by on the dimension). */
193 +export function dimAggCols(sql: Sql): Fragment {
194 + const known = knownMwAgg(sql), pipe = pipelineMwAgg(sql), counted = countedAgg(sql), hasMw = hasMwAgg(sql);
195 + return sql`
196 + count(f.id) filter (where ${counted})::int as facilities,
197 + count(f.id) filter (where ${counted} and f.status = any(${OPERATIONAL_SET}))::int as operational,
198 + count(f.id) filter (where ${counted} and f.status = any(${CONSTRUCTION_SET}))::int as construction,
199 + count(f.id) filter (where ${counted} and f.status = any(${PLANNED_SET}))::int as planned,
200 + sum(${known}) filter (where f.status = any(${OPERATIONAL_SET}))::float as known_mw,
201 + sum(${pipe}) filter (where f.status = any(${CONSTRUCTION_SET}))::float as construction_mw,
202 + sum(${pipe}) filter (where f.status = any(${PLANNED_SET}))::float as planned_mw,
203 + count(f.id) filter (where ${counted} and ${hasMw})::int as with_mw,
204 + count(distinct f.operator_id)::int as operators,
205 + count(f.id) filter (where ${counted} and (f.ai_evidence in ('confirmed', 'likely') or f.is_ai or f.facility_type = 'ai'))::int as ai,
206 + count(f.id) filter (where ${counted} and (f.is_hyperscale or f.facility_type = 'hyperscale'))::int as hyperscale`;
207 +}
modified apps/api/src/lib/quality.ts +1 −3
@@ -6,7 +6,7 @@
6 6 import type { ClaimDTO, DataQualitySummary, EntityHistory, EventDTO, HistoryPoint, ProvenanceDTO } from "@dci/core";
7 7 import { CAPACITY_COLUMN } from "@dci/core";
8 8 import { pg, claimCols, eventCols, eventJoins } from "./sql.js";
9 −import { bool, int, iso, num, reqStr, str, type Row } from "./rows.js";
9 +import { int, iso, num, reqStr, str, type Row } from "./rows.js";
10 10 import { asSeverity, claimDto, eventDto, provenanceDto } from "./dto.js";
11 11
12 12 const MW_FIELDS = new Set(["itCapacityMw", "totalPowerMw", "plannedPowerMw", "utilityCapacityMw", "gridConnectionMw", "ultimateCampusMw", "plannedMw"]);
@@ -144,5 +144,3 @@ export async function dataQualityFor(entityType: string, entityId: string, compl
144 144 pendingDuplicate: dup.length > 0,
145 145 };
146 146 }
147 −
148 −export { bool as _bool };
modified apps/worker/probe-fixture.ts +4 −2
@@ -19,8 +19,10 @@ const ctx = createTestContext(cfg);
19 19 const recs = await connector.extract(ctx, doc);
20 20 const ents = await connector.normalize(ctx, recs);
21 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));
22 +const pick = (e: Record<string, unknown>) => { const keep = ["entityType","key","name","operatorName","campusName","address","city","regionName","countryIso2","postalCode","geo","status","facilityType","itCapacityMw","totalPowerMw","plannedPowerMw","buildingSqm","siteAreaHa","pue","rackCount","tier","renewableClaim","coolingType","certifications","plannedMw","investmentUsd","projectClass","title","publishedAt","eventType","pageType","summary","mentions"]; const o: Record<string, unknown> = {}; for (const k of keep) if (e[k] != null && !(Array.isArray(e[k]) && !(e[k] as unknown[]).length)) o[k] = e[k]; return o; };
23 +console.log(`records=${recs.length} entities=${ents.length} valid=${report.valid} rejected=${report.rejected}`);
24 +for (const e of ents) console.log(JSON.stringify(pick(e as unknown as Record<string, unknown>)));
25 +if (report.issues.length) console.log("ISSUES", JSON.stringify(report.issues));
24 26 if (cfg.kind === "news" || cfg.kind === "government") {
25 27 const c = await articleContent(doc);
26 28 const a = extractAnnouncement(c.title, c.text, { publishedAt: c.published });
modified packages/connectors/src/classify.ts +37 −18
@@ -1,38 +1,56 @@
1 1 import type { EventType, PageType } from "@dci/core";
2 −import { parseAllMw, inferStatus } from "@dci/core";
2 +import { inferStatus } from "@dci/core";
3 +import { projectMwFigures } from "./mw-figures.js";
3 4
4 5 /**
5 − * Deterministic page classification. URL rules first, then title/body keyword rules.
6 + * Deterministic page classification. URL rules first (a facility URL wins over whatever the body says — a facility
7 + * page mentioning "acquired in 2019" stays a facility page), then title/body keyword rules, first match wins.
6 8 * An LLM may be plugged in later as an assistive classifier for "unknown" only.
7 9 */
10 +
11 +/**
12 + * Data-center vocabulary in the languages of the markets we index: English, German (Rechenzentrum/-zentren),
13 + * French (centre de données), Spanish (centro de datos), Portuguese (centro de dados), Dutch (datacentrum/-centra),
14 + * Nordic (datacenter, datasenter, datacentral, serverhall). Used by the relevance gate and the text rules.
15 + */
16 +export const DC_WORD_SRC = "data ?cent(?:er|re)s?|rechenzentr(?:um|en|ums)|centres? de donn[ée]es|centros? de datos|centros? de dados|datacentr(?:um|a)|datacenters?|datasent(?:er|re)|datacentral(?:er)?|serverhal(?:l|len|lar)";
17 +const DC_WORD_RE = new RegExp(`\\b(?:${DC_WORD_SRC})\\b`, "i");
18 +
8 19 const URL_RULES: Array<[RegExp, PageType]> = [
9 − [/\/(data-?cent(er|re)s?|facilities|locations|colocation|ibx|campus(es)?)\/[a-z0-9-]+\/?$/i, "facility_page"],
10 − [/\/(data-?cent(er|re)s?|facilities|locations|colocation)\/?$/i, "facility_index"],
11 − [/\/(news(room)?|press(-releases?)?|media|announcements|blog)\/?$/i, "news_index"],
12 − [/\/(news|press|media|announcements|newsroom)\/./i, "press_release"],
20 + [/\/(data-?cent(er|re)s?|facilities|locations|colocation|ibx|campus(es)?|rechenzentr(um|en)|datacenters?|centres?-de-donnees|centros?-de-datos)\/[a-z0-9-]+\/?$/i, "facility_page"],
21 + [/\/(data-?cent(er|re)s?|facilities|locations|colocation|rechenzentr(um|en)|datacenters?)\/?$/i, "facility_index"],
22 + [/\/(news(room)?|press(-releases?)?|media|announcements|blog|presse|actualites|noticias)\/?$/i, "news_index"],
23 + [/\/(news|press|media|announcements|newsroom|presse|actualites|noticias)\/./i, "press_release"],
13 24 [/\/(regions?|global-infrastructure|locations|about\/locations|infrastructure\/regions)/i, "cloud_region"],
14 − [/\.pdf(\?|$)/i, "planning_document"],
15 25 [/sitemap.*\.xml/i, "sitemap"],
16 26 [/\/(investors?|ir|financials?|sec-filings|annual-report|10-?k|quarterly)/i, "financial_disclosure"],
17 27 [/\/(planning|permits?|zoning|development-applications?|applications?)/i, "planning_document"],
18 28 ];
19 29
30 +/** A PDF is a planning document only when its URL or title says so; other PDFs stay `unknown` for the news parser. */
31 +export const PDF_URL_RE = /\.pdf(\?|$)/i;
32 +export const PLANNING_VOCAB_RE = /\b(planning|permit(s|ting)?|zoning|rezon\w*|application|environmental|eia|site[- ]plan|council|committee|agenda|bebauungsplan|plan local d'urbanisme|permis de construire|licencia de obras)\b/i;
33 +
20 34 const TEXT_RULES: Array<[RegExp, PageType]> = [
21 − [/\b(acquires|acquisition of|to acquire|completes acquisition|agreed to acquire|acquired)\b/i, "acquisition"],
22 − [/\b(expansion|expands|expanding|adds? \d+ ?MW|additional (capacity|phase)|new phase)\b/i, "expansion"],
35 + [/\b(acquires|acquisition of|to acquire|completes acquisition|agreed to acquire|acquired|übernimmt|übernahme|rachète|acquisition de|adquiere|adquisición)\b/i, "acquisition"],
36 + [/\b(expansion|expands|expanding|adds? \d+ ?MW|additional (capacity|phase)|new phase|erweitert|erweiterung|agrandit|extension|ampliación|amplía|uitbreiding)\b/i, "expansion"],
23 37 [/\b(closes?|closure|closing|decommission|shut(ting)? down|exit(s|ing)? (the )?market)\b/i, "closure"],
24 − [/\b(planning (application|permission|commission|board)|rezon(e|ing)|site plan|conditional use|environmental (assessment|impact)|permit(ting)?)\b/i, "planning_document"],
25 − [/\b(under construction|breaks? ground|groundbreaking|topping out|construction (update|progress|milestone)|shell (is )?complete)\b/i, "construction_update"],
26 − [/\b(substation|transmission line|power (purchase|supply) agreement|PPA|grid connection|interconnection agreement|utility|megawatts of power|gigawatt|nuclear|SMR|gas turbine|solar farm|wind farm)\b/i, "power_infrastructure"],
38 + [/\b(planning (application|permission|commission|board)|rezon(e|ing)|site plan|conditional use|environmental (assessment|impact)|permit(ting)?|bebauungsplan|permis de construire)\b/i, "planning_document"],
39 + [/\b(under construction|breaks? ground|groundbreaking|topping out|construction (update|progress|milestone)|shell (is )?complete|spatenstich|im bau|en construction|en construcción|em construção|in aanbouw)\b/i, "construction_update"],
40 + [/\b(substation|transmission line|power (purchase|supply) agreement|PPA|grid connection|interconnection agreement|utility|megawatts of power|gigawatt|nuclear|SMR|gas turbine|solar farm|wind farm|umspannwerk|netzanschluss|raccordement au réseau)\b/i, "power_infrastructure"],
27 41 [/\b(cloud region|availability zones?|launches? (a )?new region|region (now )?(available|generally available|GA)|sovereign cloud)\b/i, "cloud_region"],
28 − [/\b(announces?|unveils?|plans (to build|for)|to build|will build|new (data ?cent(er|re)|campus|facility)|investment of|to invest|\$\d+(\.\d+)? ?(billion|million)|€\d+|£\d+)\b/i, "project_announcement"],
29 − [/\b(press release|for immediate release|newsroom|media contact)\b/i, "press_release"],
42 + [new RegExp(`\\b(announces?|unveils?|plans (to build|for)|to build|will build|new (${DC_WORD_SRC}|campus|facility)|neues? (${DC_WORD_SRC})|nouveau (${DC_WORD_SRC})|nuevo (${DC_WORD_SRC})|novo (${DC_WORD_SRC})|nieuw (${DC_WORD_SRC})|investment of|to invest|investiert|investit|invierte|\\$\\d+(\\.\\d+)? ?(billion|million)|€\\d+|£\\d+)\\b`, "i"), "project_announcement"],
43 + [/\b(press release|for immediate release|newsroom|media contact|pressemitteilung|communiqué de presse|nota de prensa|comunicado de imprensa|persbericht)\b/i, "press_release"],
30 44 [/\b(colocation|carrier-neutral|interconnection|meet-me room|rack|cabinet|cage|cross[- ]connect|N\+1|2N|redundan|PUE|Tier (III|IV|3|4)|uptime institute)\b/i, "facility_page"],
31 45 ];
32 46
33 47 export interface Classification { pageType: PageType; rule: string; eventType?: EventType; mw: number[]; status?: string | null }
34 48
35 49 export function classifyPage(url: string, title: string | null, text: string): Classification {
50 + if (PDF_URL_RE.test(url)) {
51 + const isPlanning = PLANNING_VOCAB_RE.test(url) || (title != null && PLANNING_VOCAB_RE.test(title));
52 + return finish(isPlanning ? "planning_document" : "unknown", isPlanning ? "url:pdf+planning" : "url:pdf", title, text);
53 + }
36 54 for (const [re, pt] of URL_RULES) if (re.test(url)) return finish(pt, `url:${re.source.slice(0, 30)}`, title, text);
37 55 const head = `${title ?? ""}\n${text.slice(0, 4000)}`;
38 56 for (const [re, pt] of TEXT_RULES) if (re.test(head)) return finish(pt, `text:${re.source.slice(0, 30)}`, title, text);
@@ -41,7 +59,8 @@ export function classifyPage(url: string, title: string | null, text: string): C
41 59
42 60 function finish(pageType: PageType, rule: string, title: string | null, text: string): Classification {
43 61 const head = `${title ?? ""}\n${text.slice(0, 6000)}`;
44 − const mw = parseAllMw(head);
62 + // site figures only: "brings the portfolio to 1.4 GW" is not a capacity of this page's subject
63 + const mw = projectMwFigures(head);
45 64 const status = inferStatus(head);
46 65 const eventType = eventTypeFor(pageType, head);
47 66 return { pageType, rule, eventType, mw, status };
@@ -56,14 +75,14 @@ export function eventTypeFor(pageType: PageType, text: string): EventType | unde
56 75 case "planning_document": return /\b(approved|approval granted|granted permission)\b/i.test(text) ? "planning_approved" : "planning_filed";
57 76 case "power_infrastructure": return "power_agreement";
58 77 case "cloud_region": return /\b(announc|plans?|upcoming|coming)/i.test(text) && !/\b(now available|generally available|launched|is live|opens)\b/i.test(text) ? "cloud_region_announced" : "cloud_region_launched";
59 − case "project_announcement": return /\b(invest(ment|s)? of|to invest|\$\d)/i.test(text) && !/\b(data ?cent|campus|MW)\b/i.test(text) ? "investment_announced" : "project_announced";
78 + case "project_announcement": return /\b(invest(ment|s)? of|to invest|\$\d)/i.test(text) && !new RegExp(`\\b(${DC_WORD_SRC}|campus|MW)\\b`, "i").test(text) ? "investment_announced" : "project_announced";
60 79 case "press_release": return "news";
61 80 case "facility_page": return "facility_updated";
62 81 default: return undefined;
63 82 }
64 83 }
65 84
66 −/** Is this text about data centers at all? Used to drop irrelevant newsroom items early. */
85 +/** Is this text about data centers at all (any of the indexed languages)? Used to drop irrelevant newsroom items early. */
67 86 export function isDataCenterRelevant(text: string): boolean {
68 − return /\b(data ?cent(er|re)s?|colocation|hyperscale|cloud region|availability zone|MW|megawatt|campus|interconnection|IBX|server farm|AI infrastructure|GPU cluster|compute capacity|digital infrastructure)\b/i.test(text);
87 + return DC_WORD_RE.test(text) || /\b(colocation|hyperscale|cloud region|availability zone|MW|megawatt|campus|interconnection|IBX|server farm|AI infrastructure|GPU cluster|compute capacity|digital infrastructure|colocación|colocatie|hyperscaler)\b/i.test(text);
69 88 }
modified packages/connectors/src/extract.ts +38 −8
@@ -1,5 +1,5 @@
1 1 import { load, type CheerioAPI } from "cheerio";
2 −import { extractText as pdfExtractText, getDocumentProxy } from "unpdf";
2 +import { getDocumentProxy } from "unpdf";
3 3 import { parseMw, parseAreaSqm, parseAreaHa, parseMoney, parsePartialDate, normalizeStatus, inferFacilityType, countryFromText, cleanText, htmlToText } from "@dci/core";
4 4 import type { SelectorRule } from "./config.js";
5 5 import type { RawDocument } from "./types.js";
@@ -146,13 +146,43 @@ export function publishedDate(html: string): string | null {
146 146 return null;
147 147 }
148 148
149 −export async function pdfText(buf: Buffer, maxPages = 60): Promise<{ text: string; pages: number; title: string | null }> {
150 − const pdf = await getDocumentProxy(new Uint8Array(buf));
151 − const { text, totalPages } = await pdfExtractText(pdf, { mergePages: true });
152 − const meta = await pdf.getMetadata().catch(() => null);
153 − const title = (meta?.info as { Title?: string } | undefined)?.Title ?? null;
154 − const body = Array.isArray(text) ? text.join("\n") : String(text);
155 − return { text: body.split("\f").slice(0, maxPages).join("\n").replace(/[ \t]+/g, " ").trim(), pages: totalPages, title: cleanText(title) };
149 +/** PDFs above this size are not parsed at all (planning documents are a few MB; bigger files are scans or bombs). */
150 +export const PDF_MAX_BYTES = 15 * 1024 * 1024;
151 +/** Hard ceiling on PDF text extraction; a pathological file must never stall a run. */
152 +export const PDF_TIMEOUT_MS = 20_000;
153 +
154 +/**
155 + * Text of the first `maxPages` pages of a PDF (page by page, so a 3 000-page file costs 60 pages), bounded by
156 + * `maxBytes` on input and `timeoutMs` on wall time. Throws on oversize / timeout — callers log and skip the document.
157 + */
158 +export async function pdfText(buf: Buffer, maxPages = 60, opts: { timeoutMs?: number; maxBytes?: number } = {}): Promise<{ text: string; pages: number; title: string | null; truncated: boolean }> {
159 + const maxBytes = opts.maxBytes ?? PDF_MAX_BYTES;
160 + const timeoutMs = opts.timeoutMs ?? PDF_TIMEOUT_MS;
161 + if (buf.length > maxBytes) throw new Error(`pdf too large: ${buf.length} bytes > ${maxBytes}`);
162 + let pdf: Awaited<ReturnType<typeof getDocumentProxy>> | null = null;
163 + const work = async () => {
164 + pdf = await getDocumentProxy(new Uint8Array(buf));
165 + const pages = pdf.numPages;
166 + const n = Math.min(pages, Math.max(1, maxPages));
167 + const parts: string[] = [];
168 + for (let i = 1; i <= n; i++) {
169 + const page = await pdf.getPage(i);
170 + const content = await page.getTextContent();
171 + parts.push(content.items.map((it) => ("str" in it ? it.str + (it.hasEOL ? "\n" : " ") : "")).join(""));
172 + page.cleanup();
173 + }
174 + const meta = await pdf.getMetadata().catch(() => null);
175 + const title = (meta?.info as { Title?: string } | undefined)?.Title ?? null;
176 + return { text: parts.join("\n").replace(/[ \t]+/g, " ").replace(/\n\s*\n+/g, "\n").trim(), pages, title: cleanText(title), truncated: pages > n };
177 + };
178 + let timer: ReturnType<typeof setTimeout> | undefined;
179 + const timeout = new Promise<never>((_, reject) => { timer = setTimeout(() => reject(new Error(`pdf extraction timed out after ${timeoutMs} ms`)), timeoutMs); });
180 + try {
181 + return await Promise.race([work(), timeout]);
182 + } finally {
183 + clearTimeout(timer);
184 + void (pdf as Awaited<ReturnType<typeof getDocumentProxy>> | null)?.destroy().catch(() => undefined);
185 + }
156 186 }
157 187
158 188 export function isPdf(doc: RawDocument): boolean { return /pdf/i.test(doc.contentType ?? "") || doc.body.subarray(0, 5).toString("latin1") === "%PDF-"; }
modified packages/connectors/src/fetchers.ts +81 −8
@@ -112,14 +112,67 @@ export class DirectFetcher implements Fetcher {
112 112 }
113 113 clearTimeout(timer);
114 114 let body = Buffer.concat(chunks);
115 − if (/\.gz$/i.test(current) || (h["content-type"] ?? "").includes("gzip")) { try { body = gunzipSync(body); } catch { /* not gzip */ } }
115 + if (/\.gz$/i.test(current) || (h["content-type"] ?? "").includes("gzip")) {
116 + const r = safeGunzip(body, maxBytes);
117 + if (r.error) return failDoc(urlStr, level, "direct", "too_large", r.error, started, res.status);
118 + body = r.body;
119 + }
116 120 const ct = base.contentType;
117 121 const isText = !ct || /text|json|xml|javascript|html|csv|markdown/i.test(ct);
118 − return { ...base, body, text: isText ? decodeText(body, ct) : "", notModified: false, durationMs: Date.now() - started };
122 + return applyRobotsDirectives({ ...base, body, text: isText ? decodeText(body, ct) : "", notModified: false, durationMs: Date.now() - started });
119 123 }
120 124 }
121 125 }
122 126
127 +/** Decompression cap: an archive may expand to at most max(4 × maxBytes, 64 MB) — anything bigger is a bomb. */
128 +export function gunzipLimit(maxBytes: number): number { return Math.max(maxBytes * 4, 64 * 1024 * 1024); }
129 +export function safeGunzip(body: Buffer, maxBytes = DEFAULT_MAX_BYTES): { body: Buffer; error?: string } {
130 + const maxOutputLength = gunzipLimit(maxBytes);
131 + try {
132 + return { body: gunzipSync(body, { maxOutputLength }) };
133 + } catch (e) {
134 + const code = (e as { code?: string }).code;
135 + if (code === "ERR_BUFFER_TOO_LARGE" || e instanceof RangeError) return { body, error: `gzip output exceeds ${maxOutputLength} bytes` };
136 + return { body }; // not gzip after all: keep the raw bytes
137 + }
138 +}
139 +
140 +/**
141 + * `X-Robots-Tag` / `<meta name="robots">` directives we honour: `noarchive` and `noindex` both mean "do not keep a
142 + * copy of this page" — the pipeline still extracts facts (facts are not copyrightable) but skips the raw archive.
143 + * Sets `doc.meta.noarchive = true` (and `doc.meta.robotsDirective` with the matched token). Idempotent.
144 + */
145 +export function applyRobotsDirectives(doc: RawDocument): RawDocument {
146 + const dir = robotsDirective(doc);
147 + if (dir) doc.meta = { ...(doc.meta ?? {}), noarchive: true, robotsDirective: dir };
148 + return doc;
149 +}
150 +export function robotsDirective(doc: RawDocument): string | null {
151 + const header = doc.headers["x-robots-tag"];
152 + if (header && /\b(noarchive|noindex|none)\b/i.test(header)) return `header:${header.match(/\b(noarchive|noindex|none)\b/i)![1]!.toLowerCase()}`;
153 + if (doc.text && /html/i.test(doc.contentType ?? "")) {
154 + const head = doc.text.slice(0, 200_000);
155 + const re = /<meta\b[^>]*\bname\s*=\s*["']?(?:robots|datacenterindexbot)["']?[^>]*>/gi;
156 + for (const m of head.matchAll(re)) {
157 + const content = m[0].match(/\bcontent\s*=\s*["']([^"']*)["']/i)?.[1] ?? "";
158 + const tok = content.match(/\b(noarchive|noindex|none)\b/i);
159 + if (tok) return `meta:${tok[1]!.toLowerCase()}`;
160 + }
161 + }
162 + return null;
163 +}
164 +
165 +/** `Retry-After` header → milliseconds (delta-seconds or HTTP-date), bounded to [0, 7 days]; null when absent/invalid. */
166 +export function parseRetryAfter(value: string | null | undefined, now = Date.now()): number | null {
167 + if (!value) return null;
168 + const v = value.trim();
169 + const max = 7 * 86_400_000;
170 + if (/^\d+$/.test(v)) return Math.min(max, Number(v) * 1000);
171 + const t = Date.parse(v);
172 + if (Number.isFinite(t)) return Math.min(max, Math.max(0, t - now));
173 + return null;
174 +}
175 +
123 176 export type PremiumProvider = "scrapfly" | "firecrawl";
124 177
125 178 /**
@@ -138,7 +191,7 @@ export function setBudgetStore(store: BudgetStore | null): void { budgetStore =
138 191 export function budgetDay(now = new Date()): string { return now.toISOString().slice(0, 10); }
139 192
140 193 /** Daily budget bookkeeping shared by premium fetchers. */
141 −class Budget {
194 +export class Budget {
142 195 private day = "";
143 196 private local = 0;
144 197 constructor(readonly provider: PremiumProvider, private readonly envKey: string, private readonly fallback: number) {}
@@ -161,13 +214,14 @@ export class FirecrawlFetcher implements Fetcher {
161 214 const started = Date.now();
162 215 if (!this.available()) return failDoc(url, 3, "firecrawl", "firecrawl_unavailable", "no key or daily budget exhausted", started);
163 216 try { await assertUrlAllowed(url); } catch (e) { return failDoc(url, 3, "firecrawl", "ssrf_blocked", (e as Error).message, started); }
164 − firecrawlBudget.spend(1);
165 217 const ac = new AbortController();
166 218 const timer = setTimeout(() => ac.abort(), opts.timeoutMs ?? 90_000);
167 219 try {
168 220 const res = await fetch("https://api.firecrawl.dev/v1/scrape", { method: "POST", signal: ac.signal, headers: { "content-type": "application/json", authorization: `Bearer ${process.env.FIRECRAWL_API_KEY}` }, body: JSON.stringify({ url, formats: ["markdown", "html"], onlyMainContent: false, waitFor: opts.waitForSelector ? 3000 : 0 }) });
169 221 const json = (await res.json()) as { success?: boolean; data?: { markdown?: string; html?: string; metadata?: { statusCode?: number; sourceURL?: string; url?: string; title?: string } }; error?: string };
222 + // the daily budget is debited only when Firecrawl actually delivered a document (failed calls are not billed as a scrape)
170 223 if (!res.ok || !json.success || !json.data) return failDoc(url, 3, "firecrawl", "firecrawl_error", `${res.status} ${json.error ?? "no data"}`, started);
224 + firecrawlBudget.spend(1);
171 225 const html = json.data.html ?? "";
172 226 const body = Buffer.from(html || json.data.markdown || "", "utf8");
173 227 return { url, finalUrl: json.data.metadata?.url ?? json.data.metadata?.sourceURL ?? url, fetchedAt: new Date().toISOString(), status: json.data.metadata?.statusCode ?? 200, contentType: html ? "text/html; charset=utf-8" : "text/markdown", body, text: body.toString("utf8"), headers: {}, etag: null, lastModified: null, notModified: false, fetcher: "firecrawl", level: 3, durationMs: Date.now() - started, credits: 1, markdown: json.data.markdown ?? null };
@@ -201,7 +255,7 @@ export class ScrapflyFetcher implements Fetcher {
201 255 if (!r.success) return { ...failDoc(url, 4, "scrapfly", "scrapfly_failed", `${r.status_code ?? 0} ${r.reason ?? r.error?.message ?? "upstream failed"}`, started, r.status_code ?? 0), credits: cost };
202 256 const h = Object.fromEntries(Object.entries(r.response_headers ?? {}).map(([k, v]) => [k.toLowerCase(), String(v)]));
203 257 const body = Buffer.from(r.content ?? "", "utf8");
204 − return { url, finalUrl: r.url ?? url, fetchedAt: new Date().toISOString(), status: r.status_code ?? 200, contentType: h["content-type"] ?? "text/html; charset=utf-8", body, text: body.toString("utf8"), headers: h, etag: null, lastModified: h["last-modified"] ?? null, notModified: false, fetcher: "scrapfly", level: 4, durationMs: Date.now() - started, credits: cost };
258 + return applyRobotsDirectives({ url, finalUrl: r.url ?? url, fetchedAt: new Date().toISOString(), status: r.status_code ?? 200, contentType: h["content-type"] ?? "text/html; charset=utf-8", body, text: body.toString("utf8"), headers: h, etag: null, lastModified: h["last-modified"] ?? null, notModified: false, fetcher: "scrapfly", level: 4, durationMs: Date.now() - started, credits: cost });
205 259 } catch (e) {
206 260 return failDoc(url, 4, "scrapfly", ac.signal.aborted ? "timeout" : "scrapfly_error", (e as Error).message, started);
207 261 } finally { clearTimeout(timer); }
@@ -223,21 +277,40 @@ export function looksBlockedOrEmpty(doc: RawDocument): boolean {
223 277 return textLen < 400 && /<(div id="(root|app|__next|__nuxt)"|app-root|noscript)/i.test(t);
224 278 }
225 279
280 +/** Expected premium cost of one fetch at a level (credits): Firecrawl 1 scrape; Scrapfly ≈1 with ASP only, ≈6 with JS rendering. */
281 +export function expectedCredits(level: FetchLevel, opts: Pick<FetchOptions, "renderJs"> = {}): number {
282 + if (level === 3) return 1;
283 + if (level === 4) return opts.renderJs === false ? 1 : 6;
284 + return 0;
285 +}
286 +
226 287 /**
227 288 * Escalating fetch: start at `level`, escalate up to `maxLevel` when blocked/empty. Premium fetchers are
228 − * skipped when unavailable. Returns the best document obtained (never throws).
289 + * skipped when unavailable, and when their expected cost exceeds `creditsLeft` (the run's remaining premium budget —
290 + * so one call never spends on L3 and L4 together when the run budget cannot cover both). HTTP 429 never escalates:
291 + * the server asked us to slow down, the document is returned with `meta.retryAfterMs` from `Retry-After` when sent.
292 + * Returns the best document obtained (never throws).
229 293 */
230 −export async function fetchWithEscalation(url: string, opts: FetchOptions & { maxLevel?: FetchLevel } = {}): Promise<RawDocument> {
294 +export async function fetchWithEscalation(url: string, opts: FetchOptions & { maxLevel?: FetchLevel; creditsLeft?: number } = {}): Promise<RawDocument> {
231 295 let level: FetchLevel = opts.level ?? 1;
232 296 const max: FetchLevel = opts.maxLevel ?? 2;
297 + let creditsLeft = opts.creditsLeft;
233 298 let best: RawDocument | null = null;
234 299 const attempts: string[] = [];
235 300 while (level <= max) {
236 301 const f = fetchers[level];
237 302 if (!f.available()) { attempts.push(`L${level}:unavailable`); level = (level + 1) as FetchLevel; continue; }
238 − const doc = await f.fetch(url, { ...opts, level });
303 + const cost = expectedCredits(level, opts);
304 + if (creditsLeft !== undefined && cost > 0 && cost > creditsLeft) { attempts.push(`L${level}:over_budget(${cost}>${creditsLeft})`); level = (level + 1) as FetchLevel; continue; }
305 + const doc = applyRobotsDirectives(await f.fetch(url, { ...opts, level }));
306 + if (creditsLeft !== undefined) creditsLeft = Math.max(0, creditsLeft - doc.credits);
239 307 doc.meta = { ...(doc.meta ?? {}), attempts: [...attempts, `L${level}:${doc.error?.code ?? doc.status}`] };
240 308 if (doc.notModified) return doc;
309 + if (doc.status === 429) {
310 + const retryAfterMs = parseRetryAfter(doc.headers["retry-after"]);
311 + doc.meta = { ...doc.meta, rateLimited: true, ...(retryAfterMs != null ? { retryAfterMs } : {}) };
312 + return doc; // never buy our way past a rate limit
313 + }
241 314 if (!doc.error && !looksBlockedOrEmpty(doc)) return doc;
242 315 // hard failures that escalation cannot fix
243 316 if (doc.error && ["ssrf_blocked", "dns", "too_large"].includes(doc.error.code)) return doc;
modified packages/connectors/src/index.ts +1 −0
@@ -6,6 +6,7 @@ export * from "./ratelimit.js";
6 6 export * from "./discovery.js";
7 7 export * from "./extract.js";
8 8 export * from "./classify.js";
9 +export * from "./mw-figures.js";
9 10 export * from "./registry.js";
10 11 export * from "./generic.js";
11 12 export * from "./testing.js";
added packages/connectors/src/mw-figures.ts +38 −0
@@ -0,0 +1,38 @@
1 +import { parseAllMw } from "@dci/core";
2 +
3 +/**
4 + * MW figures that describe the announced SITE, as opposed to a company's portfolio, a market statistic or a regional
5 + * total. Shared by the page classifier (`classifyPage().mw`) and the news announcement extractor, so a press release
6 + * saying "brings the company's footprint to 1.4 GW" never puts 1 400 MW on a project or a facility.
7 + */
8 +
9 +/** Portfolio / company-wide context BEFORE a MW figure — not the size of the announced site. */
10 +export const PORTFOLIO_BEFORE_RE = /\b(total(?:l?ing)?|across|portfolio|combined|footprint|pipeline|capacity to(?: over| more than)?|to over|to more than|brings?(?: its| the| total)?[^.]{0,60}(?:to|capacity)|bringing[^.]{0,60}to|taking[^.]{0,60}to|now (?:has|operates|owns|manages)|operates?(?: over| more than)?|owns?(?: over| more than)?|manages?(?: over| more than)?|(?:more than|over) \d+ (?:data cent(?:er|re)s?|facilities|sites|campuses)[^.]{0,20}|globally|worldwide|company-?wide|nationwide|in (?:asia|apj|japan|europe|emea|the region|australia|the us|the u\.s\.|the country|india|latin america|the americas|the uk|the middle east)|of (?:total )?capacity (?:in|across))\s*(?:of\s*)?$/i;
11 +/** Portfolio / company-wide context AFTER a MW figure. */
12 +export const PORTFOLIO_AFTER_RE = /^\s*(?:MW|GW|megawatts?|gigawatts?)?\s*(?:across|portfolio|globally|worldwide|company-?wide|nationwide|in (?:asia|apj|japan|europe|emea|the region|australia|the us|the country|india|latin america|the americas|the uk|the middle east)|of (?:total )?(?:capacity |power )?(?:across|in operation|under (?:management|development)|planned)|of (?:contracted|committed|operational|total) (?:capacity|power))\b/i;
13 +const MW_TOKEN_RE = /(\d{1,3}(?:[.,]\d{3})+|\d+(?:[.,]\d+)?)\+?\s?(gigawatts?|gw|megawatts?|mw)\b/gi;
14 +
15 +/** Is the MW token at `start`…`end` of `scope` a portfolio / company-wide figure rather than the site's? */
16 +export function isPortfolioFigure(scope: string, start: number, end: number): boolean {
17 + const before = scope.slice(Math.max(0, start - 100), start);
18 + const after = scope.slice(end, end + 40);
19 + return PORTFOLIO_BEFORE_RE.test(before) || PORTFOLIO_AFTER_RE.test(after);
20 +}
21 +
22 +/**
23 + * MW figures in a scope that describe the announced site (largest first, deduplicated, 0 < MW ≤ 20 000), skipping
24 + * company-wide totals and market statistics. When every figure in the scope is a portfolio total, none is returned —
25 + * the caller moves to the next scope.
26 + */
27 +export function projectMwFigures(scope: string): number[] {
28 + const site: number[] = [];
29 + MW_TOKEN_RE.lastIndex = 0;
30 + for (const m of scope.matchAll(MW_TOKEN_RE)) {
31 + const start = m.index ?? 0;
32 + const value = parseAllMw(m[0])[0];
33 + if (value == null || value <= 0 || value > 20_000) continue;
34 + if (isPortfolioFigure(scope, start, start + m[0].length)) continue;
35 + site.push(value);
36 + }
37 + return [...new Set(site)].sort((a, b) => b - a);
38 +}
modified packages/core/src/claims.ts +14 −9
@@ -110,8 +110,8 @@ export function authorityTier(i: AuthorityInput): AuthorityTier {
110 110 }
111 111
112 112 // ─── scope classification ──────────────────────────────────────────────────────────────────────────
113 −const PORTFOLIO_RE = /\b(across (?:its|our|their|the company'?s?|a) (?:global |entire |whole |growing |existing )?(?:portfolio|footprint|platform|network|fleet|estate)|portfolio|footprint|under management|under development globally|(?:global|total|combined|aggregate|cumulative) (?:capacity|footprint|pipeline|portfolio)|pipeline of|development pipeline|(?:brings?|bringing|takes?|taking|lifts?|lifting|grows?|growing) (?:its|the company'?s?|our|their|total) [^.]{0,60}(?:capacity |footprint |portfolio )?to (?:over |more than |nearly |approximately |about |roughly )?\d|now (?:operates|owns|manages|has) (?:over |more than |nearly )?\d|operates (?:over |more than |nearly )?\d+[\d.,]* ?(?:mw|gw|megawatts?|gigawatts?) (?:across|in|of)|(?:more than|over) \d+ (?:data cent(?:er|re)s|facilities|sites|campuses|locations)|globally|worldwide|company-?wide|enterprise-?wide|group-?wide)\b/i;
114 −const COMPANY_RE = /\b(capex|capital expenditure|will invest [^.]{0,40}(?:over the next|through|by) (?:the next )?(?:\d+ years|20\d\d)|(?:through|by|over the next|over the coming) (?:20\d\d|\d+ years|the decade)|(?:total|planned|annual) investment (?:of|in) [^.]{0,30}(?:across|in) (?:the (?:us|uk|country|region)|[A-Z][a-z]+ and [A-Z][a-z]+)|(?:in|for) (?:its|our|their) (?:fiscal|financial) (?:year|quarter)|(?:fiscal|fy) ?20\d\d|guidance|earnings|revenue|quarterly results|annual report|10-k|market (?:to|will|set to|expected to|projected to) (?:surpass|reach|hit|grow|exceed|top)|market size|cagr|forecast period)\b/i;
113 +const PORTFOLIO_RE = /\b(across (?:its|our|their|the company'?s?|a) (?:global |entire |whole |growing |existing )?(?:portfolio|footprint|platform|network|fleet|estate)|portfolio|footprint|pipeline|(?:contracted|committed|signed|operating|existing|total) (?:large[- ])?(?:data cent(?:er|re) )?load|large[- ]load agreements?|(?:agreements?|contracts?) total(?:l?ing)?|under management|under development globally|(?:global|total|combined|aggregate|cumulative) (?:capacity|footprint|pipeline|portfolio)|pipeline of|development pipeline|(?:brings?|bringing|takes?|taking|lifts?|lifting|grows?|growing) (?:its|the company'?s?|our|their|total) [^.]{0,60}(?:capacity |footprint |portfolio )?to (?:over |more than |nearly |approximately |about |roughly )?\d|now (?:operates|owns|manages|has) (?:over |more than |nearly )?\d|operates (?:over |more than |nearly )?\d+[\d.,]* ?(?:mw|gw|megawatts?|gigawatts?) (?:across|in|of)|(?:more than|over) \d+ (?:data cent(?:er|re)s|facilities|sites|campuses|locations)|globally|worldwide|company-?wide|enterprise-?wide|group-?wide)\b/i;
114 +const COMPANY_RE = /\b(capex|capital (?:expenditure|plan|program|programme|spending)|generation investments?|(?:potential |planned |total )?spending (?:exceeding|of|plan)|rate base|will invest [^.]{0,40}(?:over the next|through|by) (?:the next )?(?:\d+ years|20\d\d)|(?:through|by|over the next|over the coming) (?:20\d\d|\d+ years|the decade)|(?:total|planned|annual) investment (?:of|in) [^.]{0,30}(?:across|in) (?:the (?:us|uk|country|region)|[A-Z][a-z]+ and [A-Z][a-z]+)|(?:in|for) (?:its|our|their) (?:fiscal|financial) (?:year|quarter)|(?:fiscal|fy) ?20\d\d|guidance|earnings|revenue|quarterly results|annual report|10-k|market (?:to|will|set to|expected to|projected to) (?:surpass|reach|hit|grow|exceed|top)|market size|cagr|forecast period)\b/i;
115 115 const COUNTRY_RE = /\b(nationwide|nation-?wide|national (?:program|programme|plan|strategy|investment|pipeline)|across the country|the country'?s|country-?wide|(?:in|across) (?:the )?(?:united states|u\.s\.|us|uk|united kingdom|europe|emea|apac|apj|asia|asia-pacific|latin america|the americas|the middle east|africa|india|japan|germany|france|australia|canada|brazil|malaysia|singapore|the netherlands|ireland|spain|italy|the nordics|scandinavia|the region)\b)/i;
116 116 const METRO_RE = /\b((?:in|across) the (?:[A-Z][\w.]+ )?(?:market|metro|region|area|corridor|cluster)|(?:the )?(?:northern virginia|dfw|dallas-fort worth|silicon valley|greater [A-Z][a-z]+|the bay area) (?:market|region|area)|market'?s (?:total|inventory|capacity)|inventory|absorption|vacancy)\b/i;
117 117 const CAMPUS_RE = /\b(campus|campuses|park|complex|estate|hub|gigafactory|(?:at|when|upon|once) (?:full(?:y)? )?(?:built[- ]out|build[- ]out|buildout|complete|completed|completion)|full build-?out|ultimate(?:ly)?|(?:total|overall|eventual) (?:site|campus|planned) capacity|master-?plan(?:ned)?|(?:across|over|in) (?:\w+ )?(?:phases|buildings|data halls)|multi-?building|multi-?phase|acre (?:site|campus|development)|the (?:site|development|project) (?:will|would|could) (?:deliver|provide|offer|support|house|host|total|reach)|entire (?:site|development|project))\b/i;
@@ -171,7 +171,7 @@ export function classifyCapacitySemantics(context: string | null | undefined): C
171 171 }
172 172
173 173 // ─── investment semantics ──────────────────────────────────────────────────────────────────────────
174 −const MULTI_YEAR_RE = /\b((?:through|by|over the next|over the coming|until|to) (?:20\d\d|\d+ years|the decade|the end of the decade)|multi-?year|\d+-year (?:plan|investment|programme|program|commitment)|(?:annual|yearly) (?:capex|capital expenditure|investment)|capex|capital expenditure)\b/i;
174 +const MULTI_YEAR_RE = /\b((?:through|by|over the next|over the coming|until|to) (?:20\d\d|\d+ years|the decade|the end of the decade|the mid-?20\d\ds)|multi-?year|generation investments?|capital (?:plan|program|programme)|(?:potential |planned )?spending exceeding|\d+-year (?:plan|investment|programme|program|commitment)|(?:annual|yearly) (?:capex|capital expenditure|investment)|capex|capital expenditure)\b/i;
175 175 const COUNTRY_PROGRAM_RE = /\b((?:national|country|sovereign|government|federal|state) (?:program|programme|plan|initiative|strategy|fund|investment)|(?:invest|investment|commit\w*|pledge\w*) [^.]{0,60}(?:in|across) (?:the )?(?:united states|u\.s\.|us|uk|united kingdom|europe|india|japan|germany|france|australia|canada|brazil|malaysia|singapore|the country|the nation|the region|emea|apac)\b)/i;
176 176 const DEAL_RE = /\b(acqui(?:re|res|red|sition)|takeover|buyout|merger|stake|valuation|valued at|purchase price|sale of|sells?|sold|buys?|bought|financing|refinanc\w+|loan|credit facility|bond|notes|debt|equity|raises?|raised|funding round|series [a-e]|lease(?:d|s)? (?:valued|worth)|contract (?:valued|worth))\b/i;
177 177 const CAMPUS_INV_RE = /\b(campus|site|park|complex|development|project|phase|build-?out|the facility|the data cent(?:er|re))\b/i;
@@ -395,12 +395,17 @@ export function classifyProjectEvent(i: ProjectClassInput): ProjectClassificatio
395 395 if (re.test(lead)) { cls = k; reasons.push(`lead:${k}`); break; }
396 396 }
397 397 }
398 − if (cls === "UNKNOWN" && i.status && /^(announced|proposed|rumored)$/.test(i.status) && /\b(data ?cent(?:er|re)|campus|facility)\b/i.test(`${title} ${lead}`)) { cls = "NEW_BUILD"; reasons.push("status:announced→NEW_BUILD"); }
399 − if (cls === "UNKNOWN" && i.status === "under_construction") { cls = "CONSTRUCTION_START"; reasons.push("status:construction"); }
400 − if (cls === "UNKNOWN" && (i.status === "permitting" || i.status === "approved")) { cls = "PERMIT"; reasons.push(`status:${i.status}`); }
401 − if (cls === "UNKNOWN" && i.status === "expansion") { cls = "EXPANSION"; reasons.push("status:expansion"); }
402 −
403 − const verb = DEVELOPMENT_VERB_RE.test(title) || (!!lead && DEVELOPMENT_VERB_RE.test(lead));
398 + // status-derived fallbacks need the headline itself to talk about building / a site ("Southern's 17 GW pipeline puts AI power
399 + // demand into utility math" has an "announced" status word somewhere in the body but is an analysis piece)
400 + const titleVerb = DEVELOPMENT_VERB_RE.test(title);
401 + const titleSite = /\b(data ?cent(?:er|re)s?|campus(?:es)?|facility|facilities|ai factory|hub|site)\b/i.test(title);
402 + if (cls === "UNKNOWN" && i.status && /^(announced|proposed|rumored)$/.test(i.status) && titleVerb && titleSite) { cls = "NEW_BUILD"; reasons.push("status:announced→NEW_BUILD"); }
403 + if (cls === "UNKNOWN" && i.status === "under_construction" && titleSite) { cls = "CONSTRUCTION_START"; reasons.push("status:construction"); }
404 + if (cls === "UNKNOWN" && (i.status === "permitting" || i.status === "approved") && titleSite) { cls = "PERMIT"; reasons.push(`status:${i.status}`); }
405 + if (cls === "UNKNOWN" && i.status === "expansion" && titleSite) { cls = "EXPANSION"; reasons.push("status:expansion"); }
406 +
407 + const classFromTitle = reasons.some((r) => r.startsWith("title:"));
408 + const verb = titleVerb || (classFromTitle && !!lead && DEVELOPMENT_VERB_RE.test(lead));
404 409 const name = !!i.hasExplicitName;
405 410 const operator = !!i.hasOperator;
406 411 const location = !!i.hasLocation;
407 412