/** * `pnpm dci ` — operator CLI for the crawl runtime. * * connectors table of connectors (id, kind, mode, enabled, health, last/next run, docs) * sync YAML configs → connectors + sources tables * run [--task full|discover|crawl|reprocess] [--group g] [--limit n] [--dry-run] [--url u]… [--force] [--json] * run-all [--dry-run] [--task t] [--limit n] * discover [--dry-run] * docs [--due] [--limit n] [--group g] * inspect document row + latest version + provenance rows referencing it * reprocess [--limit n] [--group g] [--stale] re-extract from archived bodies (no network); --stale = older parser version only * trace [--live] [--json] extraction debugger: every pipeline stage for one document, nothing persisted * quarantine [on|off] preview-only mode for a connector (ingest rolled back) · quality · snapshot · gaps * stats global counts (entities, documents, runs, events, queues) * doctor Postgres / Redis / ClickHouse / MinIO, budgets, stale connectors, error rates → system_alerts * rank | metrics | refresh-stats maintenance jobs, run inline * enqueue [--task t] [--group g] push a crawl job to the queue (worker must be running) * pause | resume * scheduler run only the scheduler loop (foreground) * worker run the full worker (= main.ts) */ import { existsSync, readFileSync } from "node:fs"; import { getDb, closeDb, sql, provenance as provenanceTable, connectors as connectorsTable, eq } from "@dci/db"; import { parseConnectorConfig, validateAgainstRegistry } from "@dci/connectors"; import { registerAllConnectors } from "./connectors/index.js"; import { loadAllConnectors, requireConnector, syncConnectorsToDb } from "./configs.js"; import { connectorDocStats, dueDocuments, getDocument, latestVersion } from "./documents.js"; import { getEnv, loadEnvFile } from "./env.js"; import { closeAll, doctor } from "./maintenance.js"; import { computeDailyMetrics, refreshStats } from "./metrics.js"; import { requestAbort, runAll, runConnector, type RunResult, type RunTask } from "./pipeline.js"; import { computeRankings } from "./rankings.js"; import { closeScheduler, enqueueRun, pauseConnector, queueSnapshot, schedulerTick, startSchedulerLoop } from "./scheduler.js"; import { parseDiscoveredFrom } from "./scheduling.js"; import { getRaw } from "./storage.js"; import { traceDocument } from "./trace.js"; import { dataGaps, qualitySweep, snapshotAndCheck } from "./quality.js"; import { ensureClickHouse } from "@dci/db/clickhouse"; /* ---------- arg parsing ---------- */ interface Args { cmd: string; positional: string[]; flags: Record } function parseArgs(argv: string[]): Args { const [cmd = "help", ...rest] = argv; const positional: string[] = []; const flags: Args["flags"] = {}; for (let i = 0; i < rest.length; i++) { const a = rest[i]!; if (a.startsWith("--")) { const [k, inline] = a.slice(2).split("=", 2); const key = k!; let v: string | boolean = true; if (inline !== undefined) v = inline; else if (rest[i + 1] !== undefined && !rest[i + 1]!.startsWith("--")) v = rest[++i]!; const prev = flags[key]; if (prev === undefined) flags[key] = v; else if (Array.isArray(prev)) prev.push(String(v)); else flags[key] = [String(prev), String(v)]; } else positional.push(a); } return { cmd, positional, flags }; } const 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; }; const num = (a: Args, k: string): number | undefined => { const v = str(a, k); return v === undefined ? undefined : Number(v); }; const bool = (a: Args, k: string): boolean => a.flags[k] !== undefined && a.flags[k] !== "false"; const list = (a: Args, k: string): string[] => { const v = a.flags[k]; return v === undefined || typeof v === "boolean" ? [] : Array.isArray(v) ? v : [v]; }; /* ---------- output ---------- */ const C = { reset: "\x1b[0m", dim: "\x1b[2m", bold: "\x1b[1m", red: "\x1b[31m", green: "\x1b[32m", yellow: "\x1b[33m", cyan: "\x1b[36m" }; const color = process.stdout.isTTY && !process.env.NO_COLOR; const paint = (c: keyof typeof C, s: string) => (color ? `${C[c]}${s}${C.reset}` : s); function table(rows: Array>, cols: string[]): string { const w = cols.map((c) => Math.max(c.length, ...rows.map((r) => String(r[c] ?? "").length))); const line = (vals: string[]) => vals.map((v, i) => v.padEnd(w[i]!)).join(" "); 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"); } const 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}`; }; const healthColor = (h: string) => (h === "ok" ? paint("green", h) : h === "degraded" ? paint("yellow", h) : h === "failing" ? paint("red", h) : paint("dim", h)); function printRun(r: RunResult, json: boolean): void { if (json) { console.log(JSON.stringify(r, null, 2)); return; } const s = r.stats; const st = r.status === "ok" ? paint("green", r.status) : r.status === "partial" ? paint("yellow", r.status) : paint("red", r.status); console.log(`\n${paint("bold", r.connectorId)} ${r.task} → ${st} in ${r.durationMs} ms (run ${r.runId})`); 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}` : ""}`); 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`); if (r.error) console.log(` ${paint("red", "error")}: ${r.error.split("\n")[0]}`); const errs = r.issues.filter((i) => i.level === "error"); 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}`); } if (r.samples.length) { console.log(` ${paint("cyan", `${r.samples.length} sample entit${r.samples.length > 1 ? "ies" : "y"}`)}:`); for (const e of r.samples.slice(0, 10)) { const o = e as unknown as Record; 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)); console.log(` ${JSON.stringify(brief).slice(0, 600)}`); } } } /* ---------- commands ---------- */ async function cmdConnectors(a: Args): Promise { await registerAllConnectors(); const loaded = loadAllConnectors(); 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` select c.id, c.kind, c.mode, c.enabled, c.paused, c.health, c.last_run_at, c.last_status, c.next_run_at, (select count(*)::int from documents d where d.connector_id = c.id) as docs, (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 due from connectors c order by c.id`); const inDb = new Set(rows.map((r) => r.id)); 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") })); 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" }); 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`)); return 0; } /** ClickHouse is optional analytics: make sure its tables exist, never fail the CLI on it. */ async function ensureClickHouseQuiet(): Promise { try { await ensureClickHouse(); } catch (e) { console.error(paint("dim", `clickhouse unavailable (analytics disabled): ${(e as Error).message.split("\n")[0]}`)); } } async function cmdSync(): Promise { await registerAllConnectors(); await ensureClickHouseQuiet(); const loaded = loadAllConnectors({ reload: true }); const r = await syncConnectorsToDb(loaded); console.log(`synced ${r.connectors} connector(s), ${r.sources} source(s) from ${getEnv().configDir}${r.disabledInDb.length ? `; disabled (no YAML): ${r.disabledInDb.join(", ")}` : ""}`); 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(" ")}`); return 0; } function taskOf(a: Args, def: RunTask): RunTask { const t = str(a, "task") ?? def; if (!["full", "discover", "crawl", "reprocess"].includes(t)) throw new Error(`bad --task ${t}`); return t as RunTask; } async function cmdRun(a: Args): Promise { const id = a.positional[0]; if (!id) throw new Error("usage: dci run [--task full|discover|crawl|reprocess] [--group g] [--limit n] [--dry-run] [--url u] [--force] [--json]"); await registerAllConnectors(); const dryRun = bool(a, "dry-run"); await ensureClickHouseQuiet(); if (!dryRun) await syncConnectorsToDb([requireConnector(id)]); const urls = list(a, "url"); 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 }); printRun(r, bool(a, "json")); return r.status === "failed" ? 1 : 0; } async function cmdRunAll(a: Args): Promise { await registerAllConnectors(); const dryRun = bool(a, "dry-run"); await ensureClickHouseQuiet(); if (!dryRun) await syncConnectorsToDb(); const results = await runAll({ dryRun, task: taskOf(a, "full"), limit: num(a, "limit"), force: bool(a, "force"), ids: a.positional.length ? a.positional : undefined }); for (const r of results) printRun(r, false); 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`); return results.some((r) => r.status === "failed") ? 1 : 0; } async function cmdDocs(a: Args): Promise { const id = a.positional[0]; if (!id) throw new Error("usage: dci docs [--due] [--limit n] [--group g]"); const docs = await dueDocuments(id, { group: str(a, "group"), limit: num(a, "limit") ?? 50, force: !bool(a, "due"), includeQuarantined: bool(a, "all") }); const st = await connectorDocStats(id); console.log(paint("dim", `${st.total} documents · ${st.due} due · ${st.quarantined} quarantined · ${st.errors} with errors · ${st.extracted} extracted`)); if (bool(a, "json")) { console.log(JSON.stringify(docs, null, 2)); return 0; } 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"])); return 0; } async function cmdInspect(a: Args): Promise { const key = a.positional[0]; if (!key) throw new Error("usage: dci inspect [--body]"); const doc = await getDocument(key); if (!doc) { console.error(paint("red", `no document for ${key}`)); return 1; } const version = await latestVersion(doc.id); const prov = await getDb().select().from(provenanceTable).where(eq(provenanceTable.documentId, doc.id)).limit(200); const versions = await getDb().execute<{ n: number }>(sql`select count(*)::int as n from document_versions where document_id = ${doc.id}`); if (bool(a, "json")) { console.log(JSON.stringify({ document: doc, latestVersion: version, provenance: prov }, null, 2)); return 0; } console.log(paint("bold", `document ${doc.id}`)); 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)}`); console.log(`\n${paint("bold", `versions: ${Number(versions[0]?.n ?? 0)}`)}`); if (version) { console.log(` latest ${version.id} · ${version.fetchedAt} · L${version.fetchLevel} · ${version.statusCode} · ${version.sizeBytes} B · significance ${version.significance}`); 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)}`)); } for (const c of version.detectedChanges.slice(0, 10)) console.log(` change ${JSON.stringify(c).slice(0, 200)}`); } console.log(`\n${paint("bold", `provenance rows referencing this document: ${prov.length}`)}`); 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"])); 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)}`); } return 0; } async function cmdStats(a: Args): Promise { const db = getDb(); const one = async (q: ReturnType) => Number((await db.execute<{ n: number }>(q))[0]?.n ?? 0); const s = { facilities: await one(sql`select count(*)::int as n from facilities where merged_into is null`), operators: await one(sql`select count(*)::int as n from operators`), campuses: await one(sql`select count(*)::int as n from campuses`), cloudRegions: await one(sql`select count(*)::int as n from cloud_regions`), ixps: await one(sql`select count(*)::int as n from ixps`), projects: await one(sql`select count(*)::int as n from projects`), events: await one(sql`select count(*)::int as n from events`), events24h: await one(sql`select count(*)::int as n from events where detected_at > now() - interval '24 hours'`), newsItems: await one(sql`select count(*)::int as n from news_items`), provenance: await one(sql`select count(*)::int as n from provenance`), documents: await one(sql`select count(*)::int as n from documents`), documentsDue: await one(sql`select count(*)::int as n from documents where quarantined = false and (next_check is null or next_check <= now())`), documentsQuarantined: await one(sql`select count(*)::int as n from documents where quarantined`), versions: await one(sql`select count(*)::int as n from document_versions`), connectors: await one(sql`select count(*)::int as n from connectors`), connectorsEnabled: await one(sql`select count(*)::int as n from connectors where enabled and not paused`), runs24h: await one(sql`select count(*)::int as n from connector_runs where started_at > now() - interval '24 hours'`), credits24h: await one(sql`select coalesce(sum((stats->>'credits')::float),0) as n from connector_runs where started_at > now() - interval '24 hours'`), pendingMatches: await one(sql`select count(*)::int as n from entity_matches where status = 'pending'`), openAlerts: await one(sql`select count(*)::int as n from system_alerts where resolved_at is null`), }; let queues: Awaited> = []; try { queues = await queueSnapshot(); } catch { /* redis down */ } if (bool(a, "json")) { console.log(JSON.stringify({ ...s, queues }, null, 2)); return 0; } 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}`); 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}`); 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}`); 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}`); } return 0; } async function cmdDoctor(a: Args): Promise { const r = await doctor({ dryRun: bool(a, "dry-run") }); if (bool(a, "json")) { console.log(JSON.stringify(r, null, 2)); return r.ok ? 0 : 1; } 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)`) : ""}`); 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)`)); return r.ok ? 0 : 1; } async function cmdEnqueue(a: Args): Promise { const id = a.positional[0]; if (!id) throw new Error("usage: dci enqueue [--task full|discover|crawl|reprocess] [--group g] [--limit n] [--force]"); 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 }); console.log(r.queued ? `queued job ${r.jobId}` : `job ${r.jobId} already pending`); return 0; } async function cmdScheduler(a: Args): Promise { if (bool(a, "once")) { const r = await schedulerTick({ force: true }); console.log(JSON.stringify(r, null, 2)); return 0; } console.log(`scheduler loop every ${getEnv().schedulerIntervalMs} ms (Ctrl-C to stop)`); 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}`) }); await new Promise((resolve) => { const stop = () => { loop.stop(); resolve(); }; process.once("SIGINT", stop); process.once("SIGTERM", stop); }); return 0; } async function cmdValidateYaml(a: Args): Promise { const file = a.positional[0]; if (!file || !existsSync(file)) throw new Error("usage: dci validate "); const cfg = parseConnectorConfig(readFileSync(file, "utf8"), file); // schema + cross-field rules passed; now check the parsers / implementation it references against the registry await registerAllConnectors(); const problems = validateAgainstRegistry(cfg); if (problems.length) { for (const p of problems) console.error(paint("red", `${cfg.id}: ${p}`)); return 1; } console.log(paint("green", `valid: ${cfg.id} (${cfg.kind}/${cfg.mode}) — ${Object.keys(cfg.extractors).length} extractor(s), schedule ${JSON.stringify(cfg.schedule)}`)); return 0; } function printTrace(t: Awaited>): void { const d = t.document ?? {}; 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 ?? [])}`); 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}`); 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"}`); if (t.announcement) { const a = t.announcement as Record; const c = a.classification as { class: string; mayCreateProject: boolean; evidence: Record }; 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)}"`)); } 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(", ")}`); 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}`)); 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")}`); 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(" ")}`)); } 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}`)); } if (t.error) console.log(paint("red", `\nERROR ${t.error}`)); } function help(): number { console.log(readFileSync(new URL(import.meta.url), "utf8").split("\n").slice(1, 22).map((l) => l.replace(/^ \*\s?/, "")).join("\n")); return 0; } /* ---------- main ---------- */ async function main(): Promise { loadEnvFile(); const a = parseArgs(process.argv.slice(2)); process.on("SIGINT", () => { requestAbort("SIGINT"); console.error(paint("yellow", "\nabort requested — finishing current documents…")); setTimeout(() => process.exit(130), 60_000).unref(); }); switch (a.cmd) { case "connectors": case "ls": return cmdConnectors(a); case "sync": return cmdSync(); case "run": return cmdRun(a); case "run-all": return cmdRunAll(a); case "discover": return cmdRun({ ...a, flags: { ...a.flags, task: "discover" } }); case "reprocess": return cmdRun({ ...a, flags: { ...a.flags, task: "reprocess" } }); case "docs": return cmdDocs(a); case "inspect": return cmdInspect(a); case "stats": return cmdStats(a); case "doctor": return cmdDoctor(a); case "rank": case "rankings": { const r = await computeRankings(); console.log(JSON.stringify(r)); return 0; } case "metrics": { const r = await computeDailyMetrics(); console.log(JSON.stringify(r)); return 0; } case "refresh-stats": { const r = await refreshStats(); console.log(JSON.stringify(r)); return 0; } case "enqueue": return cmdEnqueue(a); case "pause": { if (!a.positional[0]) throw new Error("usage: dci pause "); await pauseConnector(a.positional[0], true); console.log(`paused ${a.positional[0]}`); return 0; } case "resume": { if (!a.positional[0]) throw new Error("usage: dci resume "); await pauseConnector(a.positional[0], false); console.log(`resumed ${a.positional[0]}`); return 0; } case "scheduler": return cmdScheduler(a); case "clickhouse-init": await ensureClickHouse(); console.log("clickhouse tables ensured"); return 0; case "validate": return cmdValidateYaml(a); case "trace": { const key = a.positional[0]; if (!key) throw new Error("usage: dci trace [--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; } case "quarantine": { const id = a.positional[0]; const on = (a.positional[1] ?? "on") !== "off"; if (!id) throw new Error("usage: dci quarantine [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; } case "quality": { const r = await qualitySweep(); console.log(JSON.stringify(r, null, 2)); return 0; } case "snapshot": { const r = await snapshotAndCheck(); console.log(JSON.stringify(r, null, 2)); return 0; } case "gaps": { const r = await dataGaps(); console.log(table(Object.entries(r).map(([k, v]) => ({ gap: k, count: v })), ["gap", "count"])); return 0; } case "worker": { const { startWorker } = await import("./main.js"); const w = await startWorker(); await new Promise((resolve) => { const stop = () => void w.stop().then(resolve); process.once("SIGINT", stop); process.once("SIGTERM", stop); }); return 0; } case "help": case "--help": case "-h": return help(); default: console.error(paint("red", `unknown command "${a.cmd}"`)); help(); return 2; } } main() .then(async (code) => { await closeScheduler().catch(() => undefined); await closeAll(); await closeDb().catch(() => undefined); process.exit(code); }) .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); });