diff --git a/apps/web/src/lib/search/server/meili-sink.integration.test.ts b/apps/web/src/lib/search/server/meili-sink.integration.test.ts index 86e4ada..791ce20 100644 --- a/apps/web/src/lib/search/server/meili-sink.integration.test.ts +++ b/apps/web/src/lib/search/server/meili-sink.integration.test.ts @@ -7,7 +7,7 @@ // getmeili/meilisearch:v1.46 // MEILI_TEST_URL=http://localhost:7700 MEILI_TEST_KEY=masterKey \ // pnpm vitest run src/lib/search/server/meili-sink.integration.test.ts -import { describe, it, expect, beforeAll } from 'vitest'; +import { describe, it, expect, beforeAll, afterAll } from 'vitest'; import { createMeiliSink, applyMeiliSettings, type MeiliSinkBackend } from './meili-sink'; import { searchEvents, nearMeEvents, type SearchBackend } from './meili'; @@ -136,12 +136,27 @@ run('contrail.config sink populates Meili on backfill (CLI path)', () => { await applyMeiliSettings({ url: URL!, apiKey: KEY, indexUid: INDEX }); }); + // Restore in afterAll, not inline: a failing assertion in the `it` must not + // leak SEARCH_SINK_* into the rest of the process. + afterAll(() => { + for (const k of ['SEARCH_SINK_URL', 'SEARCH_SINK_API_KEY', 'SEARCH_INDEX'] as const) { + if (savedEnv[k] === undefined) delete process.env[k]; + else process.env[k] = savedEnv[k]; + } + }); + it('indexes a backfilled event through the config sink so the read path finds it', async () => { const { config } = await import('../../contrail.config'); expect(config.sinks && config.sinks.length).toBeTruthy(); await config.sinks![0].onRecords( - [created(uri, { name: 'Backfilled Jazz Night', description: 'via cli backfill', startsAt: FUTURE })], + [ + created(uri, { + name: 'Backfilled Jazz Night', + description: 'via cli backfill', + startsAt: FUTURE + }) + ], { phase: 'backfill' } ); @@ -150,10 +165,5 @@ run('contrail.config sink populates Meili on backfill (CLI path)', () => { (r) => r.hits.some((h) => h.uri === uri) ); expect(text.hits.map((h) => h.uri)).toContain(uri); - - for (const k of ['SEARCH_SINK_URL', 'SEARCH_SINK_API_KEY', 'SEARCH_INDEX'] as const) { - if (savedEnv[k] === undefined) delete process.env[k]; - else process.env[k] = savedEnv[k]; - } }); }); diff --git a/apps/web/src/lib/search/server/reindex-cli.ts b/apps/web/src/lib/search/server/reindex-cli.ts index 6aaaf4d..fcec4df 100644 --- a/apps/web/src/lib/search/server/reindex-cli.ts +++ b/apps/web/src/lib/search/server/reindex-cli.ts @@ -5,28 +5,42 @@ // The Meili backend is resolved from the same env the sink uses, injected by the // operator (so no secret is committed): // SEARCH_SINK_URL / SEARCH_SINK_API_KEY / SEARCH_INDEX -// Reads records_event (SELECT only) and feeds the config sink. Zero D1 writes. +// Reads records_event (SELECT only) and feeds the sink. Zero D1 writes. import { getPlatformProxy } from 'wrangler'; -import { config } from '../../contrail.config'; +import { applyMeiliSettings, createMeiliSink, meiliSinkBackendFromEnv } from './meili-sink'; import { reindexEventsToSink, type ReindexDb } from './reindex'; const remote = process.argv.includes('--remote'); const binding = 'DB'; -const sink = config.sinks?.[0]; -if (!sink) { +// Resolve the backend ourselves rather than borrowing config.sinks[0]: that +// path is gated only on SEARCH_SINK_URL, so a missing/invalid SEARCH_SINK_API_KEY +// leaves the sink a silent no-op while we'd still report "N rows reindexed". +// Here a half-configured env fails loudly before we touch D1. +const backend = meiliSinkBackendFromEnv(process.env); +if (!backend) { console.error( 'No search sink configured. Export SEARCH_SINK_URL and SEARCH_SINK_API_KEY (and optionally SEARCH_INDEX) before running.' ); process.exit(1); } +// Apply index settings up front. It's idempotent, ensures the read path's +// filters (_geo / startsAt / endsAt) resolve even on a never-armed index, and +// doubles as an auth/connectivity check: a bad admin key or unreachable Meili +// throws here instead of letting every per-batch upsert silently fail. +await applyMeiliSettings(backend); +const sink = createMeiliSink(() => backend); + const { env, dispose } = await getPlatformProxy({ environment: remote ? 'production' : undefined }); try { const db = (env as Record)[binding] as ReindexDb | undefined; - if (!db) throw new Error(`No "${binding}" binding in wrangler env (${remote ? 'production' : 'default'}).`); + if (!db) + throw new Error( + `No "${binding}" binding in wrangler env (${remote ? 'production' : 'default'}).` + ); const total = await reindexEventsToSink({ db, diff --git a/apps/web/src/lib/search/server/reindex.test.ts b/apps/web/src/lib/search/server/reindex.test.ts index 1ca3ed1..213a260 100644 --- a/apps/web/src/lib/search/server/reindex.test.ts +++ b/apps/web/src/lib/search/server/reindex.test.ts @@ -101,4 +101,38 @@ describe('reindexEventsToSink', () => { expect(batches[0].records[0].cid).toBe(''); }); + + it('skips a row whose record is missing or not valid JSON without aborting the run', async () => { + const { db } = fakeDb([ + // malformed JSON — the live path uses safeParseJson, so reindex must not + // throw and abort coverage on one poison row. + { + uri: 'at://did:plc:a/community.lexicon.calendar.event/bad', + did: 'did:plc:a', + rkey: 'bad', + cid: 'c', + record: '{not json', + time_us: 1 + }, + // NULL record column. + { + uri: 'at://did:plc:a/community.lexicon.calendar.event/null', + did: 'did:plc:a', + rkey: 'null', + cid: 'c', + record: null, + time_us: 1 + }, + row('at://did:plc:a/community.lexicon.calendar.event/ok', { name: 'OK' }) + ]); + const { sink, batches } = recordingSink(); + + const total = await reindexEventsToSink({ db, sink }); + + // only the valid row is fed; the two bad rows are skipped, not thrown. + expect(total).toBe(1); + expect(batches.flatMap((b) => b.records.map((r) => r.uri))).toEqual([ + 'at://did:plc:a/community.lexicon.calendar.event/ok' + ]); + }); }); diff --git a/apps/web/src/lib/search/server/reindex.ts b/apps/web/src/lib/search/server/reindex.ts index d32b439..79d28b8 100644 --- a/apps/web/src/lib/search/server/reindex.ts +++ b/apps/web/src/lib/search/server/reindex.ts @@ -31,6 +31,21 @@ export interface ReindexOptions { onProgress?: (total: number) => void; } +/** Parse a stored `record` cell into a plain object, or null if it's missing / + * not valid JSON / not an object. The live ingest path uses safeParseJson, so a + * single poison row must be skipped, not allowed to abort the whole reindex. */ +function parseRecordObject(raw: unknown): Record | null { + if (raw == null) return null; + try { + const parsed = JSON.parse(String(raw)); + return parsed && typeof parsed === 'object' && !Array.isArray(parsed) + ? (parsed as Record) + : null; + } catch { + return null; + } +} + /** Pages `records_event` and feeds each batch to `sink.onRecords(..., {phase: * 'backfill'})`. Returns the number of event rows fed. */ export async function reindexEventsToSink(opts: ReindexOptions): Promise { @@ -49,8 +64,17 @@ export async function reindexEventsToSink(opts: ReindexOptions): Promise const rows = page.results ?? []; if (rows.length === 0) break; - const records = rows.map( - (r): RecordEvent => ({ + const records: RecordEvent[] = []; + for (const r of rows) { + // D1 stores `record` as a JSON string; the sink expects a parsed object. + // Skip (don't throw on) a missing/corrupt cell so one poison row can't + // abort coverage for the rest of the table. + const record = parseRecordObject(r.record); + if (!record) { + console.warn(`[reindex] skipping ${String(r.uri)}: record is missing or not valid JSON`); + continue; + } + records.push({ kind: 'created', uri: String(r.uri), did: String(r.did), @@ -59,16 +83,15 @@ export async function reindexEventsToSink(opts: ReindexOptions): Promise // records_event has no `collection` column (one table per collection) // and may store a null cid; the sink wants a string. cid: r.cid == null ? '' : String(r.cid), - // D1 stores `record` as a JSON string; the sink expects a parsed - // object (applyEvents passes safeParseJson(record) on the live path). - record: JSON.parse(String(r.record)) as Record, + record, time_us: Number(r.time_us) - }) - ); + }); + } - await sink.onRecords(records, { phase: 'backfill' }); - total += rows.length; - offset += batchSize; + if (records.length > 0) await sink.onRecords(records, { phase: 'backfill' }); + total += records.length; + // Advance by rows read (not records fed) so skipped rows don't stall paging. + offset += rows.length; opts.onProgress?.(total); if (rows.length < batchSize) break; }