Skip to content

Commit 16c7e1a

Browse files
fix(search): retire legacy embeddings through scoped backfill (#8430)
* fix(search): retire legacy embeddings through scoped backfill * fix(search): run retirement without cleanup flags and close ingestion gaps * chore(tests): scope migration recovery journal assertion * chore(tests): align document dispatch billing and quota fixtures
1 parent a6c9b03 commit 16c7e1a

14 files changed

Lines changed: 525 additions & 53 deletions
Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,98 @@
1+
import { db } from '@sim/db'
2+
import { document, knowledgeBase, organization, user, workspace } from '@sim/db/schema'
3+
import { generateId } from '@sim/utils/id'
4+
import { eq, inArray } from 'drizzle-orm'
5+
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
6+
import {
7+
createKnowledgeAclFixtureIds,
8+
seedKnowledgeAclFixture,
9+
} from '@/lib/knowledge/__integration__/seed-source-access-fixture'
10+
import {
11+
createDocumentRecords,
12+
createSingleDocument,
13+
processDocumentAsync,
14+
processDocumentsWithQueue,
15+
} from '@/lib/knowledge/documents/service'
16+
17+
vi.mock('@/lib/core/config/env-flags', async (importOriginal) => ({
18+
...(await importOriginal<Record<string, unknown>>()),
19+
isLiveEnterpriseSearchEnabled: true,
20+
}))
21+
22+
/** Queued work cannot revive dormant Search before retirement reaches its documents. */
23+
describe('dormant Search document processing', () => {
24+
const ids = createKnowledgeAclFixtureIds()
25+
const documentId = generateId()
26+
const source = {
27+
filename: 'fixture.txt',
28+
fileUrl: 'data:text/plain,fixture',
29+
fileSize: 7,
30+
mimeType: 'text/plain',
31+
}
32+
33+
beforeAll(async () => {
34+
await seedKnowledgeAclFixture(ids, { connectorType: 'google_drive' })
35+
await db
36+
.update(knowledgeBase)
37+
.set({ isSearchIndex: true })
38+
.where(eq(knowledgeBase.id, ids.knowledgeBaseId))
39+
await db
40+
.insert(document)
41+
.values({ id: documentId, knowledgeBaseId: ids.knowledgeBaseId, ...source })
42+
})
43+
44+
afterAll(async () => {
45+
await db.delete(knowledgeBase).where(eq(knowledgeBase.id, ids.knowledgeBaseId))
46+
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
47+
await db.delete(organization).where(eq(organization.id, ids.organizationId))
48+
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
49+
})
50+
51+
it('rejects uploads before creating documents in a dormant Search KB', async () => {
52+
await expect(
53+
createDocumentRecords([source], ids.knowledgeBaseId, generateId())
54+
).rejects.toThrow('inactive')
55+
await expect(createSingleDocument(source, ids.knowledgeBaseId, generateId())).rejects.toThrow(
56+
'inactive'
57+
)
58+
const documents = await db
59+
.select({ id: document.id })
60+
.from(document)
61+
.where(eq(document.knowledgeBaseId, ids.knowledgeBaseId))
62+
expect(documents).toEqual([{ id: documentId }])
63+
})
64+
65+
it('refuses dispatch before stamping or charging a pending Search document', async () => {
66+
const result = await processDocumentsWithQueue(
67+
[{ documentId, ...source }],
68+
ids.knowledgeBaseId,
69+
{},
70+
generateId(),
71+
undefined,
72+
'backfill'
73+
)
74+
expect(result).toEqual({
75+
requested: 1,
76+
accepted: 0,
77+
failed: 1,
78+
failedDocumentIds: [documentId],
79+
})
80+
const [stored] = await db
81+
.select({ token: document.processingQueueToken, queuedAt: document.processingQueuedAt })
82+
.from(document)
83+
.where(eq(document.id, documentId))
84+
expect(stored).toEqual({ token: null, queuedAt: null })
85+
})
86+
87+
it('skips an old worker payload before claiming the document or requiring billing', async () => {
88+
expect(await processDocumentAsync(ids.knowledgeBaseId, documentId, source)).toEqual({
89+
outcome: 'skipped',
90+
reason: 'unavailable',
91+
})
92+
const [stored] = await db
93+
.select({ status: document.processingStatus })
94+
.from(document)
95+
.where(eq(document.id, documentId))
96+
expect(stored.status).toBe('pending')
97+
})
98+
})

‎apps/sim/lib/knowledge/__integration__/embedding-insert-batches.integration.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,13 @@ import { generateId } from '@sim/utils/id'
1717
import { eq, inArray } from 'drizzle-orm'
1818
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
1919

20+
/** These transaction checks exercise indexed Search, which Live Search normally disables. */
21+
vi.mock('@/lib/core/config/env-flags', async (importOriginal) =>
22+
(await import('@sim/testing/mocks/indexed-org-search.mock')).indexedOrgSearchEnvFlags(
23+
importOriginal
24+
)
25+
)
26+
2027
const fixtures = vi.hoisted(() => ({ root: '', process: vi.fn(), embeddings: vi.fn() }))
2128
vi.mock('@/lib/uploads/core/setup.server', () => ({
2229
get UPLOAD_DIR_SERVER() {

‎apps/sim/lib/knowledge/__integration__/processing-lock-scope.integration.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,13 @@ import { eq, inArray } from 'drizzle-orm'
1313
import postgres from 'postgres'
1414
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest'
1515

16+
/** These transaction checks exercise indexed Search, which Live Search normally disables. */
17+
vi.mock('@/lib/core/config/env-flags', async (importOriginal) =>
18+
(await import('@sim/testing/mocks/indexed-org-search.mock')).indexedOrgSearchEnvFlags(
19+
importOriginal
20+
)
21+
)
22+
1623
const fixtures = vi.hoisted(() => ({ root: '', process: vi.fn(), embeddings: vi.fn() }))
1724
vi.mock('@/lib/uploads/core/setup.server', () => ({
1825
get UPLOAD_DIR_SERVER() {

‎apps/sim/lib/knowledge/documents/document-processing-source.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1084,6 +1084,7 @@ describe('in-process quota continuation dispatch', () => {
10841084
Object.assign(env, { ...defaultMockEnv, TRIGGER_SECRET_KEY: undefined })
10851085
dbChainMockFns.returning.mockResolvedValue([{ id: 'document-1' }])
10861086
dbChainMockFns.limit
1087+
.mockResolvedValueOnce([{ isSearchIndex: false }])
10871088
.mockResolvedValueOnce([{ userId: 'knowledge-owner', workspaceId: 'workspace-1' }])
10881089
.mockResolvedValueOnce([PERSISTED_CONTEXT])
10891090
.mockResolvedValueOnce([PERSISTED_PROVENANCE_ROW])

‎apps/sim/lib/knowledge/documents/processing-queue.test.ts‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -89,8 +89,8 @@ describe('processDocumentsWithQueue billing attribution', () => {
8989
})
9090

9191
it('rejects missing workspace attribution without enqueueing', async () => {
92-
dbChainMockFns.limit.mockResolvedValueOnce([
93-
{ userId: 'knowledge-owner', workspaceId: 'workspace-1' },
92+
dbChainMockFns.limit.mockResolvedValue([
93+
{ isSearchIndex: false, userId: 'knowledge-owner', workspaceId: 'workspace-1' },
9494
])
9595

9696
await expect(
@@ -112,8 +112,8 @@ describe('processDocumentsWithQueue billing attribution', () => {
112112
})
113113

114114
it('rejects mismatched workspace attribution without enqueueing', async () => {
115-
dbChainMockFns.limit.mockResolvedValueOnce([
116-
{ userId: 'knowledge-owner', workspaceId: 'workspace-2' },
115+
dbChainMockFns.limit.mockResolvedValue([
116+
{ isSearchIndex: false, userId: 'knowledge-owner', workspaceId: 'workspace-2' },
117117
])
118118

119119
await expect(
@@ -130,8 +130,8 @@ describe('processDocumentsWithQueue billing attribution', () => {
130130
})
131131

132132
it('rejects a knowledge base without a workspace or organization owner', async () => {
133-
dbChainMockFns.limit.mockResolvedValueOnce([
134-
{ userId: 'legacy-owner', workspaceId: null, organizationId: null },
133+
dbChainMockFns.limit.mockResolvedValue([
134+
{ isSearchIndex: false, userId: 'legacy-owner', workspaceId: null, organizationId: null },
135135
])
136136

137137
await expect(

‎apps/sim/lib/knowledge/documents/service.ts‎

Lines changed: 35 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,10 @@ import {
8181
SYSTEM_ACCESS_SCOPE,
8282
} from '@/lib/knowledge/access/types'
8383
import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error'
84+
import {
85+
connectorIndexingCondition,
86+
requiresConnectorIndexing,
87+
} from '@/lib/knowledge/connectors/indexing-policy'
8488
import { assertSyncLeaseHeldInTx, type SyncWriteLease } from '@/lib/knowledge/connectors/sync-lock'
8589
import { documentConnectorIsActive } from '@/lib/knowledge/documents/connector-lifecycle'
8690
import {
@@ -229,12 +233,13 @@ class SupersededProcessingOutput extends Error {
229233
}
230234
}
231235

232-
/** The document's knowledge base has not been deleted. */
236+
/** The document's knowledge base is active and its backend still accepts indexed content. */
233237
function knowledgeBaseIsActive() {
234238
return sql`EXISTS (
235239
SELECT 1 FROM ${knowledgeBase}
236240
WHERE ${knowledgeBase.id} = ${document.knowledgeBaseId}
237241
AND ${knowledgeBase.deletedAt} IS NULL
242+
AND ${connectorIndexingCondition() ?? sql`true`}
238243
)`
239244
}
240245

@@ -1186,6 +1191,19 @@ export async function processDocumentsWithQueue(
11861191
}
11871192

11881193
const requested = uniqueDocuments.length
1194+
const [indexingTarget] = await db
1195+
.select({ isSearchIndex: knowledgeBase.isSearchIndex })
1196+
.from(knowledgeBase)
1197+
.where(eq(knowledgeBase.id, knowledgeBaseId))
1198+
.limit(1)
1199+
if (indexingTarget && !requiresConnectorIndexing(indexingTarget.isSearchIndex)) {
1200+
return {
1201+
requested,
1202+
accepted: 0,
1203+
failed: requested,
1204+
failedDocumentIds: uniqueDocuments.map((doc) => doc.documentId),
1205+
}
1206+
}
11891207
const queuedAt = new Date()
11901208
const documentIds = uniqueDocuments.map((doc) => doc.documentId)
11911209
const {
@@ -1609,6 +1627,7 @@ export async function processDocumentAsync(
16091627

16101628
const contextRows = await db
16111629
.select({
1630+
isSearchIndex: knowledgeBase.isSearchIndex,
16121631
workspaceId: knowledgeBase.workspaceId,
16131632
organizationId: knowledgeBase.organizationId,
16141633
chunkingConfig: knowledgeBase.chunkingConfig,
@@ -1653,6 +1672,9 @@ export async function processDocumentAsync(
16531672
)
16541673
.limit(1)
16551674

1675+
if (contextRows[0] && !requiresConnectorIndexing(contextRows[0].isSearchIndex)) {
1676+
return { outcome: 'skipped', reason: 'unavailable' }
1677+
}
16561678
if (contextRows.length === 0) {
16571679
logger.warn(
16581680
`[${documentId}] Skipping document processing: document or knowledge base ${knowledgeBaseId} no longer exists`
@@ -2421,6 +2443,7 @@ async function resolveDocumentStorageAdmission(
24212443
): Promise<DocumentStorageAdmission> {
24222444
const [kb] = await db
24232445
.select({
2446+
isSearchIndex: knowledgeBase.isSearchIndex,
24242447
workspaceId: knowledgeBase.workspaceId,
24252448
organizationId: knowledgeBase.organizationId,
24262449
userId: knowledgeBase.userId,
@@ -2432,6 +2455,9 @@ async function resolveDocumentStorageAdmission(
24322455
throw new OrchestrationError('not_found', 'Knowledge base not found')
24332456
}
24342457

2458+
if (!requiresConnectorIndexing(kb.isSearchIndex)) {
2459+
throw new OrchestrationError('validation', 'This search index is inactive; use Sim Search.')
2460+
}
24352461
if (kb.organizationId)
24362462
throw new OrchestrationError(
24372463
'validation',
@@ -2490,6 +2516,7 @@ export async function createDocumentRecords(
24902516
const kb = await tx
24912517
.select({
24922518
id: knowledgeBase.id,
2519+
isSearchIndex: knowledgeBase.isSearchIndex,
24932520
workspaceId: knowledgeBase.workspaceId,
24942521
organizationId: knowledgeBase.organizationId,
24952522
userId: knowledgeBase.userId,
@@ -2501,6 +2528,9 @@ export async function createDocumentRecords(
25012528
if (kb.length === 0) {
25022529
throw new OrchestrationError('not_found', 'Knowledge base not found')
25032530
}
2531+
if (!requiresConnectorIndexing(kb[0].isSearchIndex)) {
2532+
throw new OrchestrationError('validation', 'This search index is inactive; use Sim Search.')
2533+
}
25042534

25052535
if (kb[0].workspaceId !== admission.workspaceId) {
25062536
throw new Error(
@@ -3158,6 +3188,7 @@ export async function createSingleDocument(
31583188
const kb = await tx
31593189
.select({
31603190
id: knowledgeBase.id,
3191+
isSearchIndex: knowledgeBase.isSearchIndex,
31613192
workspaceId: knowledgeBase.workspaceId,
31623193
organizationId: knowledgeBase.organizationId,
31633194
userId: knowledgeBase.userId,
@@ -3169,6 +3200,9 @@ export async function createSingleDocument(
31693200
if (kb.length === 0) {
31703201
throw new OrchestrationError('not_found', 'Knowledge base not found')
31713202
}
3203+
if (!requiresConnectorIndexing(kb[0].isSearchIndex)) {
3204+
throw new OrchestrationError('validation', 'This search index is inactive; use Sim Search.')
3205+
}
31723206

31733207
if (
31743208
options?.expectedWorkspaceId !== undefined &&

‎apps/sim/lib/sim-search/indexed/README.md‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ While the gate is off:
1313
- 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.
1414
- 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.
1515
- 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.
16-
- Nothing crawls into search indexes: content syncs, member syncs, and processing recovery skip them (`lib/knowledge/connectors/indexing-policy.ts`).
16+
- 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`).
1717
- 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.
1818

1919
## Layout
@@ -29,7 +29,8 @@ The dormant UI sits in `indexed/` folders next to the component that picks it fr
2929

3030
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.
3131
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.
32-
3. Resume and fully resync the connectors of search-index knowledge bases, so content that went stale while dormant is indexed again.
32+
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.
33+
4. Resume and fully resync the connectors of search-index knowledge bases, so content that went stale while dormant is indexed again.
3334

3435
Projection rows written before projections carried their document's source and ACL are decided on their document until they are rewritten.
3536

‎packages/db/drizzle.config.ts‎

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -15,9 +15,6 @@ export default {
1515
dbCredentials: {
1616
url: process.env.DATABASE_URL!,
1717
},
18-
/* script_migrations is the one-off-script ledger (script-migrations/index.ts) —
19-
deliberately managed outside drizzle. Without this filter, dev's `db:push` sees an
20-
unknown table: it prompts "created or renamed?" (no TTY in CI → red) and would DROP
21-
the ledger, making every applied one-off script re-run. */
22-
tablesFilter: ['!script_migrations'],
18+
/** Runner-owned journals and resumable cleanup cursors must survive a development schema push. */
19+
tablesFilter: ['!script_migrations', '!search_embedding_cleanup_progress'],
2320
} satisfies Config

‎packages/db/script-migrations-paused-billing-attribution.test.ts‎

Lines changed: 0 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,6 @@ import {
2323
type SubscriptionCandidate,
2424
selectFrozenPersonalSubscription,
2525
} from './script-migrations/0002_backfill_paused_billing_attribution'
26-
import { scriptMigrations } from './script-migrations/index'
2726

2827
const ORGANIZATION_ATTRIBUTION: BillingAttributionSnapshot = {
2928
actorUserId: 'actor-1',
@@ -354,31 +353,3 @@ describe('paused billing attribution safety', () => {
354353
)
355354
})
356355
})
357-
358-
describe('script migration registry', () => {
359-
it('keeps script migrations in append-only order', () => {
360-
expect(scriptMigrations.map((migration) => migration.name)).toEqual([
361-
'0001_backfill_table_order_keys',
362-
'0002_backfill_paused_billing_attribution',
363-
'0003_backfill_workspace_storage_usage',
364-
'0004_backfill_fork_kb_file_ownership',
365-
'0005_repair_unknown_table_row_provenance',
366-
'0006_repair_unknown_table_row_provenance_second_pass',
367-
'0007_repair_unknown_workspace_file_provenance',
368-
'0010_backfill_credential_group_resource_policies',
369-
'0011_remap_legacy_knowledge_connector_credentials',
370-
'0012_reconcile_oauth_provider_lifecycle',
371-
'0013_backfill_legacy_knowledge_base_workspaces',
372-
'0014_require_knowledge_base_owner',
373-
'0016_backfill_search_vectors',
374-
'0017_index_search_documents',
375-
'0018_repair_workspace_file_content_revision',
376-
'0019_tin_keyword_projection',
377-
'0022_projection_source_acl_backfill',
378-
'0023_projection_acl_skip_unfilled',
379-
'0024_knowledge_projection_async',
380-
'0025_scope_keyword_projections',
381-
'0026_user_table_schema_for_write',
382-
])
383-
})
384-
})

‎packages/db/script-migrations/0016_backfill_search_vectors.integration.ts‎

Lines changed: 2 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -384,19 +384,11 @@ describe('search projection upgrade in PostgreSQL', () => {
384384
}
385385
await runScriptMigrations(sql)
386386
expect(
387-
await sql`SELECT name FROM script_migrations WHERE name >= '0015' ORDER BY name`
387+
await sql`SELECT name FROM script_migrations
388+
WHERE name IN ('0015_backfill_embedding_search', '0016_backfill_search_vectors') ORDER BY name`
388389
).toEqual([
389390
{ name: '0015_backfill_embedding_search' },
390391
{ name: '0016_backfill_search_vectors' },
391-
{ name: '0017_index_search_documents' },
392-
{ name: '0018_repair_workspace_file_content_revision' },
393-
{ name: '0019_tin_keyword_projection' },
394-
{ name: '0021_embedding_search_connector' },
395-
{ name: '0022_projection_source_acl_backfill' },
396-
{ name: '0023_projection_acl_skip_unfilled' },
397-
{ name: '0024_knowledge_projection_async' },
398-
{ name: '0025_scope_keyword_projections' },
399-
{ name: '0026_user_table_schema_for_write' },
400392
])
401393
const [{ complete }] = await sql`SELECT count(*)::int AS complete FROM embedding e
402394
JOIN embedding_search s ON s.id = e.id JOIN embedding_keyword_search k ON k.id = e.id

0 commit comments

Comments
 (0)