spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Project reconciliation + persistence (docs/PROJECT-EXTRACTION.md, docs/CLAIMS.md).3 *4 * - The announcement CLASS decides first: appointments, financing, PPAs, partnerships, customer deals, market research5 * never create a project. Non-physical classes attach a timeline row + event to an existing, identifiable project.6 * - Evidence threshold: a new project needs (explicit name or operator) + location + a development verb.7 * - Capacity and investment go through the claim store (scope + sanity engine); a company-wide or portfolio figure is8 * kept as a claim and never written to the project's columns.9 * - Lifecycle transitions follow the state machine; a backward move needs a source that outranks the stored one.10 * - Projects without coordinates are geocoded at CITY / METRO level (never more precise) from the metro seeds.11 */12import { sql } from "@dci/db";13import {14 ASSOCIATED_CLASSES,15 PHYSICAL_CLASSES,16 cleanText,17 classifyAiEvidence,18 classifyInvestmentSemantics,19 classifyScope,20 findEvidence,21 newId,22 normalizeName,23 projectTransition,24 sha256,25 validLatLng,26 type CapacityPredicate,27 type ClaimScope,28 type ConfidenceLevel,29 type EventType,30 type InvestmentPredicate,31 type NormalizedProject,32 type ProjectClass,33} from "@dci/core";34import { addRef, authority, bump, provenanceFor, safeCountry, uniqueSlug, type IngestContext, type Tx } from "./common.js";35import { emitDiffEvents, recordEvent, TRACKED_PROJECT_FIELDS } from "./events.js";36import { validGeo } from "./geo.js";37import { geocodeCity } from "./geocode.js";38import { entityIdForKey, facilityIdForKey, upsertKey } from "./keys.js";39import { isPipelineStatus, shouldReplace, shouldReplaceGeo, type FieldObservation } from "./match.js";40import { assignMetro } from "./metros.js";41import { operatorNames, resolveOperator } from "./operators.js";42import { backingObservation, loadCurrentProvenance, writeProvenance, type CurrentProvenance, type ObservedField } from "./provenance.js";43import { markWinners, recordCapacityClaim, recordInvestmentClaim, reviewPriority, writeQualityFlags, writeTextClaim } from "./claims.js";44import { resolveCampus } from "./campuses.js";45import { maybeOperatorExpansion } from "./facilities.js";4647type Row = Record<string, unknown>;4849const COLS: Record<string, string> = {50 name: "name",51 operatorId: "operator_id",52 facilityId: "facility_id",53 metroId: "metro_id",54 countryIso2: "country_iso2",55 city: "city",56 regionName: "region_name",57 lat: "lat",58 lng: "lng",59 geoPrecision: "geo_precision",60 status: "status",61 announcedOn: "announced_on",62 expectedOpening: "expected_opening",63 plannedMw: "planned_mw",64 investmentUsd: "investment_usd",65 investmentCurrency: "investment_currency",66 investmentOriginal: "investment_original",67 acreage: "acreage",68 phaseCount: "phase_count",69 isAi: "is_ai",70 description: "description",71 sourceUrl: "source_url",72 confidence: "confidence",73 externalIds: "external_ids",74 projectClass: "project_class",75 evidenceLevel: "evidence_level",76 aiEvidence: "ai_evidence",77 capacityScope: "capacity_scope",78 capacitySemantics: "capacity_semantics",79 investmentScope: "investment_scope",80 investmentSemantics: "investment_semantics",81 developerId: "developer_id",82 tenantId: "tenant_id",83 campusId: "campus_id",84 constructionStartedOn: "construction_started_on",85 approvedOn: "approved_on",86 permitFiledOn: "permit_filed_on",87 openedOn: "opened_on",88 reviewPriority: "review_priority",89 hidden: "hidden",90};9192/** External-id namespaces that identify ONE project (allowlist). */93export const PROJECT_IDENTIFYING_KEYS: ReadonlySet<string> = new Set(["wikidata", "planning_ref", "permit_id", "case_number", "application_id", "docket", "planning_application", "rezoning_case", "eia_ref"]);9495/** Timeline / event type for a non-physical (associated) class. */96const ASSOCIATED_EVENT: Record<string, EventType> = { POWER_AGREEMENT: "power_agreement", FINANCING: "investment_announced", ACQUISITION: "acquisition", PARTNERSHIP: "partnership", CUSTOMER_AGREEMENT: "customer_agreement", GRID_CONNECTION: "grid_connection", LAND_ACQUISITION: "land_acquired", PERMIT: "planning_filed", CONSTRUCTION_START: "construction_started", EXPANSION: "expansion_announced", NEW_BUILD: "project_announced" };9798async function loadProject(tx: Tx, id: string): Promise<Row | null> {99 const r = (await tx.execute(sql`select * from projects where id = ${id}`))[0];100 if (!r) return null;101 const out: Row = { id: r.id, slug: r.slug };102 for (const [camel, col] of Object.entries(COLS)) out[camel] = r[col] ?? null;103 return out;104}105106async function followMerged(tx: Tx, id: string): Promise<string> {107 let cur = id;108 for (let i = 0; i < 5; i++) {109 const r = await tx.execute(sql`select merged_into from projects where id = ${cur}`);110 const m = r[0]?.merged_into;111 if (!m) break;112 cur = String(m);113 }114 return cur;115}116117async function resolveProject(tx: Tx, ctx: IngestContext, p: NormalizedProject, operatorId: string | null, country: string | null): Promise<string | null> {118 const byKey = await entityIdForKey(tx, p.key, "project");119 if (byKey) return followMerged(tx, byKey);120 if (p.externalIds) {121 for (const [k, v] of Object.entries(p.externalIds)) {122 if (v == null || v === "" || !PROJECT_IDENTIFYING_KEYS.has(k)) continue;123 const r = await tx.execute(sql`select id from projects where merged_into is null and external_ids @> ${JSON.stringify({ [k]: v })}::jsonb limit 1`);124 if (r[0]) return String(r[0].id);125 }126 }127 const src = p.sourceUrl ?? null;128 if (src) {129 const r = await tx.execute(sql`select id from projects where merged_into is null and source_url = ${src} limit 1`);130 if (r[0]) return followMerged(tx, String(r[0].id));131 }132 const norm = normalizeName(p.name);133 if (!norm) return null;134 const rows = await tx.execute(sql`135 select id, similarity(normalized_name, ${norm}) as sim from projects136 where merged_into is null137 and (${country}::text is null or country_iso2 is null or country_iso2 = ${country})138 and (${operatorId}::text is null or operator_id is null or operator_id = ${operatorId})139 and (normalized_name = ${norm} or similarity(normalized_name, ${norm}) >= 0.85)140 order by (operator_id = ${operatorId}) desc nulls last, sim desc limit 1`);141 if (rows[0]) return String(rows[0].id);142 return sameAnnouncement(tx, ctx, p, operatorId, country);143}144145/** Announcement window: two outlets reporting the same project within this many days are one project. */146export const ANNOUNCEMENT_WINDOW_DAYS = 30;147148/**149 * The same announcement covered by several outlets: same operator (or, without an operator, the same city), planned MW150 * within ±10 % and announced within 30 days → one project. Without an operator the city is mandatory.151 */152async function sameAnnouncement(tx: Tx, ctx: IngestContext, p: NormalizedProject, operatorId: string | null, country: string | null): Promise<string | null> {153 if (p.plannedMw == null || p.plannedMw <= 0 || !p.announcedOn) return null;154 const day = p.announcedOn.length >= 10 ? p.announcedOn.slice(0, 10) : null;155 if (!day) return null;156 const city = p.city ? normalizeName(p.city) : null;157 if (!operatorId && !city) return null;158 const rows = await tx.execute(sql`159 select id from projects160 where merged_into is null and planned_mw is not null161 and planned_mw between ${p.plannedMw * 0.9} and ${p.plannedMw * 1.1}162 and announced_on is not null and length(announced_on) >= 10163 and abs(announced_on::date - ${day}::date) <= ${ANNOUNCEMENT_WINDOW_DAYS}164 and (${country}::text is null or country_iso2 is null or country_iso2 = ${country})165 and (${operatorId}::text is null or operator_id is null or operator_id = ${operatorId})166 and (${city}::text is null or city is null or lower(regexp_replace(city, '[^A-Za-z0-9]+', ' ', 'g')) = ${city})167 and (${operatorId}::text is not null or (${city}::text is not null and city is not null))168 order by created_at asc limit 1`);169 if (!rows[0]) return null;170 ctx.stats.projectDedup = (ctx.stats.projectDedup ?? 0) + 1;171 return String(rows[0].id);172}173174function obs(value: unknown, field: string, p: NormalizedProject, ctx: IngestContext): FieldObservation {175 const pv = provenanceFor(p.provenance, field, undefined);176 return { value, sourceKind: ctx.run.sourceKind, confidence: pv.confidence, isEstimate: !!pv.isEstimate, observedAt: pv.lastObserved || ctx.now, sourceId: pv.sourceId || ctx.run.sourceId, url: pv.url || ctx.doc?.url || null };177}178function stored(prov: CurrentProvenance[], field: string, current: unknown): FieldObservation | null {179 if (current == null) return null;180 const b = backingObservation(prov, field, current);181 return { value: current, sourceKind: b?.sourceKind ?? null, confidence: b?.confidence ?? null, isEstimate: b?.isEstimate ?? false, observedAt: b?.lastObserved ?? null, sourceId: b?.sourceId ?? null, url: b?.url ?? null };182}183184const isCampusName = (s: string | null | undefined) => /\b(campus|park|complex|hub|cluster|gigafactory|estate)\b/i.test(s ?? "");185186export async function ingestProject(tx: Tx, ctx: IngestContext, p: NormalizedProject): Promise<void> {187 const name = cleanText(p.name);188 if (!name) throw new Error(`project ${p.key}: name is required`);189 const url = p.provenance.url || p.sourceUrl || ctx.doc?.url || "";190 if (!url) throw new Error(`project ${p.key}: provenance url is required`);191192 const cls = (p.projectClass ?? null) as ProjectClass | null;193 const physical = cls == null || PHYSICAL_CLASSES.has(cls);194 const associated = cls != null && ASSOCIATED_CLASSES.has(cls);195196 const operator = p.operatorName ? await resolveOperator(tx, ctx, { name: p.operatorName }) : null;197 const country0 = await safeCountry(tx, ctx, p.countryIso2);198 let geo = validGeo(p.geo) ? p.geo : null;199 let geoMethod = geo ? "source" : null;200 if (!geo && p.city) {201 const g = await geocodeCity(tx, { city: p.city, regionName: p.regionName, countryIso2: country0 });202 if (g) { geo = { lat: g.lat, lng: g.lng, precision: g.precision, source: g.source }; geoMethod = g.source; }203 }204 const metro = await assignMetro(tx, { lat: geo?.lat, lng: geo?.lng, city: p.city, countryIso2: country0 });205 const country = country0 ?? metro.countryIso2;206 const facilityId = p.facilityKey ? await facilityIdForKey(tx, ctx, p.facilityKey) : null;207208 const existingId = await resolveProject(tx, ctx, p, operator?.id ?? null, country);209 const existing = existingId ? await loadProject(tx, existingId) : null;210211 // ─── veto: non-physical announcements never create a project; weak evidence never creates a project ─────────────212 if (!existing) {213 if (!physical) { ctx.stats.projectsVetoed = (ctx.stats.projectsVetoed ?? 0) + 1; return; }214 if (p.evidenceLevel === "none") { ctx.stats.projectsVetoed = (ctx.stats.projectsVetoed ?? 0) + 1; return; }215 }216 // ─── associated class on an existing project: timeline row + event only, no field changes ─────────────────────────217 if (existing && !physical) {218 const evType = (cls && ASSOCIATED_EVENT[cls]) || "news";219 const date = p.announcedOn ?? ctx.day;220 const desc = cleanText(p.description) ?? name;221 const tid = `ptl_${sha256(`${existing.id}|${date}|${evType}|${desc.toLowerCase()}`).slice(0, 20)}`;222 const ins = await tx.execute(sql`insert into project_timeline (id, project_id, event_date, event_type, description, source_id, document_id, url)223 values (${tid}, ${existing.id}, ${date}, ${evType}, ${desc.slice(0, 400)}, ${ctx.run.sourceId}, ${ctx.doc?.documentId ?? null}, ${url}) on conflict (id) do nothing returning id`);224 if (ins.length && associated) {225 await recordEvent(tx, ctx, { entityType: "project", entityId: String(existing.id), eventType: evType, title: `${String(existing.name)}: ${cls!.toLowerCase().replace(/_/g, " ")} — ${name}`.slice(0, 300), summary: desc.slice(0, 400), newValue: { class: cls, title: name, mw: p.plannedMw ?? null, investmentUsd: p.investmentUsd ?? null }, significance: cls === "POWER_AGREEMENT" || cls === "GRID_CONNECTION" ? 60 : cls === "FINANCING" ? 55 : 45, confidence: p.provenance.confidence, effectiveDate: p.announcedOn ?? null, url, countryIso2: existing.countryIso2 as string | null, operatorId: existing.operatorId as string | null, metroId: existing.metroId as string | null, projectId: String(existing.id) });226 await tx.execute(sql`update projects set last_update = ${ctx.now} where id = ${existing.id}`);227 }228 // money attached to a deal / financing headline is a claim about the project, never its investment column229 if (p.investmentUsd != null && p.investmentUsd > 0) {230 const sem = classifyInvestmentSemantics(p.claimContext?.investmentUsd ?? name);231 await recordInvestmentClaim(tx, ctx, "project", String(existing.id), { value: p.investmentUsd, currency: p.investmentCurrency ?? "USD", predicate: cls === "FINANCING" || cls === "ACQUISITION" ? "deal_value_usd" : (sem.predicate as InvestmentPredicate), scope: cls === "FINANCING" || cls === "ACQUISITION" ? "company" : sem.scope, evidence: p.claimContext?.investmentUsd ? findEvidence(p.claimContext.investmentUsd, p.investmentUsd, "usd") : null, recordScope: "project", publishedAt: p.announcedOn ?? null, provenance: p.provenance });232 }233 addRef(ctx, "project", String(existing.id));234 bump(ctx, "project");235 return;236 }237238 const prov = existing ? await loadCurrentProvenance(tx, "project", existingId!) : [];239 const id = existing ? String(existing.id) : newId("project");240 // a project hidden by the automated repair re-qualifies when a re-extraction yields a physical class with strong evidence241 // (human decisions — resolved_by <> 'quality-apply' / 'system' — are never reversed)242 if (existing && existing.hidden && physical && p.evidenceLevel === "strong" && !ctx.run.dryRun) {243 const auto = await tx.execute(sql`select 1 from quality_flags where entity_type = 'project' and entity_id = ${id} and code = 'project_false_positive' and resolved_by in ('quality-apply', 'system') limit 1`);244 const human = await tx.execute(sql`select 1 from quality_flags where entity_type = 'project' and entity_id = ${id} and code = 'project_false_positive' and resolved_by not in ('quality-apply', 'system') limit 1`);245 if (auto.length && !human.length) {246 await tx.execute(sql`update projects set hidden = false, updated_at = now() where id = ${id}`);247 await tx.execute(sql`update quality_flags set resolution = 'auto: re-qualified by re-extraction (physical class, strong evidence)', resolved_at = now() where entity_type = 'project' and entity_id = ${id} and code = 'project_false_positive'`);248 existing.hidden = false;249 ctx.stats.projectsRequalified = (ctx.stats.projectsRequalified ?? 0) + 1;250 }251 }252 const before: Row = existing ? { ...existing } : {};253 const next: Row = existing ? { ...existing } : { id, status: "announced", geoPrecision: "unknown", isAi: false, confidence: "moderate", externalIds: {}, aiEvidence: "unknown", hidden: false, reviewPriority: 0 };254 const observed: ObservedField[] = [];255 const campusDesignation = isCampusName(name) || isCampusName(p.campusName);256257 const scalars: Record<string, unknown> = {258 name,259 city: cleanText(p.city),260 regionName: cleanText(p.regionName),261 announcedOn: p.announcedOn ?? null,262 expectedOpening: p.expectedOpening ?? null,263 acreage: p.acreage ?? null,264 phaseCount: p.phaseCount ?? null,265 description: cleanText(p.description),266 sourceUrl: p.sourceUrl ?? null,267 constructionStartedOn: p.constructionStartedOn ?? null,268 approvedOn: p.approvedOn ?? null,269 permitFiledOn: p.permitFiledOn ?? null,270 countryIso2: country,271 };272 for (const [field, v] of Object.entries(scalars)) {273 if (v == null) continue;274 observed.push({ field, value: v, provenance: p.provenance });275 if (shouldReplace(obs(v, field, p, ctx), stored(prov, field, existing?.[field])).replace) next[field] = v;276 }277278 // ─── status: lifecycle state machine ────────────────────────────────────────────────────────────────────────────────279 const incomingStatus = p.status && p.status !== "unknown" ? p.status : null;280 if (incomingStatus) {281 observed.push({ field: "status", value: incomingStatus, provenance: p.provenance });282 const verdict = projectTransition(existing?.status as string | null, incomingStatus);283 if (verdict === "forward" || verdict === "side" || verdict === "resume" || (verdict === "same" && !existing)) next.status = incomingStatus;284 else if (verdict === "backward") {285 // a backward move is accepted only from a source that outranks the stored one (a government filing correcting a news story)286 const st = stored(prov, "status", existing?.status);287 if (authority(ctx.run.sourceKind, p.provenance.confidence, false) > authority(st?.sourceKind, st?.confidence, false)) next.status = incomingStatus;288 else await writeQualityFlags(tx, ctx, "project", id, [{ code: "status_backward", severity: "warn", message: `source says ${incomingStatus} but project is ${String(existing?.status)} — kept, lower-authority source`, field: "status" }]);289 } else if (verdict === "invalid") {290 await writeQualityFlags(tx, ctx, "project", id, [{ code: "status_invalid_transition", severity: "warn", message: `invalid lifecycle transition ${String(existing?.status)} → ${incomingStatus}`, field: "status" }]);291 }292 // stage dates derived from the announcement that moved the stage293 const day = p.announcedOn ?? null;294 if (day) {295 if (incomingStatus === "under_construction" && !next.constructionStartedOn) next.constructionStartedOn = day;296 if (incomingStatus === "approved" && !next.approvedOn) next.approvedOn = day;297 if (incomingStatus === "permitting" && !next.permitFiledOn) next.permitFiledOn = day;298 if ((incomingStatus === "operational" || incomingStatus === "partially_operational") && !next.openedOn) next.openedOn = day;299 }300 await writeTextClaim(tx, ctx, "project", id, "status", incomingStatus, p.provenance, { scope: campusDesignation ? "campus" : "facility", publishedAt: p.announcedOn ?? null });301 }302303 // ─── capacity through the claim store ───────────────────────────────────────────────────────────────────────────────304 if (p.plannedMw != null && p.plannedMw > 0) {305 const context = p.claimContext?.plannedMw ?? null;306 const scopeHint: ClaimScope = campusDesignation ? "campus" : "facility";307 const sc = p.capacityScope ? { scope: p.capacityScope as ClaimScope, reason: "extractor" } : classifyScope(context, context ? null : scopeHint);308 const predicate = (p.capacitySemantics ?? "planned_power_mw") as CapacityPredicate;309 const evidence = context ? findEvidence(context, p.plannedMw, "mw") ?? { text: context.slice(0, 600), start: 0, end: Math.min(600, context.length) } : null;310 const decision = await recordCapacityClaim(tx, ctx, "project", id, { field: "plannedMw", predicate, value: p.plannedMw, scope: sc.scope, scopeReason: sc.reason, evidence, context, previous: (existing?.plannedMw as number | null) ?? null, recordScope: "project", campusDesignation, publishedAt: p.announcedOn ?? null, provenance: p.provenance, semanticsDefaulted: !p.capacitySemantics });311 if (decision.assign) {312 observed.push({ field: "plannedMw", value: p.plannedMw, provenance: p.provenance, scope: sc.scope });313 if (shouldReplace(obs(p.plannedMw, "plannedMw", p, ctx), stored(prov, "plannedMw", existing?.plannedMw)).replace) { next.plannedMw = p.plannedMw; next.capacityScope = sc.scope; next.capacitySemantics = predicate; }314 } else ctx.stats.unscopedClaims = (ctx.stats.unscopedClaims ?? 0) + 1;315 }316 // ─── investment through the claim store ─────────────────────────────────────────────────────────────────────────────317 const money = p.investmentOriginal != null && p.investmentOriginal > 0 ? { amount: p.investmentOriginal, currency: p.investmentCurrency ?? "USD" } : p.investmentUsd != null && p.investmentUsd > 0 ? { amount: p.investmentUsd, currency: "USD" } : null;318 if (money) {319 const context = p.claimContext?.investmentUsd ?? null;320 const sem = p.investmentSemantics && p.investmentScope ? { predicate: p.investmentSemantics as InvestmentPredicate, scope: p.investmentScope as ClaimScope, reason: "extractor" } : classifyInvestmentSemantics(context ?? name, campusDesignation ? "campus" : "facility");321 const evidence = context ? findEvidence(context, money.amount, "usd") ?? { text: context.slice(0, 600), start: 0, end: Math.min(600, context.length) } : null;322 const decision = await recordInvestmentClaim(tx, ctx, "project", id, { value: money.amount, currency: money.currency, predicate: sem.predicate, scope: sem.scope, scopeReason: sem.reason, evidence, context, previous: (existing?.investmentUsd as number | null) ?? null, recordScope: "project", publishedAt: p.announcedOn ?? null, provenance: p.provenance });323 if (decision.assign && p.investmentUsd != null) {324 observed.push({ field: "investmentUsd", value: p.investmentUsd, provenance: p.provenance, scope: sem.scope });325 if (shouldReplace(obs(p.investmentUsd, "investmentUsd", p, ctx), stored(prov, "investmentUsd", existing?.investmentUsd)).replace) { next.investmentUsd = p.investmentUsd; next.investmentScope = sem.scope; next.investmentSemantics = sem.predicate; }326 } else ctx.stats.unscopedClaims = (ctx.stats.unscopedClaims ?? 0) + 1;327 if (money.currency !== "USD") { next.investmentCurrency = money.currency; next.investmentOriginal = money.amount; }328 }329330 // ─── related entities ───────────────────────────────────────────────────────────────────────────────────────────────331 if (operator) {332 observed.push({ field: "operatorName", value: operator.name, provenance: p.provenance });333 if (shouldReplace(obs(operator.id, "operatorName", p, ctx), stored(prov, "operatorId", existing?.operatorId)).replace) next.operatorId = operator.id;334 }335 if (p.developerName) { const d = await resolveOperator(tx, ctx, { name: p.developerName }); if (d) { next.developerId = d.id; observed.push({ field: "developerName", value: d.name, provenance: p.provenance }); } }336 if (p.tenantName) { const t = await resolveOperator(tx, ctx, { name: p.tenantName }); if (t) { next.tenantId = t.id; observed.push({ field: "tenantName", value: t.name, provenance: p.provenance }); } }337 if (facilityId) next.facilityId = facilityId;338 if (p.campusName) { const c = await resolveCampus(tx, ctx, { name: p.campusName, operatorId: operator?.id ?? null, countryIso2: country, city: p.city ?? null, lat: geo?.lat ?? null, lng: geo?.lng ?? null }); if (c) next.campusId = c.id; }339 if (geo) {340 observed.push({ field: "geo", value: { lat: geo.lat, lng: geo.lng, precision: geo.precision, source: geo.source }, provenance: { ...p.provenance, method: geoMethod ?? p.provenance.method } });341 if (shouldReplaceGeo(geo.precision, existing?.geoPrecision as string | null, validLatLng(existing?.lat, existing?.lng))) { next.lat = geo.lat; next.lng = geo.lng; next.geoPrecision = geo.precision; }342 }343 {344 const m = validLatLng(next.lat, next.lng) ? await assignMetro(tx, { lat: next.lat as number, lng: next.lng as number, city: next.city as string | null, countryIso2: next.countryIso2 as string | null }) : metro;345 if (m.metroId) next.metroId = m.metroId;346 if (!next.countryIso2 && m.countryIso2) next.countryIso2 = m.countryIso2;347 }348 // ─── AI evidence: graded, never from one keyword ──────────────────────────────────────────────────────────────────349 {350 const rank: Record<string, number> = { unknown: 0, associated: 1, likely: 2, confirmed: 3 };351 const graded = classifyAiEvidence(`${name} ${p.description ?? ""}`);352 const level = p.aiEvidence && p.aiEvidence !== "unknown" ? p.aiEvidence : graded.level;353 if ((rank[level] ?? 0) > (rank[String(next.aiEvidence ?? "unknown")] ?? 0)) next.aiEvidence = level;354 next.isAi = next.aiEvidence === "confirmed" || next.aiEvidence === "likely";355 }356 if (cls) { next.projectClass = cls; if (p.evidenceLevel) next.evidenceLevel = p.evidenceLevel; }357 if (p.externalIds && Object.keys(p.externalIds).length) next.externalIds = { ...((next.externalIds as Record<string, unknown>) ?? {}), ...p.externalIds };358359 // ─── persist ────────────────────────────────────────────────────────────────────────────────────────────────────────360 const normalizedName = normalizeName(String(next.name));361 if (!existing) {362 const slug = await uniqueSlug(tx, "projects", `${operator && !normalizeName(name).includes(normalizeName(operator.name)) ? `${operator.name} ` : ""}${name}`, (next.city as string | null) ?? (next.countryIso2 as string | null));363 next.confidence = p.evidenceLevel === "weak" ? "unverified" : (p.provenance.confidence ?? "moderate");364 await tx.execute(sql`insert into projects (id, slug, name, normalized_name, operator_id, facility_id, metro_id, country_iso2, city, region_name, lat, lng, geo_precision, status, announced_on, expected_opening, planned_mw,365 investment_usd, investment_currency, investment_original, acreage, phase_count, is_ai, description, source_url, confidence, external_ids, last_update,366 project_class, evidence_level, ai_evidence, capacity_scope, capacity_semantics, investment_scope, investment_semantics, developer_id, tenant_id, campus_id, construction_started_on, approved_on, permit_filed_on, opened_on, hidden)367 values (${id}, ${slug}, ${next.name}, ${normalizedName}, ${next.operatorId ?? null}, ${next.facilityId ?? null}, ${next.metroId ?? null}, ${next.countryIso2 ?? null}, ${next.city ?? null}, ${next.regionName ?? null},368 ${next.lat ?? null}, ${next.lng ?? null}, ${next.geoPrecision ?? "unknown"}, ${next.status ?? "announced"}, ${next.announcedOn ?? null}, ${next.expectedOpening ?? null}, ${next.plannedMw ?? null},369 ${next.investmentUsd ?? null}, ${next.investmentCurrency ?? null}, ${next.investmentOriginal ?? null}, ${next.acreage ?? null}, ${next.phaseCount ?? null}, ${!!next.isAi}, ${next.description ?? null}, ${next.sourceUrl ?? null}, ${next.confidence}, ${JSON.stringify(next.externalIds ?? {})}::jsonb, ${ctx.now},370 ${next.projectClass ?? null}, ${next.evidenceLevel ?? null}, ${next.aiEvidence ?? "unknown"}, ${next.capacityScope ?? null}, ${next.capacitySemantics ?? null}, ${next.investmentScope ?? null}, ${next.investmentSemantics ?? null}, ${next.developerId ?? null}, ${next.tenantId ?? null}, ${next.campusId ?? null},371 ${next.constructionStartedOn ?? null}, ${next.approvedOn ?? null}, ${next.permitFiledOn ?? null}, ${next.openedOn ?? null}, false)`);372 ctx.stats.created++;373 if (p.evidenceLevel === "weak") await writeQualityFlags(tx, ctx, "project", id, [{ code: "project_weak_evidence", severity: "warn", message: "created from weak evidence (missing operator/name or location) — review", field: "name" }]);374 if (/\s[|]\s|\s[–—]\s| - /.test(name) || name.length > 90) await writeQualityFlags(tx, ctx, "project", id, [{ code: "project_title_like_name", severity: "warn", message: "project name looks like an article headline", field: "name" }]);375 } else {376 const sets = [];377 for (const [camel, col] of Object.entries(COLS)) {378 if (camel === "confidence" || camel === "reviewPriority" || camel === "hidden") continue;379 if (JSON.stringify(before[camel] ?? null) === JSON.stringify(next[camel] ?? null)) continue;380 if (camel === "externalIds") sets.push(sql`${sql.identifier(col)} = ${JSON.stringify(next[camel] ?? {})}::jsonb`);381 else sets.push(sql`${sql.identifier(col)} = ${next[camel] as string | number | boolean | null}`);382 }383 if (String(before.name) !== String(next.name)) sets.push(sql`normalized_name = ${normalizedName}`);384 if (sets.length) {385 sets.push(sql`last_update = ${ctx.now}`, sql`updated_at = now()`);386 await tx.execute(sql`update projects set ${sql.join(sets, sql`, `)} where id = ${id}`);387 ctx.stats.updated++;388 } else ctx.stats.unchanged++;389 }390 await upsertKey(tx, ctx, p.key, "project", id);391 addRef(ctx, "project", id);392 bump(ctx, "project");393 await writeProvenance(tx, ctx, "project", id, observed, p.key);394 await markWinners(tx, ctx, "project", id, ["status", "plannedMw", "investmentUsd", "expectedOpening", "announcedOn", "countryIso2", "city", "name"].map((f) => ({ field: f, value: next[f] })).concat([{ field: "operatorId", value: next.operatorId }]));395 // review priority from open flags + size396 {397 const r = (await tx.execute(sql`select coalesce(max(priority), 0) as p from quality_flags where entity_type = 'project' and entity_id = ${id} and status = 'open'`))[0];398 const rp = Math.max(Number(r?.p ?? 0), reviewPriority({ mw: next.plannedMw as number | null, investmentUsd: next.investmentUsd as number | null, confidence: String(next.confidence ?? ""), ai: !!next.isAi }) - 30);399 if (!ctx.run.dryRun) await tx.execute(sql`update projects set review_priority = ${Math.max(0, rp)} where id = ${id}`);400 }401402 // ─── timeline (deduplicated) ────────────────────────────────────────────────────────────────────────────────────────403 let timelineAdded = 0;404 const timeline = [...(p.timeline ?? [])];405 if (!timeline.length && p.announcedOn && cls) {406 const evType = ASSOCIATED_EVENT[cls] ?? "project_announced";407 timeline.push({ date: p.announcedOn, type: incomingStatus === "under_construction" ? "construction_started" : incomingStatus === "approved" ? "planning_approved" : incomingStatus === "permitting" ? "planning_filed" : evType, description: cleanText(p.description)?.slice(0, 300) ?? name, url });408 }409 if (existing && p.sourceUrl && existing.sourceUrl && p.sourceUrl !== existing.sourceUrl && p.announcedOn && p.description) {410 timeline.push({ date: p.announcedOn, type: "reported", description: `Also reported: ${String(p.description).slice(0, 300)}`, url: p.sourceUrl });411 }412 for (const t of timeline) {413 if (!t?.date || !t.description) continue;414 const tid = `ptl_${sha256(`${id}|${t.date}|${t.type}|${t.description.trim().toLowerCase()}`).slice(0, 20)}`;415 const r = await tx.execute(sql`insert into project_timeline (id, project_id, event_date, event_type, description, source_id, document_id, url)416 values (${tid}, ${id}, ${t.date}, ${String(t.type)}, ${t.description.trim()}, ${ctx.run.sourceId}, ${ctx.doc?.documentId ?? null}, ${t.url ?? url}) on conflict (id) do nothing returning id`);417 if (r.length) timelineAdded++;418 }419 if (timelineAdded && existing) await tx.execute(sql`update projects set last_update = ${ctx.now} where id = ${id}`);420421 // ─── events ─────────────────────────────────────────────────────────────────────────────────────────────────────────422 const confidence = (existing ? String(existing.confidence) : String(next.confidence ?? "moderate")) as ConfidenceLevel;423 if (!existing) {424 const mw = next.plannedMw as number | null;425 const where = [next.city, next.countryIso2].filter(Boolean).join(", ");426 const evType: EventType = incomingStatus === "under_construction" ? "construction_started" : incomingStatus === "approved" ? "planning_approved" : incomingStatus === "permitting" ? "planning_filed" : cls === "EXPANSION" ? "expansion_announced" : cls === "LAND_ACQUISITION" ? "land_acquired" : cls === "GRID_CONNECTION" ? "grid_connection" : "project_announced";427 await recordEvent(tx, ctx, {428 entityType: "project",429 entityId: id,430 eventType: evType,431 title: `New project: ${next.name}${operator && !String(next.name).toLowerCase().includes(operator.name.toLowerCase()) ? ` (${operator.name})` : ""}${where ? ` — ${where}` : ""}`,432 summary: [mw != null ? `${mw} MW planned (${String(next.capacityScope ?? "site")} scope).` : null, next.status ? `Status: ${String(next.status).replace(/_/g, " ")}.` : null, next.expectedOpening ? `Expected ${String(next.expectedOpening)}.` : null].filter(Boolean).join(" ") || null,433 newValue: { name: next.name, status: next.status, plannedMw: mw, class: cls },434 significance: mw != null && mw >= 100 ? 75 : isPipelineStatus(next.status as string) ? 60 : 50,435 confidence,436 effectiveDate: (next.announcedOn as string | null) ?? null,437 url,438 countryIso2: next.countryIso2 as string | null,439 operatorId: next.operatorId as string | null,440 metroId: next.metroId as string | null,441 projectId: id,442 isAi: !!next.isAi,443 reviewStatus: p.evidenceLevel === "weak" ? "pending" : "auto",444 });445 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 });446 } else {447 const names = await operatorNames(tx, ctx, [before.operatorId as string | null, next.operatorId as string | null]);448 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 });449 // delayed / cancelled get their own, louder event types on top of the generic status change450 if (before.status !== next.status && (next.status === "delayed" || next.status === "cancelled")) {451 await recordEvent(tx, ctx, { entityType: "project", entityId: id, eventType: next.status === "delayed" ? "project_delayed" : "project_cancelled", title: `${String(next.name)}: ${next.status === "delayed" ? "delayed" : "cancelled"}`, summary: cleanText(p.description)?.slice(0, 400) ?? null, oldValue: before.status, newValue: next.status, significance: next.status === "cancelled" ? 85 : 70, confidence, effectiveDate: p.announcedOn ?? null, url, countryIso2: next.countryIso2 as string | null, operatorId: next.operatorId as string | null, metroId: next.metroId as string | null, projectId: id, isAi: !!next.isAi });452 }453 }454}455456/** Admin: hide a false-positive project (kept for audit, removed from every listing and aggregate). */457export async function hideProject(tx: Tx, id: string, reason: string, decidedBy = "admin"): Promise<void> {458 await tx.execute(sql`update projects set hidden = true, updated_at = now() where id = ${id}`);459 await tx.execute(sql`update events set review_status = 'rejected' where project_id = ${id}`);460 await tx.execute(sql`insert into quality_flags (id, entity_type, entity_id, code, severity, field, message, priority, status, resolution, resolved_by, resolved_at, dedupe_key)461 values (${`flg_${sha256(`project|${id}|hidden`).slice(0, 16)}`}, 'project', ${id}, 'project_false_positive', 'critical', 'name', ${reason.slice(0, 1000)}, 0, 'resolved', 'hidden', ${decidedBy}, now(), ${`project|${id}|project_false_positive|`})462 on conflict (dedupe_key) do update set status = 'resolved', resolution = 'hidden', resolved_by = ${decidedBy}, resolved_at = now(), message = excluded.message`);463}464465/** Admin: fold `fromId` into `intoId` (keys, provenance, claims, timeline, events, news). */466export async function mergeProjects(tx: Tx, fromId: string, intoId: string, decidedBy = "admin"): Promise<void> {467 if (fromId === intoId) return;468 await tx.execute(sql`update entity_keys set entity_id = ${intoId} where entity_type = 'project' and entity_id = ${fromId}`);469 await tx.execute(sql`update provenance set entity_id = ${intoId} where entity_type = 'project' and entity_id = ${fromId} and not exists (select 1 from provenance q where q.entity_type = 'project' and q.entity_id = ${intoId} and q.field = provenance.field and q.source_id = provenance.source_id and q.url = provenance.url)`);470 await tx.execute(sql`delete from provenance where entity_type = 'project' and entity_id = ${fromId}`);471 await tx.execute(sql`update claims set subject_id = ${intoId} where subject_type = 'project' and subject_id = ${fromId}`);472 await tx.execute(sql`update project_timeline set project_id = ${intoId} where project_id = ${fromId}`);473 await tx.execute(sql`update events set project_id = ${intoId}, entity_id = case when entity_type = 'project' and entity_id = ${fromId} then ${intoId} else entity_id end where project_id = ${fromId} or (entity_type = 'project' and entity_id = ${fromId})`);474 await tx.execute(sql`update news_items set project_id = ${intoId} where project_id = ${fromId}`);475 await tx.execute(sql`update projects set merged_into = ${intoId}, hidden = true, updated_at = now() where id = ${fromId}`);476 await tx.execute(sql`update quality_flags set status = 'resolved', resolution = ${`merged into ${intoId}`}, resolved_by = ${decidedBy}, resolved_at = now() where entity_type = 'project' and entity_id = ${fromId} and status = 'open'`);477}478