spb/cancerindex
Public
TypeScript 97.2%
SQL 1.5%
CSS 0.6%
JavaScript 0.5%
1import { sql } from 'drizzle-orm';2import type { FastifyPluginAsyncZod } from 'fastify-type-provider-zod';3import { z } from 'zod';4import { traceValue } from '@cancerindex/ranking';5import { normalizeLabel } from '@cancerindex/shared';6import { BadRequest, NotFound } from '../lib/errors.js';7import { JOBS, enqueue } from '../lib/queue.js';8import { requireAdmin } from '../plugins/auth.js';9import { AnyList, AnyRecord, camel, camelRows, ok, respond } from '../lib/respond.js';1011const ACTOR = 'admin-api';1213// Optional JSON bodies: Fastify hands the validator `null` when a POST carries no body, hence nullish().14const ReasonBody = z.object({ reason: z.string().max(500).optional() }).nullish();15const RunBody = z.object({ mode: z.enum(['full', 'incremental', 'backfill', 'dry_run']).optional(), maxMinutes: z.number().int().min(1).max(600).optional(), maxRecords: z.number().int().min(1).optional(), resetCursor: z.boolean().optional(), reason: z.string().max(500).optional() }).nullish();1617/**18 * Operator endpoints (CLAUDE.md §348: every mutation is audited). Guarded by `x-admin-token`.19 * Runs are never executed in the API process: they are enqueued to pg-boss and executed by the worker.20 */21export const adminRoutes: FastifyPluginAsyncZod = async (app) => {22 app.addHook('onRequest', requireAdmin);2324 const audit = async (action: string, entityType: string | null, entityId: string | null, before: unknown, after: unknown, reason?: string) => {25 await app.db.execute(sql`INSERT INTO audit_log (actor, action, entity_type, entity_id, before, after, reason) VALUES (${ACTOR}, ${action}, ${entityType}, ${entityId}, ${before === undefined ? null : JSON.stringify(before)}::jsonb, ${after === undefined ? null : JSON.stringify(after)}::jsonb, ${reason ?? null})`);26 };2728 const connectorExists = async (id: string) => {29 const rows = await app.db.execute<{ slug: string; status: string }>(sql`SELECT slug, status FROM sources WHERE slug = ${id}`);30 if (!rows[0]) throw new NotFound('connector', id);31 return rows[0];32 };3334 app.get('/connectors', { schema: { tags: ['admin'], summary: 'Connector health, cursors and the last 5 runs each', security: [{ adminToken: [] }], response: ok(AnyList) } }, async () => {35 const [sources, runs] = await Promise.all([36 app.db.execute<Record<string, unknown>>(sql`37 SELECT s.id AS source_id, s.slug AS connector_id, s.name, s.category, s.tier, s.status, s.license_status, s.manifest->>'schedule' AS schedule, s.manifest->>'documentationVerifiedAt' AS documentation_verified_at,38 cc.health, cc.health_detail, cc.paused, cc.cursor, cc.last_success_at, cc.last_attempt_at,39 (SELECT count(*) FROM source_records sr WHERE sr.source_id = s.id) AS record_count,40 (SELECT count(*) FROM unresolved_labels u WHERE u.source_id = s.id AND u.status = 'open') AS unresolved_open41 FROM sources s LEFT JOIN connector_cursors cc ON cc.connector_id = s.slug ORDER BY s.tier, s.slug`),42 app.db.execute<Record<string, unknown>>(sql`43 SELECT * FROM (44 SELECT id, connector_id, mode, status, started_at, finished_at, duration_ms, records_fetched, records_created, records_updated, records_unchanged, records_rejected, http_requests, http_failures, rate_limit_events, validation_failures,45 jsonb_array_length(schema_drift) AS drift_signals, error, anomaly, dataset_version, row_number() OVER (PARTITION BY connector_id ORDER BY started_at DESC) AS rn46 FROM ingest_runs) x WHERE rn <= 5 ORDER BY connector_id, started_at DESC`),47 ]);48 const byConnector = new Map<string, Array<Record<string, unknown>>>();49 for (const r of runs) {50 const { rn: _rn, ...rest } = r;51 const list = byConnector.get(r.connector_id as string) ?? [];52 list.push(camel(rest));53 byConnector.set(r.connector_id as string, list);54 }55 const data = sources.map((s) => ({ ...camel(s), health: s.health ?? 'never_run', recentRuns: byConnector.get(s.connector_id as string) ?? [] }));56 return respond(app, data, []);57 });5859 app.post('/connectors/:id/run', { schema: { tags: ['admin'], summary: 'Enqueue a connector run (pg-boss connector.run, singleton per connector)', security: [{ adminToken: [] }], params: z.object({ id: z.string() }), body: RunBody, response: ok(AnyRecord) } }, async (req) => {60 const src = await connectorExists(req.params.id);61 const body = req.body ?? {};62 // pg-boss refuses undefined option values — build the payload with defined keys only.63 const payload: Record<string, unknown> = { id: req.params.id, requestedBy: ACTOR };64 if (body.mode) payload.mode = body.mode;65 if (body.maxMinutes) payload.maxMinutes = body.maxMinutes;66 if (body.maxRecords) payload.maxRecords = body.maxRecords;67 if (body.resetCursor) payload.resetCursor = true;68 const jobId = await enqueue(JOBS.connectorRun, payload, req.params.id);69 await audit('connector.run.enqueue', 'connector', req.params.id, null, { jobId, payload, sourceStatus: src.status }, body.reason);70 return respond(app, { connectorId: req.params.id, jobId, queued: jobId !== null, note: jobId === null ? 'a run for this connector is already queued or active (singleton)' : 'queued; the worker will pick it up' }, []);71 });7273 for (const action of ['pause', 'resume'] as const) {74 app.post(`/connectors/:id/${action}`, { schema: { tags: ['admin'], summary: `${action === 'pause' ? 'Pause' : 'Resume'} a connector (worker skips paused connectors)`, security: [{ adminToken: [] }], params: z.object({ id: z.string() }), body: ReasonBody, response: ok(AnyRecord) } }, async (req) => {75 await connectorExists(req.params.id);76 const before = await app.db.execute<Record<string, unknown>>(sql`SELECT paused, health FROM connector_cursors WHERE connector_id = ${req.params.id}`);77 const paused = action === 'pause';78 await app.db.execute(sql`INSERT INTO connector_cursors (connector_id, paused, updated_at) VALUES (${req.params.id}, ${paused}, now()) ON CONFLICT (connector_id) DO UPDATE SET paused = ${paused}, updated_at = now()`);79 await app.db.execute(sql`UPDATE sources SET status = ${paused ? 'paused' : 'active'}, updated_at = now() WHERE slug = ${req.params.id} AND status IN ('active','paused','degraded')`);80 await audit(`connector.${action}`, 'connector', req.params.id, before[0] ?? null, { paused }, req.body?.reason);81 return respond(app, { connectorId: req.params.id, paused }, []);82 });83 }8485 app.get('/runs/:runId', { schema: { tags: ['admin'], summary: 'Full ingest run record: counters, log, schema drift, cursors', security: [{ adminToken: [] }], params: z.object({ runId: z.string() }), response: ok(AnyRecord) } }, async (req) => {86 const rows = await app.db.execute<Record<string, unknown>>(sql`SELECT * FROM ingest_runs WHERE id = ${req.params.runId}`);87 if (!rows[0]) throw new NotFound('run', req.params.runId);88 return respond(app, camel(rows[0]), [rows[0].source_id as string]);89 });9091 app.get('/unresolved', { schema: { tags: ['admin'], summary: 'Curation queue: labels no connector could reconcile (never dropped, §222)', security: [{ adminToken: [] }], querystring: z.object({ entityKind: z.string().optional(), status: z.enum(['open', 'mapped', 'rejected', 'ignored']).default('open'), sourceId: z.string().optional(), limit: z.coerce.number().int().min(1).max(500).default(100), offset: z.coerce.number().int().min(0).default(0) }), response: ok(AnyList, true) } }, async (req) => {92 const q = req.query;93 const conds = [sql`u.status = ${q.status}`];94 if (q.entityKind) conds.push(sql`u.entity_kind = ${q.entityKind}`);95 if (q.sourceId) conds.push(sql`(u.source_id = ${q.sourceId} OR s.slug = ${q.sourceId})`);96 const rows = await app.db.execute<Record<string, unknown> & { total: string }>(sql`97 SELECT u.*, s.slug AS source_slug, c.canonical_name AS suggested_name, count(*) OVER() AS total98 FROM unresolved_labels u JOIN sources s ON s.id = u.source_id LEFT JOIN cancers c ON c.id = u.suggested_id99 WHERE ${sql.join(conds, sql` AND `)} ORDER BY u.count DESC, u.id LIMIT ${q.limit} OFFSET ${q.offset}`);100 const total = rows.length ? Number(rows[0]!.total) : 0;101 const data = rows.map((r) => {102 const { total: _t, ...rest } = r;103 return camel(rest);104 });105 return respond(app, data, rows.map((r) => r.source_id as string), { total, limit: q.limit, offset: q.offset, hasMore: q.offset + q.limit < total });106 });107108 app.post('/unresolved/:id/resolve', { schema: { tags: ['admin'], summary: 'Resolve an unresolved label: map to a cancer (adds a curated alias; connectors pick it up on the next run) or reject', security: [{ adminToken: [] }], params: z.object({ id: z.coerce.number().int() }), body: z.object({ cancerId: z.string().optional(), rejected: z.boolean().optional(), reason: z.string().max(500).optional() }), response: ok(AnyRecord) } }, async (req) => {109 const rows = await app.db.execute<Record<string, unknown>>(sql`SELECT * FROM unresolved_labels WHERE id = ${req.params.id}`);110 const u = rows[0];111 if (!u) throw new NotFound('unresolved label', String(req.params.id));112 const { cancerId, rejected, reason } = req.body;113 if (!cancerId && !rejected) throw new BadRequest('provide cancerId or rejected=true');114 if (cancerId && rejected) throw new BadRequest('cancerId and rejected are mutually exclusive');115 let after: Record<string, unknown>;116 if (rejected) {117 await app.db.execute(sql`UPDATE unresolved_labels SET status = 'rejected', resolved_by = ${ACTOR}, updated_at = now() WHERE id = ${req.params.id}`);118 after = { status: 'rejected' };119 } else {120 const cancer = await app.db.execute<{ id: string; canonical_name: string }>(sql`SELECT id, canonical_name FROM cancers WHERE id = ${cancerId!} AND status = 'active'`);121 if (!cancer[0]) throw new NotFound('cancer', cancerId!);122 if (u.entity_kind !== 'cancer') throw new BadRequest(`only cancer labels can be mapped here (entity_kind=${u.entity_kind as string})`);123 const normalized = (u.normalized as string) || normalizeLabel(u.source_text as string);124 await app.db.transaction(async (tx) => {125 await tx.execute(sql`INSERT INTO cancer_aliases (cancer_id, alias, normalized, alias_type, source_id, source_terminology, language)126 VALUES (${cancer[0]!.id}, ${u.source_text as string}, ${normalized}, 'synonym', ${u.source_id as string}, 'curation', 'en') ON CONFLICT DO NOTHING`);127 await tx.execute(sql`UPDATE unresolved_labels SET status = 'mapped', resolved_id = ${cancer[0]!.id}, resolved_by = ${ACTOR}, updated_at = now() WHERE id = ${req.params.id}`);128 // Immediate effect on trial conditions that carry the same normalized label; other tables are re-reconciled by their connector's next run.129 await tx.execute(sql`UPDATE trial_conditions SET cancer_id = ${cancer[0]!.id}, match_type = 'CURATED_EXACT', confidence = 1 WHERE cancer_id IS NULL AND normalized = ${normalized}`);130 await tx.execute(sql`INSERT INTO change_events (entity_type, entity_id, kind, summary, after) VALUES ('cancer', ${cancer[0]!.id}, 'alias_added', ${`Curated alias "${u.source_text as string}" added from unresolved label #${req.params.id}`}, ${JSON.stringify({ alias: u.source_text, normalized, sourceId: u.source_id })}::jsonb)`);131 });132 after = { status: 'mapped', resolvedId: cancer[0].id, cancerName: cancer[0].canonical_name, aliasAdded: u.source_text, normalized };133 }134 await audit('unresolved.resolve', 'unresolved_label', String(req.params.id), camel(u), after, reason);135 return respond(app, { id: req.params.id, ...after }, [u.source_id as string]);136 });137138 app.get('/trace/:table/:id', { schema: { tags: ['admin'], summary: 'Lineage trace: ranked value → inputs → observation → provenance → raw record (§252-253)', security: [{ adminToken: [] }], params: z.object({ table: z.enum(['rankings', 'epidemiology_observations', 'cancer_gene_frequencies', 'literature_counts', 'survival_observations']), id: z.string() }), response: ok(AnyRecord) } }, async (req) => {139 const trace = await traceValue(app.db, req.params.table, req.params.id);140 return respond(app, trace, []);141 });142143 for (const [path, job, key] of [144 ['/jobs/counters', JOBS.counters, 'counters'],145 ['/jobs/rank', JOBS.rank, 'rank'],146 ] as const) {147 app.post(path, { schema: { tags: ['admin'], summary: `Enqueue ${job}`, security: [{ adminToken: [] }], body: ReasonBody, response: ok(AnyRecord) } }, async (req) => {148 const jobId = await enqueue(job, { requestedBy: ACTOR, thenRank: job === JOBS.counters ? false : undefined } as Record<string, unknown>, key);149 await audit(`${job}.enqueue`, 'job', job, null, { jobId }, req.body?.reason);150 return respond(app, { job, jobId, queued: jobId !== null }, []);151 });152 }153154 app.get('/alerts', { schema: { tags: ['admin'], summary: 'Internal alerts (§170): connector failures, anomalies, schema drift, stale/failing sources', security: [{ adminToken: [] }], querystring: z.object({ status: z.enum(['open', 'acknowledged', 'resolved', 'active']).default('active'), connectorId: z.string().optional(), limit: z.coerce.number().int().min(1).max(500).default(100) }), response: ok(AnyList) } }, async (req) => {155 const q = req.query;156 const conds = [q.status === 'active' ? sql`status IN ('open','acknowledged')` : sql`status = ${q.status}`];157 if (q.connectorId) conds.push(sql`connector_id = ${q.connectorId}`);158 const rows = await app.db.execute<Record<string, unknown>>(sql`SELECT * FROM system_alerts WHERE ${sql.join(conds, sql` AND `)} ORDER BY CASE severity WHEN 'critical' THEN 0 WHEN 'warn' THEN 1 ELSE 2 END, last_seen_at DESC LIMIT ${q.limit}`);159 return respond(app, camelRows(rows), []);160 });161162 app.get('/audit', { schema: { tags: ['admin'], summary: 'Recent audit log entries', security: [{ adminToken: [] }], querystring: z.object({ limit: z.coerce.number().int().min(1).max(500).default(100) }), response: ok(AnyList) } }, async (req) => {163 const rows = await app.db.execute<Record<string, unknown>>(sql`SELECT * FROM audit_log ORDER BY created_at DESC, id DESC LIMIT ${req.query.limit}`);164 return respond(app, camelRows(rows), []);165 });166};167