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%
4.5 KB · 55 lines typescript
Raw Blame History
1import { sql } from 'drizzle-orm';2import type { FastifyPluginAsyncZod } from 'fastify-type-provider-zod';3import { z } from 'zod';4import { NotFound } from '../lib/errors.js';5import { AnyList, AnyRecord, camel, camelRows, ok, respond } from '../lib/respond.js';67const SOURCE_COLS = sql`8  s.id, s.slug, s.name, s.organization, s.category, s.description, s.homepage, s.docs_url, s.terms_url, s.access_type, s.access_auth, s.license, s.license_status, s.commercial_use, s.redistribution,9  s.attribution, s.license_reviewed_at, s.approved_for_production, s.update_frequency, s.supports_incremental, s.entities, s.metrics, s.rate_limit, s.status, s.tier,10  s.manifest->>'documentationVerifiedAt' AS documentation_verified_at, s.manifest->>'schedule' AS schedule, s.manifest->>'termsNotes' AS terms_notes,11  cc.health, cc.health_detail, cc.paused, cc.last_success_at, cc.last_attempt_at, cc.cursor,12  (SELECT count(*) FROM source_records sr WHERE sr.source_id = s.id) AS record_count,13  (SELECT count(*) FROM provenance p WHERE p.source_id = s.id) AS provenance_count,14  (SELECT count(*) FROM unresolved_labels u WHERE u.source_id = s.id AND u.status = 'open') AS unresolved_open`;1516function shape(r: Record<string, unknown>) {17  const c = camel<Record<string, unknown>>(r);18  const { health, healthDetail, paused, lastSuccessAt, lastAttemptAt, cursor, recordCount, provenanceCount, unresolvedOpen, ...rest } = c;19  return {20    ...rest,21    connector: { health: health ?? 'never_run', healthDetail: healthDetail ?? null, paused: paused ?? false, lastSuccessAt: lastSuccessAt ?? null, lastAttemptAt: lastAttemptAt ?? null, cursor: cursor ?? null },22    counts: { sourceRecords: Number(recordCount ?? 0), provenanceRows: Number(provenanceCount ?? 0), unresolvedLabelsOpen: Number(unresolvedOpen ?? 0) },23  };24}2526export const sourceRoutes: FastifyPluginAsyncZod = async (app) => {27  app.get('/sources', { schema: { tags: ['sources'], summary: 'Source registry with license status, connector health, last runs and record counts', response: ok(AnyList) } }, async () => {28    const [rows, lastRuns] = await Promise.all([29      app.db.execute<Record<string, unknown>>(sql`SELECT ${SOURCE_COLS} FROM sources s LEFT JOIN connector_cursors cc ON cc.connector_id = s.slug ORDER BY s.tier, s.slug`),30      app.db.execute<Record<string, unknown>>(sql`31        SELECT DISTINCT ON (connector_id) connector_id, id, mode, status, started_at, finished_at, duration_ms, records_fetched, records_created, records_updated, records_unchanged, records_rejected, error, anomaly, dataset_version32        FROM ingest_runs ORDER BY connector_id, started_at DESC`),33    ]);34    const lastBy = new Map(lastRuns.map((r) => [r.connector_id as string, camel(r)]));35    const data = rows.map((r) => ({ ...shape(r), lastRun: lastBy.get(r.slug as string) ?? null }));36    return respond(app, data, rows.map((r) => r.id as string));37  });3839  app.get('/sources/:slug', { schema: { tags: ['sources'], summary: 'Source detail: manifest, license, health, last 10 runs, record counts by entity', params: z.object({ slug: z.string() }), response: ok(AnyRecord) } }, async (req) => {40    const rows = await app.db.execute<Record<string, unknown>>(sql`SELECT ${SOURCE_COLS}, s.manifest FROM sources s LEFT JOIN connector_cursors cc ON cc.connector_id = s.slug WHERE s.slug = ${req.params.slug} OR s.id = ${req.params.slug}`);41    const src = rows[0];42    if (!src) throw new NotFound('source', req.params.slug);43    const [runs, byEntity, fieldStats] = await Promise.all([44      app.db.execute<Record<string, unknown>>(sql`45        SELECT 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,46          jsonb_array_length(schema_drift) AS drift_signals, error, anomaly, dataset_version47        FROM ingest_runs WHERE connector_id = ${src.slug as string} ORDER BY started_at DESC LIMIT 10`),48      app.db.execute<Record<string, unknown>>(sql`SELECT entity_kind, count(*) AS n, max(retrieved_at) AS last_retrieved_at FROM source_records WHERE source_id = ${src.id as string} GROUP BY entity_kind ORDER BY entity_kind`),49      app.db.execute<{ n: string }>(sql`SELECT count(*) AS n FROM connector_field_stats WHERE connector_id = ${src.slug as string}`),50    ]);51    const data = { ...shape(src), recentRuns: camelRows(runs), recordsByEntity: camelRows(byEntity), observedFields: Number(fieldStats[0]?.n ?? 0) };52    return respond(app, data, [src.id as string]);53  });54};55