diff --git a/apps/sim/lib/knowledge/__integration__/embedding-insert-batches.integration.ts b/apps/sim/lib/knowledge/__integration__/embedding-insert-batches.integration.ts index 0761840834d..bffc2212d1e 100644 --- a/apps/sim/lib/knowledge/__integration__/embedding-insert-batches.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/embedding-insert-batches.integration.ts @@ -49,6 +49,11 @@ describe('bounded embedding insert transactions', () => { beforeAll(async () => { fixtures.root = mkdtempSync(path.join(tmpdir(), 'sim-embedding-batches-')) await seedKnowledgeAclFixture(ids, { connectorType: 'google_drive' }) + /** A search index, the only kind of base whose chunks the keyword projection holds. */ + await db + .update(knowledgeBase) + .set({ isSearchIndex: true }) + .where(eq(knowledgeBase.id, ids.knowledgeBaseId)) vi.spyOn(embeddingClient, 'assertKnowledgeEmbeddingCapacity').mockResolvedValue(undefined) }) diff --git a/apps/sim/lib/knowledge/__integration__/knowledge-projection.integration.ts b/apps/sim/lib/knowledge/__integration__/knowledge-projection.integration.ts index 37f9abf9a5e..56ff0027e65 100644 --- a/apps/sim/lib/knowledge/__integration__/knowledge-projection.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/knowledge-projection.integration.ts @@ -24,6 +24,7 @@ import { credentialGroupEnrollment, document, embedding, + embeddingKeywordSearch, embeddingKeywordTin, embeddingSearch, knowledgeBase, @@ -383,7 +384,6 @@ beforeAll(async () => { sql.raw(`CREATE SCHEMA tin; CREATE FUNCTION tin.full_score(tid) RETURNS double precision LANGUAGE sql IMMUTABLE AS 'SELECT 1.0::float8'; CREATE FUNCTION knowledge_tin_base_token(text) RETURNS text LANGUAGE sql IMMUTABLE AS $$SELECT 'kb'$$; - CREATE FUNCTION knowledge_tin_membership_key(text) RETURNS bigint LANGUAGE sql IMMUTABLE AS $$SELECT hashtextextended('embedding_keyword_tin:' || $1, 0)$$; CREATE FUNCTION knowledge_tin_stream(vector tsvector) RETURNS text LANGUAGE sql IMMUTABLE AS $$ SELECT coalesce(string_agg(entry.lexeme, ' ' ORDER BY position), '') FROM unnest(vector) AS entry(lexeme, positions, weights), unnest(entry.positions) AS position @@ -400,7 +400,6 @@ afterAll(async () => { sql.raw(`DROP OPERATOR IF EXISTS ==> (text, text); DROP FUNCTION IF EXISTS tin_fixture_match(text, text); DROP FUNCTION IF EXISTS knowledge_tin_base_token(text); - DROP FUNCTION IF EXISTS knowledge_tin_membership_key(text); DROP FUNCTION IF EXISTS knowledge_tin_stream(tsvector); DROP SCHEMA IF EXISTS tin CASCADE;`) ) @@ -730,6 +729,46 @@ describe('the projector', () => { expect(await markOf()).toBeUndefined() }) + it('writes keyword rows for search-index knowledge bases only', async () => { + const workspaceBaseId = generateId() + const workspaceDocument = generateId() + const workspaceChunk = generateId() + await db.insert(knowledgeBase).values({ + id: workspaceBaseId, + userId: ids.aliceId, + workspaceId: ids.workspaceId, + name: 'Workspace keyword fixture', + chunkingConfig: { maxSize: 1024, minSize: 1, overlap: 20 }, + }) + try { + await db.insert(document).values({ + id: workspaceDocument, + knowledgeBaseId: workspaceBaseId, + filename: 'keyword.md', + fileUrl: 'https://fixture.test/keyword', + fileSize: 12, + mimeType: 'text/plain', + processingStatus: 'completed', + }) + await write('async', (tx) => + tx.insert(embedding).values({ + ...chunkRow(workspaceChunk, 0), + documentId: workspaceDocument, + knowledgeBaseId: workspaceBaseId, + }) + ) + await project() + const keywordRows = await db + .select({ id: embeddingKeywordSearch.id }) + .from(embeddingKeywordSearch) + .where(inArray(embeddingKeywordSearch.id, [chunkId, workspaceChunk])) + expect(keywordRows).toEqual([{ id: chunkId }]) + expect(await rowAcl(embeddingSearch, workspaceChunk)).toBeDefined() + } finally { + await db.delete(knowledgeBase).where(eq(knowledgeBase.id, workspaceBaseId)) + } + }) + it.each(['sync', 'async'] as const)( 'writes %s projection rows from a chunk commit only when the writer did not defer them', async (mode) => { diff --git a/apps/sim/lib/knowledge/__integration__/processing-lock-scope.integration.ts b/apps/sim/lib/knowledge/__integration__/processing-lock-scope.integration.ts index 0fe2f9d88e1..5c39be414e8 100644 --- a/apps/sim/lib/knowledge/__integration__/processing-lock-scope.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/processing-lock-scope.integration.ts @@ -54,6 +54,11 @@ describe('document processing commit lock scope', () => { beforeAll(async () => { fixtures.root = mkdtempSync(path.join(tmpdir(), 'sim-processing-lock-scope-')) await seedKnowledgeAclFixture(ids, { connectorType: 'google_drive' }) + /** A search index, the only kind of base whose chunks the keyword projection holds. */ + await db + .update(knowledgeBase) + .set({ isSearchIndex: true }) + .where(eq(knowledgeBase.id, ids.knowledgeBaseId)) vi.spyOn(embeddingClient, 'assertKnowledgeEmbeddingCapacity').mockResolvedValue(undefined) fixtures.process.mockResolvedValue({ chunks, @@ -152,6 +157,14 @@ describe('document processing commit lock scope', () => { expect(await db.select().from(document).where(eq(document.id, file.documentId))).toMatchObject([ { processingStatus: 'completed', chunkCount: 3 }, ]) + for (const projection of ['embedding_search', 'embedding_keyword_search']) { + expect( + await db.$client.unsafe( + `SELECT count(*)::int AS count FROM ${projection} WHERE document_id = $1`, + [file.documentId] + ) + ).toEqual([{ count: 3 }]) + } }) it.each([ diff --git a/apps/sim/lib/sim-search/indexed/README.md b/apps/sim/lib/sim-search/indexed/README.md index de0358f05cc..a392625fb7f 100644 --- a/apps/sim/lib/sim-search/indexed/README.md +++ b/apps/sim/lib/sim-search/indexed/README.md @@ -14,7 +14,7 @@ While the gate is off: - Indexed-only surfaces refuse with `SearchIndexDormantError` (a `409`): the Stats report and connecting a source that crawls into a search index. The indexed document page is not found. - A knowledge search that names a search-index knowledge base (the Knowledge block, v1, v2, Sim's knowledge tool) still answers from the documents it already holds, decided on each document exactly as a workspace knowledge base is. - Nothing crawls into search indexes: content syncs, member syncs, and processing recovery skip them (`lib/knowledge/connectors/indexing-policy.ts`). -- The projector owes search-index documents nothing: their marks are released with the rest, and it writes no Tin keyword rows. +- The projector owes search-index documents nothing: their marks are released with the rest, and it writes no Tin keyword rows. The GIN keyword projection follows `is_search_index` alone, so it keeps its search-index rows either way. ## Layout @@ -28,7 +28,7 @@ The dormant UI sits in `indexed/` folders next to the component that picks it fr ## Re-enabling 1. Set `SIM_SEARCH_LIVE=false` in both the app and the Trigger.dev environment, and deploy. The container entrypoint (`apps/sim/bootstrap.ts`) mirrors it to `NEXT_PUBLIC_SIM_SEARCH_LIVE` for the client; crawling, processing, and projection read it in whichever process runs them. -2. Confirm the Tin objects exist (`0019_tin_keyword_projection`, `0024_knowledge_projection_async`), backfill `embedding_keyword_tin` for every search-index knowledge base, and build its index. +2. Confirm the keyword projection objects exist (`0019_tin_keyword_projection`, `0024_knowledge_projection_async`, `0025_scope_keyword_projections`). Both keyword projections, `embedding_keyword_search` and `embedding_keyword_tin`, hold only search-index rows, written by the chunk triggers and by the trigger on `knowledge_base.is_search_index`. Backfill both for every search-index knowledge base whose rows were removed while dormant, and build the Tin index. 3. Resume and fully resync the connectors of search-index knowledge bases, so content that went stale while dormant is indexed again. Projection rows written before projections carried their document's source and ACL are decided on their document until they are rewritten. diff --git a/packages/db/knowledge-projection.ts b/packages/db/knowledge-projection.ts index dc8cfcb6398..fa55824873f 100644 --- a/packages/db/knowledge-projection.ts +++ b/packages/db/knowledge-projection.ts @@ -106,6 +106,10 @@ export type SourceAclProjection = (typeof SOURCE_ACL_PROJECTIONS)[number] const mirrorsSourceAcl = (projection: KnowledgeProjection): projection is SourceAclProjection => (SOURCE_ACL_PROJECTIONS as readonly string[]).includes(projection) +/** The keyword projections, which hold only the rows of search-index knowledge bases. */ +const holdsSearchIndexesOnly = (projection: KnowledgeProjection) => + projection === 'embedding_keyword_search' || projection === 'embedding_keyword_tin' + /** * The chunks one page covers: the next {@link PROJECTION_ROW_BATCH_SIZE} of the document in * chunk order, read off `emb_doc_chunk_idx`. Every projection pages the same way, so a page is @@ -163,10 +167,16 @@ function contentPageStatement(projection: KnowledgeProjection): string { } if (projection === 'embedding_keyword_search') { const compared = ['knowledge_base_id', 'document_id', 'enabled', 'content_tsv'] - return `WITH ${PAGE}, written AS ( - INSERT INTO embedding_keyword_search AS s (${['id', ...compared].join(', ')}) - SELECT e.id, e.knowledge_base_id, e.document_id, e.enabled, e.content_tsv + return `WITH ${PAGE}, source AS MATERIALIZED ( + SELECT e.id, ${compared.map((column) => `e.${column}`).join(', ')}, k.is_search_index FROM page p JOIN embedding e ON e.id = p.id + JOIN knowledge_base k ON k.id = e.knowledge_base_id + ), removed AS ( + DELETE FROM embedding_keyword_search s USING source + WHERE s.id = source.id AND NOT source.is_search_index + ), written AS ( + INSERT INTO embedding_keyword_search AS s (${['id', ...compared].join(', ')}) + SELECT id, ${compared.join(', ')} FROM source WHERE is_search_index ON CONFLICT (id) DO UPDATE SET ${compared.map((column) => `${column} = EXCLUDED.${column}`).join(', ')} WHERE (${compared.map((column) => `s.${column}`).join(', ')}) @@ -249,8 +259,9 @@ export interface KnowledgeProjectionOptions { /** * Whether search-index rows are read, that is whether indexed organization search is on (see * {@link MarkScope}). Only then is the Tin keyword projection written, since it holds only - * search-index rows; the other projections are written either way, and Tin is still skipped - * where it is not installed. + * search-index rows; the other projections are written either way, the GIN keyword projection + * for search-index bases alone whatever this says, and Tin is still skipped where it is not + * installed. */ searchIndexes: boolean /** Called after each page commits, for tests that interleave writes with a run. */ @@ -375,8 +386,8 @@ async function projectDocumentRows( if (Date.now() >= deadline) return { pages, written, finished: false } const page = await sql.begin(async (tx) => { await enterProjectorTransaction(tx, PROJECTION_PAGE_LOCK_TIMEOUT_MS) - /** Shares the base's Tin membership lock, as the embedding trigger does, so a flip of its marker waits. */ - if (projection === 'embedding_keyword_tin') + /** Shares the base's membership lock, as the embedding triggers do, so a flip of its marker waits. */ + if (holdsSearchIndexesOnly(projection)) await tx`SELECT pg_advisory_xact_lock_shared(knowledge_tin_membership_key(${mark.knowledgeBaseId}))` const [row] = await tx.unsafe< Array<{ scanned: number; written: number; last_chunk: number | null }> diff --git a/packages/db/script-migrations-paused-billing-attribution.test.ts b/packages/db/script-migrations-paused-billing-attribution.test.ts index 6b067c5f2b7..f7901d6c65a 100644 --- a/packages/db/script-migrations-paused-billing-attribution.test.ts +++ b/packages/db/script-migrations-paused-billing-attribution.test.ts @@ -377,6 +377,7 @@ describe('script migration registry', () => { '0022_projection_source_acl_backfill', '0023_projection_acl_skip_unfilled', '0024_knowledge_projection_async', + '0025_scope_keyword_projections', ]) }) }) diff --git a/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts b/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts index 9ee2f81c346..d9972d14d3d 100644 --- a/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts +++ b/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts @@ -395,6 +395,7 @@ describe('search projection upgrade in PostgreSQL', () => { { name: '0022_projection_source_acl_backfill' }, { name: '0023_projection_acl_skip_unfilled' }, { name: '0024_knowledge_projection_async' }, + { name: '0025_scope_keyword_projections' }, ]) const [{ complete }] = await sql`SELECT count(*)::int AS complete FROM embedding e JOIN embedding_search s ON s.id = e.id JOIN embedding_keyword_search k ON k.id = e.id diff --git a/packages/db/script-migrations/0019_tin_keyword_projection.ts b/packages/db/script-migrations/0019_tin_keyword_projection.ts index 5bcd6ac7eb7..99b46ca0a26 100644 --- a/packages/db/script-migrations/0019_tin_keyword_projection.ts +++ b/packages/db/script-migrations/0019_tin_keyword_projection.ts @@ -1,8 +1,9 @@ +import { SYNCHRONOUS_PROJECTION_WHEN } from '@sim/db/knowledge-projection' import { EMBEDDING_KEYWORD_TIN_INDEX } from '@sim/db/schema' import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url' import { type ScriptMigration, ScriptMigrationDeferred } from '@sim/db/script-migrations/types' import { createLogger } from '@sim/logger' -import postgres, { type Sql } from 'postgres' +import postgres, { type Sql, type TransactionSql } from 'postgres' const logger = createLogger('TinKeywordProjection') const BATCH_SIZE = 500 @@ -39,6 +40,75 @@ async function tinAvailable(sql: Sql): Promise { return rows.length > 0 } +/** + * The advisory lock key that serializes a knowledge base's keyword projection writers with a change + * of its search-index marker (see {@link installProjection}). + */ +export async function installMembershipKey(tx: Sql | TransactionSql): Promise { + await tx.unsafe(`CREATE OR REPLACE FUNCTION knowledge_tin_membership_key(knowledge_base_id text) + RETURNS bigint LANGUAGE sql IMMUTABLE PARALLEL SAFE AS $$ + SELECT hashtextextended('embedding_keyword_tin:' || knowledge_base_id, 0) + $$`) +} + +/** + * The chunk trigger's body. A chunk outside a search index writes nothing; an update removes the + * row a chunk left behind by moving out of one. An insert has nothing to remove, since a new chunk + * id has no row and a base adopted meanwhile waits on the membership lock the insert holds. + */ +export async function installTinChunkSync(tx: Sql | TransactionSql): Promise { + await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_embedding_keyword_tin() + RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + PERFORM pg_advisory_xact_lock_shared(knowledge_tin_membership_key(NEW.knowledge_base_id)); + IF NOT EXISTS ( + SELECT 1 FROM knowledge_base WHERE id = NEW.knowledge_base_id AND is_search_index + ) THEN + IF TG_OP = 'UPDATE' THEN + DELETE FROM embedding_keyword_tin WHERE id = NEW.id; + END IF; + RETURN NEW; + END IF; + INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content) + VALUES (NEW.id, NEW.knowledge_base_id, NEW.document_id, NEW.enabled, + knowledge_tin_base_token(NEW.knowledge_base_id) || ' ' || knowledge_tin_stream(NEW.content_tsv)) + ON CONFLICT (id) DO UPDATE SET + knowledge_base_id = EXCLUDED.knowledge_base_id, document_id = EXCLUDED.document_id, + enabled = EXCLUDED.enabled, content = EXCLUDED.content; + RETURN NEW; + END; + $$`) +} + +/** + * The knowledge base trigger's body: projects or removes a whole base when its search-index marker + * changes. The chunks are key-share locked as they are read, as the backfill locks its own, so a + * chunk delete in flight is waited for and its chunk skipped rather than failing the adoption on + * the row's foreign key. + */ +export async function installTinMembershipSync(tx: Sql | TransactionSql): Promise { + await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_knowledge_base_keyword_tin() + RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + PERFORM pg_advisory_xact_lock(knowledge_tin_membership_key(NEW.id)); + IF NOT NEW.is_search_index THEN + DELETE FROM embedding_keyword_tin t USING embedding e + WHERE e.knowledge_base_id = NEW.id AND t.id = e.id; + RETURN NEW; + END IF; + INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content) + SELECT id, knowledge_base_id, document_id, enabled, + knowledge_tin_base_token(knowledge_base_id) || ' ' || knowledge_tin_stream(content_tsv) + FROM embedding WHERE knowledge_base_id = NEW.id + ORDER BY id FOR KEY SHARE + ON CONFLICT (id) DO UPDATE SET + knowledge_base_id = EXCLUDED.knowledge_base_id, document_id = EXCLUDED.document_id, + enabled = EXCLUDED.enabled, content = EXCLUDED.content; + RETURN NEW; + END; + $$`) +} + /** * Installs the stream and base-token functions shared by the triggers and the query path, and the * embedding and knowledge base triggers, atomically with respect to embedding writers. @@ -60,7 +130,9 @@ async function tinAvailable(sql: Sql): Promise { * other's rows. A row lock cannot do this: a key-share lock on the marker row is compatible with * the uncommitted marker update, so a writer would read the old marker without waiting. A whole * base is reached through `embedding`'s knowledge base index, since a projected row always carries - * its chunk's base and the projection keeps no index but Tin's. + * its chunk's base and the projection keeps no index but Tin's. The embedding trigger carries the + * guard `0024_knowledge_projection_async` added, so a database adopting Tin later ends where a full + * migration run does. */ export async function installProjection(sql: Sql): Promise { await sql.begin(async (tx) => { @@ -75,51 +147,12 @@ export async function installProjection(sql: Sql): Promise { RETURNS text LANGUAGE sql IMMUTABLE PARALLEL SAFE AS $$ SELECT 'zkb' || md5(knowledge_base_id) $$`) - await tx.unsafe(`CREATE OR REPLACE FUNCTION knowledge_tin_membership_key(knowledge_base_id text) - RETURNS bigint LANGUAGE sql IMMUTABLE PARALLEL SAFE AS $$ - SELECT hashtextextended('embedding_keyword_tin:' || knowledge_base_id, 0) - $$`) - await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_embedding_keyword_tin() - RETURNS trigger LANGUAGE plpgsql AS $$ - BEGIN - PERFORM pg_advisory_xact_lock_shared(knowledge_tin_membership_key(NEW.knowledge_base_id)); - IF NOT EXISTS ( - SELECT 1 FROM knowledge_base WHERE id = NEW.knowledge_base_id AND is_search_index - ) THEN - DELETE FROM embedding_keyword_tin WHERE id = NEW.id; - RETURN NEW; - END IF; - INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content) - VALUES (NEW.id, NEW.knowledge_base_id, NEW.document_id, NEW.enabled, - knowledge_tin_base_token(NEW.knowledge_base_id) || ' ' || knowledge_tin_stream(NEW.content_tsv)) - ON CONFLICT (id) DO UPDATE SET - knowledge_base_id = EXCLUDED.knowledge_base_id, document_id = EXCLUDED.document_id, - enabled = EXCLUDED.enabled, content = EXCLUDED.content; - RETURN NEW; - END; - $$`) + await installMembershipKey(tx) + await installTinChunkSync(tx) await tx.unsafe(`CREATE OR REPLACE TRIGGER embedding_keyword_tin_sync AFTER INSERT OR UPDATE OF knowledge_base_id, document_id, enabled, content ON embedding - FOR EACH ROW EXECUTE FUNCTION sync_embedding_keyword_tin()`) - await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_knowledge_base_keyword_tin() - RETURNS trigger LANGUAGE plpgsql AS $$ - BEGIN - PERFORM pg_advisory_xact_lock(knowledge_tin_membership_key(NEW.id)); - IF NOT NEW.is_search_index THEN - DELETE FROM embedding_keyword_tin t USING embedding e - WHERE e.knowledge_base_id = NEW.id AND t.id = e.id; - RETURN NEW; - END IF; - INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content) - SELECT id, knowledge_base_id, document_id, enabled, - knowledge_tin_base_token(knowledge_base_id) || ' ' || knowledge_tin_stream(content_tsv) - FROM embedding WHERE knowledge_base_id = NEW.id - ON CONFLICT (id) DO UPDATE SET - knowledge_base_id = EXCLUDED.knowledge_base_id, document_id = EXCLUDED.document_id, - enabled = EXCLUDED.enabled, content = EXCLUDED.content; - RETURN NEW; - END; - $$`) + FOR EACH ROW WHEN (${SYNCHRONOUS_PROJECTION_WHEN}) EXECUTE FUNCTION sync_embedding_keyword_tin()`) + await installTinMembershipSync(tx) await tx.unsafe(`CREATE OR REPLACE TRIGGER knowledge_base_keyword_tin_sync AFTER UPDATE OF is_search_index ON knowledge_base FOR EACH ROW WHEN (OLD.is_search_index IS DISTINCT FROM NEW.is_search_index) diff --git a/packages/db/script-migrations/0025_scope_keyword_projections.integration.ts b/packages/db/script-migrations/0025_scope_keyword_projections.integration.ts new file mode 100644 index 00000000000..2b3e738c1d6 --- /dev/null +++ b/packages/db/script-migrations/0025_scope_keyword_projections.integration.ts @@ -0,0 +1,310 @@ +import { backfillSearchKeywords } from '@sim/db/script-migrations/0016_backfill_search_vectors' +import { installProjection } from '@sim/db/script-migrations/0019_tin_keyword_projection' +import { installProjectionSourceAcl } from '@sim/db/script-migrations/0021_embedding_search_connector' +import { installKnowledgeProjectionAsync } from '@sim/db/script-migrations/0024_knowledge_projection_async' +import { scopeKeywordProjections } from '@sim/db/script-migrations/0025_scope_keyword_projections' +import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure' +import { generateId } from '@sim/utils/id' +import postgres, { type Sql, type TransactionSql } from 'postgres' +import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest' + +const databaseUrl = readTestDatabaseUrl() + +/** Resolves once backend `pid` waits on a lock, so each race runs in a fixed order. */ +async function waitUntilBlocked(observer: Sql, pid: number): Promise { + for (let attempt = 0; attempt < 200; attempt++) { + const [row] = await observer`SELECT wait_event_type FROM pg_stat_activity WHERE pid = ${pid}` + if (row?.wait_event_type === 'Lock') return + await observer`SELECT pg_sleep(0.01)` + } + throw new Error(`Backend ${pid} never waited on a lock`) +} + +/** Scans this transaction has run on `tables`, sequential and index alike. */ +async function scansOf(tx: TransactionSql, tables: readonly string[]): Promise { + const [row] = await tx>` + SELECT coalesce(sum(seq_scan + coalesce(idx_scan, 0)), 0)::int AS scans + FROM pg_stat_xact_user_tables WHERE relid = ANY(${tables as string[]}::regclass[])` + return row.scans +} + +/** + * The tables carry only the columns the triggers read and write. The Tin extension is not + * available here, but its triggers are plain SQL, so each projection is installed as the + * migrations before `0025` leave it and then scoped. + */ +describe('keyword projections scoped to search indexes in PostgreSQL', () => { + let admin: Sql + let sql: Sql + let other: Sql + const schemaName = `keyword_scope_${generateId().replaceAll('-', '')}` + + const keywordIds = () => + sql<{ id: string }[]>`SELECT id FROM embedding_keyword_search ORDER BY id`.then((rows) => + rows.map((row) => row.id) + ) + const tinIds = () => + sql<{ id: string }[]>`SELECT id FROM embedding_keyword_tin ORDER BY id`.then((rows) => + rows.map((row) => row.id) + ) + + beforeAll(async () => { + admin = postgres(databaseUrl, { max: 1, onnotice: () => undefined }) + await admin.unsafe(`CREATE SCHEMA "${schemaName}"`) + const connect = () => + postgres(databaseUrl, { + max: 1, + onnotice: () => undefined, + connection: { search_path: schemaName }, + }) + sql = connect() + other = connect() + await sql`CREATE TABLE knowledge_base ( + id text PRIMARY KEY, is_search_index boolean NOT NULL DEFAULT false + )` + await sql`CREATE TABLE document ( + id text PRIMARY KEY, connector_id text, acl text[] NOT NULL DEFAULT '{ws}' + )` + await sql`CREATE TABLE knowledge_projection_dirty ( + document_id text PRIMARY KEY REFERENCES document (id) ON DELETE CASCADE, + generation bigint NOT NULL DEFAULT 1, content boolean NOT NULL DEFAULT false, + marked_at timestamptz NOT NULL DEFAULT now() + )` + await sql`CREATE TABLE embedding ( + id text PRIMARY KEY, knowledge_base_id text NOT NULL REFERENCES knowledge_base (id), + document_id text NOT NULL REFERENCES document (id), enabled boolean NOT NULL DEFAULT true, + content text, content_tsv tsvector NOT NULL, embedding text, embedding_384 text, + embedding_768 text, embedding_1024 text, embedding_3072 text + )` + await sql`CREATE INDEX ON embedding (knowledge_base_id)` + await sql`CREATE TABLE embedding_search ( + id text PRIMARY KEY REFERENCES embedding (id) ON DELETE CASCADE, document_id text NOT NULL, + enabled boolean NOT NULL DEFAULT true, connector_id text, acl text[] + )` + await sql`CREATE TABLE embedding_keyword_search ( + id text PRIMARY KEY REFERENCES embedding (id) ON DELETE CASCADE, + knowledge_base_id text NOT NULL, document_id text NOT NULL, enabled boolean NOT NULL, + content_tsv tsvector NOT NULL + )` + await sql`CREATE INDEX ON embedding_keyword_search (knowledge_base_id)` + await sql`CREATE TABLE embedding_keyword_tin ( + id text PRIMARY KEY REFERENCES embedding (id) ON DELETE CASCADE, + knowledge_base_id text NOT NULL, document_id text NOT NULL, enabled boolean NOT NULL, + content text NOT NULL, connector_id text, acl text[] + )` + await sql.unsafe(`CREATE FUNCTION sync_embedding_search() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN RETURN NULL; END; $$`) + await backfillSearchKeywords(sql) + await installProjection(sql) + await installProjectionSourceAcl(sql) + await installKnowledgeProjectionAsync(sql) + await scopeKeywordProjections(sql) + }, 60_000) + + afterAll(async () => { + await sql?.end() + await other?.end() + await admin?.unsafe(`DROP SCHEMA IF EXISTS "${schemaName}" CASCADE`) + await admin?.end() + }) + + beforeEach(async () => { + await sql`TRUNCATE embedding_keyword_search, embedding_keyword_tin, embedding_search, + knowledge_projection_dirty, embedding, document, knowledge_base` + await sql`INSERT INTO knowledge_base (id, is_search_index) VALUES ('legacy', false), ('index', true)` + await sql`INSERT INTO document (id) VALUES ('doc')` + }) + + it('projects keyword rows only for chunks of search indexes', async () => { + await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) VALUES + ('in-index', 'index', 'doc', to_tsvector('english', 'Release notes')), + ('in-legacy', 'legacy', 'doc', to_tsvector('english', 'Release notes'))` + expect(await keywordIds()).toEqual(['in-index']) + expect(await tinIds()).toEqual(['in-index']) + }) + + it('never reads the Tin projection for a chunk inserted outside a search index', async () => { + const scans = await sql.begin(async (tx) => { + const before = await scansOf(tx, ['embedding_keyword_tin']) + await tx`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) + VALUES ('in-legacy', 'legacy', 'doc', to_tsvector('english', 'Release notes'))` + return (await scansOf(tx, ['embedding_keyword_tin'])) - before + }) + expect(scans).toBe(0) + }) + + it('removes the keyword rows of a chunk that moves out of a search index', async () => { + await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) + VALUES ('moving', 'index', 'doc', to_tsvector('english', 'Moving chunk'))` + await sql`UPDATE embedding SET knowledge_base_id = 'legacy' WHERE id = 'moving'` + expect(await keywordIds()).toEqual([]) + expect(await tinIds()).toEqual([]) + }) + + it('keeps the scoped and guarded Tin trigger when Tin is adopted after this migration', async () => { + await installProjection(sql) + const [trigger] = await sql<{ definition: string }[]>` + SELECT pg_get_triggerdef(oid) AS definition FROM pg_trigger + WHERE tgname = 'embedding_keyword_tin_sync' AND tgrelid = 'embedding'::regclass` + expect(trigger?.definition).toContain('sim.projection_mode') + const scans = await sql.begin(async (tx) => { + const before = await scansOf(tx, ['embedding_keyword_tin']) + await tx`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) + VALUES ('in-legacy', 'legacy', 'doc', to_tsvector('english', 'Release notes'))` + return (await scansOf(tx, ['embedding_keyword_tin'])) - before + }) + expect(scans).toBe(0) + }) + + it('projects a base adopted as a search index, and removes it when it is not', async () => { + await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, enabled, content_tsv) VALUES + ('legacy-1', 'legacy', 'doc', true, to_tsvector('english', 'Quarterly planning')), + ('legacy-2', 'legacy', 'doc', false, to_tsvector('english', 'Draft agenda'))` + await sql`UPDATE knowledge_base SET is_search_index = true WHERE id = 'legacy'` + expect(await sql`SELECT id, enabled FROM embedding_keyword_search ORDER BY id`).toMatchObject([ + { id: 'legacy-1', enabled: true }, + { id: 'legacy-2', enabled: false }, + ]) + await sql`UPDATE knowledge_base SET is_search_index = false WHERE id = 'legacy'` + expect(await keywordIds()).toEqual([]) + }) + + it('leaves a current keyword row unwritten when its base is adopted', async () => { + await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) + VALUES ('kept', 'legacy', 'doc', to_tsvector('english', 'Kept chunk'))` + /** A row written before the projection was scoped, already current. */ + await sql`INSERT INTO embedding_keyword_search (id, knowledge_base_id, document_id, enabled, content_tsv) + SELECT id, knowledge_base_id, document_id, enabled, content_tsv FROM embedding WHERE id = 'kept'` + const updated = await sql.begin(async (tx) => { + const read = async () => { + const [row] = await tx>` + SELECT n_tup_upd::int AS updated FROM pg_stat_xact_user_tables + WHERE relid = 'embedding_keyword_search'::regclass` + return row.updated + } + const before = await read() + await tx`UPDATE knowledge_base SET is_search_index = true WHERE id = 'legacy'` + return (await read()) - before + }) + expect(updated).toBe(0) + expect(await keywordIds()).toEqual(['kept']) + }) + + it('projects a chunk whose insert commits while its base is being adopted', async () => { + let inserted!: () => void + const insertedSignal = new Promise((resolve) => { + inserted = resolve + }) + let release!: () => void + const released = new Promise((resolve) => { + release = resolve + }) + const writer = sql.begin(async (tx) => { + await tx`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) + VALUES ('racing', 'legacy', 'doc', to_tsvector('english', 'Racing chunk'))` + inserted() + await released + }) + await insertedSignal + const [{ pid }] = await other`SELECT pg_backend_pid() AS pid` + const adoption = + other`UPDATE knowledge_base SET is_search_index = true WHERE id = 'legacy'`.execute() + await waitUntilBlocked(admin, pid) + release() + await writer + await adoption + expect(await keywordIds()).toEqual(['racing']) + }) + + it('projects a chunk inserted while the adoption of its base is uncommitted', async () => { + let adopted!: () => void + const adoptedSignal = new Promise((resolve) => { + adopted = resolve + }) + let release!: () => void + const released = new Promise((resolve) => { + release = resolve + }) + const adoption = other.begin(async (tx) => { + await tx`UPDATE knowledge_base SET is_search_index = true WHERE id = 'legacy'` + adopted() + await released + }) + await adoptedSignal + const [{ pid }] = await sql`SELECT pg_backend_pid() AS pid` + const insert = sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) + VALUES ('waiting', 'legacy', 'doc', to_tsvector('english', 'Waiting chunk'))`.execute() + await waitUntilBlocked(admin, pid) + release() + await adoption + await insert + expect(await keywordIds()).toEqual(['waiting']) + }) + + /** + * Adopts `legacy` while a second connection holds a delete of one of its chunks open, and resolves + * once both have committed. + */ + async function adoptWhileDeleting(): Promise { + await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) VALUES + ('kept', 'legacy', 'doc', to_tsvector('english', 'Kept chunk')), + ('deleted', 'legacy', 'doc', to_tsvector('english', 'Deleted chunk'))` + let deleted!: () => void + const deletedSignal = new Promise((resolve) => { + deleted = resolve + }) + let release!: () => void + const released = new Promise((resolve) => { + release = resolve + }) + const deleter = sql.begin(async (tx) => { + await tx`DELETE FROM embedding WHERE id = 'deleted'` + deleted() + await released + }) + await deletedSignal + const [{ pid }] = await other`SELECT pg_backend_pid() AS pid` + const adoption = + other`UPDATE knowledge_base SET is_search_index = true WHERE id = 'legacy'`.execute() + await waitUntilBlocked(admin, pid) + release() + await deleter + await adoption + } + + it('adopts a base while one of its chunks is being deleted', async () => { + await adoptWhileDeleting() + expect(await keywordIds()).toEqual(['kept']) + expect(await tinIds()).toEqual(['kept']) + }) + + it('adopts a base into Tin while a chunk is being deleted, with no keyword trigger ahead of it', async () => { + await sql`ALTER TABLE knowledge_base DISABLE TRIGGER knowledge_base_keyword_search_sync` + try { + await adoptWhileDeleting() + } finally { + await sql`ALTER TABLE knowledge_base ENABLE TRIGGER knowledge_base_keyword_search_sync` + } + expect(await tinIds()).toEqual(['kept']) + }) + + it('runs no fan-out for an inserted document, and still fans out a changed ACL', async () => { + const projections = ['embedding_search', 'embedding_keyword_tin'] as const + const scans = await sql.begin(async (tx) => { + const before = await scansOf(tx, projections) + await tx`INSERT INTO document (id, connector_id, acl) VALUES ('new', 'src', ARRAY['u:alice'])` + return (await scansOf(tx, projections)) - before + }) + expect(scans).toBe(0) + + await sql`INSERT INTO embedding (id, knowledge_base_id, document_id, content_tsv) + VALUES ('chunk', 'index', 'new', to_tsvector('english', 'Shared chunk'))` + await sql`INSERT INTO embedding_search (id, document_id) VALUES ('chunk', 'new')` + await sql`UPDATE document SET acl = ARRAY['u:bob'] WHERE id = 'new'` + for (const projection of projections) { + expect(await sql`SELECT acl FROM ${sql(projection)} WHERE id = 'chunk'`).toEqual([ + { acl: ['u:bob'] }, + ]) + } + }) +}) diff --git a/packages/db/script-migrations/0025_scope_keyword_projections.ts b/packages/db/script-migrations/0025_scope_keyword_projections.ts new file mode 100644 index 00000000000..cac41532ca3 --- /dev/null +++ b/packages/db/script-migrations/0025_scope_keyword_projections.ts @@ -0,0 +1,172 @@ +import { + installMembershipKey, + installTinChunkSync, + installTinMembershipSync, +} from '@sim/db/script-migrations/0019_tin_keyword_projection' +import { resolveMigrationDatabaseUrl } from '@sim/db/script-migrations/database-url' +import type { ScriptMigration } from '@sim/db/script-migrations/types' +import { createLogger } from '@sim/logger' +import postgres, { type Sql, type TransactionSql } from 'postgres' +import { retryOnLockTimeout } from '../scripts/lock-timeout-retry' + +const logger = createLogger('ScopeKeywordProjections') + +/** Each attempt's wait for the table locks the triggers need; the retries find a quiet moment. */ +const TRIGGER_LOCK_TIMEOUT = '5s' +const TRIGGER_LOCK_RETRY_BUDGET_MS = 20 * 60_000 +const TRIGGER_LOCK_RETRY_BACKOFF = { baseMs: 2_000, maxMs: 30_000 } as const + +/** The columns a GIN keyword row copies from its chunk. */ +const KEYWORD_SEARCH_COLUMNS = [ + 'knowledge_base_id', + 'document_id', + 'enabled', + 'content_tsv', +] as const + +/** + * How both GIN keyword writers bring an existing row (aliased `s`) up to its chunk. A row already + * current is left unwritten, as the projector's upserts leave it, so adopting a base whose rows + * survived rewrites only what changed while it holds the membership lock. + */ +const KEYWORD_SEARCH_UPSERT = `ON CONFLICT (id) DO UPDATE SET + ${KEYWORD_SEARCH_COLUMNS.map((column) => `${column} = EXCLUDED.${column}`).join(', ')} + WHERE (${KEYWORD_SEARCH_COLUMNS.map((column) => `s.${column}`).join(', ')}) + IS DISTINCT FROM (${KEYWORD_SEARCH_COLUMNS.map((column) => `EXCLUDED.${column}`).join(', ')})` + +/** + * The GIN keyword projection's chunk trigger, scoped to search indexes as the Tin projection is: + * it shares the base's membership lock before reading the marker, so a marker change in flight is + * waited for and read under a fresh snapshot. A chunk outside a search index writes nothing, and an + * update removes a row it left behind by moving out of one. + */ +async function installKeywordSearchSync(tx: TransactionSql): Promise { + await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_embedding_keyword_search() + RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + PERFORM pg_advisory_xact_lock_shared(knowledge_tin_membership_key(NEW.knowledge_base_id)); + IF NOT EXISTS ( + SELECT 1 FROM knowledge_base WHERE id = NEW.knowledge_base_id AND is_search_index + ) THEN + IF TG_OP = 'UPDATE' THEN + DELETE FROM embedding_keyword_search WHERE id = NEW.id; + END IF; + RETURN NEW; + END IF; + INSERT INTO embedding_keyword_search AS s (id, knowledge_base_id, document_id, enabled, content_tsv) + VALUES (NEW.id, NEW.knowledge_base_id, NEW.document_id, NEW.enabled, NEW.content_tsv) + ${KEYWORD_SEARCH_UPSERT}; + RETURN NEW; + END; + $$`) +} + +/** + * Projects or removes a whole base's GIN keyword rows when its search-index marker changes, under + * the membership lock taken exclusively, as the Tin projection's knowledge base trigger does. It is + * its own trigger because that one exists only where Tin is installed. The chunks are key-share + * locked as they are read, as the projection backfills lock theirs: a chunk delete in flight is + * waited for and its chunk skipped, rather than failing the adoption on the row's foreign key. + */ +async function installKeywordSearchMembership(tx: TransactionSql): Promise { + await tx.unsafe(`CREATE OR REPLACE FUNCTION sync_knowledge_base_keyword_search() + RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN + PERFORM pg_advisory_xact_lock(knowledge_tin_membership_key(NEW.id)); + IF NOT NEW.is_search_index THEN + DELETE FROM embedding_keyword_search WHERE knowledge_base_id = NEW.id; + RETURN NEW; + END IF; + INSERT INTO embedding_keyword_search AS s (id, knowledge_base_id, document_id, enabled, content_tsv) + SELECT id, knowledge_base_id, document_id, enabled, content_tsv + FROM embedding WHERE knowledge_base_id = NEW.id + ORDER BY id FOR KEY SHARE + ${KEYWORD_SEARCH_UPSERT}; + RETURN NEW; + END; + $$`) + await tx.unsafe(`CREATE OR REPLACE TRIGGER knowledge_base_keyword_search_sync + AFTER UPDATE OF is_search_index ON knowledge_base + FOR EACH ROW WHEN (OLD.is_search_index IS DISTINCT FROM NEW.is_search_index) + EXECUTE FUNCTION sync_knowledge_base_keyword_search()`) +} + +/** + * The Tin trigger bodies as `0019_tin_keyword_projection` now installs them, where that migration + * installed them: the chunk trigger without the delete an insert outside a search index ran, and + * the knowledge base trigger key-share locking the chunks it projects. + */ +async function installTinSync(tx: TransactionSql): Promise { + const [row] = await tx>` + SELECT to_regprocedure('sync_embedding_keyword_tin()') IS NOT NULL AS installed` + if (!row?.installed) return + await installTinChunkSync(tx) + await installTinMembershipSync(tx) +} + +/** + * The document trigger, firing only on an update that changes the source or ACL. An inserted + * document has no chunks yet, since a chunk's foreign key needs its document, so the fan-out it + * ran matched no row and it marked nothing; an upsert that updates still fires the update arm. Its + * body is the one `0024_knowledge_projection_async` installed. + */ +async function installDocumentTrigger(tx: TransactionSql): Promise { + await tx.unsafe(`CREATE OR REPLACE TRIGGER projection_source_acl_sync + AFTER UPDATE OF connector_id, acl ON document + FOR EACH ROW WHEN (OLD.connector_id IS DISTINCT FROM NEW.connector_id OR OLD.acl IS DISTINCT FROM NEW.acl) + EXECUTE FUNCTION sync_projection_source_acl()`) +} + +/** + * Keeps the keyword projections for search indexes only, the bases whose keyword rows indexed + * organization search ranks; every other search ranks keywords on `embedding`. Membership follows + * `knowledge_base.is_search_index` rather than the deployment's search mode, so the schema holds + * whichever mode runs. The document trigger stops firing on inserts and on updates that change + * nothing it carries. + * + * Rows already written for other bases, and those a rerun of `0016_backfill_search_vectors` writes + * before this reruns after it, are left in place: nothing reads them, and removing them is a table + * owner's maintenance, not a deploy's. All of it is installed in one transaction, so no marker change sees + * the scoped chunk trigger without the base trigger that backfills it. Each attempt waits at most + * {@link TRIGGER_LOCK_TIMEOUT} for the trigger DDL's locks and is retried within the budget. + * Idempotent. + */ +export async function scopeKeywordProjections(sql: Sql): Promise { + await retryOnLockTimeout( + () => + sql.begin(async (tx) => { + await tx.unsafe(`SET LOCAL lock_timeout = '${TRIGGER_LOCK_TIMEOUT}'`) + await installMembershipKey(tx) + await installKeywordSearchSync(tx) + await installKeywordSearchMembership(tx) + await installTinSync(tx) + await installDocumentTrigger(tx) + }), + { + budgetMs: TRIGGER_LOCK_RETRY_BUDGET_MS, + backoff: TRIGGER_LOCK_RETRY_BACKOFF, + onRetry: ({ attempt, delayMs }) => + logger.warn('Keyword projection triggers waited out their lock timeout; retrying', { + attempt, + retryInMs: Math.round(delayMs), + }), + } + ) +} + +export const scopeKeywordProjectionsMigration: ScriptMigration = { + name: '0025_scope_keyword_projections', + up: scopeKeywordProjections, +} + +/** Run directly by `db:push`, after the projection migrations it builds on. */ +if (import.meta.main) { + const url = resolveMigrationDatabaseUrl() + if (!url) throw new Error('DATABASE_URL is required to scope the keyword projections') + const sql = postgres(url, { max: 1, max_lifetime: null, onnotice: () => undefined }) + try { + await scopeKeywordProjections(sql) + } finally { + await sql.end() + } +} diff --git a/packages/db/script-migrations/index.ts b/packages/db/script-migrations/index.ts index c38b01a7464..b469e505669 100644 --- a/packages/db/script-migrations/index.ts +++ b/packages/db/script-migrations/index.ts @@ -8,6 +8,7 @@ import { tinKeywordProjectionMigration } from '@sim/db/script-migrations/0019_ti import { projectionSourceAclBackfillMigration } from '@sim/db/script-migrations/0022_projection_source_acl_backfill' import { projectionAclSkipUnfilledMigration } from '@sim/db/script-migrations/0023_projection_acl_skip_unfilled' import { knowledgeProjectionAsyncMigration } from '@sim/db/script-migrations/0024_knowledge_projection_async' +import { scopeKeywordProjectionsMigration } from '@sim/db/script-migrations/0025_scope_keyword_projections' import type { Sql } from 'postgres' import { backfillTableOrderKeys } from './0001_backfill_table_order_keys' import { backfillPausedBillingAttribution } from './0002_backfill_paused_billing_attribution' @@ -52,6 +53,8 @@ export const scriptMigrations: readonly ScriptMigration[] = [ projectionAclSkipUnfilledMigration, /** 0024 marks changed documents for the knowledge projector and lets writers defer to it. */ knowledgeProjectionAsyncMigration, + /** 0025 keeps the keyword projections for search indexes only. */ + scopeKeywordProjectionsMigration, ] /** diff --git a/packages/db/scripts/push.ts b/packages/db/scripts/push.ts index cd47c9c635b..90b2288954c 100644 --- a/packages/db/scripts/push.ts +++ b/packages/db/scripts/push.ts @@ -8,6 +8,7 @@ const RECONCILIATION_COMMANDS = [ ['bun', '--env-file=.env', 'run', './script-migrations/0019_tin_keyword_projection.ts'], ['bun', '--env-file=.env', 'run', './script-migrations/0021_embedding_search_connector.ts'], ['bun', '--env-file=.env', 'run', './script-migrations/0024_knowledge_projection_async.ts'], + ['bun', '--env-file=.env', 'run', './script-migrations/0025_scope_keyword_projections.ts'], ] /**