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%
15.0 KB · 219 lines typescript
Raw Blame History
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