SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
13.9 KB · 150 lines typescript
Raw Blame History
1/**2 * Extraction debugger — the full pipeline for ONE document, stage by stage, without publishing anything:3 *4 *   SOURCE DOCUMENT → RAW FETCH (archived body or live) → PARSED TEXT → STRUCTURED DATA (extracted records) →5 *   EXTRACTED ENTITIES (normalized) → EXTRACTED CLAIMS (figures with scope / semantics / evidence sentence) →6 *   NORMALIZED VALUES → MATCH CANDIDATES (facility resolution preview) → RECONCILIATION (dry-run ingest: created /7 *   updated / events / rejected) → RESULTING DATABASE CHANGES (what a real run would write).8 *9 * Exposed by `dci trace <doc-id|url>` and by the worker HTTP endpoint `GET /trace/<doc-id>` (internal network; the API10 * proxies it as `/api/admin/documents/:id/trace`).11 */12import type { RawDocument } from "@dci/connectors";13import { entityIsValid, mainText, pageTitle, isPdf, pdfText, classifyPage } from "@dci/connectors";14import type { NormalizedEntity, NormalizedFacility, NormalizedProject, ValidationIssue } from "@dci/core";15import { classifyAiEvidence, classifyCapacitySemantics, classifyInvestmentSemantics, classifyProjectEvent, classifyScope, findEvidence, newId, parseAllMw } from "@dci/core";16import { getDb } from "@dci/db";17import { requireConnector } from "./configs.js";18import { createRunContext } from "./context.js";19import { getDocument } from "./documents.js";20import { parseDiscoveredFrom } from "./scheduling.js";21import { ingestEntities } from "./ingest/index.js";22import { newContext } from "./ingest/common.js";23import { previewFacilityResolution } from "./ingest/facilities.js";24import { extractAnnouncement } from "./connectors/news/extract-project.js";25import { getRaw } from "./storage.js";2627export interface TraceClaim { entityKey: string; field: string; value: number; unit: "MW" | "USD"; scope: string; scopeReason: string; semantics: string | null; evidence: { text: string; start: number; end: number } | null; }2829export interface TraceResult {30  document: Record<string, unknown> | null;31  fetch: { source: "archive" | "live"; status: number; contentType: string | null; bytes: number; storageKey: string | null; level: number; fetcher: string } | null;32  text: { title: string | null; length: number; excerpt: string; classification: { pageType: string; rule: string; eventType: string | null; mw: number[] } | null };33  announcement: Record<string, unknown> | null;34  records: Array<{ kind: string; key: string; certainty: number | null; pageType: string | null; data: Record<string, unknown>; methods: Record<string, string> }>;35  entities: NormalizedEntity[];36  validation: { total: number; valid: number; rejected: number; issues: ValidationIssue[] };37  claims: TraceClaim[];38  matches: Array<{ entityKey: string; how: string; matchedId: string | null; candidates: Array<{ id: string; name: string; operatorName: string | null; city: string | null; score: number; reasons: string[]; distanceKm: number | null }> }>;39  reconciliation: { created: number; updated: number; unchanged: number; merged: number; pendingMatches: number; rejected: number; events: number; provenanceRows: number; claims: number; qualityFlags: number; unscopedClaims: number; projectsVetoed: number; changes: unknown[]; refs: Array<{ type: string; id: string }> } | null;40  logs: string[];41  error: string | null;42}4344/** Highlight offsets of every MW / money figure in the text (for the debugger UI). */45export function highlightFigures(text: string): Array<{ start: number; end: number; kind: "mw" | "usd"; raw: string }> {46  const out: Array<{ start: number; end: number; kind: "mw" | "usd"; raw: string }> = [];47  for (const m of text.matchAll(/(\d{1,3}(?:[.,]\d{3})+|\d+(?:[.,]\d+)?)\+?\s?(gigawatts?|gw|megawatts?|mw)\b/gi)) out.push({ start: m.index ?? 0, end: (m.index ?? 0) + m[0].length, kind: "mw", raw: m[0] });48  for (const m of text.matchAll(/(?:US\$|USD|\$|€|£|A\$|C\$|S\$)\s?\d[\d.,]*\s*(?:trillion|billion|million|bn|m\b|b\b|k\b)?/gi)) out.push({ start: m.index ?? 0, end: (m.index ?? 0) + m[0].length, kind: "usd", raw: m[0] });49  return out.sort((a, b) => a.start - b.start);50}5152export async function traceDocument(idOrUrl: string, opts: { live?: boolean } = {}): Promise<TraceResult> {53  const logs: string[] = [];54  const result: TraceResult = { document: null, fetch: null, text: { title: null, length: 0, excerpt: "", classification: null }, announcement: null, records: [], entities: [], validation: { total: 0, valid: 0, rejected: 0, issues: [] }, claims: [], matches: [], reconciliation: null, logs, error: null };55  const doc = await getDocument(idOrUrl);56  if (!doc) { result.error = `no document for ${idOrUrl}`; return result; }57  result.document = { id: doc.id, connectorId: doc.connectorId, url: doc.url, pageType: doc.pageType, classifier: doc.classifier, contentHash: doc.contentHash, extractorVersion: doc.extractorVersion, extractOk: doc.extractOk, extractCount: doc.extractCount, error: doc.error, lastFetched: doc.lastFetched, lastChanged: doc.lastChanged, fetchLevel: doc.fetchLevel, storageKey: doc.storageKey, title: doc.title, entityRefs: doc.entityRefs, quarantined: doc.quarantined };58  const loaded = requireConnector(doc.connectorId);59  const runId = newId("run");60  const ctx = createRunContext(loaded, { runId, dryRun: true, onLog: (l) => logs.push(`${l.level} ${l.msg}`), logLevel: "debug" });61  try {62    // ── raw63    let raw: RawDocument | null = null;64    if (!opts.live && doc.storageKey) {65      const stored = await getRaw(doc.storageKey);66      if (stored) {67        const contentType = stored.contentType ?? doc.contentType;68        const isText = !contentType || /text|json|xml|javascript|html|csv|markdown/i.test(contentType);69        const { group } = parseDiscoveredFrom(doc.discoveredFrom);70        raw = { url: doc.url, finalUrl: doc.url, fetchedAt: doc.lastFetched ?? new Date().toISOString(), status: doc.statusCode ?? 200, contentType, body: stored.body, text: isText ? stored.body.toString("utf8") : "", headers: {}, etag: doc.etag, lastModified: doc.lastModified, notModified: false, fetcher: "cache", level: doc.fetchLevel as RawDocument["level"], durationMs: 0, credits: 0, group, pageType: doc.pageType !== "unknown" ? (doc.pageType as RawDocument["pageType"]) : undefined };71        result.fetch = { source: "archive", status: raw.status, contentType, bytes: stored.body.length, storageKey: doc.storageKey, level: doc.fetchLevel, fetcher: "cache" };72      }73    }74    if (!raw) {75      const { group } = parseDiscoveredFrom(doc.discoveredFrom);76      raw = await loaded.connector.fetch(ctx, { url: doc.url, group, pageType: doc.pageType !== "unknown" ? (doc.pageType as RawDocument["pageType"]) : undefined, priority: 100 });77      result.fetch = { source: "live", status: raw.status, contentType: raw.contentType, bytes: raw.body.length, storageKey: null, level: raw.level, fetcher: raw.fetcher };78      if (raw.error) { result.error = `fetch: ${raw.error.code} ${raw.error.message}`; return result; }79    }80    // ── text81    let title: string | null = null;82    let text = "";83    if (isPdf(raw)) { const p = await pdfText(raw.body); text = p.text; title = p.title ?? null; }84    else if (/html|xml/.test(raw.contentType ?? "") || /<html/i.test(raw.text.slice(0, 2000))) { title = pageTitle(raw.text); text = mainText(raw.text); }85    else text = raw.markdown ?? raw.text;86    const cls = classifyPage(raw.finalUrl, title, text);87    result.text = { title, length: text.length, excerpt: text.slice(0, 6000), classification: { pageType: cls.pageType, rule: cls.rule, eventType: cls.eventType ?? null, mw: cls.mw ?? [] } };88    // ── announcement analysis (news-like pages)89    if (text.length > 80 && /news|press|announcement|planning|construction|expansion|acquisition|power|financial|project|unknown/.test(cls.pageType)) {90      try {91        const a = extractAnnouncement(title, text);92        result.announcement = { title: a.title, pageType: a.pageType, eventType: a.eventType ?? null, relevance: a.relevance, leadRelevance: a.leadRelevance, titleSignal: a.titleSignal, nonProjectTitle: a.nonProjectTitle, classification: a.classification, status: a.status, headlineMw: a.headlineMw, mwAll: a.mwAll, money: a.money, investmentUsd: a.investmentUsd, acreage: a.acreage, phaseCount: a.phaseCount, expectedOpening: a.expectedOpening, location: a.location, operator: a.operator?.name ?? null, operators: a.operators.map((o) => o.name), explicitName: a.explicitName, projectName: a.projectName, capacityScope: a.capacityScope, capacitySemantics: a.capacitySemantics, investmentScope: a.investmentScope, investmentSemantics: a.investmentSemantics, aiEvidence: a.aiEvidence, hqGuarded: a.hqGuarded, evidence: a.evidence, methods: a.methods, figures: highlightFigures(`${a.title}\n${text.slice(0, 6000)}`) };93      } catch (e) { logs.push(`warn announcement analysis failed: ${(e as Error).message}`); }94    }95    // ── structured data96    const records = await loaded.connector.extract(ctx, raw);97    result.records = records.map((r) => ({ kind: r.kind, key: r.key, certainty: r.certainty ?? null, pageType: r.pageType ?? null, data: r.data, methods: r.methods ?? {} }));98    // ── entities + validation99    const entities = await loaded.connector.normalize(ctx, records);100    const report = await loaded.connector.validate(ctx, entities);101    result.entities = entities;102    result.validation = { total: report.total, valid: report.valid, rejected: report.rejected, issues: report.issues };103    // ── claims (what the claim store would receive)104    for (const e of entities) {105      if (e.entityType === "facility") {106        const f = e as NormalizedFacility;107        for (const field of ["itCapacityMw", "totalPowerMw", "plannedPowerMw"] as const) {108          const v = f[field];109          if (v == null || v <= 0) continue;110          const context = f.claimContext?.[field] ?? null;111          const sc = context ? classifyScope(context, "facility") : { scope: "facility", reason: "structured:record" };112          const sem = context ? classifyCapacitySemantics(context) : { predicate: field === "itCapacityMw" ? "it_capacity_mw" : field === "totalPowerMw" ? "current_power_mw" : "planned_power_mw", reason: "field" };113          result.claims.push({ entityKey: f.key, field, value: v, unit: "MW", scope: sc.scope, scopeReason: sc.reason, semantics: sem.predicate, evidence: context ? findEvidence(context, v, "mw") : null });114        }115      } else if (e.entityType === "project") {116        const p = e as NormalizedProject;117        if (p.plannedMw != null && p.plannedMw > 0) { const c = p.claimContext?.plannedMw ?? null; result.claims.push({ entityKey: p.key, field: "plannedMw", value: p.plannedMw, unit: "MW", scope: p.capacityScope ?? classifyScope(c, "facility").scope, scopeReason: p.capacityScope ? "extractor" : classifyScope(c, "facility").reason, semantics: p.capacitySemantics ?? null, evidence: c ? findEvidence(c, p.plannedMw, "mw") : null }); }118        if (p.investmentUsd != null && p.investmentUsd > 0) { const c = p.claimContext?.investmentUsd ?? null; const sem = classifyInvestmentSemantics(c ?? p.name); result.claims.push({ entityKey: p.key, field: "investmentUsd", value: p.investmentUsd, unit: "USD", scope: p.investmentScope ?? sem.scope, scopeReason: p.investmentScope ? "extractor" : sem.reason, semantics: p.investmentSemantics ?? sem.predicate, evidence: c ? findEvidence(c, p.investmentUsd, "usd") : null }); }119      } else if (e.entityType === "news_event") {120        for (const mw of e.mentions?.mw ?? []) result.claims.push({ entityKey: e.key, field: "mw", value: mw, unit: "MW", scope: "unknown", scopeReason: "news mention (never assigned)", semantics: null, evidence: findEvidence(text, mw, "mw") });121      }122    }123    // ── match candidates (facilities) — read-only, inside a rolled-back transaction124    const valid = entities.filter((x) => entityIsValid(report, x.key));125    const db = getDb();126    const run = { connectorId: loaded.cfg.id, sourceId: loaded.sourceId, runId, sourceKind: loaded.cfg.kind, sourcePriority: loaded.cfg.priority, dryRun: true };127    await db.transaction(async (tx) => {128      const ictx = newContext(run, { documentId: doc.id, url: raw!.finalUrl, pageType: cls.pageType, fetchedAt: raw!.fetchedAt });129      for (const e of valid) {130        if (e.entityType !== "facility") continue;131        try { const m = await previewFacilityResolution(tx, ictx, e as NormalizedFacility); result.matches.push({ entityKey: e.key, ...m }); } catch (err) { logs.push(`warn match preview ${e.key}: ${(err as Error).message}`); }132      }133      throw new Error("__trace_rollback__");134    }).catch((e: Error) => { if (e.message !== "__trace_rollback__") throw e; });135    // ── reconciliation (dry-run ingest: everything runs, nothing is committed)136    if (valid.length) {137      const st = await ingestEntities(run, valid, { documentId: doc.id, url: raw.finalUrl, pageType: cls.pageType, fetchedAt: raw.fetchedAt });138      result.reconciliation = { created: st.created, updated: st.updated, unchanged: st.unchanged, merged: st.merged, pendingMatches: st.pendingMatches, rejected: st.rejected, events: st.events, provenanceRows: st.provenanceRows, claims: st.claims ?? 0, qualityFlags: st.qualityFlags ?? 0, unscopedClaims: st.unscopedClaims ?? 0, projectsVetoed: st.projectsVetoed ?? 0, changes: st.changes, refs: st.refs };139    } else result.reconciliation = { created: 0, updated: 0, unchanged: 0, merged: 0, pendingMatches: 0, rejected: report.rejected, events: 0, provenanceRows: 0, claims: 0, qualityFlags: 0, unscopedClaims: 0, projectsVetoed: 0, changes: [], refs: [] };140    // AI grading of the whole text for the debugger141    logs.push(`debug ai-evidence: ${JSON.stringify(classifyAiEvidence(`${title ?? ""} ${text.slice(0, 3000)}`))}`);142    logs.push(`debug mw figures in text: ${parseAllMw(text.slice(0, 6000)).join(", ") || "none"}`);143    // the announcement classifier verdict for non-news pages too144    if (!result.announcement && title) logs.push(`debug classifier: ${JSON.stringify(classifyProjectEvent({ title, lead: text.slice(0, 600) }))}`);145  } catch (e) {146    result.error = (e as Error).stack ?? (e as Error).message;147  }148  return result;149}150