diff --git a/.changeset/indexed-d1-fts.md b/.changeset/indexed-d1-fts.md new file mode 100644 index 0000000..1cd8b82 --- /dev/null +++ b/.changeset/indexed-d1-fts.md @@ -0,0 +1,5 @@ +--- +"@atmo-dev/contrail": patch +--- + +Replace SQLite and D1 full-text-search URI scans with an ordinary unique URI-to-rowid mapping and direct FTS5 rowid mutations. Existing URI-bearing FTS tables are rebuilt transactionally from canonical records, stale-fingerprint projections are rebuilt and verified before acceptance, duplicate search rows are removed, and incremental/rebuild whitespace normalization now agrees. PostgreSQL search remains unchanged. diff --git a/packages/contrail/src/core/db/records.ts b/packages/contrail/src/core/db/records.ts index eda15c1..e06cac0 100644 --- a/packages/contrail/src/core/db/records.ts +++ b/packages/contrail/src/core/db/records.ts @@ -23,8 +23,17 @@ import { feedTargetMaxItems, DEFAULT_FOLLOW_SHORT, } from "../types"; -import { getSearchableFields, ftsTableName, buildFtsContent } from "../search"; -import { ftsQueryClause, getDialect } from "../dialect"; +import { + getSearchableFields, + ftsRowTableName, + ftsTableName, + buildFtsContent, +} from "../search"; +import { + ftsQueryClause, + getDialect, + sqliteFtsContentExpression, +} from "../dialect"; import type { SourcePosition } from "../sources"; // --- Counts --- @@ -215,32 +224,40 @@ function buildFtsStatements( if (!fields || fields.length === 0) return []; const table = ftsTableName(short); - const stmts: Statement[] = []; + const rowsTable = ftsRowTableName(short); + const deleteFtsRow = () => + db + .prepare( + `DELETE FROM ${table} WHERE rowid = (SELECT id FROM ${rowsTable} WHERE uri = ?)`, + ) + .bind(event.uri); + const deleteMapping = () => + db.prepare(`DELETE FROM ${rowsTable} WHERE uri = ?`).bind(event.uri); if (event.operation === "delete") { - stmts.push(db.prepare(`DELETE FROM ${table} WHERE uri = ?`).bind(event.uri)); - } else { - const record = event.record ? JSON.parse(event.record) : null; - if (!record) return []; - - // Always delete first so FTS sync is idempotent. The FTS virtual table has no - // uniqueness constraint, so a bare insert appends a duplicate row when one - // already exists. existingMap is unreliable here: backfill runs with - // skipReplayDetection, leaving it empty, so a re-applied record would look new - // and accumulate duplicate rows that fan out the search JOIN. The delete is - // unconditional so it also evicts a stale row when an update clears all - // searchable fields (content is null); only the re-insert is gated on content. - stmts.push(db.prepare(`DELETE FROM ${table} WHERE uri = ?`).bind(event.uri)); - - const content = buildFtsContent(record, fields); - if (content) { - stmts.push( - db.prepare(`INSERT INTO ${table} (uri, content) VALUES (?, ?)`).bind(event.uri, content) - ); - } + return [deleteFtsRow(), deleteMapping()]; } - return stmts; + const record = event.record ? JSON.parse(event.record) : null; + if (!record) return []; + const content = buildFtsContent(record, fields); + if (!content) { + return [deleteFtsRow(), deleteMapping()]; + } + + return [ + db + .prepare( + `INSERT INTO ${rowsTable} (uri) VALUES (?) ON CONFLICT(uri) DO NOTHING`, + ) + .bind(event.uri), + deleteFtsRow(), + db + .prepare( + `INSERT INTO ${table} (rowid, content) SELECT id, ? FROM ${rowsTable} WHERE uri = ?`, + ) + .bind(content, event.uri), + ]; } // --- Feeds --- @@ -1137,21 +1154,29 @@ export async function rebuildDerivedProjections( if (!fields || fields.length === 0) continue; const ftsTable = ftsTableName(short); + const rowsTable = ftsRowTableName(short); const recordsTable = recordsTableName(short); - const terms = fields.map((field) => { - const path = `$.${field}`; - return `CASE WHEN json_type(record, '${path}') = 'text' THEN json_extract(record, '${path}') ELSE '' END`; - }); - const content = `trim(${terms.join(" || ' ' || ")})`; + const content = sqliteFtsContentExpression(fields); try { await db.batch([ db.prepare(`DELETE FROM ${ftsTable}`), + db.prepare(`DELETE FROM ${rowsTable}`), + db.prepare( + `INSERT INTO ${rowsTable} (uri) + SELECT uri FROM ( + SELECT uri, ${content} AS content FROM ${recordsTable} + ) rebuilt + WHERE content <> '' + ORDER BY uri`, + ), db.prepare( - `INSERT INTO ${ftsTable} (uri, content) - SELECT uri, content FROM ( + `INSERT INTO ${ftsTable} (rowid, content) + SELECT fts_rows.id, rebuilt.content + FROM ( SELECT uri, ${content} AS content FROM ${recordsTable} ) rebuilt - WHERE content <> ''` + JOIN ${rowsTable} fts_rows ON fts_rows.uri = rebuilt.uri + WHERE rebuilt.content <> ''`, ), ]); } catch { diff --git a/packages/contrail/src/core/db/schema.ts b/packages/contrail/src/core/db/schema.ts index fde2f3e..95fa42d 100644 --- a/packages/contrail/src/core/db/schema.ts +++ b/packages/contrail/src/core/db/schema.ts @@ -2,10 +2,16 @@ import type { ContrailConfig, Database, ResolvedContrailConfig, + Statement, ResolvedMaps, } from "../types"; import type { SqlDialect } from "../dialect"; -import { buildFtsSchema, getDialect, postgresDialect } from "../dialect"; +import { + buildFtsSchema, + getDialect, + postgresDialect, + sqliteFtsContentExpression, +} from "../dialect"; import { countColumnName, getRelationField, @@ -13,11 +19,15 @@ import { recordsTableName, resolveConfig, } from "../types"; -import { getSearchableFields } from "../search"; +import { + ftsRowTableName, + ftsTableName, + getSearchableFields, +} from "../search"; import { buildLabelsSchema } from "../labels/schema"; import { getMeta, setMeta } from "./meta"; -export const CONTRAIL_SCHEMA_VERSION = 10; +export const CONTRAIL_SCHEMA_VERSION = 11; const SCHEMA_FINGERPRINT_KEY = "schema_fingerprint"; function getResolved(config: ContrailConfig): ResolvedMaps { @@ -360,6 +370,181 @@ export function buildFtsTables( return statements; } +async function virtualFtsColumns( + db: Database, + table: string, +): Promise { + const rows = await db + .prepare(`PRAGMA table_info(${table})`) + .all<{ name: string }>(); + return (rows.results ?? []).map((row) => row.name); +} + +function populateVirtualFtsStatements( + db: Database, + collection: string, + fields: string[], +): Statement[] { + const recordsTable = recordsTableName(collection); + const ftsTable = ftsTableName(collection); + const rowsTable = ftsRowTableName(collection); + const content = sqliteFtsContentExpression(fields); + return [ + db.prepare( + `INSERT INTO ${rowsTable} (uri) + SELECT uri FROM ( + SELECT uri, ${content} AS content FROM ${recordsTable} + ) rebuilt + WHERE content <> '' + ORDER BY uri`, + ), + db.prepare( + `INSERT INTO ${ftsTable} (rowid, content) + SELECT fts_rows.id, rebuilt.content + FROM ( + SELECT uri, ${content} AS content FROM ${recordsTable} + ) rebuilt + JOIN ${rowsTable} fts_rows ON fts_rows.uri = rebuilt.uri + WHERE rebuilt.content <> ''`, + ), + ]; +} + +async function verifyVirtualFtsProjection( + db: Database, + collection: string, + fields: string[], +): Promise { + const recordsTable = recordsTableName(collection); + const ftsTable = ftsTableName(collection); + const rowsTable = ftsRowTableName(collection); + const content = sqliteFtsContentExpression(fields); + const row = await db + .prepare( + `SELECT + (SELECT COUNT(*) FROM ( + SELECT uri, ${content} AS content FROM ${recordsTable} + ) expected WHERE content <> '') AS expected_rows, + (SELECT COUNT(*) FROM ${rowsTable}) AS mapping_rows, + (SELECT COUNT(*) FROM ${ftsTable}) AS fts_rows, + (SELECT COUNT(*) + FROM ${rowsTable} mapped + JOIN ${ftsTable} fts ON fts.rowid = mapped.id + JOIN ( + SELECT uri FROM ( + SELECT uri, ${content} AS content FROM ${recordsTable} + ) searchable + WHERE content <> '' + ) expected ON expected.uri = mapped.uri) AS joined_rows`, + ) + .first<{ + expected_rows: number | string; + mapping_rows: number | string; + fts_rows: number | string; + joined_rows: number | string; + }>(); + const expected = Number(row?.expected_rows ?? 0); + const mapping = Number(row?.mapping_rows ?? 0); + const fts = Number(row?.fts_rows ?? 0); + const joined = Number(row?.joined_rows ?? 0); + if (mapping !== expected || fts !== expected || joined !== expected) { + throw new Error( + `FTS migration verification failed for ${collection}: ` + + `expected=${expected}, mapping=${mapping}, fts=${fts}, joined=${joined}`, + ); + } +} + +async function migrateLegacyVirtualFts( + db: Database, + config: ContrailConfig, + dialect: SqlDialect, + collection: string, + fields: string[], +): Promise { + const recordsTable = recordsTableName(collection); + const ftsTable = ftsTableName(collection); + const rowsTable = ftsRowTableName(collection); + const schema = buildFtsSchema(dialect, recordsTable, fields); + await db.batch([ + db.prepare(`DROP TABLE IF EXISTS ${ftsTable}`), + db.prepare(`DROP TABLE IF EXISTS ${rowsTable}`), + ...schema.map((statement) => db.prepare(statement)), + ...populateVirtualFtsStatements(db, collection, fields), + ]); + await verifyVirtualFtsProjection(db, collection, fields); + (config.logger ?? console).log( + `[schema] migrated ${ftsTable} to indexed rowid maintenance`, + ); +} + +async function rebuildVirtualFtsProjection( + db: Database, + collection: string, + fields: string[], +): Promise { + const ftsTable = ftsTableName(collection); + const rowsTable = ftsRowTableName(collection); + await db.batch([ + db.prepare(`DELETE FROM ${ftsTable}`), + db.prepare(`DELETE FROM ${rowsTable}`), + ...populateVirtualFtsStatements(db, collection, fields), + ]); + await verifyVirtualFtsProjection(db, collection, fields); +} + +async function applyFtsTables( + db: Database, + config: ContrailConfig, + dialect: SqlDialect, +): Promise { + if (dialect.ftsStrategy === "generated-column") { + for (const statement of buildFtsTables(config, dialect)) { + await db.prepare(statement).run(); + } + return; + } + + for (const [collection, colConfig] of Object.entries(config.collections)) { + const fields = getSearchableFields(collection, colConfig); + if (!fields || fields.length === 0) continue; + const ftsTable = ftsTableName(collection); + const columns = await virtualFtsColumns(db, ftsTable); + if (columns.includes("uri")) { + await migrateLegacyVirtualFts( + db, + config, + dialect, + collection, + fields, + ); + continue; + } + + const schema = buildFtsSchema( + dialect, + recordsTableName(collection), + fields, + ); + await db.prepare(schema[0]!).run(); + try { + await db.prepare(schema[1]!).run(); + } catch { + // FTS5 is optional in some SQLite builds. The ordinary mapping table is + // harmless when the virtual table cannot be created. + } + const currentColumns = + columns.length > 0 ? columns : await virtualFtsColumns(db, ftsTable); + if (currentColumns.length > 0) { + // applyFtsTables only runs while the global schema fingerprint is stale. + // Rebuild even an already content-only table before accepting the new + // fingerprint: a prior migration may have committed and then failed + // verification, or the configured searchable fields may have changed. + await rebuildVirtualFtsProjection(db, collection, fields); + } + } +} + interface MigrationOp { table: string; column: string; @@ -540,13 +725,7 @@ export async function initSchema( } } - for (const statement of fts) { - try { - await db.prepare(statement).run(); - } catch { - // FTS5 is not available in every SQLite build. - } - } + await applyFtsTables(db, config, dialect); const hasFeeds = !!(config.feeds && Object.keys(config.feeds).length > 0); await runMigrations(db, hasFeeds); diff --git a/packages/contrail/src/core/dialect.ts b/packages/contrail/src/core/dialect.ts index f5c3371..f03bad2 100644 --- a/packages/contrail/src/core/dialect.ts +++ b/packages/contrail/src/core/dialect.ts @@ -1,3 +1,5 @@ +import { FTS_TRIM_CODE_POINTS } from "./search"; + /** Get the dialect from a Database, defaulting to SQLite (for D1 compatibility) */ export function getDialect(db: { dialect?: SqlDialect }): SqlDialect { return db.dialect ?? sqliteDialect; @@ -11,6 +13,23 @@ function assertSafeField(field: string): void { } } +/** Build the SQLite expression used by both incremental and set-based FTS + * projection. Only JSON string values participate, matching buildFtsContent(). */ +export function sqliteFtsContentExpression( + fields: string[], + recordColumn = "record", +): string { + assertSafeField(recordColumn); + if (fields.length === 0) throw new Error("FTS fields must not be empty"); + const terms = fields.map((field) => { + assertSafeField(field); + const path = `$.${field}`; + return `CASE WHEN json_type(${recordColumn}, '${path}') = 'text' THEN json_extract(${recordColumn}, '${path}') ELSE '' END`; + }); + const trimCharacters = `char(${FTS_TRIM_CODE_POINTS.join(", ")})`; + return `trim(${terms.join(" || ' ' || ")}, ${trimCharacters})`; +} + export interface SqlDialect { /** json_extract(col, '$.field') or col->>'field' */ jsonExtract(column: string, field: string): string; @@ -91,8 +110,13 @@ export function buildFtsSchema( ): string[] { if (dialect.ftsStrategy === "virtual-table") { const ftsTable = recordsTable.replace("records_", "fts_"); + const rowsTable = `${ftsTable}_rows`; return [ - `CREATE VIRTUAL TABLE IF NOT EXISTS ${ftsTable} USING fts5(uri UNINDEXED, content)` + `CREATE TABLE IF NOT EXISTS ${rowsTable} ( + id INTEGER PRIMARY KEY, + uri TEXT NOT NULL UNIQUE + )`, + `CREATE VIRTUAL TABLE IF NOT EXISTS ${ftsTable} USING fts5(content)`, ]; } else { const concatExpr = fields @@ -117,8 +141,9 @@ export function ftsQueryClause( } { if (dialect.ftsStrategy === "virtual-table") { const ftsTable = recordsTable.replace("records_", "fts_"); + const rowsTable = `${ftsTable}_rows`; return { - join: `JOIN ${ftsTable} fts ON fts.uri = r.uri`, + join: `JOIN ${rowsTable} fts_rows ON fts_rows.uri = r.uri JOIN ${ftsTable} fts ON fts.rowid = fts_rows.id`, condition: "fts.content MATCH ?", orderExpr: "fts.rank", orderDirection: "asc", diff --git a/packages/contrail/src/core/search.ts b/packages/contrail/src/core/search.ts index a33dc5e..462b7d6 100644 --- a/packages/contrail/src/core/search.ts +++ b/packages/contrail/src/core/search.ts @@ -1,6 +1,57 @@ import type { CollectionConfig } from "./types"; import { getNestedValue } from "./types"; +/** ECMAScript whitespace and line-terminator code points trimmed from the + * combined FTS document. SQLite receives this same explicit set via char(). */ +export const FTS_TRIM_CODE_POINTS = [ + 0x0009, + 0x000a, + 0x000b, + 0x000c, + 0x000d, + 0x0020, + 0x00a0, + 0x1680, + 0x2000, + 0x2001, + 0x2002, + 0x2003, + 0x2004, + 0x2005, + 0x2006, + 0x2007, + 0x2008, + 0x2009, + 0x200a, + 0x2028, + 0x2029, + 0x202f, + 0x205f, + 0x3000, + 0xfeff, +] as const; + +const FTS_TRIM_CHARACTERS = new Set(FTS_TRIM_CODE_POINTS); + +export function trimFtsWhitespace(value: string): string { + const codePoints = [...value]; + let start = 0; + let end = codePoints.length; + while ( + start < end && + FTS_TRIM_CHARACTERS.has(codePoints[start]!.codePointAt(0)!) + ) { + start++; + } + while ( + end > start && + FTS_TRIM_CHARACTERS.has(codePoints[end - 1]!.codePointAt(0)!) + ) { + end--; + } + return codePoints.slice(start, end).join(""); +} + /** * Resolve which fields are searchable for a collection. * Returns null if search is disabled or no fields found. @@ -13,9 +64,18 @@ export function getSearchableFields( return colConfig.searchable.length > 0 ? colConfig.searchable : null; } -/** Sanitized FTS table name for a collection. */ +function sanitizedCollectionName(collection: string): string { + return collection.replace(/[^a-zA-Z0-9]/g, "_"); +} + +/** Sanitized FTS virtual-table name for a collection. */ export function ftsTableName(collection: string): string { - return `fts_${collection.replace(/[^a-zA-Z0-9]/g, "_")}`; + return `fts_${sanitizedCollectionName(collection)}`; +} + +/** Ordinary unique URI-to-FTS-rowid mapping table for a collection. */ +export function ftsRowTableName(collection: string): string { + return `${ftsTableName(collection)}_rows`; } /** Extract searchable field values from a record and join them into a single string. */ @@ -27,5 +87,6 @@ export function buildFtsContent(record: unknown, fields: string[]): string | nul parts.push(value); } } - return parts.length > 0 ? parts.join(" ") : null; + const content = trimFtsWhitespace(parts.join(" ")); + return content.length > 0 ? content : null; } diff --git a/packages/contrail/src/core/verification.ts b/packages/contrail/src/core/verification.ts index 8dfaa01..44c520c 100644 --- a/packages/contrail/src/core/verification.ts +++ b/packages/contrail/src/core/verification.ts @@ -1,5 +1,11 @@ import type { ContrailConfig, Database } from "./types"; import { recordsTableName } from "./types"; +import { getDialect, sqliteFtsContentExpression } from "./dialect"; +import { + ftsRowTableName, + ftsTableName, + getSearchableFields, +} from "./search"; export const BOOTSTRAP_VERIFICATION_META_KEY = "bootstrap_verification"; @@ -86,6 +92,49 @@ export async function verifyBootstrapCandidate( ), ), ); + + const fields = getSearchableFields(shortName, collectionConfig); + if (getDialect(db).ftsStrategy === "virtual-table" && fields) { + const ftsTable = ftsTableName(shortName); + const rowsTable = ftsRowTableName(shortName); + const tables = await count( + db, + "SELECT COUNT(*) AS count FROM sqlite_master WHERE type = 'table' AND name IN (?, ?)", + [ftsTable, rowsTable], + ); + checks.push(check(`fts-tables:${shortName}`, Math.max(0, 2 - tables))); + if (tables === 2) { + const content = sqliteFtsContentExpression(fields); + const expected = await count( + db, + `SELECT COUNT(*) AS count FROM ( + SELECT ${content} AS content FROM ${table} + ) searchable WHERE content <> ''`, + ); + const mapping = await count( + db, + `SELECT COUNT(*) AS count FROM ${rowsTable}`, + ); + const fts = await count(db, `SELECT COUNT(*) AS count FROM ${ftsTable}`); + const joined = await count( + db, + `SELECT COUNT(*) AS count + FROM ${rowsTable} mapped + JOIN ${ftsTable} fts ON fts.rowid = mapped.id + JOIN ( + SELECT uri FROM ( + SELECT uri, ${content} AS content FROM ${table} + ) searchable + WHERE content <> '' + ) expected ON expected.uri = mapped.uri`, + ); + checks.push( + check(`fts-mapping:${shortName}`, Math.abs(expected - mapping)), + check(`fts-rows:${shortName}`, Math.abs(expected - fts)), + check(`fts-joined:${shortName}`, Math.abs(expected - joined)), + ); + } + } } return { diff --git a/packages/contrail/tests/dialect.test.ts b/packages/contrail/tests/dialect.test.ts index 209d728..1a8fea1 100644 --- a/packages/contrail/tests/dialect.test.ts +++ b/packages/contrail/tests/dialect.test.ts @@ -140,11 +140,14 @@ describe("indexExpression", () => { }); describe("FTS schema generation", () => { - it("sqlite generates virtual table", () => { + it("sqlite generates a unique URI map and content-only virtual table", () => { const stmts = buildFtsSchema(sqliteDialect, "records_community_lexicon_calendar_event", ["name", "description"]); - expect(stmts).toHaveLength(1); - expect(stmts[0]).toContain("CREATE VIRTUAL TABLE"); - expect(stmts[0]).toContain("USING fts5"); + expect(stmts).toHaveLength(2); + expect(stmts[0]).toContain("fts_community_lexicon_calendar_event_rows"); + expect(stmts[0]).toContain("uri TEXT NOT NULL UNIQUE"); + expect(stmts[1]).toContain("CREATE VIRTUAL TABLE"); + expect(stmts[1]).toContain("USING fts5(content)"); + expect(stmts[1]).not.toContain("uri UNINDEXED"); }); it("postgres generates tsvector column + GIN index", () => { @@ -158,7 +161,8 @@ describe("FTS schema generation", () => { describe("FTS query clause", () => { it("sqlite uses FTS5 join and MATCH", () => { const clause = ftsQueryClause(sqliteDialect, "records_community_lexicon_calendar_event"); - expect(clause.join).toContain("JOIN fts_"); + expect(clause.join).toContain("JOIN fts_community_lexicon_calendar_event_rows"); + expect(clause.join).toContain("fts.rowid = fts_rows.id"); expect(clause.condition).toContain("MATCH"); expect(clause.orderExpr).toBe("fts.rank"); }); diff --git a/packages/contrail/tests/search.test.ts b/packages/contrail/tests/search.test.ts index 3bc2c2e..f305c23 100644 --- a/packages/contrail/tests/search.test.ts +++ b/packages/contrail/tests/search.test.ts @@ -2,7 +2,7 @@ import { describe, it, expect, beforeEach } from "vitest"; import type { Database } from "../src/index"; import { resolveConfig } from "../src/index"; import { createTestDb, makeEvent } from "./helpers"; -import { initSchema } from "../src/index"; +import { initSchema, verifyBootstrapCandidate } from "../src/index"; import { ingestRecords, queryRecords } from "../src/index"; import { rebuildDerivedProjections } from "../src/core/db/records"; @@ -199,10 +199,9 @@ describe.skipIf(!hasFts)("FTS sync", () => { }); it("does not duplicate FTS rows when the same record is re-applied during backfill", async () => { - // Backfill passes call ingestRecords with skipReplayDetection: true, which leaves - // existingMap empty so every record looks brand-new. Re-backfilling the same - // record must not append a second FTS row; otherwise the search JOIN fans out - // and returns the event more than once (which crashes keyed lists downstream). + // The ordinary mapping table owns URI uniqueness even when a caller bypasses + // existing-record replay detection. Search can therefore never fan out one + // canonical URI through duplicate virtual-table rows. const event = makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", collection, @@ -215,6 +214,166 @@ describe.skipIf(!hasFts)("FTS sync", () => { const result = await queryRecords(db, SEARCH_CONFIG, { collection, search: "Meetup" }); expect(result.records).toHaveLength(1); + expect( + await db + .prepare( + "SELECT COUNT(*) AS count FROM fts_community_lexicon_calendar_event_rows WHERE uri = ?", + ) + .bind(event.uri) + .first<{ count: number }>(), + ).toEqual({ count: 1 }); + expect( + await db + .prepare("SELECT COUNT(*) AS count FROM fts_community_lexicon_calendar_event") + .first<{ count: number }>(), + ).toEqual({ count: 1 }); + }); + + it("maps only searchable records and supports absent delete plus recreate", async () => { + const uri = "at://did:plc:a/community.lexicon.calendar.event/lifecycle"; + await ingestRecords(db, [ + makeEvent({ + uri, + collection, + rkey: "lifecycle", + operation: "delete", + time_us: 500, + }), + ], SEARCH_CONFIG); + await ingestRecords(db, [ + makeEvent({ + uri, + collection, + rkey: "lifecycle", + record: { name: " ", startsAt: "2026-01-01T00:00:00Z" }, + time_us: 1000, + }), + ], SEARCH_CONFIG); + expect( + await db + .prepare( + "SELECT COUNT(*) AS count FROM fts_community_lexicon_calendar_event_rows WHERE uri = ?", + ) + .bind(uri) + .first<{ count: number }>(), + ).toEqual({ count: 0 }); + + await ingestRecords(db, [ + makeEvent({ + uri, + collection, + rkey: "lifecycle", + operation: "update", + record: { name: "MappedOnce" }, + time_us: 2000, + }), + ], SEARCH_CONFIG); + const firstMapping = await db + .prepare( + "SELECT id FROM fts_community_lexicon_calendar_event_rows WHERE uri = ?", + ) + .bind(uri) + .first<{ id: number }>(); + expect(firstMapping).toBeTruthy(); + + await ingestRecords(db, [ + makeEvent({ + uri, + collection, + rkey: "lifecycle", + operation: "delete", + time_us: 3000, + }), + ], SEARCH_CONFIG); + expect( + await db + .prepare( + "SELECT COUNT(*) AS count FROM fts_community_lexicon_calendar_event_rows WHERE uri = ?", + ) + .bind(uri) + .first<{ count: number }>(), + ).toEqual({ count: 0 }); + + await ingestRecords(db, [ + makeEvent({ + uri, + collection, + rkey: "lifecycle", + operation: "create", + record: { name: "MappedAgain" }, + time_us: 4000, + }), + ], SEARCH_CONFIG); + expect((await queryRecords(db, SEARCH_CONFIG, { + collection, + search: "MappedAgain", + })).records).toHaveLength(1); + expect( + await db + .prepare( + "SELECT COUNT(*) AS count FROM fts_community_lexicon_calendar_event_rows WHERE uri = ?", + ) + .bind(uri) + .first<{ count: number }>(), + ).toEqual({ count: 1 }); + }); + + it("normalizes tab, newline, and NBSP identically during writes and rebuilds", async () => { + const prefix = "at://did:plc:a/community.lexicon.calendar.event/whitespace-"; + await ingestRecords(db, [ + makeEvent({ + uri: `${prefix}tab`, + collection, + rkey: "whitespace-tab", + record: { name: "\t" }, + time_us: 1000, + }), + makeEvent({ + uri: `${prefix}newline`, + collection, + rkey: "whitespace-newline", + record: { name: "\n" }, + time_us: 1100, + }), + makeEvent({ + uri: `${prefix}nbsp`, + collection, + rkey: "whitespace-nbsp", + record: { name: "\u00a0" }, + time_us: 1200, + }), + makeEvent({ + uri: `${prefix}wrapped`, + collection, + rkey: "whitespace-wrapped", + record: { name: "\t\n\u00a0WhitespaceNeedle\u00a0\n" }, + time_us: 1300, + }), + ], SEARCH_CONFIG); + + const projected = async () => + db + .prepare( + `SELECT mapped.uri, fts.content + FROM fts_community_lexicon_calendar_event_rows mapped + JOIN fts_community_lexicon_calendar_event fts ON fts.rowid = mapped.id + WHERE mapped.uri LIKE ? + ORDER BY mapped.uri`, + ) + .bind(`${prefix}%`) + .all<{ uri: string; content: string }>(); + + expect(await projected()).toEqual({ + results: [{ uri: `${prefix}wrapped`, content: "WhitespaceNeedle" }], + }); + await rebuildDerivedProjections(db, SEARCH_CONFIG); + expect(await projected()).toEqual({ + results: [{ uri: `${prefix}wrapped`, content: "WhitespaceNeedle" }], + }); + expect((await queryRecords(db, SEARCH_CONFIG, { + collection, + search: "WhitespaceNeedle", + })).records).toHaveLength(1); }); it("evicts the stale FTS row when an update clears all searchable fields", async () => { @@ -226,10 +385,269 @@ describe.skipIf(!hasFts)("FTS sync", () => { ], SEARCH_CONFIG); expect((await queryRecords(db, SEARCH_CONFIG, { collection, search: "Searchable" })).records).toHaveLength(1); + const uri = "at://did:plc:a/community.lexicon.calendar.event/1"; await ingestRecords(db, [ - makeEvent({ uri: "at://did:plc:a/community.lexicon.calendar.event/1", collection, rkey: "1", record: { startsAt: "2026-01-01T00:00:00Z" }, operation: "update", time_us: 2000 }), + makeEvent({ uri, collection, rkey: "1", record: { startsAt: "2026-01-01T00:00:00Z" }, operation: "update", time_us: 2000 }), ], SEARCH_CONFIG); expect((await queryRecords(db, SEARCH_CONFIG, { collection, search: "Searchable" })).records).toHaveLength(0); + expect( + await db + .prepare( + "SELECT COUNT(*) AS count FROM fts_community_lexicon_calendar_event_rows WHERE uri = ?", + ) + .bind(uri) + .first<{ count: number }>(), + ).toEqual({ count: 0 }); + }); +}); + +describe.skipIf(!hasFts)("FTS rowid schema migration", () => { + it("rebuilds a duplicate legacy URI table from canonical records", async () => { + const legacy = createTestDb(); + const recordsTable = "records_community_lexicon_calendar_event"; + const ftsTable = "fts_community_lexicon_calendar_event"; + const rowsTable = `${ftsTable}_rows`; + const uri = "at://did:plc:legacy/community.lexicon.calendar.event/one"; + + await legacy + .prepare( + `CREATE TABLE ${recordsTable} ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + cid TEXT, + record TEXT, + time_us INTEGER NOT NULL, + indexed_at INTEGER NOT NULL + )`, + ) + .run(); + await legacy + .prepare( + `INSERT INTO ${recordsTable} + (uri, did, rkey, cid, record, time_us, indexed_at) + VALUES (?, ?, ?, ?, ?, ?, ?)`, + ) + .bind( + uri, + "did:plc:legacy", + "one", + "cid-legacy", + JSON.stringify({ name: "MigrationNeedle" }), + 1, + 1, + ) + .run(); + await legacy + .prepare( + `CREATE VIRTUAL TABLE ${ftsTable} USING fts5(uri UNINDEXED, content)`, + ) + .run(); + await legacy + .prepare(`INSERT INTO ${ftsTable} (uri, content) VALUES (?, ?), (?, ?)`) + .bind(uri, "stale", uri, "duplicate") + .run(); + + await initSchema(legacy, SEARCH_CONFIG); + + const columns = await legacy + .prepare(`PRAGMA table_info(${ftsTable})`) + .all<{ name: string }>(); + expect(columns.results.map(({ name }) => name)).toEqual(["content"]); + expect( + await legacy + .prepare(`SELECT uri FROM ${rowsTable}`) + .all<{ uri: string }>(), + ).toEqual({ results: [{ uri }] }); + expect( + await legacy + .prepare(`SELECT COUNT(*) AS count FROM ${ftsTable}`) + .first<{ count: number }>(), + ).toEqual({ count: 1 }); + expect((await queryRecords(legacy, SEARCH_CONFIG, { + collection: "community.lexicon.calendar.event", + search: "MigrationNeedle", + })).records).toHaveLength(1); + const verification = await verifyBootstrapCandidate(legacy, SEARCH_CONFIG); + const ftsChecks = verification.checks.filter( + ({ name }) => + name.includes("fts-") && + name.endsWith("community.lexicon.calendar.event"), + ); + expect(ftsChecks).toHaveLength(4); + expect(ftsChecks.every(({ ok }) => ok)).toBe(true); + + const plan = await legacy + .prepare( + `EXPLAIN QUERY PLAN DELETE FROM ${ftsTable} + WHERE rowid = (SELECT id FROM ${rowsTable} WHERE uri = ?)`, + ) + .bind(uri) + .all<{ detail: string }>(); + expect(plan.results.some(({ detail }) => + detail.includes(rowsTable) && detail.includes("uri=?") + )).toBe(true); + expect(plan.results.some(({ detail }) => + detail === `SCAN ${ftsTable} VIRTUAL TABLE INDEX 0:` + )).toBe(false); + expect(plan.results.some(({ detail }) => + detail === `SCAN ${ftsTable} VIRTUAL TABLE INDEX 0:=` + )).toBe(true); + + await initSchema(legacy, SEARCH_CONFIG); + expect( + await legacy + .prepare(`SELECT COUNT(*) AS count FROM ${ftsTable}`) + .first<{ count: number }>(), + ).toEqual({ count: 1 }); + }); + + it("does not accept a stale content-only schema until rebuilding verifies", async () => { + const stale = createTestDb(); + const migrationConfig = resolveConfig({ + namespace: "com.example", + profiles: [], + collections: { + event: { + collection: "community.lexicon.calendar.event", + searchable: ["name"], + }, + }, + }); + const uri = "at://did:plc:legacy/community.lexicon.calendar.event/retry"; + await stale + .prepare("CREATE TABLE _contrail_meta (key TEXT PRIMARY KEY, value TEXT NOT NULL)") + .run(); + await stale + .prepare("INSERT INTO _contrail_meta (key, value) VALUES ('schema_fingerprint', 'stale')") + .run(); + await stale + .prepare( + `CREATE TABLE records_event ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + cid TEXT, + record TEXT, + time_us INTEGER NOT NULL, + indexed_at INTEGER NOT NULL + )`, + ) + .run(); + await stale + .prepare( + `INSERT INTO records_event + (uri, did, rkey, cid, record, time_us, indexed_at) + VALUES (?, 'did:plc:legacy', 'retry', 'cid-retry', 'not-json', 1, 1)`, + ) + .bind(uri) + .run(); + await stale + .prepare( + "CREATE TABLE fts_event_rows (id INTEGER PRIMARY KEY, uri TEXT NOT NULL UNIQUE)", + ) + .run(); + await stale + .prepare("INSERT INTO fts_event_rows (uri) VALUES ('at://stale')") + .run(); + await stale + .prepare("CREATE VIRTUAL TABLE fts_event USING fts5(content)") + .run(); + await stale + .prepare("INSERT INTO fts_event (rowid, content) VALUES (1, 'stale')") + .run(); + + await expect(initSchema(stale, migrationConfig)).rejects.toThrow(); + expect( + await stale + .prepare("SELECT value FROM _contrail_meta WHERE key = 'schema_fingerprint'") + .first<{ value: string }>(), + ).toEqual({ value: "stale" }); + expect( + await stale.prepare("SELECT uri FROM fts_event_rows").all<{ uri: string }>(), + ).toEqual({ results: [{ uri: "at://stale" }] }); + + await stale + .prepare("UPDATE records_event SET record = ? WHERE uri = ?") + .bind(JSON.stringify({ name: "RecoveryNeedle" }), uri) + .run(); + await initSchema(stale, migrationConfig); + + expect( + await stale.prepare("SELECT uri FROM fts_event_rows").all<{ uri: string }>(), + ).toEqual({ results: [{ uri }] }); + expect( + await stale.prepare("SELECT content FROM fts_event").all<{ content: string }>(), + ).toEqual({ results: [{ content: "RecoveryNeedle" }] }); + expect( + await stale + .prepare("SELECT value FROM _contrail_meta WHERE key = 'schema_fingerprint'") + .first<{ value: string }>(), + ).not.toEqual({ value: "stale" }); + }); + + it("rolls the legacy table back when canonical rebuilding fails", async () => { + const legacy = createTestDb(); + const migrationConfig = resolveConfig({ + namespace: "com.example", + profiles: [], + collections: { + event: { + collection: "community.lexicon.calendar.event", + searchable: ["name"], + }, + }, + }); + await legacy + .prepare( + `CREATE TABLE records_event ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + cid TEXT, + record TEXT, + time_us INTEGER NOT NULL, + indexed_at INTEGER NOT NULL + )`, + ) + .run(); + await legacy + .prepare( + `INSERT INTO records_event + (uri, did, rkey, cid, record, time_us, indexed_at) + VALUES ('at://did:plc:legacy/community.lexicon.calendar.event/bad', + 'did:plc:legacy', 'bad', 'cid-bad', 'not-json', 1, 1)`, + ) + .run(); + await legacy + .prepare( + "CREATE VIRTUAL TABLE fts_event USING fts5(uri UNINDEXED, content)", + ) + .run(); + await legacy + .prepare( + "INSERT INTO fts_event (uri, content) VALUES ('at://legacy', 'keep-old')", + ) + .run(); + + await expect(initSchema(legacy, migrationConfig)).rejects.toThrow(); + + const columns = await legacy + .prepare("PRAGMA table_info(fts_event)") + .all<{ name: string }>(); + expect(columns.results.map(({ name }) => name)).toEqual(["uri", "content"]); + expect( + await legacy + .prepare("SELECT uri, content FROM fts_event") + .all<{ uri: string; content: string }>(), + ).toEqual({ results: [{ uri: "at://legacy", content: "keep-old" }] }); + expect( + await legacy + .prepare( + "SELECT COUNT(*) AS count FROM sqlite_master WHERE type = 'table' AND name = 'fts_event_rows'", + ) + .first<{ count: number }>(), + ).toEqual({ count: 0 }); }); });