SPB Git forge

spb/cancerindex

Public
37commits 1branches 0releases
2.9 MBsize
maindefault branch
10 days agolast push
TypeScript 97.2% SQL 1.5% CSS 0.6% JavaScript 0.5%
18.6 KB · 335 lines typescript
Raw Blame History
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