/** * Project reconciliation + persistence (docs/PROJECT-EXTRACTION.md, docs/CLAIMS.md). * * - The announcement CLASS decides first: appointments, financing, PPAs, partnerships, customer deals, market research * never create a project. Non-physical classes attach a timeline row + event to an existing, identifiable project. * - Evidence threshold: a new project needs (explicit name or operator) + location + a development verb. * - Capacity and investment go through the claim store (scope + sanity engine); a company-wide or portfolio figure is * kept as a claim and never written to the project's columns. * - Lifecycle transitions follow the state machine; a backward move needs a source that outranks the stored one. * - Projects without coordinates are geocoded at CITY / METRO level (never more precise) from the metro seeds. */ import { sql } from "@dci/db"; import { ASSOCIATED_CLASSES, PHYSICAL_CLASSES, cleanText, classifyAiEvidence, classifyInvestmentSemantics, classifyScope, findEvidence, newId, normalizeName, projectTransition, sha256, validLatLng, type CapacityPredicate, type ClaimScope, type ConfidenceLevel, type EventType, type InvestmentPredicate, type NormalizedProject, type ProjectClass, } from "@dci/core"; import { addRef, authority, bump, provenanceFor, safeCountry, uniqueSlug, type IngestContext, type Tx } from "./common.js"; import { emitDiffEvents, recordEvent, TRACKED_PROJECT_FIELDS } from "./events.js"; import { validGeo } from "./geo.js"; import { geocodeCity } from "./geocode.js"; import { entityIdForKey, facilityIdForKey, upsertKey } from "./keys.js"; import { isPipelineStatus, shouldReplace, shouldReplaceGeo, type FieldObservation } from "./match.js"; import { assignMetro } from "./metros.js"; import { operatorNames, resolveOperator } from "./operators.js"; import { backingObservation, loadCurrentProvenance, writeProvenance, type CurrentProvenance, type ObservedField } from "./provenance.js"; import { markWinners, recordCapacityClaim, recordInvestmentClaim, reviewPriority, writeQualityFlags, writeTextClaim } from "./claims.js"; import { resolveCampus } from "./campuses.js"; import { maybeOperatorExpansion } from "./facilities.js"; type Row = Record; const COLS: Record = { name: "name", operatorId: "operator_id", facilityId: "facility_id", metroId: "metro_id", countryIso2: "country_iso2", city: "city", regionName: "region_name", lat: "lat", lng: "lng", geoPrecision: "geo_precision", status: "status", announcedOn: "announced_on", expectedOpening: "expected_opening", plannedMw: "planned_mw", investmentUsd: "investment_usd", investmentCurrency: "investment_currency", investmentOriginal: "investment_original", acreage: "acreage", phaseCount: "phase_count", isAi: "is_ai", description: "description", sourceUrl: "source_url", confidence: "confidence", externalIds: "external_ids", projectClass: "project_class", evidenceLevel: "evidence_level", aiEvidence: "ai_evidence", capacityScope: "capacity_scope", capacitySemantics: "capacity_semantics", investmentScope: "investment_scope", investmentSemantics: "investment_semantics", developerId: "developer_id", tenantId: "tenant_id", campusId: "campus_id", constructionStartedOn: "construction_started_on", approvedOn: "approved_on", permitFiledOn: "permit_filed_on", openedOn: "opened_on", reviewPriority: "review_priority", hidden: "hidden", }; /** External-id namespaces that identify ONE project (allowlist). */ export const PROJECT_IDENTIFYING_KEYS: ReadonlySet = new Set(["wikidata", "planning_ref", "permit_id", "case_number", "application_id", "docket", "planning_application", "rezoning_case", "eia_ref"]); /** Timeline / event type for a non-physical (associated) class. */ const ASSOCIATED_EVENT: Record = { 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" }; async function loadProject(tx: Tx, id: string): Promise { const r = (await tx.execute(sql`select * from projects where id = ${id}`))[0]; if (!r) return null; const out: Row = { id: r.id, slug: r.slug }; for (const [camel, col] of Object.entries(COLS)) out[camel] = r[col] ?? null; return out; } async function followMerged(tx: Tx, id: string): Promise { let cur = id; for (let i = 0; i < 5; i++) { const r = await tx.execute(sql`select merged_into from projects where id = ${cur}`); const m = r[0]?.merged_into; if (!m) break; cur = String(m); } return cur; } async function resolveProject(tx: Tx, ctx: IngestContext, p: NormalizedProject, operatorId: string | null, country: string | null): Promise { const byKey = await entityIdForKey(tx, p.key, "project"); if (byKey) return followMerged(tx, byKey); if (p.externalIds) { for (const [k, v] of Object.entries(p.externalIds)) { if (v == null || v === "" || !PROJECT_IDENTIFYING_KEYS.has(k)) continue; 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`); if (r[0]) return String(r[0].id); } } const src = p.sourceUrl ?? null; if (src) { const r = await tx.execute(sql`select id from projects where merged_into is null and source_url = ${src} limit 1`); if (r[0]) return followMerged(tx, String(r[0].id)); } const norm = normalizeName(p.name); if (!norm) return null; const rows = await tx.execute(sql` select id, similarity(normalized_name, ${norm}) as sim from projects where merged_into is null and (${country}::text is null or country_iso2 is null or country_iso2 = ${country}) and (${operatorId}::text is null or operator_id is null or operator_id = ${operatorId}) and (normalized_name = ${norm} or similarity(normalized_name, ${norm}) >= 0.85) order by (operator_id = ${operatorId}) desc nulls last, sim desc limit 1`); if (rows[0]) return String(rows[0].id); return sameAnnouncement(tx, ctx, p, operatorId, country); } /** Announcement window: two outlets reporting the same project within this many days are one project. */ export const ANNOUNCEMENT_WINDOW_DAYS = 30; /** * The same announcement covered by several outlets: same operator (or, without an operator, the same city), planned MW * within ±10 % and announced within 30 days → one project. Without an operator the city is mandatory. */ async function sameAnnouncement(tx: Tx, ctx: IngestContext, p: NormalizedProject, operatorId: string | null, country: string | null): Promise { if (p.plannedMw == null || p.plannedMw <= 0 || !p.announcedOn) return null; const day = p.announcedOn.length >= 10 ? p.announcedOn.slice(0, 10) : null; if (!day) return null; const city = p.city ? normalizeName(p.city) : null; if (!operatorId && !city) return null; const rows = await tx.execute(sql` select id from projects where merged_into is null and planned_mw is not null and planned_mw between ${p.plannedMw * 0.9} and ${p.plannedMw * 1.1} and announced_on is not null and length(announced_on) >= 10 and abs(announced_on::date - ${day}::date) <= ${ANNOUNCEMENT_WINDOW_DAYS} and (${country}::text is null or country_iso2 is null or country_iso2 = ${country}) and (${operatorId}::text is null or operator_id is null or operator_id = ${operatorId}) and (${city}::text is null or city is null or lower(regexp_replace(city, '[^A-Za-z0-9]+', ' ', 'g')) = ${city}) and (${operatorId}::text is not null or (${city}::text is not null and city is not null)) order by created_at asc limit 1`); if (!rows[0]) return null; ctx.stats.projectDedup = (ctx.stats.projectDedup ?? 0) + 1; return String(rows[0].id); } function obs(value: unknown, field: string, p: NormalizedProject, ctx: IngestContext): FieldObservation { const pv = provenanceFor(p.provenance, field, undefined); 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 }; } function stored(prov: CurrentProvenance[], field: string, current: unknown): FieldObservation | null { if (current == null) return null; const b = backingObservation(prov, field, current); 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 }; } const isCampusName = (s: string | null | undefined) => /\b(campus|park|complex|hub|cluster|gigafactory|estate)\b/i.test(s ?? ""); export async function ingestProject(tx: Tx, ctx: IngestContext, p: NormalizedProject): Promise { const name = cleanText(p.name); if (!name) throw new Error(`project ${p.key}: name is required`); const url = p.provenance.url || p.sourceUrl || ctx.doc?.url || ""; if (!url) throw new Error(`project ${p.key}: provenance url is required`); const cls = (p.projectClass ?? null) as ProjectClass | null; const physical = cls == null || PHYSICAL_CLASSES.has(cls); const associated = cls != null && ASSOCIATED_CLASSES.has(cls); const operator = p.operatorName ? await resolveOperator(tx, ctx, { name: p.operatorName }) : null; const country0 = await safeCountry(tx, ctx, p.countryIso2); let geo = validGeo(p.geo) ? p.geo : null; let geoMethod = geo ? "source" : null; if (!geo && p.city) { const g = await geocodeCity(tx, { city: p.city, regionName: p.regionName, countryIso2: country0 }); if (g) { geo = { lat: g.lat, lng: g.lng, precision: g.precision, source: g.source }; geoMethod = g.source; } } const metro = await assignMetro(tx, { lat: geo?.lat, lng: geo?.lng, city: p.city, countryIso2: country0 }); const country = country0 ?? metro.countryIso2; const facilityId = p.facilityKey ? await facilityIdForKey(tx, ctx, p.facilityKey) : null; const existingId = await resolveProject(tx, ctx, p, operator?.id ?? null, country); const existing = existingId ? await loadProject(tx, existingId) : null; // ─── veto: non-physical announcements never create a project; weak evidence never creates a project ───────────── if (!existing) { if (!physical) { ctx.stats.projectsVetoed = (ctx.stats.projectsVetoed ?? 0) + 1; return; } if (p.evidenceLevel === "none") { ctx.stats.projectsVetoed = (ctx.stats.projectsVetoed ?? 0) + 1; return; } } // ─── associated class on an existing project: timeline row + event only, no field changes ───────────────────────── if (existing && !physical) { const evType = (cls && ASSOCIATED_EVENT[cls]) || "news"; const date = p.announcedOn ?? ctx.day; const desc = cleanText(p.description) ?? name; const tid = `ptl_${sha256(`${existing.id}|${date}|${evType}|${desc.toLowerCase()}`).slice(0, 20)}`; const ins = await tx.execute(sql`insert into project_timeline (id, project_id, event_date, event_type, description, source_id, document_id, url) 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`); if (ins.length && associated) { 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) }); await tx.execute(sql`update projects set last_update = ${ctx.now} where id = ${existing.id}`); } // money attached to a deal / financing headline is a claim about the project, never its investment column if (p.investmentUsd != null && p.investmentUsd > 0) { const sem = classifyInvestmentSemantics(p.claimContext?.investmentUsd ?? name); 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 }); } addRef(ctx, "project", String(existing.id)); bump(ctx, "project"); return; } const prov = existing ? await loadCurrentProvenance(tx, "project", existingId!) : []; const id = existing ? String(existing.id) : newId("project"); // a project hidden by the automated repair re-qualifies when a re-extraction yields a physical class with strong evidence // (human decisions — resolved_by <> 'quality-apply' / 'system' — are never reversed) if (existing && existing.hidden && physical && p.evidenceLevel === "strong" && !ctx.run.dryRun) { 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`); 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`); if (auto.length && !human.length) { await tx.execute(sql`update projects set hidden = false, updated_at = now() where id = ${id}`); 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'`); existing.hidden = false; ctx.stats.projectsRequalified = (ctx.stats.projectsRequalified ?? 0) + 1; } } const before: Row = existing ? { ...existing } : {}; const next: Row = existing ? { ...existing } : { id, status: "announced", geoPrecision: "unknown", isAi: false, confidence: "moderate", externalIds: {}, aiEvidence: "unknown", hidden: false, reviewPriority: 0 }; const observed: ObservedField[] = []; const campusDesignation = isCampusName(name) || isCampusName(p.campusName); const scalars: Record = { name, city: cleanText(p.city), regionName: cleanText(p.regionName), announcedOn: p.announcedOn ?? null, expectedOpening: p.expectedOpening ?? null, acreage: p.acreage ?? null, phaseCount: p.phaseCount ?? null, description: cleanText(p.description), sourceUrl: p.sourceUrl ?? null, constructionStartedOn: p.constructionStartedOn ?? null, approvedOn: p.approvedOn ?? null, permitFiledOn: p.permitFiledOn ?? null, countryIso2: country, }; for (const [field, v] of Object.entries(scalars)) { if (v == null) continue; observed.push({ field, value: v, provenance: p.provenance }); if (shouldReplace(obs(v, field, p, ctx), stored(prov, field, existing?.[field])).replace) next[field] = v; } // ─── status: lifecycle state machine ──────────────────────────────────────────────────────────────────────────────── const incomingStatus = p.status && p.status !== "unknown" ? p.status : null; if (incomingStatus) { observed.push({ field: "status", value: incomingStatus, provenance: p.provenance }); const verdict = projectTransition(existing?.status as string | null, incomingStatus); if (verdict === "forward" || verdict === "side" || verdict === "resume" || (verdict === "same" && !existing)) next.status = incomingStatus; else if (verdict === "backward") { // a backward move is accepted only from a source that outranks the stored one (a government filing correcting a news story) const st = stored(prov, "status", existing?.status); if (authority(ctx.run.sourceKind, p.provenance.confidence, false) > authority(st?.sourceKind, st?.confidence, false)) next.status = incomingStatus; 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" }]); } else if (verdict === "invalid") { await writeQualityFlags(tx, ctx, "project", id, [{ code: "status_invalid_transition", severity: "warn", message: `invalid lifecycle transition ${String(existing?.status)} → ${incomingStatus}`, field: "status" }]); } // stage dates derived from the announcement that moved the stage const day = p.announcedOn ?? null; if (day) { if (incomingStatus === "under_construction" && !next.constructionStartedOn) next.constructionStartedOn = day; if (incomingStatus === "approved" && !next.approvedOn) next.approvedOn = day; if (incomingStatus === "permitting" && !next.permitFiledOn) next.permitFiledOn = day; if ((incomingStatus === "operational" || incomingStatus === "partially_operational") && !next.openedOn) next.openedOn = day; } await writeTextClaim(tx, ctx, "project", id, "status", incomingStatus, p.provenance, { scope: campusDesignation ? "campus" : "facility", publishedAt: p.announcedOn ?? null }); } // ─── capacity through the claim store ─────────────────────────────────────────────────────────────────────────────── if (p.plannedMw != null && p.plannedMw > 0) { const context = p.claimContext?.plannedMw ?? null; const scopeHint: ClaimScope = campusDesignation ? "campus" : "facility"; const sc = p.capacityScope ? { scope: p.capacityScope as ClaimScope, reason: "extractor" } : classifyScope(context, context ? null : scopeHint); const predicate = (p.capacitySemantics ?? "planned_power_mw") as CapacityPredicate; const evidence = context ? findEvidence(context, p.plannedMw, "mw") ?? { text: context.slice(0, 600), start: 0, end: Math.min(600, context.length) } : null; 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 }); if (decision.assign) { observed.push({ field: "plannedMw", value: p.plannedMw, provenance: p.provenance, scope: sc.scope }); 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; } } else ctx.stats.unscopedClaims = (ctx.stats.unscopedClaims ?? 0) + 1; } // ─── investment through the claim store ───────────────────────────────────────────────────────────────────────────── 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; if (money) { const context = p.claimContext?.investmentUsd ?? null; const sem = p.investmentSemantics && p.investmentScope ? { predicate: p.investmentSemantics as InvestmentPredicate, scope: p.investmentScope as ClaimScope, reason: "extractor" } : classifyInvestmentSemantics(context ?? name, campusDesignation ? "campus" : "facility"); const evidence = context ? findEvidence(context, money.amount, "usd") ?? { text: context.slice(0, 600), start: 0, end: Math.min(600, context.length) } : null; 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 }); if (decision.assign && p.investmentUsd != null) { observed.push({ field: "investmentUsd", value: p.investmentUsd, provenance: p.provenance, scope: sem.scope }); 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; } } else ctx.stats.unscopedClaims = (ctx.stats.unscopedClaims ?? 0) + 1; if (money.currency !== "USD") { next.investmentCurrency = money.currency; next.investmentOriginal = money.amount; } } // ─── related entities ─────────────────────────────────────────────────────────────────────────────────────────────── if (operator) { observed.push({ field: "operatorName", value: operator.name, provenance: p.provenance }); if (shouldReplace(obs(operator.id, "operatorName", p, ctx), stored(prov, "operatorId", existing?.operatorId)).replace) next.operatorId = operator.id; } 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 }); } } 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 }); } } if (facilityId) next.facilityId = facilityId; 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; } if (geo) { 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 } }); 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; } } { 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; if (m.metroId) next.metroId = m.metroId; if (!next.countryIso2 && m.countryIso2) next.countryIso2 = m.countryIso2; } // ─── AI evidence: graded, never from one keyword ────────────────────────────────────────────────────────────────── { const rank: Record = { unknown: 0, associated: 1, likely: 2, confirmed: 3 }; const graded = classifyAiEvidence(`${name} ${p.description ?? ""}`); const level = p.aiEvidence && p.aiEvidence !== "unknown" ? p.aiEvidence : graded.level; if ((rank[level] ?? 0) > (rank[String(next.aiEvidence ?? "unknown")] ?? 0)) next.aiEvidence = level; next.isAi = next.aiEvidence === "confirmed" || next.aiEvidence === "likely"; } if (cls) { next.projectClass = cls; if (p.evidenceLevel) next.evidenceLevel = p.evidenceLevel; } if (p.externalIds && Object.keys(p.externalIds).length) next.externalIds = { ...((next.externalIds as Record) ?? {}), ...p.externalIds }; // ─── persist ──────────────────────────────────────────────────────────────────────────────────────────────────────── const normalizedName = normalizeName(String(next.name)); if (!existing) { 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)); next.confidence = p.evidenceLevel === "weak" ? "unverified" : (p.provenance.confidence ?? "moderate"); 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, investment_usd, investment_currency, investment_original, acreage, phase_count, is_ai, description, source_url, confidence, external_ids, last_update, 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) values (${id}, ${slug}, ${next.name}, ${normalizedName}, ${next.operatorId ?? null}, ${next.facilityId ?? null}, ${next.metroId ?? null}, ${next.countryIso2 ?? null}, ${next.city ?? null}, ${next.regionName ?? null}, ${next.lat ?? null}, ${next.lng ?? null}, ${next.geoPrecision ?? "unknown"}, ${next.status ?? "announced"}, ${next.announcedOn ?? null}, ${next.expectedOpening ?? null}, ${next.plannedMw ?? null}, ${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}, ${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}, ${next.constructionStartedOn ?? null}, ${next.approvedOn ?? null}, ${next.permitFiledOn ?? null}, ${next.openedOn ?? null}, false)`); ctx.stats.created++; 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" }]); 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" }]); } else { const sets = []; for (const [camel, col] of Object.entries(COLS)) { if (camel === "confidence" || camel === "reviewPriority" || camel === "hidden") continue; if (JSON.stringify(before[camel] ?? null) === JSON.stringify(next[camel] ?? null)) continue; if (camel === "externalIds") sets.push(sql`${sql.identifier(col)} = ${JSON.stringify(next[camel] ?? {})}::jsonb`); else sets.push(sql`${sql.identifier(col)} = ${next[camel] as string | number | boolean | null}`); } if (String(before.name) !== String(next.name)) sets.push(sql`normalized_name = ${normalizedName}`); if (sets.length) { sets.push(sql`last_update = ${ctx.now}`, sql`updated_at = now()`); await tx.execute(sql`update projects set ${sql.join(sets, sql`, `)} where id = ${id}`); ctx.stats.updated++; } else ctx.stats.unchanged++; } await upsertKey(tx, ctx, p.key, "project", id); addRef(ctx, "project", id); bump(ctx, "project"); await writeProvenance(tx, ctx, "project", id, observed, p.key); 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 }])); // review priority from open flags + size { 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]; 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); if (!ctx.run.dryRun) await tx.execute(sql`update projects set review_priority = ${Math.max(0, rp)} where id = ${id}`); } // ─── timeline (deduplicated) ──────────────────────────────────────────────────────────────────────────────────────── let timelineAdded = 0; const timeline = [...(p.timeline ?? [])]; if (!timeline.length && p.announcedOn && cls) { const evType = ASSOCIATED_EVENT[cls] ?? "project_announced"; 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 }); } if (existing && p.sourceUrl && existing.sourceUrl && p.sourceUrl !== existing.sourceUrl && p.announcedOn && p.description) { timeline.push({ date: p.announcedOn, type: "reported", description: `Also reported: ${String(p.description).slice(0, 300)}`, url: p.sourceUrl }); } for (const t of timeline) { if (!t?.date || !t.description) continue; const tid = `ptl_${sha256(`${id}|${t.date}|${t.type}|${t.description.trim().toLowerCase()}`).slice(0, 20)}`; const r = await tx.execute(sql`insert into project_timeline (id, project_id, event_date, event_type, description, source_id, document_id, url) 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`); if (r.length) timelineAdded++; } if (timelineAdded && existing) await tx.execute(sql`update projects set last_update = ${ctx.now} where id = ${id}`); // ─── events ───────────────────────────────────────────────────────────────────────────────────────────────────────── const confidence = (existing ? String(existing.confidence) : String(next.confidence ?? "moderate")) as ConfidenceLevel; if (!existing) { const mw = next.plannedMw as number | null; const where = [next.city, next.countryIso2].filter(Boolean).join(", "); 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"; await recordEvent(tx, ctx, { entityType: "project", entityId: id, eventType: evType, title: `New project: ${next.name}${operator && !String(next.name).toLowerCase().includes(operator.name.toLowerCase()) ? ` (${operator.name})` : ""}${where ? ` — ${where}` : ""}`, 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, newValue: { name: next.name, status: next.status, plannedMw: mw, class: cls }, significance: mw != null && mw >= 100 ? 75 : isPipelineStatus(next.status as string) ? 60 : 50, confidence, effectiveDate: (next.announcedOn as string | null) ?? 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, reviewStatus: p.evidenceLevel === "weak" ? "pending" : "auto", }); 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 }); } else { const names = await operatorNames(tx, ctx, [before.operatorId as string | null, next.operatorId as string | null]); 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 }); // delayed / cancelled get their own, louder event types on top of the generic status change if (before.status !== next.status && (next.status === "delayed" || next.status === "cancelled")) { 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 }); } } } /** Admin: hide a false-positive project (kept for audit, removed from every listing and aggregate). */ export async function hideProject(tx: Tx, id: string, reason: string, decidedBy = "admin"): Promise { await tx.execute(sql`update projects set hidden = true, updated_at = now() where id = ${id}`); await tx.execute(sql`update events set review_status = 'rejected' where project_id = ${id}`); 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) 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|`}) on conflict (dedupe_key) do update set status = 'resolved', resolution = 'hidden', resolved_by = ${decidedBy}, resolved_at = now(), message = excluded.message`); } /** Admin: fold `fromId` into `intoId` (keys, provenance, claims, timeline, events, news). */ export async function mergeProjects(tx: Tx, fromId: string, intoId: string, decidedBy = "admin"): Promise { if (fromId === intoId) return; await tx.execute(sql`update entity_keys set entity_id = ${intoId} where entity_type = 'project' and entity_id = ${fromId}`); 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)`); await tx.execute(sql`delete from provenance where entity_type = 'project' and entity_id = ${fromId}`); await tx.execute(sql`update claims set subject_id = ${intoId} where subject_type = 'project' and subject_id = ${fromId}`); await tx.execute(sql`update project_timeline set project_id = ${intoId} where project_id = ${fromId}`); 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})`); await tx.execute(sql`update news_items set project_id = ${intoId} where project_id = ${fromId}`); await tx.execute(sql`update projects set merged_into = ${intoId}, hidden = true, updated_at = now() where id = ${fromId}`); 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'`); }