import { readFileSync } from 'node:fs'; import { and, desc, eq, inArray, isNull, sql } from 'drizzle-orm'; import { connectorFieldStats, connectorRuns, normalizedRecords, rawRecords, sales } from '@rareindex/database'; import { DRIFT_FIELDS, driftPresence, fieldNullDrift, loadConnector, priceDistributionAnomaly, type RareIndexConnector } from '@rareindex/connectors'; import { NormalizedRecordSchema, logger, newId, toDateOnly, type NormalizedRecord } from '@rareindex/shared'; import { db } from '../lib/db.ts'; /** * Schema-drift detection (SPEC §13): accumulate per-field presence counters for today, then compare * today's null rates with the trailing 7-day baseline. Returns anomaly strings (possibly empty). */ async function recordFieldStats(connectorId: string, records: NormalizedRecord[]): Promise { if (!records.length) return []; const counts: Record = {}; for (const f of DRIFT_FIELDS) counts[f] = { total: 0, nulls: 0 }; for (const r of records) { const p = driftPresence(r as unknown as Record); for (const f of DRIFT_FIELDS) { counts[f]!.total++; if (!p[f]) counts[f]!.nulls++; } } const day = toDateOnly(new Date()); for (const [field, c] of Object.entries(counts)) { await db() .insert(connectorFieldStats) .values({ connectorId, day, field, total: c.total, nulls: c.nulls }) .onConflictDoUpdate({ target: [connectorFieldStats.connectorId, connectorFieldStats.day, connectorFieldStats.field], set: { total: sql`${connectorFieldStats.total} + ${c.total}`, nulls: sql`${connectorFieldStats.nulls} + ${c.nulls}` } }); } const rows = (await db().execute(sql`select field, day::text as day, total, nulls from connector_field_stats where connector_id = ${connectorId} and day >= (current_date - interval '7 days')::date`)) as unknown as Array<{ field: string; day: string; total: number; nulls: number }>; const baseline: Record = {}; const today: Record = {}; for (const r of rows) { const target = r.day === day ? today : baseline; const t = (target[r.field] ??= { total: 0, nulls: 0 }); t.total += Number(r.total); t.nulls += Number(r.nulls); } return fieldNullDrift(baseline, today); } export interface NormalizeResult { processed: number; produced: number; failed: number; anomalies: string[]; } const connectorCache = new Map>(); function getConnector(id: string): Promise { let p = connectorCache.get(id); if (!p) { p = loadConnector(id); connectorCache.set(id, p); } return p; } /** * Normalizer (§108): raw_records → normalized_records. Each raw row is processed exactly once; * failures are recorded on the raw row (process_error) and counted as anomalies — never swallowed. */ export async function normalizeBatch(opts: { connectorId?: string; rawIds?: string[]; limit?: number } = {}): Promise { const log = logger.child({ component: 'normalizer', connector: opts.connectorId }); const limit = opts.limit ?? 500; const where = opts.rawIds?.length ? inArray(rawRecords.id, opts.rawIds) : and(isNull(rawRecords.processedAt), opts.connectorId ? eq(rawRecords.connectorId, opts.connectorId) : undefined); const rows = await db().select().from(rawRecords).where(where).orderBy(rawRecords.fetchedAt).limit(limit); const result: NormalizeResult = { processed: 0, produced: 0, failed: 0, anomalies: [] }; if (rows.length === 0) return result; const byConnector = new Map(); for (const r of rows) byConnector.set(r.connectorId, [...(byConnector.get(r.connectorId) ?? []), r]); for (const [connectorId, group] of byConnector) { let connector: RareIndexConnector; try { connector = await getConnector(connectorId); } catch (err) { const msg = `connector load failed: ${err instanceof Error ? err.message : String(err)}`; log.error({ err }, msg); await db().update(rawRecords).set({ processedAt: new Date(), processError: msg }).where(inArray(rawRecords.id, group.map((g) => g.id))); result.failed += group.length; result.anomalies.push(`parse_failure: ${msg}`); continue; } const inserts: Array = []; const done: Array<{ id: string; error: string | null }> = []; const batchPrices: number[] = []; const produced: NormalizedRecord[] = []; for (const raw of group) { try { let payload = raw.payload; if (raw.snapshotRef && payload && typeof payload === 'object' && !('snapshot' in (payload as object))) { try { payload = { ...(payload as object), snapshot: readFileSync(raw.snapshotRef, 'utf8') }; } catch { /* snapshot missing: proceed with payload only */ } } const out = await connector.normalize({ id: raw.id, connectorId, sourceId: raw.sourceId, url: raw.url, externalId: raw.externalId, kind: raw.kind as NormalizedRecord['kind'], payload, fetchedAt: raw.fetchedAt, engine: raw.engine as 'api' }); let n = 0; for (const rec of out) { const parsed = NormalizedRecordSchema.safeParse(rec); if (!parsed.success) { result.anomalies.push(`parse_failure: schema ${parsed.error.issues[0]?.path.join('.')} ${parsed.error.issues[0]?.message}`); result.failed++; continue; } const p = parsed.data; produced.push(p); if ((p.kind === 'sale' || p.kind === 'price_observation') && p.price > 0) batchPrices.push(p.price); inserts.push({ id: newId('raw'), rawRecordId: raw.id, connectorId, sourceId: raw.sourceId, kind: p.kind, seq: n, payload: p, status: 'pending' }); n++; } done.push({ id: raw.id, error: n === 0 && out.length === 0 ? 'normalize produced no records' : null }); result.produced += n; } catch (err) { const msg = err instanceof Error ? err.message : String(err); log.warn({ err, raw: raw.id }, 'normalize failed'); done.push({ id: raw.id, error: `normalize: ${msg}`.slice(0, 500) }); result.failed++; result.anomalies.push(`parse_failure: ${msg.slice(0, 120)}`); } result.processed++; } // normalized_records has a unique raw_record_id index; a raw record producing several records // is common (catalog + observations) so we key on (raw, kind, position) instead. for (let i = 0; i < inserts.length; i += 500) { await db().insert(normalizedRecords).values(inserts.slice(i, i + 500)).onConflictDoNothing(); } const okIds = done.filter((d) => !d.error).map((d) => d.id); if (okIds.length) await db().update(rawRecords).set({ processedAt: new Date(), processError: null }).where(inArray(rawRecords.id, okIds)); for (const d of done.filter((d) => d.error)) await db().update(rawRecords).set({ processedAt: new Date(), processError: d.error }).where(eq(rawRecords.id, d.id)); // Schema drift: field-null rates vs the trailing week (SPEC §13) try { result.anomalies.push(...(await recordFieldStats(connectorId, produced))); } catch (err) { log.warn({ err }, 'field stats failed'); } // Price distribution anomaly vs the connector's recent accepted sales (§105) if (batchPrices.length >= 10) { const baseline = await db().select({ p: sales.priceUsd }).from(sales).where(eq(sales.connectorId, connectorId)).orderBy(desc(sales.saleDate)).limit(500); const a = priceDistributionAnomaly(baseline.map((b) => Number(b.p)), batchPrices); if (a) result.anomalies.push(a); } if (result.anomalies.length) { const uniq = [...new Set(result.anomalies)].slice(0, 20); const [lastRun] = await db().select({ id: connectorRuns.id }).from(connectorRuns).where(eq(connectorRuns.connectorId, connectorId)).orderBy(desc(connectorRuns.startedAt)).limit(1); if (lastRun) await db().update(connectorRuns).set({ anomalies: sql`(SELECT jsonb_agg(DISTINCT x) FROM jsonb_array_elements(${connectorRuns.anomalies} || ${JSON.stringify(uniq)}::jsonb) x)` }).where(eq(connectorRuns.id, lastRun.id)); } } log.info(result, 'normalize batch done'); return result; }