diff --git a/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts b/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts new file mode 100644 index 00000000000..98fc26882b8 --- /dev/null +++ b/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts @@ -0,0 +1,98 @@ +import { db } from '@sim/db' +import { document, knowledgeBase, organization, user, workspace } from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq, inArray } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' +import { + createKnowledgeAclFixtureIds, + seedKnowledgeAclFixture, +} from '@/lib/knowledge/__integration__/seed-source-access-fixture' +import { + createDocumentRecords, + createSingleDocument, + processDocumentAsync, + processDocumentsWithQueue, +} from '@/lib/knowledge/documents/service' + +vi.mock('@/lib/core/config/env-flags', async (importOriginal) => ({ + ...(await importOriginal>()), + isLiveEnterpriseSearchEnabled: true, +})) + +/** Queued work cannot revive dormant Search before retirement reaches its documents. */ +describe('dormant Search document processing', () => { + const ids = createKnowledgeAclFixtureIds() + const documentId = generateId() + const source = { + filename: 'fixture.txt', + fileUrl: 'data:text/plain,fixture', + fileSize: 7, + mimeType: 'text/plain', + } + + beforeAll(async () => { + await seedKnowledgeAclFixture(ids, { connectorType: 'google_drive' }) + await db + .update(knowledgeBase) + .set({ isSearchIndex: true }) + .where(eq(knowledgeBase.id, ids.knowledgeBaseId)) + await db + .insert(document) + .values({ id: documentId, knowledgeBaseId: ids.knowledgeBaseId, ...source }) + }) + + afterAll(async () => { + await db.delete(knowledgeBase).where(eq(knowledgeBase.id, ids.knowledgeBaseId)) + await db.delete(workspace).where(eq(workspace.id, ids.workspaceId)) + await db.delete(organization).where(eq(organization.id, ids.organizationId)) + await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId])) + }) + + it('rejects uploads before creating documents in a dormant Search KB', async () => { + await expect( + createDocumentRecords([source], ids.knowledgeBaseId, generateId()) + ).rejects.toThrow('inactive') + await expect(createSingleDocument(source, ids.knowledgeBaseId, generateId())).rejects.toThrow( + 'inactive' + ) + const documents = await db + .select({ id: document.id }) + .from(document) + .where(eq(document.knowledgeBaseId, ids.knowledgeBaseId)) + expect(documents).toEqual([{ id: documentId }]) + }) + + it('refuses dispatch before stamping or charging a pending Search document', async () => { + const result = await processDocumentsWithQueue( + [{ documentId, ...source }], + ids.knowledgeBaseId, + {}, + generateId(), + undefined, + 'backfill' + ) + expect(result).toEqual({ + requested: 1, + accepted: 0, + failed: 1, + failedDocumentIds: [documentId], + }) + const [stored] = await db + .select({ token: document.processingQueueToken, queuedAt: document.processingQueuedAt }) + .from(document) + .where(eq(document.id, documentId)) + expect(stored).toEqual({ token: null, queuedAt: null }) + }) + + it('skips an old worker payload before claiming the document or requiring billing', async () => { + expect(await processDocumentAsync(ids.knowledgeBaseId, documentId, source)).toEqual({ + outcome: 'skipped', + reason: 'unavailable', + }) + const [stored] = await db + .select({ status: document.processingStatus }) + .from(document) + .where(eq(document.id, documentId)) + expect(stored.status).toBe('pending') + }) +}) 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 bffc2212d1e..25a060da5fb 100644 --- a/apps/sim/lib/knowledge/__integration__/embedding-insert-batches.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/embedding-insert-batches.integration.ts @@ -17,6 +17,13 @@ import { generateId } from '@sim/utils/id' import { eq, inArray } from 'drizzle-orm' import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' +/** These transaction checks exercise indexed Search, which Live Search normally disables. */ +vi.mock('@/lib/core/config/env-flags', async (importOriginal) => + (await import('@sim/testing/mocks/indexed-org-search.mock')).indexedOrgSearchEnvFlags( + importOriginal + ) +) + const fixtures = vi.hoisted(() => ({ root: '', process: vi.fn(), embeddings: vi.fn() })) vi.mock('@/lib/uploads/core/setup.server', () => ({ get UPLOAD_DIR_SERVER() { 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 5c39be414e8..427f1e758de 100644 --- a/apps/sim/lib/knowledge/__integration__/processing-lock-scope.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/processing-lock-scope.integration.ts @@ -13,6 +13,13 @@ import { eq, inArray } from 'drizzle-orm' import postgres from 'postgres' import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest' +/** These transaction checks exercise indexed Search, which Live Search normally disables. */ +vi.mock('@/lib/core/config/env-flags', async (importOriginal) => + (await import('@sim/testing/mocks/indexed-org-search.mock')).indexedOrgSearchEnvFlags( + importOriginal + ) +) + const fixtures = vi.hoisted(() => ({ root: '', process: vi.fn(), embeddings: vi.fn() })) vi.mock('@/lib/uploads/core/setup.server', () => ({ get UPLOAD_DIR_SERVER() { diff --git a/apps/sim/lib/knowledge/documents/document-processing-source.test.ts b/apps/sim/lib/knowledge/documents/document-processing-source.test.ts index 9fa048e4ac4..e8068c08a7e 100644 --- a/apps/sim/lib/knowledge/documents/document-processing-source.test.ts +++ b/apps/sim/lib/knowledge/documents/document-processing-source.test.ts @@ -1084,6 +1084,7 @@ describe('in-process quota continuation dispatch', () => { Object.assign(env, { ...defaultMockEnv, TRIGGER_SECRET_KEY: undefined }) dbChainMockFns.returning.mockResolvedValue([{ id: 'document-1' }]) dbChainMockFns.limit + .mockResolvedValueOnce([{ isSearchIndex: false }]) .mockResolvedValueOnce([{ userId: 'knowledge-owner', workspaceId: 'workspace-1' }]) .mockResolvedValueOnce([PERSISTED_CONTEXT]) .mockResolvedValueOnce([PERSISTED_PROVENANCE_ROW]) diff --git a/apps/sim/lib/knowledge/documents/processing-queue.test.ts b/apps/sim/lib/knowledge/documents/processing-queue.test.ts index 9f38fd2d168..d8a0b985cf9 100644 --- a/apps/sim/lib/knowledge/documents/processing-queue.test.ts +++ b/apps/sim/lib/knowledge/documents/processing-queue.test.ts @@ -89,8 +89,8 @@ describe('processDocumentsWithQueue billing attribution', () => { }) it('rejects missing workspace attribution without enqueueing', async () => { - dbChainMockFns.limit.mockResolvedValueOnce([ - { userId: 'knowledge-owner', workspaceId: 'workspace-1' }, + dbChainMockFns.limit.mockResolvedValue([ + { isSearchIndex: false, userId: 'knowledge-owner', workspaceId: 'workspace-1' }, ]) await expect( @@ -112,8 +112,8 @@ describe('processDocumentsWithQueue billing attribution', () => { }) it('rejects mismatched workspace attribution without enqueueing', async () => { - dbChainMockFns.limit.mockResolvedValueOnce([ - { userId: 'knowledge-owner', workspaceId: 'workspace-2' }, + dbChainMockFns.limit.mockResolvedValue([ + { isSearchIndex: false, userId: 'knowledge-owner', workspaceId: 'workspace-2' }, ]) await expect( @@ -130,8 +130,8 @@ describe('processDocumentsWithQueue billing attribution', () => { }) it('rejects a knowledge base without a workspace or organization owner', async () => { - dbChainMockFns.limit.mockResolvedValueOnce([ - { userId: 'legacy-owner', workspaceId: null, organizationId: null }, + dbChainMockFns.limit.mockResolvedValue([ + { isSearchIndex: false, userId: 'legacy-owner', workspaceId: null, organizationId: null }, ]) await expect( diff --git a/apps/sim/lib/knowledge/documents/service.ts b/apps/sim/lib/knowledge/documents/service.ts index 6373e269d69..82c09359367 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -81,6 +81,10 @@ import { SYSTEM_ACCESS_SCOPE, } from '@/lib/knowledge/access/types' import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error' +import { + connectorIndexingCondition, + requiresConnectorIndexing, +} from '@/lib/knowledge/connectors/indexing-policy' import { assertSyncLeaseHeldInTx, type SyncWriteLease } from '@/lib/knowledge/connectors/sync-lock' import { documentConnectorIsActive } from '@/lib/knowledge/documents/connector-lifecycle' import { @@ -229,12 +233,13 @@ class SupersededProcessingOutput extends Error { } } -/** The document's knowledge base has not been deleted. */ +/** The document's knowledge base is active and its backend still accepts indexed content. */ function knowledgeBaseIsActive() { return sql`EXISTS ( SELECT 1 FROM ${knowledgeBase} WHERE ${knowledgeBase.id} = ${document.knowledgeBaseId} AND ${knowledgeBase.deletedAt} IS NULL + AND ${connectorIndexingCondition() ?? sql`true`} )` } @@ -1186,6 +1191,19 @@ export async function processDocumentsWithQueue( } const requested = uniqueDocuments.length + const [indexingTarget] = await db + .select({ isSearchIndex: knowledgeBase.isSearchIndex }) + .from(knowledgeBase) + .where(eq(knowledgeBase.id, knowledgeBaseId)) + .limit(1) + if (indexingTarget && !requiresConnectorIndexing(indexingTarget.isSearchIndex)) { + return { + requested, + accepted: 0, + failed: requested, + failedDocumentIds: uniqueDocuments.map((doc) => doc.documentId), + } + } const queuedAt = new Date() const documentIds = uniqueDocuments.map((doc) => doc.documentId) const { @@ -1609,6 +1627,7 @@ export async function processDocumentAsync( const contextRows = await db .select({ + isSearchIndex: knowledgeBase.isSearchIndex, workspaceId: knowledgeBase.workspaceId, organizationId: knowledgeBase.organizationId, chunkingConfig: knowledgeBase.chunkingConfig, @@ -1653,6 +1672,9 @@ export async function processDocumentAsync( ) .limit(1) + if (contextRows[0] && !requiresConnectorIndexing(contextRows[0].isSearchIndex)) { + return { outcome: 'skipped', reason: 'unavailable' } + } if (contextRows.length === 0) { logger.warn( `[${documentId}] Skipping document processing: document or knowledge base ${knowledgeBaseId} no longer exists` @@ -2421,6 +2443,7 @@ async function resolveDocumentStorageAdmission( ): Promise { const [kb] = await db .select({ + isSearchIndex: knowledgeBase.isSearchIndex, workspaceId: knowledgeBase.workspaceId, organizationId: knowledgeBase.organizationId, userId: knowledgeBase.userId, @@ -2432,6 +2455,9 @@ async function resolveDocumentStorageAdmission( throw new OrchestrationError('not_found', 'Knowledge base not found') } + if (!requiresConnectorIndexing(kb.isSearchIndex)) { + throw new OrchestrationError('validation', 'This search index is inactive; use Sim Search.') + } if (kb.organizationId) throw new OrchestrationError( 'validation', @@ -2490,6 +2516,7 @@ export async function createDocumentRecords( const kb = await tx .select({ id: knowledgeBase.id, + isSearchIndex: knowledgeBase.isSearchIndex, workspaceId: knowledgeBase.workspaceId, organizationId: knowledgeBase.organizationId, userId: knowledgeBase.userId, @@ -2501,6 +2528,9 @@ export async function createDocumentRecords( if (kb.length === 0) { throw new OrchestrationError('not_found', 'Knowledge base not found') } + if (!requiresConnectorIndexing(kb[0].isSearchIndex)) { + throw new OrchestrationError('validation', 'This search index is inactive; use Sim Search.') + } if (kb[0].workspaceId !== admission.workspaceId) { throw new Error( @@ -3158,6 +3188,7 @@ export async function createSingleDocument( const kb = await tx .select({ id: knowledgeBase.id, + isSearchIndex: knowledgeBase.isSearchIndex, workspaceId: knowledgeBase.workspaceId, organizationId: knowledgeBase.organizationId, userId: knowledgeBase.userId, @@ -3169,6 +3200,9 @@ export async function createSingleDocument( if (kb.length === 0) { throw new OrchestrationError('not_found', 'Knowledge base not found') } + if (!requiresConnectorIndexing(kb[0].isSearchIndex)) { + throw new OrchestrationError('validation', 'This search index is inactive; use Sim Search.') + } if ( options?.expectedWorkspaceId !== undefined && diff --git a/apps/sim/lib/sim-search/indexed/README.md b/apps/sim/lib/sim-search/indexed/README.md index a392625fb7f..07440bedb8e 100644 --- a/apps/sim/lib/sim-search/indexed/README.md +++ b/apps/sim/lib/sim-search/indexed/README.md @@ -13,7 +13,7 @@ While the gate is off: - Search, the MCP tools, and Sim's `search_workspace` and `read_document` tools serve Live Search, and personal Search integrations are read from live accounts. - 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`). +- Nothing crawls into search indexes: content syncs, member syncs, processing recovery, document dispatch, and queued document workers 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 GIN keyword projection follows `is_search_index` alone, so it keeps its search-index rows either way. ## Layout @@ -29,7 +29,8 @@ The dormant UI sits in `indexed/` folders next to the component that picks it fr 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 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. +3. If the `0027_retire_search_embeddings` cleanup ran for a base, deliberately restore its retired documents' eligibility before resyncing. Its deleted embeddings cannot be recovered by changing the backend flag alone. See `packages/db/script-migrations/search-embedding-retirement.md` for the cleanup lifecycle. +4. 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/drizzle.config.ts b/packages/db/drizzle.config.ts index 79dd4be671c..aead0f23910 100644 --- a/packages/db/drizzle.config.ts +++ b/packages/db/drizzle.config.ts @@ -15,9 +15,6 @@ export default { dbCredentials: { url: process.env.DATABASE_URL!, }, - /* script_migrations is the one-off-script ledger (script-migrations/index.ts) — - deliberately managed outside drizzle. Without this filter, dev's `db:push` sees an - unknown table: it prompts "created or renamed?" (no TTY in CI → red) and would DROP - the ledger, making every applied one-off script re-run. */ - tablesFilter: ['!script_migrations'], + /** Runner-owned journals and resumable cleanup cursors must survive a development schema push. */ + tablesFilter: ['!script_migrations', '!search_embedding_cleanup_progress'], } satisfies Config diff --git a/packages/db/script-migrations-paused-billing-attribution.test.ts b/packages/db/script-migrations-paused-billing-attribution.test.ts index 131364665ac..a428608db7d 100644 --- a/packages/db/script-migrations-paused-billing-attribution.test.ts +++ b/packages/db/script-migrations-paused-billing-attribution.test.ts @@ -23,7 +23,6 @@ import { type SubscriptionCandidate, selectFrozenPersonalSubscription, } from './script-migrations/0002_backfill_paused_billing_attribution' -import { scriptMigrations } from './script-migrations/index' const ORGANIZATION_ATTRIBUTION: BillingAttributionSnapshot = { actorUserId: 'actor-1', @@ -354,31 +353,3 @@ describe('paused billing attribution safety', () => { ) }) }) - -describe('script migration registry', () => { - it('keeps script migrations in append-only order', () => { - expect(scriptMigrations.map((migration) => migration.name)).toEqual([ - '0001_backfill_table_order_keys', - '0002_backfill_paused_billing_attribution', - '0003_backfill_workspace_storage_usage', - '0004_backfill_fork_kb_file_ownership', - '0005_repair_unknown_table_row_provenance', - '0006_repair_unknown_table_row_provenance_second_pass', - '0007_repair_unknown_workspace_file_provenance', - '0010_backfill_credential_group_resource_policies', - '0011_remap_legacy_knowledge_connector_credentials', - '0012_reconcile_oauth_provider_lifecycle', - '0013_backfill_legacy_knowledge_base_workspaces', - '0014_require_knowledge_base_owner', - '0016_backfill_search_vectors', - '0017_index_search_documents', - '0018_repair_workspace_file_content_revision', - '0019_tin_keyword_projection', - '0022_projection_source_acl_backfill', - '0023_projection_acl_skip_unfilled', - '0024_knowledge_projection_async', - '0025_scope_keyword_projections', - '0026_user_table_schema_for_write', - ]) - }) -}) 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 fc9cc6f741c..e3c92cf8955 100644 --- a/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts +++ b/packages/db/script-migrations/0016_backfill_search_vectors.integration.ts @@ -384,19 +384,11 @@ describe('search projection upgrade in PostgreSQL', () => { } await runScriptMigrations(sql) expect( - await sql`SELECT name FROM script_migrations WHERE name >= '0015' ORDER BY name` + await sql`SELECT name FROM script_migrations + WHERE name IN ('0015_backfill_embedding_search', '0016_backfill_search_vectors') ORDER BY name` ).toEqual([ { name: '0015_backfill_embedding_search' }, { name: '0016_backfill_search_vectors' }, - { name: '0017_index_search_documents' }, - { name: '0018_repair_workspace_file_content_revision' }, - { name: '0019_tin_keyword_projection' }, - { name: '0021_embedding_search_connector' }, - { name: '0022_projection_source_acl_backfill' }, - { name: '0023_projection_acl_skip_unfilled' }, - { name: '0024_knowledge_projection_async' }, - { name: '0025_scope_keyword_projections' }, - { name: '0026_user_table_schema_for_write' }, ]) 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/0027_retire_search_embeddings.integration.ts b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts new file mode 100644 index 00000000000..80ef1608a5c --- /dev/null +++ b/packages/db/script-migrations/0027_retire_search_embeddings.integration.ts @@ -0,0 +1,149 @@ +import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' +import { ScriptMigrationDeferred } from '@sim/db/script-migrations/types' +import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure' +import { generateId } from '@sim/utils/id' +import postgres, { type Sql } from 'postgres' +import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest' + +/** Proves destructive scope, cascading integrity, and atomic restart against real PostgreSQL. */ +describe('retiring dormant Search embeddings', () => { + const schema = `search_retirement_${generateId().replaceAll('-', '')}` + let admin: Sql + let sql: Sql + + beforeAll(async () => { + admin = postgres(readTestDatabaseUrl(), { max: 1, onnotice: () => undefined }) + await admin.unsafe(`CREATE SCHEMA "${schema}"`) + sql = postgres(readTestDatabaseUrl(), { + max: 1, + connection: { search_path: schema }, + onnotice: () => undefined, + }) + await sql`CREATE TABLE knowledge_base (id text PRIMARY KEY, is_search_index boolean NOT NULL)` + await sql`CREATE TABLE document ( + id text PRIMARY KEY, knowledge_base_id text REFERENCES knowledge_base(id), + user_excluded boolean NOT NULL DEFAULT false, enabled boolean NOT NULL DEFAULT true, + processing_queue_token text, processing_queued_at timestamp, processing_deferred_until timestamp)` + await sql`CREATE TABLE embedding ( + id text PRIMARY KEY, knowledge_base_id text REFERENCES knowledge_base(id), + document_id text REFERENCES document(id))` + await sql`CREATE TABLE embedding_search (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)` + await sql`CREATE TABLE embedding_keyword_search (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)` + await sql`CREATE TABLE embedding_keyword_tin (id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)` + await sql`CREATE TABLE embedding_secret_provenance (embedding_id text PRIMARY KEY REFERENCES embedding(id) ON DELETE CASCADE)` + }) + + afterAll(async () => { + await sql?.end() + await admin.unsafe(`DROP SCHEMA "${schema}" CASCADE`) + await admin.end() + }) + + beforeEach(async () => { + await sql`TRUNCATE knowledge_base, document, embedding, embedding_search, + embedding_keyword_search, embedding_keyword_tin, embedding_secret_provenance` + await sql`DROP TABLE IF EXISTS search_embedding_cleanup_progress` + await sql`INSERT INTO knowledge_base VALUES ('search', true), ('ordinary', false)` + await sql`INSERT INTO document (id, knowledge_base_id, processing_queue_token) + VALUES ('search-doc', 'search', 'old-dispatch'), ('ordinary-doc', 'ordinary', 'keep-dispatch')` + await sql`INSERT INTO embedding + SELECT lpad(i::text, 5, '0'), CASE WHEN i % 2 = 0 THEN 'search' ELSE 'ordinary' END, + CASE WHEN i % 2 = 0 THEN 'search-doc' ELSE 'ordinary-doc' END + FROM generate_series(1, 1002) i` + await sql`INSERT INTO embedding_search SELECT id FROM embedding` + await sql`INSERT INTO embedding_keyword_search SELECT id FROM embedding` + await sql`INSERT INTO embedding_keyword_tin SELECT id FROM embedding` + await sql`INSERT INTO embedding_secret_provenance SELECT id FROM embedding` + }) + + async function pass() { + try { + await retireSearchEmbeddingsMigration.up(sql) + return true + } catch (error) { + if (error instanceof ScriptMigrationDeferred) return false + throw error + } + } + + it('does nothing without Search data and defers an ambiguous target without changing data', async () => { + await sql`UPDATE knowledge_base SET is_search_index = false` + expect(await pass()).toBe(true) + await sql`UPDATE knowledge_base SET is_search_index = true` + expect(await pass()).toBe(false) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(1002) + expect( + (await sql`SELECT user_excluded FROM document WHERE id = 'search-doc'`)[0].user_excluded + ).toBe(false) + }) + + it('cascades only Search chunks and preserves configuration and ordinary documents', async () => { + expect(await pass()).toBe(true) + expect( + ( + await sql`SELECT user_excluded, processing_queue_token FROM document WHERE id = 'search-doc'` + )[0] + ).toEqual({ user_excluded: true, processing_queue_token: null }) + expect( + (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0].n + ).toBe(0) + expect( + (await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'ordinary'`)[0] + .n + ).toBe(501) + for (const table of [ + 'embedding_search', + 'embedding_keyword_search', + 'embedding_keyword_tin', + 'embedding_secret_provenance', + ]) { + expect((await sql.unsafe(`SELECT count(*)::int AS n FROM ${table}`))[0].n).toBe(501) + } + expect((await sql`SELECT count(*)::int AS n FROM knowledge_base`)[0].n).toBe(2) + expect( + ( + await sql`SELECT user_excluded, processing_queue_token FROM document WHERE id = 'ordinary-doc'` + )[0] + ).toEqual({ user_excluded: false, processing_queue_token: 'keep-dispatch' }) + expect(await pass()).toBe(true) + }) + + it('rolls back failed pages and resumes the frozen target, retiring documents inserted behind the cursor', async () => { + await sql`CREATE TABLE deletion_blocker (id text REFERENCES embedding(id))` + await sql`INSERT INTO deletion_blocker VALUES ('00002')` + await expect(pass()).rejects.toThrow() + const before = await sql`SELECT * FROM search_embedding_cleanup_progress` + await expect(pass()).rejects.toThrow() + expect(await sql`SELECT * FROM search_embedding_cleanup_progress`).toEqual(before) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(1002) + await sql`DROP TABLE deletion_blocker` + await sql`INSERT INTO document (id, knowledge_base_id) VALUES ('aaa-late-document', 'search')` + await sql`INSERT INTO knowledge_base VALUES ('other-search', true)` + await sql`INSERT INTO document (id, knowledge_base_id, processing_queue_token) + VALUES ('other-search-doc', 'other-search', 'keep-dispatch')` + await sql`INSERT INTO embedding VALUES ('other-search-chunk', 'other-search', 'other-search-doc')` + expect(await pass()).toBe(true) + expect( + (await sql`SELECT user_excluded FROM document WHERE id = 'aaa-late-document'`)[0] + .user_excluded + ).toBe(true) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(502) + expect( + ( + await sql`SELECT user_excluded, processing_queue_token FROM document WHERE id = 'other-search-doc'` + )[0] + ).toEqual({ user_excluded: false, processing_queue_token: 'keep-dispatch' }) + }) + + it('defers at the page budget and resumes without skipping remaining chunks', async () => { + await sql`INSERT INTO embedding + SELECT lpad(i::text, 5, '0'), 'search', 'search-doc' FROM generate_series(1003, 51002) i` + expect(await pass()).toBe(false) + const [remaining] = + await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'` + expect(remaining.n).toBeGreaterThan(0) + expect(remaining.n).toBeLessThan(50501) + expect(await pass()).toBe(true) + expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501) + }, 60_000) +}) diff --git a/packages/db/script-migrations/0027_retire_search_embeddings.ts b/packages/db/script-migrations/0027_retire_search_embeddings.ts new file mode 100644 index 00000000000..5afec02dc11 --- /dev/null +++ b/packages/db/script-migrations/0027_retire_search_embeddings.ts @@ -0,0 +1,141 @@ +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 { sleep } from '@sim/utils/helpers' +import postgres, { type Sql } from 'postgres' + +const logger = createLogger('RetireSearchEmbeddings') +const BATCH_SIZE = 500 +const MAX_BATCHES = 100 +const RUN_BUDGET_MS = 60_000 + +interface Progress { + knowledge_base_id: string + phase: 'documents' | 'embeddings' | 'done' + after_id: string +} + +/** + * Retires the sole legacy Search KB after the move to live Search. Each page commits with its + * durable cursor. Other KBs are traversed without modification. The runner owns bookkeeping, as it owns `script_migrations`. + */ +export const retireSearchEmbeddingsMigration: ScriptMigration = { + name: '0027_retire_search_embeddings', + async up(sql) { + await sql`CREATE TABLE IF NOT EXISTS search_embedding_cleanup_progress ( + id integer PRIMARY KEY CHECK (id = 1), knowledge_base_id text NOT NULL, + phase text NOT NULL CHECK (phase IN ('documents', 'embeddings', 'done')), + after_id text NOT NULL + )` + const existing = await sql< + Progress[] + >`SELECT * FROM search_embedding_cleanup_progress WHERE id = 1` + if (existing.length === 0) { + const targets = await sql<{ id: string }[]>` + SELECT id FROM knowledge_base WHERE is_search_index LIMIT 2` + if (targets.length === 0) return + if (targets.length > 1) { + throw new ScriptMigrationDeferred( + 'Multiple Search knowledge bases found; cleanup target is ambiguous' + ) + } + await sql`INSERT INTO search_embedding_cleanup_progress VALUES (1, ${targets[0].id}, 'documents', '') + ON CONFLICT (id) DO NOTHING` + } + const [progress] = await sql< + Progress[] + >`SELECT * FROM search_embedding_cleanup_progress WHERE id = 1` + const knowledgeBaseId = progress.knowledge_base_id + + const deadline = Date.now() + RUN_BUDGET_MS + for (let batch = 0; batch < MAX_BATCHES && Date.now() < deadline; batch++) { + if (await retirePage(sql, knowledgeBaseId)) { + logger.info('Selected Search knowledge base embeddings retired') + return + } + await sleep(100) + } + throw new ScriptMigrationDeferred( + 'Cleanup page budget reached; rerun to resume the saved cursor' + ) + }, +} + +async function retirePage(sql: Sql, knowledgeBaseId: string): Promise { + return sql.begin(async (tx) => { + await tx`SET LOCAL statement_timeout = '15s'` + await tx`SET LOCAL lock_timeout = '1s'` + const [progress] = await tx` + SELECT knowledge_base_id, phase, after_id FROM search_embedding_cleanup_progress WHERE id = 1 FOR UPDATE` + if (progress.knowledge_base_id !== knowledgeBaseId) { + throw new Error('Cannot change the selected Search KB during a cleanup') + } + const [target] = await tx`SELECT id FROM knowledge_base + WHERE id = ${knowledgeBaseId} AND is_search_index FOR SHARE` + if (!target) throw new Error('Cleanup target is no longer a Search knowledge base') + if (progress.phase === 'done') return true + + if (progress.phase === 'documents') { + const rows = await tx<{ id: string }[]>` + SELECT id FROM document WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE}` + if (rows.length === 0) { + await tx`UPDATE search_embedding_cleanup_progress SET phase = 'embeddings', after_id = '' WHERE id = 1` + return false + } + await tx`UPDATE document + SET user_excluded = true, enabled = false, processing_queue_token = NULL, + processing_queued_at = NULL, processing_deferred_until = NULL + WHERE id IN ${tx(rows.map((row) => row.id))} AND knowledge_base_id = ${knowledgeBaseId} + AND (NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL + OR processing_queued_at IS NOT NULL OR processing_deferred_until IS NOT NULL)` + await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${rows[rows.length - 1].id} WHERE id = 1` + return false + } + + const rows = await tx<{ id: string }[]>` + SELECT id FROM embedding WHERE id > ${progress.after_id} ORDER BY id LIMIT ${BATCH_SIZE}` + if (rows.length === 0) { + /** A late insert may sort behind either UUID cursor; completion must recheck the target. */ + const [unretired] = + await tx`SELECT id FROM document WHERE knowledge_base_id = ${knowledgeBaseId} + AND (NOT user_excluded OR enabled OR processing_queue_token IS NOT NULL + OR processing_queued_at IS NOT NULL OR processing_deferred_until IS NOT NULL) LIMIT 1` + if (unretired) { + await tx`UPDATE search_embedding_cleanup_progress SET phase = 'documents', after_id = '' WHERE id = 1` + return false + } + const [remaining] = + await tx`SELECT id FROM embedding WHERE knowledge_base_id = ${knowledgeBaseId} LIMIT 1` + if (remaining) { + await tx`UPDATE search_embedding_cleanup_progress SET after_id = '' WHERE id = 1` + return false + } + await tx`UPDATE search_embedding_cleanup_progress SET phase = 'done' WHERE id = 1` + return true + } + const ids = rows.map((row) => row.id) + const [unretired] = await tx` + SELECT e.id FROM embedding e JOIN document d ON d.id = e.document_id + WHERE e.id IN ${tx(ids)} AND e.knowledge_base_id = ${knowledgeBaseId} + AND (NOT d.user_excluded OR d.knowledge_base_id <> e.knowledge_base_id) LIMIT 1` + if (unretired) + throw new Error('Search content changed after retirement; stop writers before resuming') + /** Foreign keys cascade to vector/keyword projections and private chunk provenance. */ + await tx`DELETE FROM embedding WHERE id IN ${tx(ids)} AND knowledge_base_id = ${knowledgeBaseId}` + await tx`UPDATE search_embedding_cleanup_progress SET after_id = ${rows[rows.length - 1].id} WHERE id = 1` + return false + }) +} + +/** The standalone entry uses the deployment journal and leaves budget-deferred work unrecorded. */ +if (import.meta.main) { + const url = resolveMigrationDatabaseUrl() + if (!url) throw new Error('DATABASE_URL is required for Search retirement') + const sql = postgres(url, { max: 1, max_lifetime: null, onnotice: () => undefined }) + try { + const { runScriptMigrations } = await import('@sim/db/script-migrations/index') + await runScriptMigrations(sql, [retireSearchEmbeddingsMigration]) + } finally { + await sql.end() + } +} diff --git a/packages/db/script-migrations/index.ts b/packages/db/script-migrations/index.ts index a98a4cfcd6e..77c9cc5869d 100644 --- a/packages/db/script-migrations/index.ts +++ b/packages/db/script-migrations/index.ts @@ -10,6 +10,7 @@ import { projectionAclSkipUnfilledMigration } from '@sim/db/script-migrations/00 import { knowledgeProjectionAsyncMigration } from '@sim/db/script-migrations/0024_knowledge_projection_async' import { scopeKeywordProjectionsMigration } from '@sim/db/script-migrations/0025_scope_keyword_projections' import { userTableSchemaForWriteMigration } from '@sim/db/script-migrations/0026_user_table_schema_for_write' +import { retireSearchEmbeddingsMigration } from '@sim/db/script-migrations/0027_retire_search_embeddings' import type { Sql } from 'postgres' import { backfillTableOrderKeys } from './0001_backfill_table_order_keys' import { backfillPausedBillingAttribution } from './0002_backfill_paused_billing_attribution' @@ -58,6 +59,7 @@ export const scriptMigrations: readonly ScriptMigration[] = [ scopeKeywordProjectionsMigration, /** 0026 installs the schema guard every table row write takes before it validates. */ userTableSchemaForWriteMigration, + retireSearchEmbeddingsMigration, ] /** diff --git a/packages/db/script-migrations/search-embedding-retirement.md b/packages/db/script-migrations/search-embedding-retirement.md new file mode 100644 index 00000000000..e012eab8286 --- /dev/null +++ b/packages/db/script-migrations/search-embedding-retirement.md @@ -0,0 +1,72 @@ +# Retiring one legacy Search index + +`0027_retire_search_embeddings` runs through the existing script-migration registry without +cleanup flags. On its first invocation it discovers the sole knowledge base whose persisted +`is_search_index` marker is true and saves that target. No Search KB is a completed no-op; +multiple Search KBs defer without changing content because the target is ambiguous. Once saved, +the target stays fixed even if another Search KB is created. Ordinary KBs and the target's live +source/credential configuration, document metadata, and backing files are preserved. + +## Deployment and execution + +The app and workers must already use live Search, and older indexing jobs must be drained before +this cleanup ships: deployment migrations run before the new app switches over. `SIM_SEARCH_LIVE=true` +(the default) makes `isIndexedOrgSearchEnabled()` false. **`SIM_SEARCH_LIVE=false` enables indexed +Search again.** The cleanup does not inspect this flag. Live source setup may still create a Search +KB for configuration; it does not index content. Document uploads, dispatch and queued processing +also honor the indexed-search gate. + +The ordinary migration runner starts the cleanup automatically. For subsequent maintenance passes, +run `bun run packages/db/script-migrations/0027_retire_search_embeddings.ts` with the writer supplied +through the normal `MIGRATION_DATABASE_URL`/`DATABASE_URL` configuration. Repeat until +`script_migrations` records `0027_retire_search_embeddings`. A budget deferral exits normally without +recording completion; it is **not** a completed purge. + +Each invocation handles at most 100 pages of 500 IDs, with a 60-second budget between pages, +a 15-second statement timeout, a one-second lock timeout and 100 ms pacing. One in-flight statement +can finish after the run budget. Keep one maintenance worker and monitor primary latency, WAL, +replica lag and available disk. Lock/query failures stop the invocation; rerun after resolving them. + +The runner-owned `search_embedding_cleanup_progress` table stores the selected KB, phase and ID +cursor. Page mutations and cursor advancement commit together. The one-off migration journal +records only completion. `db:push` excludes the progress table from schema diffing. + +The documents phase fences queued and in-flight processing by marking only target documents +excluded/disabled and clearing their dispatch stamps. The embeddings phase deletes only target +chunks, and refuses a page whose target document was not retired or has inconsistent ownership. +The existing foreign keys cascade to vector/keyword projections and chunk provenance. Both phases +walk the primary key in bounded pages; unrelated rows are read only as IDs and are never updated. +This avoids sorting a whole KB or repeatedly scanning earlier pages when no suitable composite +cleanup index exists. Resumption continues the saved scan, including across pages containing only +unrelated rows. Before completion, the cleanup checks for unretired documents and remaining chunks +behind either cursor and restarts the affected phase if needed. Keep target writers stopped and +do not change its marker during the pass. + +Inspect progress with: + +```sql +SELECT knowledge_base_id, phase, after_id FROM search_embedding_cleanup_progress; +SELECT name, applied_at FROM script_migrations WHERE name = '0027_retire_search_embeddings'; +``` + +After completion, check the selected KB has no `embedding` rows, verify its live Search, and verify +ordinary KB retrieval. A zero-row absence check may still scan index entries; use an appropriate +timeout. The cleanup is destructive and not reversible by flipping the search flag. Re-enabling +indexed Search requires deliberately restoring document eligibility and fully resyncing its sources. + +## Storage maintenance + +No vacuum, reindex or table rewrite runs inside this backfill. Plan maintenance separately on the +writer, outside transactions. Ordinary `VACUUM (ANALYZE, TRUNCATE FALSE)` on affected source, +projection and document tables makes dead space reusable and refreshes planner statistics. It +generally does not shrink allocated files. `VACUUM FULL` takes an exclusive lock and requires extra +space; use a separately planned maintenance window or a provider-supported online rewrite for +physical shrinkage. HNSW may benefit from `REINDEX INDEX CONCURRENTLY` before vacuum, with enough +temporary disk and WAL capacity. See the [PostgreSQL vacuum documentation](https://www.postgresql.org/docs/17/sql-vacuum.html) +and [pgvector maintenance guidance](https://github.com/pgvector/pgvector#vacuuming). + +Document counts are historical until a later document cleanup. Do not raw-delete documents or +bucket objects: their application hard-delete path also enqueues identity-bound storage cleanup +and applies accounting. Its ordinary scoped mode excludes retired documents, so a follow-up must +explicitly support these rows while preserving those side effects. Do not delete source accounts, +integration policies or permission grants used by live Search.