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%
13.9 KB · 167 lines typescript
Raw Blame History
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