spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Claim store + quality flags + winner bookkeeping (docs/CLAIMS.md).3 *4 * Every numerical figure observed about a subject becomes a `claims` row carrying its scope, semantics, supporting5 * sentence, parser and field-level authority tier. The ingest layer then asks `resolveCapacityColumn()` which claim6 * may populate a column: only site-scoped claims (building / facility / campus) that pass the deterministic sanity7 * engine; everything else is stored with status `unscoped` / `review` and surfaces in the admin quality dashboard.8 */9import { sql } from "@dci/db";10import {11 authorityFieldFor,12 authorityTier,13 capacitySanity,14 investmentSanity,15 isSiteScope,16 newId,17 stableId,18 TIER_RANK,19 type AuthorityTier,20 type CapacityPredicate,21 type ClaimScope,22 type ClaimStatus,23 type Evidence,24 type InvestmentPredicate,25 type Provenance,26 type SanityFlag,27} from "@dci/core";28import type { IngestContext, Tx } from "./common.js";2930export interface ClaimInput {31 predicate: string;32 value?: number | null;33 valueText?: string | null;34 unit?: string | null;35 scope: ClaimScope;36 scopeReason?: string | null;37 evidence?: Evidence | null;38 publishedAt?: string | null;39 provenance: Provenance;40 parserName?: string | null;41 /** override the computed status (e.g. "review" after a sanity flag) */42 status?: ClaimStatus;43 rejectionReason?: string | null;44}4546export interface WrittenClaim { id: string; status: ClaimStatus; tier: AuthorityTier }4748export function claimId(subjectType: string, subjectId: string, c: { predicate: string; value?: number | null; valueText?: string | null; provenance: Provenance }, url: string): string {49 return stableId("claim", `${subjectType}|${subjectId}|${c.predicate}|${c.provenance.sourceId}|${url}|${c.value ?? ""}|${c.valueText ?? ""}`);50}5152/** Upsert one claim. Same source + url + predicate with a different value supersedes the earlier claim. */53export async function writeClaim(tx: Tx, ctx: IngestContext, subjectType: string, subjectId: string, c: ClaimInput): Promise<WrittenClaim> {54 const p = c.provenance;55 const url = p.url || ctx.doc?.url || "";56 const sourceId = p.sourceId || ctx.run.sourceId;57 const tier = authorityTier({ field: authorityFieldFor(c.predicate), sourceKind: ctx.run.sourceKind, isEstimate: p.isEstimate, method: p.method });58 const status: ClaimStatus = c.status ?? (isSiteScope(c.scope) ? "current" : "unscoped");59 const id = claimId(subjectType, subjectId, c, url);60 const observed = p.lastObserved || ctx.now;61 if (!ctx.run.dryRun) {62 await tx.execute(sql`63 insert into claims (id, subject_type, subject_id, predicate, value, value_text, unit, scope, scope_reason, source_id, connector_id, document_id, url, published_at, retrieved_at,64 confidence, is_estimate, authority_tier, evidence_text, evidence_start, evidence_end, parser_name, parser_version, run_id, status, rejection_reason, first_observed, last_observed)65 values (${id}, ${subjectType}, ${subjectId}, ${c.predicate}, ${c.value ?? null}, ${c.valueText ?? null}, ${c.unit ?? null}, ${c.scope}, ${c.scopeReason ?? null}, ${sourceId}, ${p.connectorId || ctx.run.connectorId},66 ${p.documentId ?? ctx.doc?.documentId ?? null}, ${url}, ${c.publishedAt ?? null}, ${p.retrievedAt || ctx.now}, ${p.confidence ?? "moderate"}, ${!!p.isEstimate}, ${tier},67 ${c.evidence?.text ?? null}, ${c.evidence?.start ?? null}, ${c.evidence?.end ?? null}, ${c.parserName ?? p.method ?? null}, ${p.extractorVersion ?? null}, ${ctx.run.runId}, ${status}, ${c.rejectionReason ?? null}, ${p.firstObserved || observed}, ${observed})68 on conflict (id) do update set last_observed = excluded.last_observed, retrieved_at = excluded.retrieved_at, confidence = excluded.confidence, authority_tier = excluded.authority_tier,69 evidence_text = coalesce(excluded.evidence_text, claims.evidence_text), evidence_start = coalesce(excluded.evidence_start, claims.evidence_start), evidence_end = coalesce(excluded.evidence_end, claims.evidence_end),70 scope = excluded.scope, scope_reason = excluded.scope_reason, run_id = excluded.run_id, parser_version = excluded.parser_version,71 status = case when claims.status in ('rejected') then claims.status else excluded.status end, rejection_reason = coalesce(excluded.rejection_reason, claims.rejection_reason)`);72 // the same source, same page, same predicate now says something else → the older claim is superseded73 await tx.execute(sql`update claims set status = 'superseded' where subject_type = ${subjectType} and subject_id = ${subjectId} and predicate = ${c.predicate} and source_id = ${sourceId} and url = ${url} and id <> ${id} and status in ('current', 'review')`);74 }75 ctx.stats.claims = (ctx.stats.claims ?? 0) + 1;76 return { id, status, tier };77}7879/** Review priority 0–100: impact-weighted (MW / money size, low confidence, new country, AI, ambiguity, unexpected change). */80export function reviewPriority(i: { mw?: number | null; investmentUsd?: number | null; confidence?: string | null; newCountry?: boolean; ai?: boolean; ambiguous?: boolean; unexpectedChange?: boolean; severity?: "info" | "warn" | "critical"; homepageVisible?: boolean }): number {81 let p = 0;82 const mw = i.mw ?? 0;83 p += mw >= 1000 ? 40 : mw >= 300 ? 30 : mw >= 100 ? 20 : mw >= 20 ? 10 : mw > 0 ? 4 : 0;84 const inv = i.investmentUsd ?? 0;85 p += inv >= 10e9 ? 25 : inv >= 1e9 ? 15 : inv >= 100e6 ? 8 : 0;86 if (i.confidence === "unverified" || i.confidence === "estimated") p += 10;87 if (i.newCountry) p += 10;88 if (i.ai) p += 5;89 if (i.ambiguous) p += 8;90 if (i.unexpectedChange) p += 12;91 if (i.severity === "critical") p += 15; else if (i.severity === "warn") p += 5;92 if (i.homepageVisible) p += 10;93 return Math.max(0, Math.min(100, Math.round(p)));94}9596/** Upsert deterministic quality flags for an entity (deduped by entity + code + field). Resolved flags are not reopened for the same key. */97export async function writeQualityFlags(tx: Tx, ctx: IngestContext, entityType: string, entityId: string, flags: SanityFlag[], opts: { claimId?: string | null; priority?: number; details?: Record<string, unknown> } = {}): Promise<number> {98 let n = 0;99 for (const f of flags) {100 if (f.severity === "info") continue;101 const dedupeKey = `${entityType}|${entityId}|${f.code}|${f.field ?? ""}`;102 const id = stableId("flag", dedupeKey);103 const priority = opts.priority ?? reviewPriority({ mw: /mw/.test(f.code) ? f.value : null, investmentUsd: /inv/.test(f.code) ? f.value : null, severity: f.severity });104 if (!ctx.run.dryRun) {105 await tx.execute(sql`106 insert into quality_flags (id, entity_type, entity_id, claim_id, code, severity, field, message, details, priority, status, run_id, dedupe_key)107 values (${id}, ${entityType}, ${entityId}, ${opts.claimId ?? null}, ${f.code}, ${f.severity}, ${f.field ?? null}, ${f.message.slice(0, 1000)}, ${JSON.stringify({ ...(opts.details ?? {}), value: f.value ?? null })}::jsonb, ${priority}, 'open', ${ctx.run.runId}, ${dedupeKey})108 on conflict (dedupe_key) do update set message = excluded.message, details = excluded.details, priority = greatest(quality_flags.priority, excluded.priority), claim_id = coalesce(excluded.claim_id, quality_flags.claim_id),109 run_id = excluded.run_id, updated_at = now(), status = case when quality_flags.status = 'dismissed' then 'dismissed' else 'open' end`);110 }111 n++;112 }113 ctx.stats.qualityFlags = (ctx.stats.qualityFlags ?? 0) + n;114 return n;115}116117/** Close open flags of the given codes for an entity (the condition no longer holds). */118export async function resolveQualityFlags(tx: Tx, ctx: IngestContext, entityType: string, entityId: string, codes: string[], resolution = "auto: condition cleared"): Promise<void> {119 if (!codes.length || ctx.run.dryRun) return;120 await tx.execute(sql`update quality_flags set status = 'resolved', resolution = ${resolution}, resolved_by = 'system', resolved_at = now(), updated_at = now()121 where entity_type = ${entityType} and entity_id = ${entityId} and status = 'open' and code in ${codes}`);122}123124/**125 * Mark, per field, the current provenance row whose value equals the stored column value as the winner (and unmark126 * the others). This is what makes "why does the page show 300 MW?" answerable without a diff.127 */128export async function markWinners(tx: Tx, ctx: IngestContext, entityType: string, entityId: string, fields: Array<{ field: string; value: unknown }>): Promise<void> {129 if (ctx.run.dryRun) return;130 for (const f of fields) {131 if (f.value == null) continue;132 const valueJson = JSON.stringify(f.value);133 await tx.execute(sql`update provenance set is_winner = (value = ${valueJson}::jsonb) where entity_type = ${entityType} and entity_id = ${entityId} and field = ${f.field} and is_current and is_winner <> (value = ${valueJson}::jsonb)`);134 }135}136137export interface CapacityDecision {138 /** write the figure into the record's column */139 assign: boolean;140 claim: WrittenClaim;141 flags: SanityFlag[];142 blockedBy: string[];143}144145export interface CapacityClaimInput {146 field: "itCapacityMw" | "totalPowerMw" | "plannedPowerMw" | "plannedMw" | "utilityCapacityMw" | "gridConnectionMw" | "ultimateCampusMw";147 predicate: CapacityPredicate;148 value: number;149 scope: ClaimScope;150 scopeReason?: string | null;151 evidence?: Evidence | null;152 context?: string | null;153 previous?: number | null;154 recordScope: "building" | "facility" | "campus" | "project";155 campusDesignation?: boolean;156 publishedAt?: string | null;157 provenance: Provenance;158 parserName?: string | null;159 /** the sentence did not say what kind of MW this is; semantics were defaulted from the record status */160 semanticsDefaulted?: boolean;161}162163/** Store a capacity claim, run the sanity engine, write flags and decide whether the column may take the value. */164export async function recordCapacityClaim(tx: Tx, ctx: IngestContext, subjectType: "facility" | "project" | "campus", subjectId: string, i: CapacityClaimInput): Promise<CapacityDecision> {165 const flags = capacitySanity({ value: i.value, predicate: i.semanticsDefaulted ? null : i.predicate, scope: i.scope, recordScope: i.recordScope, previous: i.previous ?? null, context: i.context ?? i.evidence?.text ?? null, campusDesignation: i.campusDesignation });166 const blockedBy = flags.filter((f) => f.blocks).map((f) => f.code);167 const critical = flags.some((f) => f.severity === "critical");168 const status: ClaimStatus = blockedBy.length ? (isSiteScope(i.scope) ? "review" : "unscoped") : critical ? "review" : "current";169 const claim = await writeClaim(tx, ctx, subjectType, subjectId, { predicate: i.predicate, value: i.value, unit: "MW", scope: i.scope, scopeReason: i.scopeReason ?? null, evidence: i.evidence ?? null, publishedAt: i.publishedAt ?? null, provenance: i.provenance, parserName: i.parserName ?? null, status, rejectionReason: blockedBy.length ? blockedBy.join(",") : null });170 await writeQualityFlags(tx, ctx, subjectType, subjectId, flags.map((f) => ({ ...f, field: f.field ?? i.field })), { claimId: claim.id, priority: reviewPriority({ mw: i.value, severity: flags.some((f) => f.severity === "critical") ? "critical" : flags.some((f) => f.severity === "warn") ? "warn" : "info" }) });171 if (!blockedBy.length) {172 const cleared = ["scope_unknown", "scope_company", "scope_portfolio", "scope_country", "scope_metro", "mw_market_statistic", "mw_invalid"].filter((c) => !flags.some((f) => f.code === c));173 await resolveQualityFlags(tx, ctx, subjectType, subjectId, cleared.map((c) => c));174 }175 // a critical (non-blocking) flag such as "> 1 000 MW single site" still assigns — the figure is what the source says — but176 // stays visible for review; a blocking flag never assigns177 return { assign: blockedBy.length === 0, claim, flags, blockedBy };178}179180export interface InvestmentClaimInput {181 value: number;182 currency?: string | null;183 predicate: InvestmentPredicate;184 scope: ClaimScope;185 scopeReason?: string | null;186 evidence?: Evidence | null;187 context?: string | null;188 previous?: number | null;189 recordScope: "facility" | "campus" | "project";190 publishedAt?: string | null;191 provenance: Provenance;192 parserName?: string | null;193}194195export async function recordInvestmentClaim(tx: Tx, ctx: IngestContext, subjectType: "facility" | "project" | "campus" | "operator", subjectId: string, i: InvestmentClaimInput): Promise<CapacityDecision> {196 const flags = investmentSanity({ value: i.value, scope: i.scope, predicate: i.predicate, recordScope: i.recordScope, previous: i.previous ?? null, context: i.context ?? i.evidence?.text ?? null });197 const blockedBy = flags.filter((f) => f.blocks).map((f) => f.code);198 const status: ClaimStatus = blockedBy.length ? (isSiteScope(i.scope) && i.predicate === "project_investment_usd" ? "review" : "unscoped") : flags.some((f) => f.severity === "critical") ? "review" : "current";199 const claim = await writeClaim(tx, ctx, subjectType, subjectId, { predicate: i.predicate, value: i.value, unit: i.currency ?? "USD", scope: i.scope, scopeReason: i.scopeReason ?? null, evidence: i.evidence ?? null, publishedAt: i.publishedAt ?? null, provenance: i.provenance, parserName: i.parserName ?? null, status, rejectionReason: blockedBy.length ? blockedBy.join(",") : null });200 await writeQualityFlags(tx, ctx, subjectType, subjectId, flags.map((f) => ({ ...f, field: f.field ?? "investmentUsd" })), { claimId: claim.id, priority: reviewPriority({ investmentUsd: i.value, severity: flags.some((f) => f.severity === "critical") ? "critical" : "warn" }) });201 return { assign: blockedBy.length === 0 && (i.currency ?? "USD") === "USD", claim, flags, blockedBy };202}203204/** Best current claim for a predicate (highest tier, then most recent), for the evidence drawer and reconciliation. */205export async function bestCurrentClaim(tx: Tx, subjectType: string, subjectId: string, predicate: string): Promise<{ id: string; value: number | null; tier: AuthorityTier; url: string; evidence: string | null } | null> {206 const rows = await tx.execute(sql`select id, value, authority_tier, url, evidence_text from claims where subject_type = ${subjectType} and subject_id = ${subjectId} and predicate = ${predicate} and status = 'current' order by last_observed desc limit 50`);207 if (!rows.length) return null;208 const sorted = [...rows].sort((a, b) => (TIER_RANK[String(b.authority_tier) as AuthorityTier] ?? 0) - (TIER_RANK[String(a.authority_tier) as AuthorityTier] ?? 0));209 const r = sorted[0]!;210 return { id: String(r.id), value: r.value == null ? null : Number(r.value), tier: String(r.authority_tier) as AuthorityTier, url: String(r.url), evidence: r.evidence_text == null ? null : String(r.evidence_text) };211}212213/** A generic non-numeric claim (status, operator, opening date…) — same store, text value. */214export async function writeTextClaim(tx: Tx, ctx: IngestContext, subjectType: string, subjectId: string, predicate: string, value: string, provenance: Provenance, opts: { scope?: ClaimScope; evidence?: Evidence | null; publishedAt?: string | null } = {}): Promise<WrittenClaim> {215 return writeClaim(tx, ctx, subjectType, subjectId, { predicate, valueText: value, scope: opts.scope ?? "facility", evidence: opts.evidence ?? null, publishedAt: opts.publishedAt ?? null, provenance, status: "current" });216}217218export { newId as _newClaimId };219