spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * `pnpm dci <cmd>` — operator CLI for the crawl runtime.3 *4 * connectors table of connectors (id, kind, mode, enabled, health, last/next run, docs)5 * sync YAML configs → connectors + sources tables6 * run <id> [--task full|discover|crawl|reprocess] [--group g] [--limit n] [--dry-run] [--url u]… [--force] [--json]7 * run-all [--dry-run] [--task t] [--limit n]8 * discover <id> [--dry-run]9 * docs <id> [--due] [--limit n] [--group g]10 * inspect <url-or-docId> document row + latest version + provenance rows referencing it11 * reprocess <id> [--limit n] [--group g] [--stale] re-extract from archived bodies (no network); --stale = older parser version only12 * trace <doc-id|url> [--live] [--json] extraction debugger: every pipeline stage for one document, nothing persisted13 * quarantine <id> [on|off] preview-only mode for a connector (ingest rolled back) · quality · snapshot · gaps14 * stats global counts (entities, documents, runs, events, queues)15 * doctor Postgres / Redis / ClickHouse / MinIO, budgets, stale connectors, error rates → system_alerts16 * rank | metrics | refresh-stats maintenance jobs, run inline17 * enqueue <id> [--task t] [--group g] push a crawl job to the queue (worker must be running)18 * pause <id> | resume <id>19 * scheduler run only the scheduler loop (foreground)20 * worker run the full worker (= main.ts)21 */22import { existsSync, readFileSync } from "node:fs";23import { getDb, closeDb, sql, provenance as provenanceTable, connectors as connectorsTable, eq } from "@dci/db";24import { parseConnectorConfig, validateAgainstRegistry } from "@dci/connectors";25import { registerAllConnectors } from "./connectors/index.js";26import { loadAllConnectors, requireConnector, syncConnectorsToDb } from "./configs.js";27import { connectorDocStats, dueDocuments, getDocument, latestVersion } from "./documents.js";28import { getEnv, loadEnvFile } from "./env.js";29import { closeAll, doctor } from "./maintenance.js";30import { computeDailyMetrics, refreshStats } from "./metrics.js";31import { requestAbort, runAll, runConnector, type RunResult, type RunTask } from "./pipeline.js";32import { computeRankings } from "./rankings.js";33import { closeScheduler, enqueueRun, pauseConnector, queueSnapshot, schedulerTick, startSchedulerLoop } from "./scheduler.js";34import { parseDiscoveredFrom } from "./scheduling.js";35import { getRaw } from "./storage.js";36import { traceDocument } from "./trace.js";37import { dataGaps, qualitySweep, snapshotAndCheck } from "./quality.js";38import { ensureClickHouse } from "@dci/db/clickhouse";3940/* ---------- arg parsing ---------- */41interface Args { cmd: string; positional: string[]; flags: Record<string, string | string[] | boolean> }42function parseArgs(argv: string[]): Args {43 const [cmd = "help", ...rest] = argv;44 const positional: string[] = [];45 const flags: Args["flags"] = {};46 for (let i = 0; i < rest.length; i++) {47 const a = rest[i]!;48 if (a.startsWith("--")) {49 const [k, inline] = a.slice(2).split("=", 2);50 const key = k!;51 let v: string | boolean = true;52 if (inline !== undefined) v = inline;53 else if (rest[i + 1] !== undefined && !rest[i + 1]!.startsWith("--")) v = rest[++i]!;54 const prev = flags[key];55 if (prev === undefined) flags[key] = v;56 else if (Array.isArray(prev)) prev.push(String(v));57 else flags[key] = [String(prev), String(v)];58 } else positional.push(a);59 }60 return { cmd, positional, flags };61}62const str = (a: Args, k: string): string | undefined => { const v = a.flags[k]; return v === undefined || typeof v === "boolean" ? undefined : Array.isArray(v) ? v[v.length - 1] : v; };63const num = (a: Args, k: string): number | undefined => { const v = str(a, k); return v === undefined ? undefined : Number(v); };64const bool = (a: Args, k: string): boolean => a.flags[k] !== undefined && a.flags[k] !== "false";65const list = (a: Args, k: string): string[] => { const v = a.flags[k]; return v === undefined || typeof v === "boolean" ? [] : Array.isArray(v) ? v : [v]; };6667/* ---------- output ---------- */68const C = { reset: "\x1b[0m", dim: "\x1b[2m", bold: "\x1b[1m", red: "\x1b[31m", green: "\x1b[32m", yellow: "\x1b[33m", cyan: "\x1b[36m" };69const color = process.stdout.isTTY && !process.env.NO_COLOR;70const paint = (c: keyof typeof C, s: string) => (color ? `${C[c]}${s}${C.reset}` : s);71function table(rows: Array<Record<string, string | number | null | undefined>>, cols: string[]): string {72 const w = cols.map((c) => Math.max(c.length, ...rows.map((r) => String(r[c] ?? "").length)));73 const line = (vals: string[]) => vals.map((v, i) => v.padEnd(w[i]!)).join(" ");74 return [paint("bold", line(cols)), paint("dim", line(w.map((n) => "─".repeat(n)))), ...rows.map((r) => line(cols.map((c) => String(r[c] ?? ""))))].join("\n");75}76const ago = (iso: string | null | undefined): string => { if (!iso) return "—"; const ms = Date.now() - Date.parse(iso); const abs = Math.abs(ms); const u = abs < 60_000 ? `${Math.round(abs / 1000)}s` : abs < 3_600_000 ? `${Math.round(abs / 60_000)}m` : abs < 86_400_000 ? `${Math.round(abs / 3_600_000)}h` : `${Math.round(abs / 86_400_000)}d`; return ms >= 0 ? `${u} ago` : `in ${u}`; };77const healthColor = (h: string) => (h === "ok" ? paint("green", h) : h === "degraded" ? paint("yellow", h) : h === "failing" ? paint("red", h) : paint("dim", h));7879function printRun(r: RunResult, json: boolean): void {80 if (json) { console.log(JSON.stringify(r, null, 2)); return; }81 const s = r.stats;82 const st = r.status === "ok" ? paint("green", r.status) : r.status === "partial" ? paint("yellow", r.status) : paint("red", r.status);83 console.log(`\n${paint("bold", r.connectorId)} ${r.task} → ${st} in ${r.durationMs} ms (run ${r.runId})`);84 console.log(` discovered ${s.discovered} · registered ${s.registered} · selected ${s.selected} · fetched ${s.fetched} · 304 ${s.notModified} · changed ${s.changed} · unchanged ${s.unchanged} · failed ${s.failed}${s.quarantined ? ` · quarantined ${s.quarantined}` : ""}${s.robotsBlocked ? ` · robots ${s.robotsBlocked}` : ""}`);85 console.log(` extracted ${s.extracted} docs · entities ${s.valid}/${s.entities} valid (${s.rejected} rejected) · created ${s.created} · updated ${s.updated} · events ${s.events} · provenance ${s.provenanceRows} · credits ${s.credits} · avg ${s.avgMs} ms/doc`);86 if (r.error) console.log(` ${paint("red", "error")}: ${r.error.split("\n")[0]}`);87 const errs = r.issues.filter((i) => i.level === "error");88 if (errs.length) { console.log(` ${paint("yellow", `${errs.length} validation error(s)`)}:`); for (const i of errs.slice(0, 15)) console.log(` ${i.key}${i.field ? "." + i.field : ""}: ${i.message}`); }89 if (r.samples.length) {90 console.log(` ${paint("cyan", `${r.samples.length} sample entit${r.samples.length > 1 ? "ies" : "y"}`)}:`);91 for (const e of r.samples.slice(0, 10)) {92 const o = e as unknown as Record<string, unknown>;93 const brief = Object.fromEntries(Object.entries(o).filter(([k, v]) => v != null && v !== "" && !(Array.isArray(v) && !v.length) && !["provenance", "facts", "description", "timeline"].includes(k)).slice(0, 14));94 console.log(` ${JSON.stringify(brief).slice(0, 600)}`);95 }96 }97}9899/* ---------- commands ---------- */100async function cmdConnectors(a: Args): Promise<number> {101 await registerAllConnectors();102 const loaded = loadAllConnectors();103 const rows = await getDb().execute<{ id: string; kind: string; mode: string; enabled: boolean; paused: boolean; health: string; last_run_at: string | null; last_status: string | null; next_run_at: string | null; docs: number; due: number }>(sql`104 select c.id, c.kind, c.mode, c.enabled, c.paused, c.health, c.last_run_at, c.last_status, c.next_run_at,105 (select count(*)::int from documents d where d.connector_id = c.id) as docs,106 (select count(*)::int from documents d where d.connector_id = c.id and d.quarantined = false and (d.next_check is null or d.next_check <= now())) as due107 from connectors c order by c.id`);108 const inDb = new Set(rows.map((r) => r.id));109 const out = rows.map((r) => ({ id: r.id, kind: r.kind, mode: r.mode, enabled: r.paused ? paint("yellow", "paused") : r.enabled ? paint("green", "yes") : paint("dim", "no"), health: healthColor(r.health), "last run": r.last_run_at ? `${ago(r.last_run_at)} ${r.last_status ?? ""}` : "never", "next run": r.next_run_at ? ago(r.next_run_at) : "—", docs: `${r.docs}${r.due ? ` (${r.due} due)` : ""}`, yaml: loaded.some((l) => l.cfg.id === r.id) ? "yes" : paint("red", "missing") }));110 for (const l of loaded) if (!inDb.has(l.cfg.id)) out.push({ id: l.cfg.id, kind: l.cfg.kind, mode: l.cfg.mode, enabled: paint("dim", "not synced"), health: paint("dim", "—"), "last run": "—", "next run": "—", docs: "—", yaml: "yes" });111 if (bool(a, "json")) console.log(JSON.stringify(rows, null, 2)); else console.log(out.length ? table(out, ["id", "kind", "mode", "enabled", "health", "last run", "next run", "docs", "yaml"]) : paint("dim", `no connectors (config dir: ${getEnv().configDir}) — run \`pnpm dci sync\` after adding YAML files`));112 return 0;113}114115/** ClickHouse is optional analytics: make sure its tables exist, never fail the CLI on it. */116async function ensureClickHouseQuiet(): Promise<void> {117 try { await ensureClickHouse(); } catch (e) { console.error(paint("dim", `clickhouse unavailable (analytics disabled): ${(e as Error).message.split("\n")[0]}`)); }118}119120async function cmdSync(): Promise<number> {121 await registerAllConnectors();122 await ensureClickHouseQuiet();123 const loaded = loadAllConnectors({ reload: true });124 const r = await syncConnectorsToDb(loaded);125 console.log(`synced ${r.connectors} connector(s), ${r.sources} source(s) from ${getEnv().configDir}${r.disabledInDb.length ? `; disabled (no YAML): ${r.disabledInDb.join(", ")}` : ""}`);126 for (const l of loaded) console.log(` ${l.cfg.id.padEnd(28)} ${l.cfg.kind.padEnd(14)} ${l.cfg.mode.padEnd(8)} ${l.cfg.enabled ? "enabled " : "disabled"} L${l.cfg.fetch.level}→${l.cfg.fetch.maxLevel} rpm ${l.cfg.fetch.rpm} ${Object.entries(l.cfg.schedule).map(([g, s]) => `${g}=${s}`).join(" ")}`);127 return 0;128}129130function taskOf(a: Args, def: RunTask): RunTask {131 const t = str(a, "task") ?? def;132 if (!["full", "discover", "crawl", "reprocess"].includes(t)) throw new Error(`bad --task ${t}`);133 return t as RunTask;134}135136async function cmdRun(a: Args): Promise<number> {137 const id = a.positional[0];138 if (!id) throw new Error("usage: dci run <connector-id> [--task full|discover|crawl|reprocess] [--group g] [--limit n] [--dry-run] [--url u] [--force] [--json]");139 await registerAllConnectors();140 const dryRun = bool(a, "dry-run");141 await ensureClickHouseQuiet();142 if (!dryRun) await syncConnectorsToDb([requireConnector(id)]);143 const urls = list(a, "url");144 const r = await runConnector(id, { task: urls.length ? "crawl" : taskOf(a, "full"), group: str(a, "group"), limit: num(a, "limit"), dryRun, urls: urls.length ? urls : undefined, force: bool(a, "force"), staleOnly: bool(a, "stale"), quarantine: bool(a, "quarantine"), logLevel: bool(a, "verbose") ? "debug" : undefined, sampleEntities: dryRun ? num(a, "samples") ?? 20 : num(a, "samples") ?? 0 });145 printRun(r, bool(a, "json"));146 return r.status === "failed" ? 1 : 0;147}148149async function cmdRunAll(a: Args): Promise<number> {150 await registerAllConnectors();151 const dryRun = bool(a, "dry-run");152 await ensureClickHouseQuiet();153 if (!dryRun) await syncConnectorsToDb();154 const results = await runAll({ dryRun, task: taskOf(a, "full"), limit: num(a, "limit"), force: bool(a, "force"), ids: a.positional.length ? a.positional : undefined });155 for (const r of results) printRun(r, false);156 console.log(`\n${results.length} run(s): ${results.filter((r) => r.status === "ok").length} ok, ${results.filter((r) => r.status === "partial").length} partial, ${results.filter((r) => r.status === "failed").length} failed`);157 return results.some((r) => r.status === "failed") ? 1 : 0;158}159160async function cmdDocs(a: Args): Promise<number> {161 const id = a.positional[0];162 if (!id) throw new Error("usage: dci docs <connector-id> [--due] [--limit n] [--group g]");163 const docs = await dueDocuments(id, { group: str(a, "group"), limit: num(a, "limit") ?? 50, force: !bool(a, "due"), includeQuarantined: bool(a, "all") });164 const st = await connectorDocStats(id);165 console.log(paint("dim", `${st.total} documents · ${st.due} due · ${st.quarantined} quarantined · ${st.errors} with errors · ${st.extracted} extracted`));166 if (bool(a, "json")) { console.log(JSON.stringify(docs, null, 2)); return 0; }167 console.log(table(docs.map((d) => ({ id: d.id, group: parseDiscoveredFrom(d.discoveredFrom).group, type: d.pageType, prio: d.priority, status: d.statusCode ?? "", "next check": d.nextCheck ? ago(d.nextCheck) : "—", checked: ago(d.lastChecked), chg: `${d.changeCount}/${d.fetchCount} ${d.changeFrequencyScore.toFixed(2)}`, ent: d.extractCount, err: d.quarantined ? "QUARANTINED" : d.errorCount ? `${d.errorCount} ${d.error?.slice(0, 30) ?? ""}` : "", url: d.url.slice(0, 90) })), ["id", "group", "type", "prio", "status", "next check", "checked", "chg", "ent", "err", "url"]));168 return 0;169}170171async function cmdInspect(a: Args): Promise<number> {172 const key = a.positional[0];173 if (!key) throw new Error("usage: dci inspect <url-or-docId> [--body]");174 const doc = await getDocument(key);175 if (!doc) { console.error(paint("red", `no document for ${key}`)); return 1; }176 const version = await latestVersion(doc.id);177 const prov = await getDb().select().from(provenanceTable).where(eq(provenanceTable.documentId, doc.id)).limit(200);178 const versions = await getDb().execute<{ n: number }>(sql`select count(*)::int as n from document_versions where document_id = ${doc.id}`);179 if (bool(a, "json")) { console.log(JSON.stringify({ document: doc, latestVersion: version, provenance: prov }, null, 2)); return 0; }180 console.log(paint("bold", `document ${doc.id}`));181 for (const [k, v] of Object.entries(doc)) if (v !== null && v !== undefined && !(Array.isArray(v) && !v.length)) console.log(` ${k.padEnd(22)} ${typeof v === "object" ? JSON.stringify(v).slice(0, 300) : String(v).slice(0, 300)}`);182 console.log(`\n${paint("bold", `versions: ${Number(versions[0]?.n ?? 0)}`)}`);183 if (version) {184 console.log(` latest ${version.id} · ${version.fetchedAt} · L${version.fetchLevel} · ${version.statusCode} · ${version.sizeBytes} B · significance ${version.significance}`);185 if (version.diffSummary) { console.log(` diff +${version.diffSummary.addedCount} −${version.diffSummary.removedCount} ratio ${version.diffSummary.ratio}`); for (const l of version.diffSummary.added.slice(0, 8)) console.log(paint("green", ` + ${l.slice(0, 140)}`)); for (const l of version.diffSummary.removed.slice(0, 8)) console.log(paint("red", ` − ${l.slice(0, 140)}`)); }186 for (const c of version.detectedChanges.slice(0, 10)) console.log(` change ${JSON.stringify(c).slice(0, 200)}`);187 }188 console.log(`\n${paint("bold", `provenance rows referencing this document: ${prov.length}`)}`);189 if (prov.length) console.log(table(prov.slice(0, 60).map((p) => ({ entity: `${p.entityType} ${p.entityId}`, field: p.field, value: JSON.stringify(p.value).slice(0, 40), conf: p.confidence, method: p.method ?? "", current: p.isCurrent ? "yes" : "no", observed: ago(p.lastObserved) })), ["entity", "field", "value", "conf", "method", "current", "observed"]));190 if (bool(a, "body") && doc.storageKey) { const raw = await getRaw(doc.storageKey); if (raw) console.log(`\n${paint("bold", `archived body ${doc.storageKey} (${raw.body.length} B)`)}\n${raw.body.toString("utf8").slice(0, 4000)}`); }191 return 0;192}193194async function cmdStats(a: Args): Promise<number> {195 const db = getDb();196 const one = async (q: ReturnType<typeof sql>) => Number((await db.execute<{ n: number }>(q))[0]?.n ?? 0);197 const s = {198 facilities: await one(sql`select count(*)::int as n from facilities where merged_into is null`),199 operators: await one(sql`select count(*)::int as n from operators`),200 campuses: await one(sql`select count(*)::int as n from campuses`),201 cloudRegions: await one(sql`select count(*)::int as n from cloud_regions`),202 ixps: await one(sql`select count(*)::int as n from ixps`),203 projects: await one(sql`select count(*)::int as n from projects`),204 events: await one(sql`select count(*)::int as n from events`),205 events24h: await one(sql`select count(*)::int as n from events where detected_at > now() - interval '24 hours'`),206 newsItems: await one(sql`select count(*)::int as n from news_items`),207 provenance: await one(sql`select count(*)::int as n from provenance`),208 documents: await one(sql`select count(*)::int as n from documents`),209 documentsDue: await one(sql`select count(*)::int as n from documents where quarantined = false and (next_check is null or next_check <= now())`),210 documentsQuarantined: await one(sql`select count(*)::int as n from documents where quarantined`),211 versions: await one(sql`select count(*)::int as n from document_versions`),212 connectors: await one(sql`select count(*)::int as n from connectors`),213 connectorsEnabled: await one(sql`select count(*)::int as n from connectors where enabled and not paused`),214 runs24h: await one(sql`select count(*)::int as n from connector_runs where started_at > now() - interval '24 hours'`),215 credits24h: await one(sql`select coalesce(sum((stats->>'credits')::float),0) as n from connector_runs where started_at > now() - interval '24 hours'`),216 pendingMatches: await one(sql`select count(*)::int as n from entity_matches where status = 'pending'`),217 openAlerts: await one(sql`select count(*)::int as n from system_alerts where resolved_at is null`),218 };219 let queues: Awaited<ReturnType<typeof queueSnapshot>> = [];220 try { queues = await queueSnapshot(); } catch { /* redis down */ }221 if (bool(a, "json")) { console.log(JSON.stringify({ ...s, queues }, null, 2)); return 0; }222 console.log(paint("bold", "entities")); console.log(` facilities ${s.facilities} · operators ${s.operators} · campuses ${s.campuses} · cloud regions ${s.cloudRegions} · IXPs ${s.ixps} · projects ${s.projects} · news ${s.newsItems}`);223 console.log(paint("bold", "history")); console.log(` events ${s.events} (${s.events24h} last 24 h) · provenance ${s.provenance} · versions ${s.versions} · pending matches ${s.pendingMatches} · open alerts ${s.openAlerts}`);224 console.log(paint("bold", "crawl")); console.log(` connectors ${s.connectorsEnabled}/${s.connectors} active · documents ${s.documents} (${s.documentsDue} due, ${s.documentsQuarantined} quarantined) · runs 24 h ${s.runs24h} · credits 24 h ${s.credits24h}`);225 if (queues.length) { console.log(paint("bold", "queues")); for (const q of queues) console.log(` ${q.name.padEnd(12)} waiting ${q.waiting} · prioritized ${q.prioritized} · active ${q.active} · delayed ${q.delayed} · failed ${q.failed} · completed ${q.completed}`); }226 return 0;227}228229async function cmdDoctor(a: Args): Promise<number> {230 const r = await doctor({ dryRun: bool(a, "dry-run") });231 if (bool(a, "json")) { console.log(JSON.stringify(r, null, 2)); return r.ok ? 0 : 1; }232 for (const c of r.checks) console.log(`${c.level === "ok" ? paint("green", " ok ") : c.level === "warn" ? paint("yellow", "warn") : paint("red", "FAIL")} ${c.name.padEnd(20)} ${c.message}${c.ms !== undefined ? paint("dim", ` (${c.ms} ms)`) : ""}`);233 console.log(r.ok ? paint("green", `\nall critical checks passed${r.alertsWritten ? ` · ${r.alertsWritten} alert(s) recorded` : ""}`) : paint("red", `\n${r.checks.filter((c) => c.level === "error").length} critical failure(s)`));234 return r.ok ? 0 : 1;235}236237async function cmdEnqueue(a: Args): Promise<number> {238 const id = a.positional[0];239 if (!id) throw new Error("usage: dci enqueue <connector-id> [--task full|discover|crawl|reprocess] [--group g] [--limit n] [--force]");240 const r = await enqueueRun({ connectorId: id, task: taskOf(a, "full"), group: str(a, "group"), limit: num(a, "limit"), force: bool(a, "force"), requestedBy: "cli" }, { priority: 1 });241 console.log(r.queued ? `queued job ${r.jobId}` : `job ${r.jobId} already pending`);242 return 0;243}244245async function cmdScheduler(a: Args): Promise<number> {246 if (bool(a, "once")) { const r = await schedulerTick({ force: true }); console.log(JSON.stringify(r, null, 2)); return 0; }247 console.log(`scheduler loop every ${getEnv().schedulerIntervalMs} ms (Ctrl-C to stop)`);248 const loop = startSchedulerLoop({ onTick: (r) => console.log(`${new Date().toISOString()} checked ${r.checked} · enqueued ${r.enqueued.length}${r.enqueued.length ? ": " + r.enqueued.join(", ") : ""}${r.lockHeldElsewhere ? " (lock held elsewhere)" : ""}`), onError: (e) => console.error(`tick failed: ${e.message}`) });249 await new Promise<void>((resolve) => { const stop = () => { loop.stop(); resolve(); }; process.once("SIGINT", stop); process.once("SIGTERM", stop); });250 return 0;251}252253async function cmdValidateYaml(a: Args): Promise<number> {254 const file = a.positional[0];255 if (!file || !existsSync(file)) throw new Error("usage: dci validate <path/to/connector.yaml>");256 const cfg = parseConnectorConfig(readFileSync(file, "utf8"), file);257 // schema + cross-field rules passed; now check the parsers / implementation it references against the registry258 await registerAllConnectors();259 const problems = validateAgainstRegistry(cfg);260 if (problems.length) { for (const p of problems) console.error(paint("red", `${cfg.id}: ${p}`)); return 1; }261 console.log(paint("green", `valid: ${cfg.id} (${cfg.kind}/${cfg.mode}) — ${Object.keys(cfg.extractors).length} extractor(s), schedule ${JSON.stringify(cfg.schedule)}`));262 return 0;263}264265function printTrace(t: Awaited<ReturnType<typeof traceDocument>>): void {266 const d = t.document ?? {};267 console.log(paint("bold", `SOURCE DOCUMENT ${String(d.id ?? "?")}`)); console.log(` ${String(d.url ?? "")}\n connector ${String(d.connectorId ?? "")} · pageType ${String(d.pageType ?? "")} · extractor ${String(d.extractorVersion ?? "")} · refs ${JSON.stringify(d.entityRefs ?? [])}`);268 if (t.fetch) console.log(`\n${paint("bold", "RAW FETCH")} ${t.fetch.source} · HTTP ${t.fetch.status} · ${t.fetch.contentType ?? "?"} · ${t.fetch.bytes} B · L${t.fetch.level} ${t.fetch.fetcher}`);269 console.log(`\n${paint("bold", "PARSED TEXT")} title: ${t.text.title ?? "—"} · ${t.text.length} chars · classified ${t.text.classification?.pageType ?? "?"} (${t.text.classification?.rule ?? ""}) · MW in text: ${t.text.classification?.mw.join(", ") || "none"}`);270 if (t.announcement) { const a = t.announcement as Record<string, unknown>; const c = a.classification as { class: string; mayCreateProject: boolean; evidence: Record<string, unknown> }; console.log(`\n${paint("bold", "ANNOUNCEMENT")} class ${paint(c.mayCreateProject ? "green" : "yellow", c.class)} · mayCreateProject ${c.mayCreateProject} · evidence ${JSON.stringify(c.evidence)}\n status ${String(a.status)} · headline MW ${String(a.headlineMw)} (scope ${String(a.capacityScope)}, ${String(a.capacitySemantics)}) · money ${JSON.stringify(a.money)} (${String(a.investmentScope)}/${String(a.investmentSemantics)}) · location ${JSON.stringify(a.location)} · operator ${String(a.operator)} · AI ${String(a.aiEvidence)}${a.hqGuarded ? " · HQ mention stripped" : ""}`); const ev = a.evidence as { plannedMw?: { text: string } | null; investment?: { text: string } | null }; if (ev?.plannedMw) console.log(paint("dim", ` MW evidence: "${ev.plannedMw.text.slice(0, 220)}"`)); if (ev?.investment) console.log(paint("dim", ` $ evidence: "${ev.investment.text.slice(0, 220)}"`)); }271 console.log(`\n${paint("bold", `STRUCTURED DATA ${t.records.length} record(s)`)}`); for (const r of t.records.slice(0, 20)) console.log(` ${r.kind.padEnd(11)} ${r.key} · certainty ${r.certainty ?? "—"} · ${Object.keys(r.data).filter((k) => r.data[k] != null && !k.startsWith("_")).slice(0, 12).join(", ")}`);272 console.log(`\n${paint("bold", `EXTRACTED ENTITIES ${t.entities.length} · valid ${t.validation.valid} · rejected ${t.validation.rejected}`)}`); for (const i of t.validation.issues.slice(0, 20)) console.log(paint(i.level === "error" ? "red" : "yellow", ` ${i.level} ${i.key}${i.field ? "." + i.field : ""}: ${i.message}`));273 console.log(`\n${paint("bold", `EXTRACTED CLAIMS ${t.claims.length}`)}`); for (const c of t.claims) console.log(` ${c.field.padEnd(14)} ${String(c.value).padStart(10)} ${c.unit} · scope ${paint(["building", "facility", "campus"].includes(c.scope) ? "green" : "yellow", c.scope)} (${c.scopeReason}) · ${c.semantics ?? "—"}${c.evidence ? paint("dim", `\n "${c.evidence.text.slice(0, 200)}"`) : paint("red", "\n no supporting sentence")}`);274 console.log(`\n${paint("bold", `MATCH CANDIDATES ${t.matches.length} facility record(s)`)}`); for (const m of t.matches) { console.log(` ${m.entityKey} → ${m.how}${m.matchedId ? ` ${m.matchedId}` : ""}`); for (const c of m.candidates.slice(0, 5)) console.log(paint("dim", ` ${c.score.toFixed(3)} ${c.name} (${c.operatorName ?? "—"}, ${c.city ?? "—"}) ${c.reasons.join(" ")}`)); }275 if (t.reconciliation) { const r = t.reconciliation; console.log(`\n${paint("bold", "RECONCILIATION (dry run)")} created ${r.created} · updated ${r.updated} · unchanged ${r.unchanged} · merged ${r.merged} · pending ${r.pendingMatches} · rejected ${r.rejected} · events ${r.events} · provenance ${r.provenanceRows} · claims ${r.claims} (${r.unscopedClaims} unscoped) · flags ${r.qualityFlags} · projects vetoed ${r.projectsVetoed}`); console.log(`\n${paint("bold", "RESULTING DATABASE CHANGES")}`); for (const c of r.changes.slice(0, 20)) console.log(` ${JSON.stringify(c).slice(0, 220)}`); for (const ref of r.refs.slice(0, 20)) console.log(paint("dim", ` ref ${ref.type} ${ref.id}`)); }276 if (t.error) console.log(paint("red", `\nERROR ${t.error}`));277}278279function help(): number {280 console.log(readFileSync(new URL(import.meta.url), "utf8").split("\n").slice(1, 22).map((l) => l.replace(/^ \*\s?/, "")).join("\n"));281 return 0;282}283284/* ---------- main ---------- */285async function main(): Promise<number> {286 loadEnvFile();287 const a = parseArgs(process.argv.slice(2));288 process.on("SIGINT", () => { requestAbort("SIGINT"); console.error(paint("yellow", "\nabort requested — finishing current documents…")); setTimeout(() => process.exit(130), 60_000).unref(); });289 switch (a.cmd) {290 case "connectors": case "ls": return cmdConnectors(a);291 case "sync": return cmdSync();292 case "run": return cmdRun(a);293 case "run-all": return cmdRunAll(a);294 case "discover": return cmdRun({ ...a, flags: { ...a.flags, task: "discover" } });295 case "reprocess": return cmdRun({ ...a, flags: { ...a.flags, task: "reprocess" } });296 case "docs": return cmdDocs(a);297 case "inspect": return cmdInspect(a);298 case "stats": return cmdStats(a);299 case "doctor": return cmdDoctor(a);300 case "rank": case "rankings": { const r = await computeRankings(); console.log(JSON.stringify(r)); return 0; }301 case "metrics": { const r = await computeDailyMetrics(); console.log(JSON.stringify(r)); return 0; }302 case "refresh-stats": { const r = await refreshStats(); console.log(JSON.stringify(r)); return 0; }303 case "enqueue": return cmdEnqueue(a);304 case "pause": { if (!a.positional[0]) throw new Error("usage: dci pause <id>"); await pauseConnector(a.positional[0], true); console.log(`paused ${a.positional[0]}`); return 0; }305 case "resume": { if (!a.positional[0]) throw new Error("usage: dci resume <id>"); await pauseConnector(a.positional[0], false); console.log(`resumed ${a.positional[0]}`); return 0; }306 case "scheduler": return cmdScheduler(a);307 case "clickhouse-init": await ensureClickHouse(); console.log("clickhouse tables ensured"); return 0;308 case "validate": return cmdValidateYaml(a);309 case "trace": { const key = a.positional[0]; if (!key) throw new Error("usage: dci trace <doc-id-or-url> [--live] [--json]"); await registerAllConnectors(); const t = await traceDocument(key, { live: bool(a, "live") }); if (bool(a, "json")) { console.log(JSON.stringify(t, null, 2)); return t.error ? 1 : 0; } printTrace(t); return t.error ? 1 : 0; }310 case "quarantine": { const id = a.positional[0]; const on = (a.positional[1] ?? "on") !== "off"; if (!id) throw new Error("usage: dci quarantine <connector-id> [on|off]"); await getDb().execute(sql`update connectors set quarantine = ${on}, updated_at = now() where id = ${id}`); console.log(`${id}: quarantine ${on ? "ON (ingest rolled back, preview only)" : "OFF"}`); return 0; }311 case "quality": { const r = await qualitySweep(); console.log(JSON.stringify(r, null, 2)); return 0; }312 case "snapshot": { const r = await snapshotAndCheck(); console.log(JSON.stringify(r, null, 2)); return 0; }313 case "gaps": { const r = await dataGaps(); console.log(table(Object.entries(r).map(([k, v]) => ({ gap: k, count: v })), ["gap", "count"])); return 0; }314 case "worker": { const { startWorker } = await import("./main.js"); const w = await startWorker(); await new Promise<void>((resolve) => { const stop = () => void w.stop().then(resolve); process.once("SIGINT", stop); process.once("SIGTERM", stop); }); return 0; }315 case "help": case "--help": case "-h": return help();316 default: console.error(paint("red", `unknown command "${a.cmd}"`)); help(); return 2;317 }318}319320main()321 .then(async (code) => { await closeScheduler().catch(() => undefined); await closeAll(); await closeDb().catch(() => undefined); process.exit(code); })322 .catch(async (e) => { console.error(paint("red", `error: ${(e as Error).message}`)); if (process.env.DCI_LOG_LEVEL === "debug") console.error((e as Error).stack); await closeScheduler().catch(() => undefined); await closeDb().catch(() => undefined); process.exit(1); });323