SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
4 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
37.5 KB · 478 lines typescript
Raw Blame History
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