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
@@ -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<Record<string, unknown>>()),
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')
})
})
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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])
Expand Down
12 changes: 6 additions & 6 deletions apps/sim/lib/knowledge/documents/processing-queue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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(
Expand All @@ -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(
Expand Down
36 changes: 35 additions & 1 deletion apps/sim/lib/knowledge/documents/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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`}
)`
}

Expand Down Expand Up @@ -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),
}
Comment thread
icecrasher321 marked this conversation as resolved.
}
const queuedAt = new Date()
const documentIds = uniqueDocuments.map((doc) => doc.documentId)
const {
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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`
Expand Down Expand Up @@ -2421,6 +2443,7 @@ async function resolveDocumentStorageAdmission(
): Promise<DocumentStorageAdmission> {
const [kb] = await db
.select({
isSearchIndex: knowledgeBase.isSearchIndex,
workspaceId: knowledgeBase.workspaceId,
organizationId: knowledgeBase.organizationId,
userId: knowledgeBase.userId,
Expand All @@ -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',
Expand Down Expand Up @@ -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,
Expand All @@ -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(
Expand Down Expand Up @@ -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,
Expand All @@ -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 &&
Expand Down
5 changes: 3 additions & 2 deletions apps/sim/lib/sim-search/indexed/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.

Expand Down
7 changes: 2 additions & 5 deletions packages/db/drizzle.config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
29 changes: 0 additions & 29 deletions packages/db/script-migrations-paused-billing-attribution.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down Expand Up @@ -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',
])
})
})
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
icecrasher321 marked this conversation as resolved.
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
Expand Down
Loading
Loading