/** * Extraction debugger — the full pipeline for ONE document, stage by stage, without publishing anything: * * SOURCE DOCUMENT → RAW FETCH (archived body or live) → PARSED TEXT → STRUCTURED DATA (extracted records) → * EXTRACTED ENTITIES (normalized) → EXTRACTED CLAIMS (figures with scope / semantics / evidence sentence) → * NORMALIZED VALUES → MATCH CANDIDATES (facility resolution preview) → RECONCILIATION (dry-run ingest: created / * updated / events / rejected) → RESULTING DATABASE CHANGES (what a real run would write). * * Exposed by `dci trace ` and by the worker HTTP endpoint `GET /trace/` (internal network; the API * proxies it as `/api/admin/documents/:id/trace`). */ import type { RawDocument } from "@dci/connectors"; import { entityIsValid, mainText, pageTitle, isPdf, pdfText, classifyPage } from "@dci/connectors"; import type { NormalizedEntity, NormalizedFacility, NormalizedProject, ValidationIssue } from "@dci/core"; import { classifyAiEvidence, classifyCapacitySemantics, classifyInvestmentSemantics, classifyProjectEvent, classifyScope, findEvidence, newId, parseAllMw } from "@dci/core"; import { getDb } from "@dci/db"; import { requireConnector } from "./configs.js"; import { createRunContext } from "./context.js"; import { getDocument } from "./documents.js"; import { parseDiscoveredFrom } from "./scheduling.js"; import { ingestEntities } from "./ingest/index.js"; import { newContext } from "./ingest/common.js"; import { previewFacilityResolution } from "./ingest/facilities.js"; import { extractAnnouncement } from "./connectors/news/extract-project.js"; import { getRaw } from "./storage.js"; export 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; } export interface TraceResult { document: Record | null; fetch: { source: "archive" | "live"; status: number; contentType: string | null; bytes: number; storageKey: string | null; level: number; fetcher: string } | null; text: { title: string | null; length: number; excerpt: string; classification: { pageType: string; rule: string; eventType: string | null; mw: number[] } | null }; announcement: Record | null; records: Array<{ kind: string; key: string; certainty: number | null; pageType: string | null; data: Record; methods: Record }>; entities: NormalizedEntity[]; validation: { total: number; valid: number; rejected: number; issues: ValidationIssue[] }; claims: TraceClaim[]; 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 }> }>; 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; logs: string[]; error: string | null; } /** Highlight offsets of every MW / money figure in the text (for the debugger UI). */ export function highlightFigures(text: string): Array<{ start: number; end: number; kind: "mw" | "usd"; raw: string }> { const out: Array<{ start: number; end: number; kind: "mw" | "usd"; raw: string }> = []; 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] }); 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] }); return out.sort((a, b) => a.start - b.start); } export async function traceDocument(idOrUrl: string, opts: { live?: boolean } = {}): Promise { const logs: string[] = []; 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 }; const doc = await getDocument(idOrUrl); if (!doc) { result.error = `no document for ${idOrUrl}`; return result; } 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 }; const loaded = requireConnector(doc.connectorId); const runId = newId("run"); const ctx = createRunContext(loaded, { runId, dryRun: true, onLog: (l) => logs.push(`${l.level} ${l.msg}`), logLevel: "debug" }); try { // ── raw let raw: RawDocument | null = null; if (!opts.live && doc.storageKey) { const stored = await getRaw(doc.storageKey); if (stored) { const contentType = stored.contentType ?? doc.contentType; const isText = !contentType || /text|json|xml|javascript|html|csv|markdown/i.test(contentType); const { group } = parseDiscoveredFrom(doc.discoveredFrom); 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 }; result.fetch = { source: "archive", status: raw.status, contentType, bytes: stored.body.length, storageKey: doc.storageKey, level: doc.fetchLevel, fetcher: "cache" }; } } if (!raw) { const { group } = parseDiscoveredFrom(doc.discoveredFrom); raw = await loaded.connector.fetch(ctx, { url: doc.url, group, pageType: doc.pageType !== "unknown" ? (doc.pageType as RawDocument["pageType"]) : undefined, priority: 100 }); result.fetch = { source: "live", status: raw.status, contentType: raw.contentType, bytes: raw.body.length, storageKey: null, level: raw.level, fetcher: raw.fetcher }; if (raw.error) { result.error = `fetch: ${raw.error.code} ${raw.error.message}`; return result; } } // ── text let title: string | null = null; let text = ""; if (isPdf(raw)) { const p = await pdfText(raw.body); text = p.text; title = p.title ?? null; } else if (/html|xml/.test(raw.contentType ?? "") || / 80 && /news|press|announcement|planning|construction|expansion|acquisition|power|financial|project|unknown/.test(cls.pageType)) { try { const a = extractAnnouncement(title, text); 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)}`) }; } catch (e) { logs.push(`warn announcement analysis failed: ${(e as Error).message}`); } } // ── structured data const records = await loaded.connector.extract(ctx, raw); 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 ?? {} })); // ── entities + validation const entities = await loaded.connector.normalize(ctx, records); const report = await loaded.connector.validate(ctx, entities); result.entities = entities; result.validation = { total: report.total, valid: report.valid, rejected: report.rejected, issues: report.issues }; // ── claims (what the claim store would receive) for (const e of entities) { if (e.entityType === "facility") { const f = e as NormalizedFacility; for (const field of ["itCapacityMw", "totalPowerMw", "plannedPowerMw"] as const) { const v = f[field]; if (v == null || v <= 0) continue; const context = f.claimContext?.[field] ?? null; const sc = context ? classifyScope(context, "facility") : { scope: "facility", reason: "structured:record" }; const sem = context ? classifyCapacitySemantics(context) : { predicate: field === "itCapacityMw" ? "it_capacity_mw" : field === "totalPowerMw" ? "current_power_mw" : "planned_power_mw", reason: "field" }; 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 }); } } else if (e.entityType === "project") { const p = e as NormalizedProject; 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 }); } 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 }); } } else if (e.entityType === "news_event") { 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") }); } } // ── match candidates (facilities) — read-only, inside a rolled-back transaction const valid = entities.filter((x) => entityIsValid(report, x.key)); const db = getDb(); const run = { connectorId: loaded.cfg.id, sourceId: loaded.sourceId, runId, sourceKind: loaded.cfg.kind, sourcePriority: loaded.cfg.priority, dryRun: true }; await db.transaction(async (tx) => { const ictx = newContext(run, { documentId: doc.id, url: raw!.finalUrl, pageType: cls.pageType, fetchedAt: raw!.fetchedAt }); for (const e of valid) { if (e.entityType !== "facility") continue; 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}`); } } throw new Error("__trace_rollback__"); }).catch((e: Error) => { if (e.message !== "__trace_rollback__") throw e; }); // ── reconciliation (dry-run ingest: everything runs, nothing is committed) if (valid.length) { const st = await ingestEntities(run, valid, { documentId: doc.id, url: raw.finalUrl, pageType: cls.pageType, fetchedAt: raw.fetchedAt }); 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 }; } 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: [] }; // AI grading of the whole text for the debugger logs.push(`debug ai-evidence: ${JSON.stringify(classifyAiEvidence(`${title ?? ""} ${text.slice(0, 3000)}`))}`); logs.push(`debug mw figures in text: ${parseAllMw(text.slice(0, 6000)).join(", ") || "none"}`); // the announcement classifier verdict for non-news pages too if (!result.announcement && title) logs.push(`debug classifier: ${JSON.stringify(classifyProjectEvent({ title, lead: text.slice(0, 600) }))}`); } catch (e) { result.error = (e as Error).stack ?? (e as Error).message; } return result; }