import { accessSync, constants, existsSync, mkdirSync, readFileSync, readdirSync, statSync, writeFileSync, unlinkSync } from 'node:fs'; import { statfs } from 'node:fs/promises'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; import { sql } from 'drizzle-orm'; import { dataDir } from '@cancerindex/shared'; import { listAlerts, type Database, type SystemAlert } from '@cancerindex/database'; import { CONNECTORS } from '../registry.js'; import { humanDuration, isStale } from './schedule.js'; export type CheckLevel = 'ok' | 'info' | 'warn' | 'fail'; export interface DoctorCheck { section: string; name: string; level: CheckLevel; detail: string; } export interface ConnectorReport { id: string; status: string; licenseStatus: string; health: string; paused: boolean; lastSuccessAt: Date | null; lastSuccessAge: string; schedule: string | null; stale: boolean; lastRun: { id: string; status: string; startedAt: Date; recordsFetched: number; anomaly: string | null; drift: number; error: string | null } | null; cursorSummary: string; } export interface DoctorReport { generatedAt: Date; checks: DoctorCheck[]; connectors: ConnectorReport[]; tables: Array<{ table: string; rows: number }>; unresolvedTop: Array<{ source: string; entityKind: string; sourceText: string; count: number }>; alerts: SystemAlert[]; disk: { rawDir: string; bytes: number; files: number; freeBytes: number | null } | null; rankings: { snapshots: number; current: number; latest: Date | null } | null; hardFailures: number; } export interface DoctorOptions { /** Folder holding drizzle migrations (defaults to packages/database/migrations resolved from this file, then cwd). */ migrationsDir?: string; /** Skip the data/raw walk (large lakes). */ skipDisk?: boolean; now?: Date; } const COUNTED_TABLES = ['sources', 'ingest_runs', 'source_records', 'provenance', 'cancers', 'cancer_aliases', 'cancer_hierarchy', 'cancer_codes', 'genes', 'variants', 'drugs', 'clinical_trials', 'trial_conditions', 'publications', 'literature_counts', 'civic_evidence_items', 'genomic_cohorts', 'cancer_gene_frequencies', 'epidemiology_observations', 'survival_observations', 'knowledge_edges', 'unresolved_labels', 'entity_counters', 'rankings', 'ranking_snapshots', 'system_alerts']; const OPTIONAL_KEYS = ['NCBI_API_KEY', 'SEER_API_KEY', 'ADMIN_TOKEN', 'ANTHROPIC_API_KEY', 'OPENAI_API_KEY']; /** * Readiness report for operators (`pnpm cix doctor`): environment, database + extensions + * migrations, table sizes, per-connector state (health, freshness, last run, anomaly, drift, * cursor), curation backlog, data-lake disk usage, ranking freshness and open alerts. * `hardFailures > 0` ⇒ the CLI exits 1. */ export async function runDoctor(db: Database, opts: DoctorOptions = {}): Promise { const now = opts.now ?? new Date(); const checks: DoctorCheck[] = []; const add = (section: string, name: string, level: CheckLevel, detail: string) => checks.push({ section, name, level, detail }); const report: DoctorReport = { generatedAt: now, checks, connectors: [], tables: [], unresolvedTop: [], alerts: [], disk: null, rankings: null, hardFailures: 0 }; /* ---------------------------------------------------------------- env */ const dbUrl = process.env.DATABASE_URL; if (dbUrl) add('env', 'DATABASE_URL', 'ok', redactUrl(dbUrl)); else add('env', 'DATABASE_URL', 'warn', 'not set — using default postgres://localhost:5432/cancerindex'); const raw = path.join(dataDir(), 'raw'); try { mkdirSync(raw, { recursive: true }); accessSync(raw, constants.W_OK); const probe = path.join(raw, `.doctor-${process.pid}`); writeFileSync(probe, 'ok'); unlinkSync(probe); add('env', 'CI_DATA_DIR', 'ok', `${dataDir()} (raw lake writable)`); } catch (e) { add('env', 'CI_DATA_DIR', 'fail', `${dataDir()} not writable: ${(e as Error).message}`); } if (process.env.NCBI_EMAIL) add('env', 'NCBI_EMAIL', 'ok', process.env.NCBI_EMAIL); else add('env', 'NCBI_EMAIL', 'warn', 'absent — NCBI E-utilities (PubMed) require tool + email'); for (const k of OPTIONAL_KEYS) { const v = process.env[k]; if (!v) add('env', k, 'info', 'absent'); else if (k === 'ADMIN_TOKEN' && (v === 'change-me' || v.length < 16)) add('env', k, 'warn', 'placeholder / too short — admin endpoints disabled or weak'); else add('env', k, 'ok', `present (${v.length} chars)`); } /* ---------------------------------------------------------------- database */ let reachable = false; try { const [v] = await db.execute<{ v: string; db: string }>(sql`SELECT version() AS v, current_database() AS db`); reachable = true; add('database', 'reachable', 'ok', `${v?.db} — ${(v?.v ?? '').split(' on ')[0]}`); } catch (e) { add('database', 'reachable', 'fail', (e as Error).message); } if (reachable) { const ext = await db.execute<{ name: string; installed: string | null }>(sql`SELECT name, installed_version AS installed FROM pg_available_extensions WHERE name IN ('vector','pg_trgm','unaccent')`); const byName = new Map(ext.map((r) => [r.name, r.installed])); for (const name of ['pg_trgm', 'unaccent', 'vector']) { const installed = byName.get(name); if (installed) add('database', `extension ${name}`, 'ok', `v${installed}`); else if (byName.has(name)) add('database', `extension ${name}`, name === 'vector' ? 'warn' : 'fail', 'available but not created (pnpm db:migrate creates it)'); else add('database', `extension ${name}`, name === 'vector' ? 'warn' : 'fail', 'not available on this server'); } // Pending migrations: journal entries newer than the last applied drizzle migration. const dir = opts.migrationsDir ?? findMigrationsDir(); const journalPath = dir ? path.join(dir, 'meta', '_journal.json') : null; if (!journalPath || !existsSync(journalPath)) add('database', 'migrations', 'warn', `migrations journal not found (${dir ?? 'no folder'})`); else { const journal = JSON.parse(readFileSync(journalPath, 'utf8')) as { entries: Array<{ tag: string; when: number }> }; const applied = await db .execute<{ n: string; last: string | null }>(sql`SELECT count(*)::text AS n, max(created_at)::text AS last FROM drizzle.__drizzle_migrations`) .catch(() => [] as Array<{ n: string; last: string | null }>); const row = applied[0]; if (!row) add('database', 'migrations', 'fail', 'drizzle.__drizzle_migrations missing — run pnpm db:migrate'); else { const last = row.last ? Number(row.last) : 0; const pending = journal.entries.filter((e) => e.when > last); if (pending.length) add('database', 'migrations', 'fail', `${pending.length} pending: ${pending.map((p) => p.tag).join(', ')} — run pnpm db:migrate`); else add('database', 'migrations', 'ok', `${row.n} applied, ${journal.entries.length} in folder, none pending`); } } const [ops] = await db.execute<{ ok: string | null }>(sql`SELECT to_regclass('public.system_alerts')::text AS ok`); if (ops?.ok) add('database', 'ops schema', 'ok', 'system_alerts present'); else add('database', 'ops schema', 'fail', 'system_alerts table missing — apply the ext-ops migration (docs/schema-changes-ops.md)'); /* ------------------------------------------------------------ tables */ const counts = await db.execute<{ relname: string; n: string }>(sql`SELECT relname, n_live_tup::text AS n FROM pg_stat_user_tables WHERE schemaname = 'public'`); const byTable = new Map(counts.map((r) => [r.relname, Number(r.n)])); report.tables = COUNTED_TABLES.map((t) => ({ table: t, rows: byTable.get(t) ?? -1 })); for (const t of ['cancers', 'sources']) if ((byTable.get(t) ?? 0) === 0) add('data', t, 'warn', 'empty — run db:seed / cix sources:sync / ncit-evs'); /* ------------------------------------------------------------ connectors */ const cursors = await db.execute<{ connector_id: string; health: string; paused: boolean; last_success_at: Date | null; cursor: Record }>(sql`SELECT connector_id, health, paused, last_success_at, cursor FROM connector_cursors`); const curById = new Map(cursors.map((c) => [c.connector_id, c])); const lastRuns = await db.execute<{ id: string; connector_id: string; status: string; started_at: Date; records_fetched: number; anomaly: string | null; drift: number; error: string | null }>(sql` SELECT DISTINCT ON (connector_id) id, connector_id, status, started_at, records_fetched, anomaly, jsonb_array_length(schema_drift) AS drift, error FROM ingest_runs WHERE mode NOT IN ('probe') ORDER BY connector_id, started_at DESC`); const runById = new Map(lastRuns.map((r) => [r.connector_id, r])); const srcRows = await db.execute<{ slug: string; license_status: string; status: string }>(sql`SELECT slug, license_status, status FROM sources`); const srcById = new Map(srcRows.map((s) => [s.slug, s])); for (const c of CONNECTORS) { const m = c.manifest; const cur = curById.get(m.id); const run = runById.get(m.id); const src = srcById.get(m.id); const lastSuccessAt = cur?.last_success_at ? new Date(cur.last_success_at) : null; const st = isStale(lastSuccessAt, m.schedule ?? null, now.getTime()); const active = m.status === 'active' && !cur?.paused; const stale = active && !!m.schedule && st.stale; const cursorSummary = summarizeCursor(cur?.cursor ?? {}); report.connectors.push({ id: m.id, status: m.status, licenseStatus: m.licenseStatus, health: cur?.health ?? 'never-run', paused: !!cur?.paused, lastSuccessAt, lastSuccessAge: lastSuccessAt ? humanDuration(now.getTime() - lastSuccessAt.getTime()) : 'never', schedule: m.schedule ?? null, stale, lastRun: run ? { id: run.id, status: run.status, startedAt: new Date(run.started_at), recordsFetched: Number(run.records_fetched), anomaly: run.anomaly, drift: Number(run.drift), error: run.error } : null, cursorSummary, }); if (!src) add('connectors', m.id, 'warn', 'not in sources table — run cix sources:sync'); else if (src.license_status !== m.licenseStatus) add('connectors', m.id, 'warn', `sources.license_status=${src.license_status} ≠ manifest ${m.licenseStatus} — run cix sources:sync`); if (stale) add('connectors', m.id, 'warn', lastSuccessAt ? `stale: last success ${humanDuration(st.ageMs)} ago > 2× schedule (${humanDuration(st.limitMs)})` : `never succeeded (schedule ${m.schedule})`); if (run?.anomaly) add('connectors', m.id, 'warn', `last run ${run.id} flagged an anomaly: ${run.anomaly.slice(0, 160)}`); else if (run?.status === 'failed' && active) add('connectors', m.id, 'warn', `last run ${run.id} failed: ${(run.error ?? '').split('\n')[0]?.slice(0, 160)}`); if (cur?.health === 'failing' && active) add('connectors', m.id, 'warn', `health failing`); } /* ------------------------------------------------------------ curation backlog */ const unresolved = await db.execute<{ slug: string; entity_kind: string; source_text: string; count: number }>(sql` SELECT s.slug, u.entity_kind, u.source_text, u.count FROM unresolved_labels u JOIN sources s ON s.id = u.source_id WHERE u.status = 'open' ORDER BY u.count DESC, u.id LIMIT 10`); report.unresolvedTop = unresolved.map((u) => ({ source: u.slug, entityKind: u.entity_kind, sourceText: u.source_text, count: Number(u.count) })); const [openTotal] = await db.execute<{ n: string }>(sql`SELECT count(*)::text AS n FROM unresolved_labels WHERE status = 'open'`); add('curation', 'unresolved labels (open)', Number(openTotal?.n ?? 0) > 5000 ? 'warn' : 'info', `${openTotal?.n ?? 0}`); /* ------------------------------------------------------------ rankings */ const [rk] = await db.execute<{ snapshots: string; current: string; latest: Date | null }>(sql`SELECT count(*)::text AS snapshots, count(*) FILTER (WHERE is_current)::text AS current, max(generated_at) AS latest FROM ranking_snapshots`); const latest = rk?.latest ? new Date(rk.latest) : null; report.rankings = { snapshots: Number(rk?.snapshots ?? 0), current: Number(rk?.current ?? 0), latest }; if (!latest) add('rankings', 'snapshots', 'warn', 'none — run cix counters && cix rank'); else if (now.getTime() - latest.getTime() > 3 * 24 * 3600_000) add('rankings', 'snapshots', 'warn', `latest ${humanDuration(now.getTime() - latest.getTime())} ago (${report.rankings.current} current)`); else add('rankings', 'snapshots', 'ok', `latest ${humanDuration(now.getTime() - latest.getTime())} ago, ${report.rankings.current} current / ${report.rankings.snapshots} total`); /* ------------------------------------------------------------ alerts */ if (ops?.ok) { report.alerts = await listAlerts(db, { status: 'active', limit: 100 }); const critical = report.alerts.filter((a) => a.severity === 'critical').length; add('alerts', 'open', report.alerts.length ? (critical ? 'fail' : 'warn') : 'ok', report.alerts.length ? `${report.alerts.length} open (${critical} critical)` : 'none'); } } /* ---------------------------------------------------------------- disk */ if (!opts.skipDisk) { try { const usage = walkSize(raw); let freeBytes: number | null = null; try { const fsStat = await statfs(raw); freeBytes = Number(fsStat.bavail) * Number(fsStat.bsize); } catch { /* statfs unsupported */ } report.disk = { rawDir: raw, bytes: usage.bytes, files: usage.files, freeBytes }; const low = freeBytes !== null && freeBytes < 20 * 1024 ** 3; add('disk', 'data/raw', low ? 'warn' : 'ok', `${humanBytes(usage.bytes)} in ${usage.files} files${freeBytes !== null ? `, ${humanBytes(freeBytes)} free${low ? ' (< 20 GB)' : ''}` : ''}`); } catch (e) { add('disk', 'data/raw', 'warn', (e as Error).message); } } report.hardFailures = checks.filter((c) => c.level === 'fail').length; return report; } /** Plain-text rendering for the terminal. */ export function formatDoctorReport(r: DoctorReport): string { const out: string[] = []; const tag = (l: CheckLevel) => ({ ok: '[ OK ]', info: '[INFO]', warn: '[WARN]', fail: '[FAIL]' })[l]; out.push(`CancerIndex doctor — ${r.generatedAt.toISOString()}`); for (const section of ['env', 'database', 'data', 'disk', 'rankings', 'curation', 'alerts', 'connectors']) { const rows = r.checks.filter((c) => c.section === section); if (!rows.length) continue; out.push('', `## ${section}`); for (const c of rows) out.push(`${tag(c.level)} ${c.name.padEnd(28)} ${c.detail}`); } if (r.tables.length) { out.push('', '## tables (≈ live rows)'); const line: string[] = []; for (const t of r.tables) line.push(`${t.table}=${t.rows < 0 ? '?' : t.rows}`); for (let i = 0; i < line.length; i += 4) out.push(' ' + line.slice(i, i + 4).map((s) => s.padEnd(34)).join('').trimEnd()); } if (r.connectors.length) { out.push('', '## connectors'); out.push(` ${'id'.padEnd(16)} ${'status'.padEnd(20)} ${'license'.padEnd(10)} ${'health'.padEnd(20)} ${'last ok'.padEnd(12)} ${'last run'.padEnd(34)} cursor`); for (const c of r.connectors) { const lr = c.lastRun ? `${c.lastRun.status}${c.lastRun.anomaly ? '!anomaly' : ''}${c.lastRun.drift ? ` drift=${c.lastRun.drift}` : ''} n=${c.lastRun.recordsFetched}` : '-'; out.push(` ${c.id.padEnd(16)} ${(c.status + (c.paused ? ' (paused)' : '')).padEnd(20)} ${c.licenseStatus.padEnd(10)} ${(c.health + (c.stale ? ' STALE' : '')).padEnd(20)} ${c.lastSuccessAge.padEnd(12)} ${lr.padEnd(34)} ${c.cursorSummary}`); } } if (r.unresolvedTop.length) { out.push('', '## unresolved labels — top 10 (open)'); for (const u of r.unresolvedTop) out.push(` ${String(u.count).padStart(6)} ${u.source.padEnd(16)} ${u.entityKind.padEnd(8)} ${u.sourceText.slice(0, 80)}`); } if (r.alerts.length) { out.push('', '## open alerts'); for (const a of r.alerts) out.push(` ${a.severity.padEnd(8)} ${a.kind.padEnd(18)} ${(a.connectorId ?? '-').padEnd(16)} ×${String(a.count).padEnd(4)} ${a.lastSeenAt.toISOString().slice(0, 16)} ${a.message.slice(0, 100)}`); } out.push('', r.hardFailures ? `RESULT: ${r.hardFailures} hard failure(s)` : 'RESULT: ready'); return out.join('\n'); } export function formatAlerts(alerts: SystemAlert[]): string { if (!alerts.length) return 'no open alerts'; const out = [`${'id'.padStart(5)} ${'severity'.padEnd(8)} ${'status'.padEnd(12)} ${'kind'.padEnd(18)} ${'connector'.padEnd(16)} ${'count'.padStart(5)} ${'first seen'.padEnd(16)} ${'last seen'.padEnd(16)} message`]; for (const a of alerts) out.push(`${String(a.id).padStart(5)} ${a.severity.padEnd(8)} ${a.status.padEnd(12)} ${a.kind.padEnd(18)} ${(a.connectorId ?? '-').padEnd(16)} ${String(a.count).padStart(5)} ${a.firstSeenAt.toISOString().slice(0, 16)} ${a.lastSeenAt.toISOString().slice(0, 16)} ${a.message}`); return out.join('\n'); } /* -------------------------------------------------------------------------------------------- */ function findMigrationsDir(): string | null { const here = path.dirname(fileURLToPath(import.meta.url)); for (const candidate of [path.resolve(here, '../../../database/migrations'), path.resolve(process.cwd(), 'packages/database/migrations')]) if (existsSync(candidate)) return candidate; return null; } function summarizeCursor(cursor: Record): string { const keys = Object.keys(cursor); if (!keys.length) return '{}'; const parts = keys.slice(0, 6).map((k) => { const v = cursor[k]; const s = typeof v === 'string' ? (v.length > 24 ? v.slice(0, 21) + '…' : v) : typeof v === 'object' && v !== null ? '{…}' : String(v); return `${k}=${s}`; }); return parts.join(' ') + (keys.length > 6 ? ` +${keys.length - 6}` : ''); } function walkSize(dir: string): { bytes: number; files: number } { let bytes = 0; let files = 0; const stack = [dir]; while (stack.length) { const d = stack.pop()!; let entries: string[] = []; try { entries = readdirSync(d); } catch { continue; } for (const name of entries) { const p = path.join(d, name); let st; try { st = statSync(p); } catch { continue; } if (st.isDirectory()) stack.push(p); else { bytes += st.size; files++; } } } return { bytes, files }; } export function humanBytes(n: number): string { if (n < 1024) return `${n} B`; const units = ['KB', 'MB', 'GB', 'TB']; let v = n / 1024; let i = 0; while (v >= 1024 && i < units.length - 1) { v /= 1024; i++; } return `${v.toFixed(v >= 100 ? 0 : 1)} ${units[i]}`; } function redactUrl(u: string): string { return u.replace(/\/\/([^:@/]+):([^@/]+)@/, '//$1:***@'); }