spb/cancerindex
Public
TypeScript 97.2%
SQL 1.5%
CSS 0.6%
JavaScript 0.5%
1import { accessSync, constants, existsSync, mkdirSync, readFileSync, readdirSync, statSync, writeFileSync, unlinkSync } from 'node:fs';2import { statfs } from 'node:fs/promises';3import path from 'node:path';4import { fileURLToPath } from 'node:url';5import { sql } from 'drizzle-orm';6import { dataDir } from '@cancerindex/shared';7import { listAlerts, type Database, type SystemAlert } from '@cancerindex/database';8import { CONNECTORS } from '../registry.js';9import { humanDuration, isStale } from './schedule.js';1011export type CheckLevel = 'ok' | 'info' | 'warn' | 'fail';1213export interface DoctorCheck {14 section: string;15 name: string;16 level: CheckLevel;17 detail: string;18}1920export interface ConnectorReport {21 id: string;22 status: string;23 licenseStatus: string;24 health: string;25 paused: boolean;26 lastSuccessAt: Date | null;27 lastSuccessAge: string;28 schedule: string | null;29 stale: boolean;30 lastRun: { id: string; status: string; startedAt: Date; recordsFetched: number; anomaly: string | null; drift: number; error: string | null } | null;31 cursorSummary: string;32}3334export interface DoctorReport {35 generatedAt: Date;36 checks: DoctorCheck[];37 connectors: ConnectorReport[];38 tables: Array<{ table: string; rows: number }>;39 unresolvedTop: Array<{ source: string; entityKind: string; sourceText: string; count: number }>;40 alerts: SystemAlert[];41 disk: { rawDir: string; bytes: number; files: number; freeBytes: number | null } | null;42 rankings: { snapshots: number; current: number; latest: Date | null } | null;43 hardFailures: number;44}4546export interface DoctorOptions {47 /** Folder holding drizzle migrations (defaults to packages/database/migrations resolved from this file, then cwd). */48 migrationsDir?: string;49 /** Skip the data/raw walk (large lakes). */50 skipDisk?: boolean;51 now?: Date;52}5354const 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'];55const OPTIONAL_KEYS = ['NCBI_API_KEY', 'SEER_API_KEY', 'ADMIN_TOKEN', 'ANTHROPIC_API_KEY', 'OPENAI_API_KEY'];5657/**58 * Readiness report for operators (`pnpm cix doctor`): environment, database + extensions +59 * migrations, table sizes, per-connector state (health, freshness, last run, anomaly, drift,60 * cursor), curation backlog, data-lake disk usage, ranking freshness and open alerts.61 * `hardFailures > 0` ⇒ the CLI exits 1.62 */63export async function runDoctor(db: Database, opts: DoctorOptions = {}): Promise<DoctorReport> {64 const now = opts.now ?? new Date();65 const checks: DoctorCheck[] = [];66 const add = (section: string, name: string, level: CheckLevel, detail: string) => checks.push({ section, name, level, detail });67 const report: DoctorReport = { generatedAt: now, checks, connectors: [], tables: [], unresolvedTop: [], alerts: [], disk: null, rankings: null, hardFailures: 0 };6869 /* ---------------------------------------------------------------- env */70 const dbUrl = process.env.DATABASE_URL;71 if (dbUrl) add('env', 'DATABASE_URL', 'ok', redactUrl(dbUrl));72 else add('env', 'DATABASE_URL', 'warn', 'not set — using default postgres://localhost:5432/cancerindex');73 const raw = path.join(dataDir(), 'raw');74 try {75 mkdirSync(raw, { recursive: true });76 accessSync(raw, constants.W_OK);77 const probe = path.join(raw, `.doctor-${process.pid}`);78 writeFileSync(probe, 'ok');79 unlinkSync(probe);80 add('env', 'CI_DATA_DIR', 'ok', `${dataDir()} (raw lake writable)`);81 } catch (e) {82 add('env', 'CI_DATA_DIR', 'fail', `${dataDir()} not writable: ${(e as Error).message}`);83 }84 if (process.env.NCBI_EMAIL) add('env', 'NCBI_EMAIL', 'ok', process.env.NCBI_EMAIL);85 else add('env', 'NCBI_EMAIL', 'warn', 'absent — NCBI E-utilities (PubMed) require tool + email');86 for (const k of OPTIONAL_KEYS) {87 const v = process.env[k];88 if (!v) add('env', k, 'info', 'absent');89 else if (k === 'ADMIN_TOKEN' && (v === 'change-me' || v.length < 16)) add('env', k, 'warn', 'placeholder / too short — admin endpoints disabled or weak');90 else add('env', k, 'ok', `present (${v.length} chars)`);91 }9293 /* ---------------------------------------------------------------- database */94 let reachable = false;95 try {96 const [v] = await db.execute<{ v: string; db: string }>(sql`SELECT version() AS v, current_database() AS db`);97 reachable = true;98 add('database', 'reachable', 'ok', `${v?.db} — ${(v?.v ?? '').split(' on ')[0]}`);99 } catch (e) {100 add('database', 'reachable', 'fail', (e as Error).message);101 }102 if (reachable) {103 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')`);104 const byName = new Map(ext.map((r) => [r.name, r.installed]));105 for (const name of ['pg_trgm', 'unaccent', 'vector']) {106 const installed = byName.get(name);107 if (installed) add('database', `extension ${name}`, 'ok', `v${installed}`);108 else if (byName.has(name)) add('database', `extension ${name}`, name === 'vector' ? 'warn' : 'fail', 'available but not created (pnpm db:migrate creates it)');109 else add('database', `extension ${name}`, name === 'vector' ? 'warn' : 'fail', 'not available on this server');110 }111112 // Pending migrations: journal entries newer than the last applied drizzle migration.113 const dir = opts.migrationsDir ?? findMigrationsDir();114 const journalPath = dir ? path.join(dir, 'meta', '_journal.json') : null;115 if (!journalPath || !existsSync(journalPath)) add('database', 'migrations', 'warn', `migrations journal not found (${dir ?? 'no folder'})`);116 else {117 const journal = JSON.parse(readFileSync(journalPath, 'utf8')) as { entries: Array<{ tag: string; when: number }> };118 const applied = await db119 .execute<{ n: string; last: string | null }>(sql`SELECT count(*)::text AS n, max(created_at)::text AS last FROM drizzle.__drizzle_migrations`)120 .catch(() => [] as Array<{ n: string; last: string | null }>);121 const row = applied[0];122 if (!row) add('database', 'migrations', 'fail', 'drizzle.__drizzle_migrations missing — run pnpm db:migrate');123 else {124 const last = row.last ? Number(row.last) : 0;125 const pending = journal.entries.filter((e) => e.when > last);126 if (pending.length) add('database', 'migrations', 'fail', `${pending.length} pending: ${pending.map((p) => p.tag).join(', ')} — run pnpm db:migrate`);127 else add('database', 'migrations', 'ok', `${row.n} applied, ${journal.entries.length} in folder, none pending`);128 }129 }130 const [ops] = await db.execute<{ ok: string | null }>(sql`SELECT to_regclass('public.system_alerts')::text AS ok`);131 if (ops?.ok) add('database', 'ops schema', 'ok', 'system_alerts present');132 else add('database', 'ops schema', 'fail', 'system_alerts table missing — apply the ext-ops migration (docs/schema-changes-ops.md)');133134 /* ------------------------------------------------------------ tables */135 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'`);136 const byTable = new Map(counts.map((r) => [r.relname, Number(r.n)]));137 report.tables = COUNTED_TABLES.map((t) => ({ table: t, rows: byTable.get(t) ?? -1 }));138 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');139140 /* ------------------------------------------------------------ connectors */141 const cursors = await db.execute<{ connector_id: string; health: string; paused: boolean; last_success_at: Date | null; cursor: Record<string, unknown> }>(sql`SELECT connector_id, health, paused, last_success_at, cursor FROM connector_cursors`);142 const curById = new Map(cursors.map((c) => [c.connector_id, c]));143 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`144 SELECT DISTINCT ON (connector_id) id, connector_id, status, started_at, records_fetched, anomaly, jsonb_array_length(schema_drift) AS drift, error145 FROM ingest_runs WHERE mode NOT IN ('probe') ORDER BY connector_id, started_at DESC`);146 const runById = new Map(lastRuns.map((r) => [r.connector_id, r]));147 const srcRows = await db.execute<{ slug: string; license_status: string; status: string }>(sql`SELECT slug, license_status, status FROM sources`);148 const srcById = new Map(srcRows.map((s) => [s.slug, s]));149 for (const c of CONNECTORS) {150 const m = c.manifest;151 const cur = curById.get(m.id);152 const run = runById.get(m.id);153 const src = srcById.get(m.id);154 const lastSuccessAt = cur?.last_success_at ? new Date(cur.last_success_at) : null;155 const st = isStale(lastSuccessAt, m.schedule ?? null, now.getTime());156 const active = m.status === 'active' && !cur?.paused;157 const stale = active && !!m.schedule && st.stale;158 const cursorSummary = summarizeCursor(cur?.cursor ?? {});159 report.connectors.push({160 id: m.id,161 status: m.status,162 licenseStatus: m.licenseStatus,163 health: cur?.health ?? 'never-run',164 paused: !!cur?.paused,165 lastSuccessAt,166 lastSuccessAge: lastSuccessAt ? humanDuration(now.getTime() - lastSuccessAt.getTime()) : 'never',167 schedule: m.schedule ?? null,168 stale,169 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,170 cursorSummary,171 });172 if (!src) add('connectors', m.id, 'warn', 'not in sources table — run cix sources:sync');173 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`);174 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})`);175 if (run?.anomaly) add('connectors', m.id, 'warn', `last run ${run.id} flagged an anomaly: ${run.anomaly.slice(0, 160)}`);176 else if (run?.status === 'failed' && active) add('connectors', m.id, 'warn', `last run ${run.id} failed: ${(run.error ?? '').split('\n')[0]?.slice(0, 160)}`);177 if (cur?.health === 'failing' && active) add('connectors', m.id, 'warn', `health failing`);178 }179180 /* ------------------------------------------------------------ curation backlog */181 const unresolved = await db.execute<{ slug: string; entity_kind: string; source_text: string; count: number }>(sql`182 SELECT s.slug, u.entity_kind, u.source_text, u.count FROM unresolved_labels u JOIN sources s ON s.id = u.source_id183 WHERE u.status = 'open' ORDER BY u.count DESC, u.id LIMIT 10`);184 report.unresolvedTop = unresolved.map((u) => ({ source: u.slug, entityKind: u.entity_kind, sourceText: u.source_text, count: Number(u.count) }));185 const [openTotal] = await db.execute<{ n: string }>(sql`SELECT count(*)::text AS n FROM unresolved_labels WHERE status = 'open'`);186 add('curation', 'unresolved labels (open)', Number(openTotal?.n ?? 0) > 5000 ? 'warn' : 'info', `${openTotal?.n ?? 0}`);187188 /* ------------------------------------------------------------ rankings */189 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`);190 const latest = rk?.latest ? new Date(rk.latest) : null;191 report.rankings = { snapshots: Number(rk?.snapshots ?? 0), current: Number(rk?.current ?? 0), latest };192 if (!latest) add('rankings', 'snapshots', 'warn', 'none — run cix counters && cix rank');193 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)`);194 else add('rankings', 'snapshots', 'ok', `latest ${humanDuration(now.getTime() - latest.getTime())} ago, ${report.rankings.current} current / ${report.rankings.snapshots} total`);195196 /* ------------------------------------------------------------ alerts */197 if (ops?.ok) {198 report.alerts = await listAlerts(db, { status: 'active', limit: 100 });199 const critical = report.alerts.filter((a) => a.severity === 'critical').length;200 add('alerts', 'open', report.alerts.length ? (critical ? 'fail' : 'warn') : 'ok', report.alerts.length ? `${report.alerts.length} open (${critical} critical)` : 'none');201 }202 }203204 /* ---------------------------------------------------------------- disk */205 if (!opts.skipDisk) {206 try {207 const usage = walkSize(raw);208 let freeBytes: number | null = null;209 try {210 const fsStat = await statfs(raw);211 freeBytes = Number(fsStat.bavail) * Number(fsStat.bsize);212 } catch {213 /* statfs unsupported */214 }215 report.disk = { rawDir: raw, bytes: usage.bytes, files: usage.files, freeBytes };216 const low = freeBytes !== null && freeBytes < 20 * 1024 ** 3;217 add('disk', 'data/raw', low ? 'warn' : 'ok', `${humanBytes(usage.bytes)} in ${usage.files} files${freeBytes !== null ? `, ${humanBytes(freeBytes)} free${low ? ' (< 20 GB)' : ''}` : ''}`);218 } catch (e) {219 add('disk', 'data/raw', 'warn', (e as Error).message);220 }221 }222223 report.hardFailures = checks.filter((c) => c.level === 'fail').length;224 return report;225}226227/** Plain-text rendering for the terminal. */228export function formatDoctorReport(r: DoctorReport): string {229 const out: string[] = [];230 const tag = (l: CheckLevel) => ({ ok: '[ OK ]', info: '[INFO]', warn: '[WARN]', fail: '[FAIL]' })[l];231 out.push(`CancerIndex doctor — ${r.generatedAt.toISOString()}`);232 for (const section of ['env', 'database', 'data', 'disk', 'rankings', 'curation', 'alerts', 'connectors']) {233 const rows = r.checks.filter((c) => c.section === section);234 if (!rows.length) continue;235 out.push('', `## ${section}`);236 for (const c of rows) out.push(`${tag(c.level)} ${c.name.padEnd(28)} ${c.detail}`);237 }238 if (r.tables.length) {239 out.push('', '## tables (≈ live rows)');240 const line: string[] = [];241 for (const t of r.tables) line.push(`${t.table}=${t.rows < 0 ? '?' : t.rows}`);242 for (let i = 0; i < line.length; i += 4) out.push(' ' + line.slice(i, i + 4).map((s) => s.padEnd(34)).join('').trimEnd());243 }244 if (r.connectors.length) {245 out.push('', '## connectors');246 out.push(` ${'id'.padEnd(16)} ${'status'.padEnd(20)} ${'license'.padEnd(10)} ${'health'.padEnd(20)} ${'last ok'.padEnd(12)} ${'last run'.padEnd(34)} cursor`);247 for (const c of r.connectors) {248 const lr = c.lastRun ? `${c.lastRun.status}${c.lastRun.anomaly ? '!anomaly' : ''}${c.lastRun.drift ? ` drift=${c.lastRun.drift}` : ''} n=${c.lastRun.recordsFetched}` : '-';249 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}`);250 }251 }252 if (r.unresolvedTop.length) {253 out.push('', '## unresolved labels — top 10 (open)');254 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)}`);255 }256 if (r.alerts.length) {257 out.push('', '## open alerts');258 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)}`);259 }260 out.push('', r.hardFailures ? `RESULT: ${r.hardFailures} hard failure(s)` : 'RESULT: ready');261 return out.join('\n');262}263264export function formatAlerts(alerts: SystemAlert[]): string {265 if (!alerts.length) return 'no open alerts';266 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`];267 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}`);268 return out.join('\n');269}270271/* -------------------------------------------------------------------------------------------- */272273function findMigrationsDir(): string | null {274 const here = path.dirname(fileURLToPath(import.meta.url));275 for (const candidate of [path.resolve(here, '../../../database/migrations'), path.resolve(process.cwd(), 'packages/database/migrations')]) if (existsSync(candidate)) return candidate;276 return null;277}278279function summarizeCursor(cursor: Record<string, unknown>): string {280 const keys = Object.keys(cursor);281 if (!keys.length) return '{}';282 const parts = keys.slice(0, 6).map((k) => {283 const v = cursor[k];284 const s = typeof v === 'string' ? (v.length > 24 ? v.slice(0, 21) + '…' : v) : typeof v === 'object' && v !== null ? '{…}' : String(v);285 return `${k}=${s}`;286 });287 return parts.join(' ') + (keys.length > 6 ? ` +${keys.length - 6}` : '');288}289290function walkSize(dir: string): { bytes: number; files: number } {291 let bytes = 0;292 let files = 0;293 const stack = [dir];294 while (stack.length) {295 const d = stack.pop()!;296 let entries: string[] = [];297 try {298 entries = readdirSync(d);299 } catch {300 continue;301 }302 for (const name of entries) {303 const p = path.join(d, name);304 let st;305 try {306 st = statSync(p);307 } catch {308 continue;309 }310 if (st.isDirectory()) stack.push(p);311 else {312 bytes += st.size;313 files++;314 }315 }316 }317 return { bytes, files };318}319320export function humanBytes(n: number): string {321 if (n < 1024) return `${n} B`;322 const units = ['KB', 'MB', 'GB', 'TB'];323 let v = n / 1024;324 let i = 0;325 while (v >= 1024 && i < units.length - 1) {326 v /= 1024;327 i++;328 }329 return `${v.toFixed(v >= 100 ? 0 : 1)} ${units[i]}`;330}331332function redactUrl(u: string): string {333 return u.replace(/\/\/([^:@/]+):([^@/]+)@/, '//$1:***@');334}335