import { desc, eq, gte, sql } from 'drizzle-orm'; import { connectorHealth, connectorRuns, connectors as connectorsTable } from '@rareindex/database'; import { createHealthContext, listConnectorMeta, loadConnector } from '@rareindex/connectors'; import { logger, type ConnectorHealth } from '@rareindex/shared'; import { db } from './lib/db.ts'; import { getRouter } from './lib/router.ts'; /** Compute and persist ConnectorHealth for every registered connector (§105, §144). */ export async function computeHealth(opts: { connectorId?: string; probe?: boolean } = {}): Promise { const metas = listConnectorMeta().filter((m) => !opts.connectorId || m.id === opts.connectorId); const out: ConnectorHealth[] = []; const since = new Date(Date.now() - 7 * 86_400_000); for (const meta of metas) { const runs = await db().select().from(connectorRuns).where(eq(connectorRuns.connectorId, meta.id)).orderBy(desc(connectorRuns.startedAt)).limit(50); const recentRuns = runs .filter((r) => r.startedAt >= since) .reverse() .map((r) => ({ startedAt: r.startedAt, status: r.status, pagesAttempted: r.pagesAttempted, pagesSuccess: r.pagesSuccess, recordsRaw: r.recordsRaw, recordsDuplicate: r.recordsDuplicate, engineStats: r.engineStats as Record, anomalies: (r.anomalies as string[]) ?? [], error: r.error })); let health: ConnectorHealth; try { const connector = await loadConnector(meta.id); const ctx = createHealthContext({ router: getRouter(meta.id), meta, recentRuns: opts.probe ? recentRuns : recentRuns.length ? recentRuns : [{ startedAt: new Date(0), status: 'unknown', pagesAttempted: 0, pagesSuccess: 0, recordsRaw: 0, recordsDuplicate: 0, engineStats: {}, anomalies: [], error: null }] }); health = await connector.healthCheck(ctx); if (!recentRuns.length && !opts.probe) health.status = 'unknown'; } catch (err) { health = { connector: meta.id, status: 'failing', success_rate_24h: null, pages_attempted: 0, pages_success: 0, firecrawl_success_rate: null, scrapfly_fallback_rate: null, parse_failure_rate: null, last_success: null, last_error: err instanceof Error ? err.message : String(err), schema_version: meta.schemaVersion, records_24h: 0, duplicates_24h: 0, anomalies: ['load_failure'] }; } const [state] = await db().select({ status: connectorsTable.status }).from(connectorsTable).where(eq(connectorsTable.id, meta.id)).limit(1); if (state?.status === 'paused') health.status = 'paused'; else if (state?.status === 'maintenance') health.status = 'maintenance'; else if (state?.status === 'disabled') health.status = 'disabled'; // data freshness: newest source-side date ingested by this connector (SPEC §12) const [fresh] = (await db().execute(sql`select greatest((select max(sale_date) from sales where connector_id = ${meta.id}), (select max(last_seen_at) from listings where connector_id = ${meta.id}), (select max(observation_date)::timestamptz from price_observations where connector_id = ${meta.id})) as ts`)) as unknown as Array<{ ts: Date | string | null }>; health.data_freshness = fresh?.ts ? new Date(fresh.ts).toISOString() : null; await db().insert(connectorHealth).values({ connectorId: meta.id, computedAt: new Date(), status: health.status, health }).onConflictDoUpdate({ target: connectorHealth.connectorId, set: { computedAt: new Date(), status: health.status, health } }); out.push(health); } logger.info({ connectors: out.length }, 'health computed'); return out; } export { gte };