Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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)
})

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import {
credentialGroupEnrollment,
document,
embedding,
embeddingKeywordSearch,
embeddingKeywordTin,
embeddingSearch,
knowledgeBase,
Expand Down Expand Up @@ -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
Expand All @@ -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;`)
)
Expand Down Expand Up @@ -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) => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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([
Expand Down
4 changes: 2 additions & 2 deletions apps/sim/lib/sim-search/indexed/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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.
Expand Down
25 changes: 18 additions & 7 deletions packages/db/knowledge-projection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(', ')})
Expand Down Expand Up @@ -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. */
Expand Down Expand Up @@ -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 }>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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',
])
})
})
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
123 changes: 78 additions & 45 deletions packages/db/script-migrations/0019_tin_keyword_projection.ts
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -39,6 +40,75 @@ async function tinAvailable(sql: Sql): Promise<boolean> {
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<void> {
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<void> {
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));
Comment thread
waleedlatif1 marked this conversation as resolved.
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<void> {
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.
Expand All @@ -60,7 +130,9 @@ async function tinAvailable(sql: Sql): Promise<boolean> {
* 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<void> {
await sql.begin(async (tx) => {
Expand All @@ -75,51 +147,12 @@ export async function installProjection(sql: Sql): Promise<void> {
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)
Expand Down
Loading
Loading