import { sql } from 'drizzle-orm'; import type { FastifyPluginAsyncZod } from 'fastify-type-provider-zod'; import { z } from 'zod'; import { traceValue } from '@cancerindex/ranking'; import { normalizeLabel } from '@cancerindex/shared'; import { BadRequest, NotFound } from '../lib/errors.js'; import { JOBS, enqueue } from '../lib/queue.js'; import { requireAdmin } from '../plugins/auth.js'; import { AnyList, AnyRecord, camel, camelRows, ok, respond } from '../lib/respond.js'; const ACTOR = 'admin-api'; // Optional JSON bodies: Fastify hands the validator `null` when a POST carries no body, hence nullish(). const ReasonBody = z.object({ reason: z.string().max(500).optional() }).nullish(); const 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(); /** * Operator endpoints (CLAUDE.md §348: every mutation is audited). Guarded by `x-admin-token`. * Runs are never executed in the API process: they are enqueued to pg-boss and executed by the worker. */ export const adminRoutes: FastifyPluginAsyncZod = async (app) => { app.addHook('onRequest', requireAdmin); const audit = async (action: string, entityType: string | null, entityId: string | null, before: unknown, after: unknown, reason?: string) => { 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})`); }; const connectorExists = async (id: string) => { const rows = await app.db.execute<{ slug: string; status: string }>(sql`SELECT slug, status FROM sources WHERE slug = ${id}`); if (!rows[0]) throw new NotFound('connector', id); return rows[0]; }; app.get('/connectors', { schema: { tags: ['admin'], summary: 'Connector health, cursors and the last 5 runs each', security: [{ adminToken: [] }], response: ok(AnyList) } }, async () => { const [sources, runs] = await Promise.all([ app.db.execute>(sql` 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, cc.health, cc.health_detail, cc.paused, cc.cursor, cc.last_success_at, cc.last_attempt_at, (SELECT count(*) FROM source_records sr WHERE sr.source_id = s.id) AS record_count, (SELECT count(*) FROM unresolved_labels u WHERE u.source_id = s.id AND u.status = 'open') AS unresolved_open FROM sources s LEFT JOIN connector_cursors cc ON cc.connector_id = s.slug ORDER BY s.tier, s.slug`), app.db.execute>(sql` SELECT * FROM ( 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, 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 rn FROM ingest_runs) x WHERE rn <= 5 ORDER BY connector_id, started_at DESC`), ]); const byConnector = new Map>>(); for (const r of runs) { const { rn: _rn, ...rest } = r; const list = byConnector.get(r.connector_id as string) ?? []; list.push(camel(rest)); byConnector.set(r.connector_id as string, list); } const data = sources.map((s) => ({ ...camel(s), health: s.health ?? 'never_run', recentRuns: byConnector.get(s.connector_id as string) ?? [] })); return respond(app, data, []); }); 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) => { const src = await connectorExists(req.params.id); const body = req.body ?? {}; // pg-boss refuses undefined option values — build the payload with defined keys only. const payload: Record = { id: req.params.id, requestedBy: ACTOR }; if (body.mode) payload.mode = body.mode; if (body.maxMinutes) payload.maxMinutes = body.maxMinutes; if (body.maxRecords) payload.maxRecords = body.maxRecords; if (body.resetCursor) payload.resetCursor = true; const jobId = await enqueue(JOBS.connectorRun, payload, req.params.id); await audit('connector.run.enqueue', 'connector', req.params.id, null, { jobId, payload, sourceStatus: src.status }, body.reason); 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' }, []); }); for (const action of ['pause', 'resume'] as const) { 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) => { await connectorExists(req.params.id); const before = await app.db.execute>(sql`SELECT paused, health FROM connector_cursors WHERE connector_id = ${req.params.id}`); const paused = action === 'pause'; 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()`); 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')`); await audit(`connector.${action}`, 'connector', req.params.id, before[0] ?? null, { paused }, req.body?.reason); return respond(app, { connectorId: req.params.id, paused }, []); }); } 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) => { const rows = await app.db.execute>(sql`SELECT * FROM ingest_runs WHERE id = ${req.params.runId}`); if (!rows[0]) throw new NotFound('run', req.params.runId); return respond(app, camel(rows[0]), [rows[0].source_id as string]); }); 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) => { const q = req.query; const conds = [sql`u.status = ${q.status}`]; if (q.entityKind) conds.push(sql`u.entity_kind = ${q.entityKind}`); if (q.sourceId) conds.push(sql`(u.source_id = ${q.sourceId} OR s.slug = ${q.sourceId})`); const rows = await app.db.execute & { total: string }>(sql` SELECT u.*, s.slug AS source_slug, c.canonical_name AS suggested_name, count(*) OVER() AS total FROM unresolved_labels u JOIN sources s ON s.id = u.source_id LEFT JOIN cancers c ON c.id = u.suggested_id WHERE ${sql.join(conds, sql` AND `)} ORDER BY u.count DESC, u.id LIMIT ${q.limit} OFFSET ${q.offset}`); const total = rows.length ? Number(rows[0]!.total) : 0; const data = rows.map((r) => { const { total: _t, ...rest } = r; return camel(rest); }); 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 }); }); 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) => { const rows = await app.db.execute>(sql`SELECT * FROM unresolved_labels WHERE id = ${req.params.id}`); const u = rows[0]; if (!u) throw new NotFound('unresolved label', String(req.params.id)); const { cancerId, rejected, reason } = req.body; if (!cancerId && !rejected) throw new BadRequest('provide cancerId or rejected=true'); if (cancerId && rejected) throw new BadRequest('cancerId and rejected are mutually exclusive'); let after: Record; if (rejected) { await app.db.execute(sql`UPDATE unresolved_labels SET status = 'rejected', resolved_by = ${ACTOR}, updated_at = now() WHERE id = ${req.params.id}`); after = { status: 'rejected' }; } else { const cancer = await app.db.execute<{ id: string; canonical_name: string }>(sql`SELECT id, canonical_name FROM cancers WHERE id = ${cancerId!} AND status = 'active'`); if (!cancer[0]) throw new NotFound('cancer', cancerId!); if (u.entity_kind !== 'cancer') throw new BadRequest(`only cancer labels can be mapped here (entity_kind=${u.entity_kind as string})`); const normalized = (u.normalized as string) || normalizeLabel(u.source_text as string); await app.db.transaction(async (tx) => { await tx.execute(sql`INSERT INTO cancer_aliases (cancer_id, alias, normalized, alias_type, source_id, source_terminology, language) VALUES (${cancer[0]!.id}, ${u.source_text as string}, ${normalized}, 'synonym', ${u.source_id as string}, 'curation', 'en') ON CONFLICT DO NOTHING`); 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}`); // Immediate effect on trial conditions that carry the same normalized label; other tables are re-reconciled by their connector's next run. 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}`); 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)`); }); after = { status: 'mapped', resolvedId: cancer[0].id, cancerName: cancer[0].canonical_name, aliasAdded: u.source_text, normalized }; } await audit('unresolved.resolve', 'unresolved_label', String(req.params.id), camel(u), after, reason); return respond(app, { id: req.params.id, ...after }, [u.source_id as string]); }); 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) => { const trace = await traceValue(app.db, req.params.table, req.params.id); return respond(app, trace, []); }); for (const [path, job, key] of [ ['/jobs/counters', JOBS.counters, 'counters'], ['/jobs/rank', JOBS.rank, 'rank'], ] as const) { app.post(path, { schema: { tags: ['admin'], summary: `Enqueue ${job}`, security: [{ adminToken: [] }], body: ReasonBody, response: ok(AnyRecord) } }, async (req) => { const jobId = await enqueue(job, { requestedBy: ACTOR, thenRank: job === JOBS.counters ? false : undefined } as Record, key); await audit(`${job}.enqueue`, 'job', job, null, { jobId }, req.body?.reason); return respond(app, { job, jobId, queued: jobId !== null }, []); }); } 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) => { const q = req.query; const conds = [q.status === 'active' ? sql`status IN ('open','acknowledged')` : sql`status = ${q.status}`]; if (q.connectorId) conds.push(sql`connector_id = ${q.connectorId}`); const rows = await app.db.execute>(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}`); return respond(app, camelRows(rows), []); }); 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) => { const rows = await app.db.execute>(sql`SELECT * FROM audit_log ORDER BY created_at DESC, id DESC LIMIT ${req.query.limit}`); return respond(app, camelRows(rows), []); }); };