/** * Claim store + quality flags + winner bookkeeping (docs/CLAIMS.md). * * Every numerical figure observed about a subject becomes a `claims` row carrying its scope, semantics, supporting * sentence, parser and field-level authority tier. The ingest layer then asks `resolveCapacityColumn()` which claim * may populate a column: only site-scoped claims (building / facility / campus) that pass the deterministic sanity * engine; everything else is stored with status `unscoped` / `review` and surfaces in the admin quality dashboard. */ import { sql } from "@dci/db"; import { authorityFieldFor, authorityTier, capacitySanity, investmentSanity, isSiteScope, newId, stableId, TIER_RANK, type AuthorityTier, type CapacityPredicate, type ClaimScope, type ClaimStatus, type Evidence, type InvestmentPredicate, type Provenance, type SanityFlag, } from "@dci/core"; import type { IngestContext, Tx } from "./common.js"; export interface ClaimInput { predicate: string; value?: number | null; valueText?: string | null; unit?: string | null; scope: ClaimScope; scopeReason?: string | null; evidence?: Evidence | null; publishedAt?: string | null; provenance: Provenance; parserName?: string | null; /** override the computed status (e.g. "review" after a sanity flag) */ status?: ClaimStatus; rejectionReason?: string | null; } export interface WrittenClaim { id: string; status: ClaimStatus; tier: AuthorityTier } export function claimId(subjectType: string, subjectId: string, c: { predicate: string; value?: number | null; valueText?: string | null; provenance: Provenance }, url: string): string { return stableId("claim", `${subjectType}|${subjectId}|${c.predicate}|${c.provenance.sourceId}|${url}|${c.value ?? ""}|${c.valueText ?? ""}`); } /** Upsert one claim. Same source + url + predicate with a different value supersedes the earlier claim. */ export async function writeClaim(tx: Tx, ctx: IngestContext, subjectType: string, subjectId: string, c: ClaimInput): Promise { const p = c.provenance; const url = p.url || ctx.doc?.url || ""; const sourceId = p.sourceId || ctx.run.sourceId; const tier = authorityTier({ field: authorityFieldFor(c.predicate), sourceKind: ctx.run.sourceKind, isEstimate: p.isEstimate, method: p.method }); const status: ClaimStatus = c.status ?? (isSiteScope(c.scope) ? "current" : "unscoped"); const id = claimId(subjectType, subjectId, c, url); const observed = p.lastObserved || ctx.now; if (!ctx.run.dryRun) { await tx.execute(sql` 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, confidence, is_estimate, authority_tier, evidence_text, evidence_start, evidence_end, parser_name, parser_version, run_id, status, rejection_reason, first_observed, last_observed) 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}, ${p.documentId ?? ctx.doc?.documentId ?? null}, ${url}, ${c.publishedAt ?? null}, ${p.retrievedAt || ctx.now}, ${p.confidence ?? "moderate"}, ${!!p.isEstimate}, ${tier}, ${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}) on conflict (id) do update set last_observed = excluded.last_observed, retrieved_at = excluded.retrieved_at, confidence = excluded.confidence, authority_tier = excluded.authority_tier, 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), scope = excluded.scope, scope_reason = excluded.scope_reason, run_id = excluded.run_id, parser_version = excluded.parser_version, status = case when claims.status in ('rejected') then claims.status else excluded.status end, rejection_reason = coalesce(excluded.rejection_reason, claims.rejection_reason)`); // the same source, same page, same predicate now says something else → the older claim is superseded 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')`); } ctx.stats.claims = (ctx.stats.claims ?? 0) + 1; return { id, status, tier }; } /** Review priority 0–100: impact-weighted (MW / money size, low confidence, new country, AI, ambiguity, unexpected change). */ export 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 { let p = 0; const mw = i.mw ?? 0; p += mw >= 1000 ? 40 : mw >= 300 ? 30 : mw >= 100 ? 20 : mw >= 20 ? 10 : mw > 0 ? 4 : 0; const inv = i.investmentUsd ?? 0; p += inv >= 10e9 ? 25 : inv >= 1e9 ? 15 : inv >= 100e6 ? 8 : 0; if (i.confidence === "unverified" || i.confidence === "estimated") p += 10; if (i.newCountry) p += 10; if (i.ai) p += 5; if (i.ambiguous) p += 8; if (i.unexpectedChange) p += 12; if (i.severity === "critical") p += 15; else if (i.severity === "warn") p += 5; if (i.homepageVisible) p += 10; return Math.max(0, Math.min(100, Math.round(p))); } /** Upsert deterministic quality flags for an entity (deduped by entity + code + field). Resolved flags are not reopened for the same key. */ export async function writeQualityFlags(tx: Tx, ctx: IngestContext, entityType: string, entityId: string, flags: SanityFlag[], opts: { claimId?: string | null; priority?: number; details?: Record } = {}): Promise { let n = 0; for (const f of flags) { if (f.severity === "info") continue; const dedupeKey = `${entityType}|${entityId}|${f.code}|${f.field ?? ""}`; const id = stableId("flag", dedupeKey); const priority = opts.priority ?? reviewPriority({ mw: /mw/.test(f.code) ? f.value : null, investmentUsd: /inv/.test(f.code) ? f.value : null, severity: f.severity }); if (!ctx.run.dryRun) { await tx.execute(sql` insert into quality_flags (id, entity_type, entity_id, claim_id, code, severity, field, message, details, priority, status, run_id, dedupe_key) 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}) 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), run_id = excluded.run_id, updated_at = now(), status = case when quality_flags.status = 'dismissed' then 'dismissed' else 'open' end`); } n++; } ctx.stats.qualityFlags = (ctx.stats.qualityFlags ?? 0) + n; return n; } /** Close open flags of the given codes for an entity (the condition no longer holds). */ export async function resolveQualityFlags(tx: Tx, ctx: IngestContext, entityType: string, entityId: string, codes: string[], resolution = "auto: condition cleared"): Promise { if (!codes.length || ctx.run.dryRun) return; await tx.execute(sql`update quality_flags set status = 'resolved', resolution = ${resolution}, resolved_by = 'system', resolved_at = now(), updated_at = now() where entity_type = ${entityType} and entity_id = ${entityId} and status = 'open' and code in ${codes}`); } /** * Mark, per field, the current provenance row whose value equals the stored column value as the winner (and unmark * the others). This is what makes "why does the page show 300 MW?" answerable without a diff. */ export async function markWinners(tx: Tx, ctx: IngestContext, entityType: string, entityId: string, fields: Array<{ field: string; value: unknown }>): Promise { if (ctx.run.dryRun) return; for (const f of fields) { if (f.value == null) continue; const valueJson = JSON.stringify(f.value); 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)`); } } export interface CapacityDecision { /** write the figure into the record's column */ assign: boolean; claim: WrittenClaim; flags: SanityFlag[]; blockedBy: string[]; } export interface CapacityClaimInput { field: "itCapacityMw" | "totalPowerMw" | "plannedPowerMw" | "plannedMw" | "utilityCapacityMw" | "gridConnectionMw" | "ultimateCampusMw"; predicate: CapacityPredicate; value: number; scope: ClaimScope; scopeReason?: string | null; evidence?: Evidence | null; context?: string | null; previous?: number | null; recordScope: "building" | "facility" | "campus" | "project"; campusDesignation?: boolean; publishedAt?: string | null; provenance: Provenance; parserName?: string | null; /** the sentence did not say what kind of MW this is; semantics were defaulted from the record status */ semanticsDefaulted?: boolean; } /** Store a capacity claim, run the sanity engine, write flags and decide whether the column may take the value. */ export async function recordCapacityClaim(tx: Tx, ctx: IngestContext, subjectType: "facility" | "project" | "campus", subjectId: string, i: CapacityClaimInput): Promise { 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 }); const blockedBy = flags.filter((f) => f.blocks).map((f) => f.code); const critical = flags.some((f) => f.severity === "critical"); const status: ClaimStatus = blockedBy.length ? (isSiteScope(i.scope) ? "review" : "unscoped") : critical ? "review" : "current"; 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 }); 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" }) }); if (!blockedBy.length) { 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)); await resolveQualityFlags(tx, ctx, subjectType, subjectId, cleared.map((c) => c)); } // a critical (non-blocking) flag such as "> 1 000 MW single site" still assigns — the figure is what the source says — but // stays visible for review; a blocking flag never assigns return { assign: blockedBy.length === 0, claim, flags, blockedBy }; } export interface InvestmentClaimInput { value: number; currency?: string | null; predicate: InvestmentPredicate; scope: ClaimScope; scopeReason?: string | null; evidence?: Evidence | null; context?: string | null; previous?: number | null; recordScope: "facility" | "campus" | "project"; publishedAt?: string | null; provenance: Provenance; parserName?: string | null; } export async function recordInvestmentClaim(tx: Tx, ctx: IngestContext, subjectType: "facility" | "project" | "campus" | "operator", subjectId: string, i: InvestmentClaimInput): Promise { 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 }); const blockedBy = flags.filter((f) => f.blocks).map((f) => f.code); const status: ClaimStatus = blockedBy.length ? (isSiteScope(i.scope) && i.predicate === "project_investment_usd" ? "review" : "unscoped") : flags.some((f) => f.severity === "critical") ? "review" : "current"; 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 }); 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" }) }); return { assign: blockedBy.length === 0 && (i.currency ?? "USD") === "USD", claim, flags, blockedBy }; } /** Best current claim for a predicate (highest tier, then most recent), for the evidence drawer and reconciliation. */ export 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> { 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`); if (!rows.length) return null; const sorted = [...rows].sort((a, b) => (TIER_RANK[String(b.authority_tier) as AuthorityTier] ?? 0) - (TIER_RANK[String(a.authority_tier) as AuthorityTier] ?? 0)); const r = sorted[0]!; 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) }; } /** A generic non-numeric claim (status, operator, opening date…) — same store, text value. */ export 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 { return writeClaim(tx, ctx, subjectType, subjectId, { predicate, valueText: value, scope: opts.scope ?? "facility", evidence: opts.evidence ?? null, publishedAt: opts.publishedAt ?? null, provenance, status: "current" }); } export { newId as _newClaimId };