From 8d9c765d42481d8ddb03c6e7fc8aa1deacfb87d6 Mon Sep 17 00:00:00 2001 From: Vikhyath Mondreti Date: Wed, 16 Sep 2026 16:28:10 -0700 Subject: [PATCH 1/3] fix(search): preserve indexed vector retrieval and literal source text --- apps/sim/lib/file-parsers/index.ts | 2 +- apps/sim/lib/file-parsers/sniff.test.ts | 32 ++ apps/sim/lib/file-parsers/sniff.ts | 14 +- apps/sim/lib/file-parsers/types.ts | 2 + .../search-latency.integration.ts | 343 ++++++++++++++---- .../knowledge/documents/document-processor.ts | 5 +- .../stored-artifact-extension.test.ts | 60 ++- apps/sim/lib/knowledge/search/queries.test.ts | 6 +- apps/sim/lib/knowledge/search/queries.ts | 49 ++- 9 files changed, 411 insertions(+), 102 deletions(-) diff --git a/apps/sim/lib/file-parsers/index.ts b/apps/sim/lib/file-parsers/index.ts index 7c4cc43a512..03bf72c26b8 100644 --- a/apps/sim/lib/file-parsers/index.ts +++ b/apps/sim/lib/file-parsers/index.ts @@ -176,7 +176,7 @@ export async function parseBuffer( } const kind = sniffFileKind(buffer, normalizedExtension) - const route = reconcileParserRoute(normalizedExtension, kind) + const route = reconcileParserRoute(normalizedExtension, kind, options) const parser = PARSERS.get(route.extension) if (!parser?.parseBuffer) { diff --git a/apps/sim/lib/file-parsers/sniff.test.ts b/apps/sim/lib/file-parsers/sniff.test.ts index 5e94bc89698..e79f0a575c0 100644 --- a/apps/sim/lib/file-parsers/sniff.test.ts +++ b/apps/sim/lib/file-parsers/sniff.test.ts @@ -308,6 +308,38 @@ describe('parseBuffer reconciles the extension with the sniffed bytes', () => { expect(result.metadata?.detectedType).toBe('html') }) + it.each([ + '
', + '{\\rtf1\\ansi Source-format example}', + ])( + 'preserves textual markup when the caller supplies a canonical text artifact', + async (content) => { + const result = await parseBuffer(Buffer.from(content), 'txt', { textMode: 'literal' }) + + expect(result.content).toBe(content) + expect(result.metadata?.detectedType).toBeUndefined() + } + ) + + it('keeps literal-text handling scoped to txt artifacts', async () => { + const result = await parseBuffer( + Buffer.from('

Readable page

'), + 'html', + { textMode: 'literal' } + ) + + expect(result.content).toContain('Readable page') + expect(result.content).not.toContain('') + await expect( + parseBuffer(Buffer.from('403 Forbidden'), 'json', { + textMode: 'literal', + }) + ).rejects.toMatchObject({ code: 'invalid_format' }) + await expect(parseBuffer(oleBinary(), 'txt', { textMode: 'literal' })).rejects.toMatchObject({ + code: 'invalid_format', + }) + }) + it('extracts a docx labelled .xlsx through the Word parser', async () => { const result = await parseBuffer(await buildDocx('Office Relocation'), 'xlsx') diff --git a/apps/sim/lib/file-parsers/sniff.ts b/apps/sim/lib/file-parsers/sniff.ts index 365b4c4b66f..a976b97e22b 100644 --- a/apps/sim/lib/file-parsers/sniff.ts +++ b/apps/sim/lib/file-parsers/sniff.ts @@ -1,5 +1,6 @@ import { FileParserError } from '@/lib/file-parsers/errors' import { isEncryptedOoxmlContainer } from '@/lib/file-parsers/ooxml-encryption' +import type { FileParseOptions } from '@/lib/file-parsers/types' import { decodeTextBuffer, detectBomlessUtf16 } from '@/lib/file-parsers/utils' import { isZipShaped } from '@/lib/file-parsers/zip-guard' @@ -307,7 +308,18 @@ function invalidFormat(extension: string, kind: SniffedKind): FileParserError { * text (as CSV under a spreadsheet extension), and an OLE2 file under a modern * Word extension is the legacy `.doc` parser's job. Legacy `.ppt` has no reader. */ -export function reconcileParserRoute(extension: string, kind: SniffedKind): ParserRoute { +export function reconcileParserRoute( + extension: string, + kind: SniffedKind, + options: Pick = {} +): ParserRoute { + if ( + extension === 'txt' && + options.textMode === 'literal' && + (kind === 'html' || kind === 'rtf') + ) { + return { extension } + } if (kind === 'rtf') { throw new FileParserError( 'unsupported_type', diff --git a/apps/sim/lib/file-parsers/types.ts b/apps/sim/lib/file-parsers/types.ts index 36b059a3bc7..028e2e7382b 100644 --- a/apps/sim/lib/file-parsers/types.ts +++ b/apps/sim/lib/file-parsers/types.ts @@ -33,6 +33,8 @@ export interface FileParseResult { export interface FileParseOptions { signal?: AbortSignal + /** Preserve textual markup in a canonical .txt artifact instead of interpreting it as HTML or RTF. */ + textMode?: 'literal' /** Complete PDF extraction rejects safety limits instead of returning preview text. */ pdfTextMode?: 'preview' | 'complete' } diff --git a/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts b/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts index 899982f5b78..eb0e16b9eee 100644 --- a/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts @@ -54,7 +54,6 @@ vi.hoisted(() => { if (process.env.KNOWLEDGE_SEARCH_PERFORMANCE_TEST === 'true') { Object.assign(process.env, { OPENAI_API_KEY: 'isolated-embedding-http-fixture', - GEMINI_API_KEY: 'isolated-gemini-http-fixture', CONFLUENCE_CLIENT_ID: 'isolated-confluence-fixture-client', CONFLUENCE_CLIENT_SECRET: 'isolated-confluence-fixture-secret', }) @@ -64,7 +63,12 @@ vi.hoisted(() => { const externalFetch = globalThis.fetch const enabled = process.env.KNOWLEDGE_SEARCH_PERFORMANCE_TEST === 'true' const chunkCount = Number(process.env.KNOWLEDGE_SEARCH_PERFORMANCE_CHUNKS ?? 20_000) +const unrelatedChunkCount = Number( + process.env.KNOWLEDGE_SEARCH_PERFORMANCE_UNRELATED_CHUNKS ?? chunkCount / 2 +) +const evictSharedBuffers = process.env.KNOWLEDGE_SEARCH_PERFORMANCE_EVICT_BUFFERS === 'true' const dimensions = 1536 +const candidateDimensions = 512 const chunksPerDocument = 4 const batchSize = 1000 const logger = createLogger('SearchLatencyIntegration') @@ -84,7 +88,11 @@ function readFixtureReport(file: string) { /** Captured SQL plans include repeated high-dimensional query parameters. */ if (statSync(file).size > 64 * 1024 * 1024) throw new Error('Fixture report exceeds 64 MiB') return z - .object({ fixture: fixtureSchema, unrelatedFixture: fixtureSchema }) + .object({ + fixture: fixtureSchema, + unrelatedFixture: fixtureSchema, + method: z.object({ fixtureVersion: z.literal(2) }), + }) .parse(JSON.parse(readFileSync(file, 'utf8'))) } const reused = reuseFile ? readFixtureReport(reuseFile) : undefined @@ -93,7 +101,11 @@ const unrelated = reused?.unrelatedFixture ?? createKnowledgeAclFixtureIds() const organizationChatId = generateId() function topicVector(topic = 0) { const vector = Array.from({ length: dimensions }, (_, index) => - Math.sin((index + 1) * (topic + 1) * 12.9898) + Math.sin( + (((index * 137 + Math.floor(index / candidateDimensions) * 57) % candidateDimensions) + 1) * + (topic + 1) * + 12.9898 + ) ) const magnitude = Math.hypot(...vector) return vector.map((value) => value / magnitude) @@ -104,15 +116,23 @@ const report: Record = { fixture: ids, unrelatedFixture: unrelated, method: { + fixtureVersion: 2, chunkCount, + unrelatedChunkCount, dimensions, + candidateDimensions, chunksPerDocument, sql: 'Captured from the real Assistant tool; no hand-written search query', providers: 'Embedding and source-permission HTTP responses are controlled; internal search and authorization code is real', vectors: - 'Normalized topic clusters with deterministic dense noise; not semantic-quality evaluation', - cache: 'First and repeated samples; no claim of a cold operating-system cache', + 'Normalized 512-dimensional topic/noise geometry with permuted copies across 1536 dimensions; verifies prefix candidate ranking, not semantic embedding quality', + cache: evictSharedBuffers + ? 'Organization samples evict PostgreSQL shared buffers before each request; operating-system cache is not cleared' + : 'First and repeated samples; no claim of a cold operating-system cache', + layout: reused + ? 'Reused fixture; physical layout is inherited from its original report' + : 'Tenant batches interleaved; chunks permuted across document identities', }, } let capture = false @@ -131,8 +151,13 @@ interface ExplainNode { 'Node Type': string 'Actual Rows': number 'Actual Loops': number + 'Plan Rows'?: number + 'Shared Hit Blocks'?: number + 'Shared Read Blocks'?: number 'Index Name'?: string 'Relation Name'?: string + 'Subplan Name'?: string + 'CTE Name'?: string Output?: string[] Plans?: ExplainNode[] } @@ -143,8 +168,13 @@ const explainNodeSchema: z.ZodType = z.lazy(() => 'Node Type': z.string(), 'Actual Rows': z.number(), 'Actual Loops': z.number(), + 'Plan Rows': z.number().optional(), + 'Shared Hit Blocks': z.number().optional(), + 'Shared Read Blocks': z.number().optional(), 'Index Name': z.string().optional(), 'Relation Name': z.string().optional(), + 'Subplan Name': z.string().optional(), + 'CTE Name': z.string().optional(), Output: z.array(z.string()).optional(), Plans: z.array(explainNodeSchema).optional(), }) @@ -160,6 +190,43 @@ function assertCompactCandidates(node: ExplainNode) { for (const child of node.Plans ?? []) assertCompactCandidates(child) } +function explainNodes(node: ExplainNode): ExplainNode[] { + return [node, ...(node.Plans ?? []).flatMap(explainNodes)] +} + +/** Broad ranking must stop the ordered ANN scan instead of sorting every accessible chunk. */ +function assertIndexedCandidates(plan: ExplainNode, candidateLimit: number) { + const nodes = explainNodes(plan) + const initial = nodes.find((node) => node['Subplan Name'] === 'CTE initial_candidates') + expect(initial).toBeDefined() + const candidateNodes = explainNodes(initial!) + expect( + candidateNodes.some( + (node) => + node['Index Name'] === 'embedding_search_512_cosine_hnsw_idx' && node['Actual Loops'] > 0 + ) + ).toBe(true) + expect(candidateNodes.some((node) => node['Node Type'] === 'Sort')).toBe(false) + expect( + candidateNodes.some((node) => node['Index Name'] === 'embedding_search_document_lookup_idx') + ).toBe(false) + const filtered = nodes.find((node) => node['Subplan Name'] === 'CTE filtered_scores') + expect(filtered).toBeDefined() + expect(filtered!.Output).toHaveLength(3) + expect(filtered!.Output![2]).toContain('<=>') + if (initial!['Actual Rows'] >= candidateLimit) { + for (const node of nodes.filter( + (item) => + item['Subplan Name'] === 'CTE visible_search_documents' || + item['CTE Name'] === 'visible_search_documents' || + item['Subplan Name'] === 'CTE filtered_scores' || + item['CTE Name'] === 'filtered_scores' + )) { + expect(node['Actual Loops']).toBe(0) + } + } +} + /** Small scopes must seek chunk metadata by document without reading the full vector projection. */ function assertIndexedChunkProbe(node: ExplainNode): number { let lookups = 0 @@ -186,12 +253,30 @@ function saveReport() { if (file) writeFileSync(file, JSON.stringify(report, null, 2), { mode: 0o600 }) } +/** Only the disposable fixture may evict shared buffers; the operating-system cache stays intact. */ +async function prepareOrganizationSample(label: string) { + if (!evictSharedBuffers) return + const [eviction] = await db.execute<{ buffers: number; evicted: number }>(sql` + WITH cached AS MATERIALIZED ( + SELECT bufferid FROM pg_buffercache + WHERE reldatabase = (SELECT oid FROM pg_database WHERE datname = current_database()) + ) SELECT count(*)::int AS buffers, + count(*) FILTER (WHERE pg_buffercache_evict(bufferid))::int AS evicted + FROM cached + `) + report[`${label}.sharedBufferEviction`] = eviction + saveReport() +} + const diagnosticSchema = z .object({ surface: z.enum(['dashboard', 'copilot']), outcome: z.enum(['success', 'partial']), elapsedMs: z.number(), vectorBudgetMs: z.number().positive(), + vectorCandidateDimensions: z.number().optional(), + vectorCandidateLimit: z.number().optional(), + vectorCandidateScan: z.enum(['planned', 'filtered']).optional(), retrievalStatus: z.enum(['complete', 'partial']), timedOutLegs: z.array(z.enum(['vector', 'keyword', 'tags'])), toolResultBytes: z.number().int().nonnegative().optional(), @@ -223,11 +308,12 @@ async function search( userId = ids.aliceId, query = 'Orion deployment', filters: WorkspaceSearchFilters = {}, - organizationScope = false + organizationScope = false, + topK = 15 ) { return resultSchema.parse( await searchWorkspaceServerTool.execute( - { query, topK: 15, ...filters }, + { query, topK, ...filters }, { userId, ...(organizationScope @@ -245,7 +331,12 @@ async function search( ) } -async function searchDashboard(query = 'Orion deployment', userId = ids.aliceId) { +async function searchDashboard( + query = 'Orion deployment', + userId = ids.aliceId, + organizationScope = false, + topK = 15 +) { const authenticate = vi.spyOn(internalSessionAuth, 'authenticate').mockResolvedValue({ kind: 'session', userId, @@ -257,9 +348,11 @@ async function searchDashboard(query = 'Orion deployment', userId = ids.aliceId) method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ - workspaceId: ids.workspaceId, + ...(organizationScope + ? { organizationId: ids.organizationId } + : { workspaceId: ids.workspaceId }), query, - topK: 15, + topK, }), }) ) @@ -285,7 +378,11 @@ function expectCompleteVectorSearch(diagnostics: z.infer ReturnType) { +async function sample( + label: string, + run: () => ReturnType, + options: { explain?: boolean } = {} +) { captured.length = 0 diagnosticLog?.mockClear() const start = performance.now() @@ -328,8 +425,22 @@ async function sample(label: string, run: () => ReturnType) { item.query.includes('WITH scored_search_candidates') || item.query.includes('WITH visible_keyword_documents')) ) - const plans = [] - for (const query of searches) { + const plans: Array< + CapturedQuery & { + kind: 'keyword' | 'vector' | 'rerank' | 'probe' + plan: z.infer + } + > = [] + report[label] = { + milliseconds, + diagnostics, + queryCount: captured.length, + resultCount: result.data.results.length, + explainsDeferred: options.explain === false, + plans, + } + saveReport() + for (const query of options.explain === false ? [] : searches) { const plan = await db.$client.begin(async (tx) => { await tx.unsafe("SET LOCAL statement_timeout = '45s'") await tx.unsafe('SET LOCAL jit = off') @@ -339,8 +450,7 @@ async function sample(label: string, run: () => ReturnType) { query.query.includes('WITH visible_search_documents') || (query.query.includes('from "embedding_search"') && query.query.includes('order by')) ) { - await tx.unsafe('SET LOCAL hnsw.max_scan_tuples = 1000') - await tx.unsafe('SET LOCAL hnsw.ef_search = 1000') + await tx.unsafe('SET LOCAL hnsw.ef_search = 100') await tx.unsafe('SET LOCAL hnsw.scan_mem_multiplier = 2') } return tx.unsafe( @@ -349,9 +459,6 @@ async function sample(label: string, run: () => ReturnType) { ) }) const parsedPlan = explainSchema.parse(plan[0]['QUERY PLAN']) - if (query.query.includes('WITH visible_keyword_documents')) { - assertScalarKeywordSorts(parsedPlan[0].Plan) - } plans.push({ kind: query.query.includes('keyword_rank') ? 'keyword' @@ -366,15 +473,17 @@ async function sample(label: string, run: () => ReturnType) { parameters: query.parameters, plan: parsedPlan, }) + saveReport() + if (query.query.includes('WITH visible_search_documents')) { + expect(query.query).toContain('"embedding_search"."vector_512"') + expect(diagnostics.vectorCandidateDimensions).toBe(candidateDimensions) + expect(diagnostics.vectorCandidateLimit).toBeGreaterThan(0) + assertIndexedCandidates(parsedPlan[0].Plan, diagnostics.vectorCandidateLimit!) + } + if (query.query.includes('WITH visible_keyword_documents')) { + assertScalarKeywordSorts(parsedPlan[0].Plan) + } } - report[label] = { - milliseconds, - diagnostics, - queryCount: captured.length, - resultCount: result.data.results.length, - plans, - } - saveReport() logger.info(label, { milliseconds, queryCount: captured.length, @@ -386,13 +495,13 @@ async function sample(label: string, run: () => ReturnType) { describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpus', () => { beforeAll(async () => { if ( - !Number.isInteger(chunkCount) || - chunkCount < 10_000 || - chunkCount > 200_000 || - chunkCount % batchSize !== 0 + [chunkCount, unrelatedChunkCount].some( + (count) => + !Number.isInteger(count) || count < 5_000 || count > 200_000 || count % batchSize !== 0 + ) ) throw new Error( - 'KNOWLEDGE_SEARCH_PERFORMANCE_CHUNKS must be a multiple of 1000 from 10000 to 200000' + 'Search performance chunk counts must be multiples of 1000 from 5000 to 200000' ) vi.stubGlobal('fetch', async (input: string | URL | Request, init?: RequestInit) => { const url = input instanceof Request ? input.url : String(input) @@ -405,33 +514,14 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu ? new Response(null, { status: 403 }) : Response.json({ type: 'known', accountId: ids.aliceId }) } - if ( - url === - 'https://generativelanguage.googleapis.com/v1beta/models/gemini-embedding-001:batchEmbedContents' - ) { - const body = z - .object({ - requests: z - .array( - z.object({ - content: z.object({ parts: z.array(z.object({ text: z.string() })).length(1) }), - }) - ) - .length(1), - }) - .parse(JSON.parse(String(init?.body))) - embeddingCalls++ - const text = body.requests[0].content.parts[0].text - const topic = Number(/^Topic (\d+) deployment$/.exec(text)?.[1] ?? 0) - return Response.json({ - embeddings: [{ values: topicVector(topic) }], - usageMetadata: { promptTokenCount: 4 }, - }) - } if (url !== 'https://api.openai.com/v1/embeddings') throw new Error(`Unexpected outbound request in search fixture: ${new URL(url).origin}`) const body = z - .object({ input: z.array(z.string()).length(1), encoding_format: z.literal('base64') }) + .object({ + input: z.array(z.string()).length(1), + encoding_format: z.literal('base64'), + model: z.literal('text-embedding-3-small'), + }) .parse(JSON.parse(String(init?.body))) embeddingCalls += body.input.length const bytes = Buffer.alloc(dimensions * 4) @@ -457,14 +547,14 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu const [size] = await db.execute<{ count: number }>( sql`SELECT count(*)::int AS count FROM embedding WHERE knowledge_base_id = ${fixture.knowledgeBaseId}` ) - expect(size.count).toBe(fixture === ids ? chunkCount : chunkCount / 2) + expect(size.count).toBe(fixture === ids ? chunkCount : unrelatedChunkCount) } await db .update(knowledgeBase) .set({ workspaceId: ids.workspaceId, organizationId: null, - embeddingModel: 'gemini-embedding-001', + embeddingModel: 'text-embedding-3-small', }) .where(eq(knowledgeBase.id, ids.knowledgeBaseId)) await db @@ -490,10 +580,10 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu } else { await seedKnowledgeAclFixture(ids, { connectorType: 'google_drive' }) await seedKnowledgeAclFixture(unrelated, { connectorType: 'google_drive' }) - /** These arbitrary dense vectors are not trained for prefix shortening. */ + /** The controlled geometry preserves prefix distances for the production 512-dimensional path. */ await db .update(knowledgeBase) - .set({ embeddingModel: 'gemini-embedding-001' }) + .set({ embeddingModel: 'text-embedding-3-small' }) .where(inArray(knowledgeBase.id, [ids.knowledgeBaseId, unrelated.knowledgeBaseId])) await db .update(knowledgeBase) @@ -507,7 +597,7 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu for (const index of indexes) await db.execute(sql`DROP INDEX ${sql.identifier(index.indexname)}`) for (const fixture of [ids, unrelated]) { - const count = fixture === ids ? chunkCount : chunkCount / 2 + const count = fixture === ids ? chunkCount : unrelatedChunkCount for (let first = 0; first < count / chunksPerDocument; first += batchSize) { const last = Math.min(first + batchSize, count / chunksPerDocument) - 1 await db.execute(sql`INSERT INTO document @@ -517,7 +607,11 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu ARRAY[${`u:${fixture.aliceId}@fixture.test`}]::text[], statement_timestamp() FROM generate_series(${first}::int, ${last}::int) n`) } - for (let first = 0; first < count; first += batchSize) { + } + for (let first = 0; first < Math.max(chunkCount, unrelatedChunkCount); first += batchSize) { + for (const fixture of [ids, unrelated]) { + const count = fixture === ids ? chunkCount : unrelatedChunkCount + if (first >= count) continue const last = Math.min(first + batchSize, count) - 1 await db.transaction(async (tx) => { await tx.execute(sql`SET LOCAL jit = off`) @@ -530,22 +624,37 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu 3000, 750, 0, 3000, l2_normalize(ARRAY(SELECT (sin(coordinate * (n % 32 + 1) * 12.9898) + 0.25 * sin(n::double precision * coordinate * 12.9898 + coordinate * 78.233))::real - FROM generate_series(1, ${dimensions}) coordinate)::vector(1536)) - FROM generate_series(${first}::int, ${last}::int) n`) + FROM ( + SELECT (((position - 1) * 137 + ((position - 1) / ${candidateDimensions}) * 57) + % ${candidateDimensions}) + 1 AS coordinate + FROM generate_series(1, ${dimensions}) position + ) coordinates)::vector(1536)) + FROM ( + SELECT (ordinal * 7919) % ${count} AS n + FROM generate_series(${first}::int, ${last}::int) ordinal + ) shuffled`) }) } - logger.info('Synthetic corpus loaded', { chunks: count }) } + logger.info('Synthetic corpora loaded', { chunkCount, unrelatedChunkCount }) for (const index of indexes) await db.execute(sql.raw(index.indexdef)) } await db.execute(sql`ANALYZE document`) await db.execute(sql`ANALYZE embedding`) await db.execute(sql`ANALYZE embedding_search`) await db.execute(sql`ANALYZE embedding_keyword_search`) + if (evictSharedBuffers) await db.execute(sql`CREATE EXTENSION IF NOT EXISTS pg_buffercache`) report.server = ( await db.execute(sql`SELECT version(), current_setting('work_mem') AS work_mem, (SELECT extversion FROM pg_extension WHERE extname = 'vector') AS pgvector`) )[0] + report.relations = await db.execute(sql` + SELECT relname, pg_relation_size(oid) AS bytes + FROM pg_class + WHERE relname IN ('embedding', 'embedding_search', 'document') + OR relname LIKE 'embedding_search%hnsw_idx' + ORDER BY relname + `) db.$client.options.debug = (_connection, query, parameters) => { if (capture && captured.length < 300) captured.push({ query, parameters: [...parameters] }) } @@ -900,7 +1009,7 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu } }, 180_000) - it.each([200, 396, 400])( + it.each([200, 396, 400, 1000, 2000])( 'keeps a selective scope of %s chunks within both retrieval budgets', async (count) => { const documentCount = count / chunksPerDocument @@ -927,9 +1036,30 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu true ) const probe = plans.find((plan) => plan.kind === 'probe')! - expect(probe.plan[0].Plan['Actual Rows']).toBe(count) - expect(assertIndexedChunkProbe(probe.plan[0].Plan)).toBe(documentCount) + expect(probe.plan[0].Plan['Actual Rows']).toBe(Math.min(count, 400)) + expect(assertIndexedChunkProbe(probe.plan[0].Plan)).toBe(Math.min(documentCount, 100)) expect(plans.filter((plan) => plan.kind === 'vector')).toHaveLength(count < 400 ? 0 : 1) + if (count > 400) { + const rerank = plans.find((plan) => plan.kind === 'rerank')! + const actual = await db.$client.unsafe(rerank.query, rerank.parameters).values() + const expected = await db.execute<{ id: string }>(sql`SELECT id FROM embedding + WHERE knowledge_base_id = ${ids.knowledgeBaseId} AND enabled + AND document_id IN (${sql.join( + documentIds.map((id) => sql`${id}`), + sql`, ` + )}) + ORDER BY (embedding <=> ${JSON.stringify(queryVector)}::vector) + 0, id + LIMIT ${actual.length}`) + const expectedIds = new Set(expected.map(({ id }) => id)) + const recall = actual.filter(([id]) => expectedIds.has(id)).length / expected.length + expect(recall).toBeGreaterThanOrEqual(0.95) + report[`recall.selective-${count}.${surface}`] = { + neighbors: expected.length, + recall, + candidateScan: diagnostics.vectorCandidateScan, + } + saveReport() + } } } finally { await db @@ -1083,7 +1213,7 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu expect(restored.result.data.results).toHaveLength(15) }, 180_000) - it('uses the same indexed retrieval through a persisted private organization Assistant chat', async () => { + it('keeps organization searches complete with stale ACL estimates and concurrent requests', async () => { await db.insert(member).values({ id: generateId(), organizationId: ids.organizationId, @@ -1104,17 +1234,72 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu .update(knowledgeConnector) .set({ connectorType: 'google_drive', credentialId: null, sourceConfig: {} }) .where(eq(knowledgeConnector.id, ids.connectorId)) - await db.execute( - sql`UPDATE document SET acl = ARRAY[${`u:${ids.aliceId}@fixture.test`}] WHERE knowledge_base_id = ${ids.knowledgeBaseId}` - ) - await db.execute(sql`ANALYZE document`) - const { result } = await sample('organization', () => - search(ids.aliceId, 'Orion deployment', {}, true) - ) - expect(result.data.results).toHaveLength(15) - expect(result.data.results.every((row) => row.knowledgeBaseId === ids.knowledgeBaseId)).toBe( - true - ) + /** Keep the deliberate tenfold visibility underestimate until the measured requests finish. */ + await db.execute(sql`ALTER TABLE document SET (autovacuum_enabled = false)`) + try { + await db.execute(sql`UPDATE document + SET acl = ARRAY[CASE WHEN external_id::int % 10 = 0 + THEN ${`u:${ids.aliceId}@fixture.test`} ELSE ${`u:${ids.bobId}@fixture.test`} END] + WHERE knowledge_base_id = ${ids.knowledgeBaseId}`) + await db.execute(sql`ANALYZE document`) + await db.execute(sql`UPDATE document SET acl = ARRAY[${`u:${ids.aliceId}@fixture.test`}] + WHERE knowledge_base_id = ${ids.knowledgeBaseId}`) + report.organizationVisibility = { + analyzedVisibleDocuments: chunkCount / chunksPerDocument / 10, + actualVisibleDocuments: chunkCount / chunksPerDocument, + unrelatedChunks: unrelatedChunkCount, + } + /** Capture latency samples before EXPLAIN ANALYZE can warm the candidate paths. */ + for (const surface of ['dashboard', 'copilot'] as const) { + const label = `organization.${surface}` + await prepareOrganizationSample(label) + const { result, diagnostics } = await sample( + label, + () => + surface === 'dashboard' + ? searchDashboard('Orion deployment', ids.aliceId, true, 20) + : search(ids.aliceId, 'Orion deployment', {}, true, 20), + { explain: false } + ) + expectCompleteVectorSearch(diagnostics) + expect(result.data.results).toHaveLength(20) + expect( + result.data.results.every((row) => row.knowledgeBaseId === ids.knowledgeBaseId) + ).toBe(true) + } + await prepareOrganizationSample('organization.concurrent') + diagnosticLog?.mockClear() + const started = performance.now() + const results = await Promise.all([ + search(ids.aliceId, 'Orion deployment', {}, true, 20), + search(ids.aliceId, 'Engineering operations', {}, true, 20), + ]) + const diagnostics = diagnosticLog!.mock.calls + .filter(([message]) => message === 'Knowledge search completed') + .map(([, metadata]) => diagnosticSchema.parse(metadata)) + report['organization.concurrent'] = { + milliseconds: performance.now() - started, + resultCounts: results.map((result) => result.data.results.length), + diagnostics, + } + saveReport() + expect(diagnostics).toHaveLength(2) + for (const item of diagnostics) expectCompleteVectorSearch(item) + for (const result of results) { + expect(result.data.results).toHaveLength(20) + expect( + result.data.results.every((row) => row.knowledgeBaseId === ids.knowledgeBaseId) + ).toBe(true) + } + const planned = await sample('organization.plans', () => + search(ids.aliceId, 'Orion deployment', {}, true, 20) + ) + expectCompleteVectorSearch(planned.diagnostics) + expect(planned.plans.filter((plan) => plan.kind === 'vector')).toHaveLength(1) + } finally { + await db.execute(sql`ALTER TABLE document RESET (autovacuum_enabled)`) + await db.execute(sql`ANALYZE document`) + } }, 180_000) /** Opt in with local Sim and Go URLs; uses the real configured provider, billing adapter, and async resume protocol. */ it.skipIf(!process.env.KNOWLEDGE_SEARCH_ASSISTANT_URL)( diff --git a/apps/sim/lib/knowledge/documents/document-processor.ts b/apps/sim/lib/knowledge/documents/document-processor.ts index a0691b43220..b285dd68d09 100644 --- a/apps/sim/lib/knowledge/documents/document-processor.ts +++ b/apps/sim/lib/knowledge/documents/document-processor.ts @@ -1350,10 +1350,11 @@ async function parseHttpFile( access.signal?.throwIfAborted() /** Prefer what we actually downloaded over what the document is *called*. */ - const extension = - resolveStoredArtifactExtension(fileUrl) ?? resolveParserExtension(filename, mimeType) + const storedExtension = resolveStoredArtifactExtension(fileUrl) + const extension = storedExtension ?? resolveParserExtension(filename, mimeType) const result = await parseBuffer(buffer, extension, { signal: access.signal, + textMode: storedExtension === 'txt' ? 'literal' : undefined, pdfTextMode: extension === 'pdf' ? 'complete' : undefined, }) return result diff --git a/apps/sim/lib/knowledge/documents/stored-artifact-extension.test.ts b/apps/sim/lib/knowledge/documents/stored-artifact-extension.test.ts index b30446cd547..18267da5b2a 100644 --- a/apps/sim/lib/knowledge/documents/stored-artifact-extension.test.ts +++ b/apps/sim/lib/knowledge/documents/stored-artifact-extension.test.ts @@ -9,7 +9,13 @@ * SharePoint PDFs with `Invalid PDF structure.` and silently double-wrapped * every spreadsheet, which "succeeded" because SheetJS accepts almost anything. */ -import { describe, expect, it } from 'vitest' +import { describe, expect, it, vi } from 'vitest' + +const { mockDownload } = vi.hoisted(() => ({ mockDownload: vi.fn() })) + +vi.mock('@/lib/uploads/utils/file-utils.server', () => ({ downloadFileFromUrl: mockDownload })) + +import { processDocument } from '@/lib/knowledge/documents/document-processor' import { resolveStoredArtifactExtension } from '@/lib/knowledge/documents/parser-extension' const CONNECTOR_PDF_URL = @@ -86,3 +92,55 @@ describe('resolveStoredArtifactExtension', () => { expect(resolveStoredArtifactExtension('/api/files/serve/s3/kb%2F1-a-Report.PDF')).toBe('pdf') }) }) + +describe('stored text document processing', () => { + const source = ` + + + + Application shell + + + +
+` + + it('indexes the source of an HTML shell stored as connector text', async () => { + mockDownload.mockResolvedValue(Buffer.from(source)) + + const result = await processDocument( + '/api/files/serve/s3/kb%2Ffixture-index.html.txt?context=knowledge-base', + 'index.html', + 'text/plain' + ) + + expect(result.chunks).toHaveLength(1) + expect(result.chunks[0].text).toContain('') + expect(result.chunks[0].text).toContain('
') + expect(result.metadata.characterCount).toBe(source.length) + }) + + it('keeps an actual HTML document on rendered-text extraction', async () => { + mockDownload.mockResolvedValue(Buffer.from(source)) + + await expect( + processDocument( + '/api/files/serve/s3/kb%2Ffixture-page.html?context=knowledge-base', + 'page.html', + 'text/html' + ) + ).rejects.toMatchObject({ code: 'no_extractable_text' }) + }) + + it('still rejects a stored text artifact containing only whitespace', async () => { + mockDownload.mockResolvedValue(Buffer.from(' \n\t ')) + + await expect( + processDocument( + '/api/files/serve/s3/kb%2Ffixture-blank.txt?context=knowledge-base', + 'blank.txt', + 'text/plain' + ) + ).rejects.toMatchObject({ code: 'no_extractable_text' }) + }) +}) diff --git a/apps/sim/lib/knowledge/search/queries.test.ts b/apps/sim/lib/knowledge/search/queries.test.ts index 1a6170a0a33..a5f20dc7f20 100644 --- a/apps/sim/lib/knowledge/search/queries.test.ts +++ b/apps/sim/lib/knowledge/search/queries.test.ts @@ -591,7 +591,7 @@ describe('live repository authorization follows ranked candidates', () => { rerankPages.length = 0 keywordPages.length = 0 dbChainMockFns.execute.mockImplementation(async (query) => - render(query).sql.includes('CROSS JOIN LATERAL') + render(query).sql.includes('SELECT scoped_chunk.id') ? (probePages.shift() ?? []) : render(query).sql.includes('WITH visible_search_documents') ? (candidatePages.shift() ?? []) @@ -630,6 +630,8 @@ describe('live repository authorization follows ranked candidates', () => { render(query).sql.includes('WITH visible_search_documents') )![0] expect(render(candidateQuery).sql).toContain('MATERIALIZED') + expect(render(candidateQuery).sql).toContain('CROSS JOIN LATERAL') + expect(render(candidateQuery).sql).toContain('LIMIT 1') expect(JSON.stringify(candidateQuery)).toContain('required_clause') expect(JSON.stringify(candidateQuery)).toContain('subvector') const rankQuery = dbChainMockFns.execute.mock.calls.find(([query]) => @@ -721,6 +723,8 @@ describe('live repository authorization follows ranked candidates', () => { )![0] expect(render(candidateQuery).sql).toContain('UNION ALL') expect(render(candidateQuery).sql).toContain('+ 0') + expect(render(candidateQuery).sql).toContain('filtered_scores AS MATERIALIZED') + expect(render(candidateQuery).sql).toContain('ORDER BY filtered_scores.distance + 0') expect(JSON.stringify(dbChainMockFns.where.mock.calls.at(-1)![0])).toContain( 'github_read_grant' ) diff --git a/apps/sim/lib/knowledge/search/queries.ts b/apps/sim/lib/knowledge/search/queries.ts index 46fe5069783..5cca4482d2c 100644 --- a/apps/sim/lib/knowledge/search/queries.ts +++ b/apps/sim/lib/knowledge/search/queries.ts @@ -49,9 +49,8 @@ const logger = createLogger('KnowledgeSearchQueries') const UNDEFINED_OBJECT_SQLSTATE = '42704' /** Tuples a relaxed-order scan may visit before giving up on filling the limit. */ const HNSW_MAX_SCAN_TUPLES = '20000' -/** Stop a permission-starved graph walk early enough to scan the filtered projection instead. */ -const CANDIDATE_HNSW_MAX_SCAN_TUPLES = '1000' -const CANDIDATE_HNSW_EF_SEARCH = '1000' +/** Iterative scans expand a modest initial neighborhood when permission filters reject neighbors. */ +const CANDIDATE_HNSW_EF_SEARCH = '100' const CANDIDATE_HNSW_SCAN_MEM_MULTIPLIER = '2' const MIN_VECTOR_RERANK_CANDIDATES = 400 const MAX_VECTOR_RERANK_CANDIDATES = 1600 @@ -86,7 +85,7 @@ async function withVectorScanSettings( await measureSearchStage('vector.settings', () => tx.execute( ranking === 'candidate' - ? sql`SELECT set_config('hnsw.iterative_scan', 'relaxed_order', true), set_config('hnsw.max_scan_tuples', ${CANDIDATE_HNSW_MAX_SCAN_TUPLES}, true), set_config('hnsw.ef_search', ${CANDIDATE_HNSW_EF_SEARCH}, true), set_config('hnsw.scan_mem_multiplier', ${CANDIDATE_HNSW_SCAN_MEM_MULTIPLIER}, true)` + ? sql`SELECT set_config('hnsw.iterative_scan', 'relaxed_order', true), set_config('hnsw.max_scan_tuples', ${HNSW_MAX_SCAN_TUPLES}, true), set_config('hnsw.ef_search', ${CANDIDATE_HNSW_EF_SEARCH}, true), set_config('hnsw.scan_mem_multiplier', ${CANDIDATE_HNSW_SCAN_MEM_MULTIPLIER}, true)` : sql`SELECT set_config('hnsw.iterative_scan', 'relaxed_order', true), set_config('hnsw.max_scan_tuples', ${HNSW_MAX_SCAN_TUPLES}, true)` ) ) @@ -813,7 +812,8 @@ export async function handleVectorOnlySearch(params: SearchParams): Promise executor.execute<{ id: string; initial_count: number }>(sql` @@ -934,16 +933,32 @@ async function selectLiveVectorResults( WHERE ${and(...candidateDocumentVisibility)} ), initial_candidates AS MATERIALIZED ( SELECT ${embeddingSearch.id} AS id FROM ${embeddingSearch} - WHERE ${candidateConditions} + CROSS JOIN LATERAL ( + SELECT 1 FROM ${document} + WHERE ${and(eq(document.id, embeddingSearch.documentId), ...candidateDocumentVisibility)} + LIMIT 1 + ) AS visible + WHERE ${and( + inArray(embeddingSearch.knowledgeBaseId, params.knowledgeBaseIds), + eq(embeddingSearch.enabled, true) + )} ORDER BY ${candidateDistance} LIMIT ${candidateLimit} + ), filtered_scores AS MATERIALIZED ( + SELECT ${embeddingSearch.id} AS id, ${embeddingSearch.documentId} AS document_id, + ${candidateDistance} AS distance FROM ${embeddingSearch} + WHERE ${and( + inArray(embeddingSearch.knowledgeBaseId, params.knowledgeBaseIds), + eq(embeddingSearch.enabled, true) + )} + AND (SELECT count(*) FROM initial_candidates) < ${candidateLimit} ), candidates AS ( SELECT id FROM initial_candidates WHERE (SELECT count(*) FROM initial_candidates) >= ${candidateLimit} UNION ALL ( - SELECT ${embeddingSearch.id} AS id FROM ${embeddingSearch} - WHERE ${candidateConditions} - AND (SELECT count(*) FROM initial_candidates) < ${candidateLimit} - ORDER BY (${candidateDistance}) + 0, ${embeddingSearch.id} + SELECT filtered_scores.id FROM filtered_scores + INNER JOIN visible_search_documents ON visible_search_documents.id = filtered_scores.document_id + WHERE (SELECT count(*) FROM initial_candidates) < ${candidateLimit} + ORDER BY filtered_scores.distance + 0, filtered_scores.id LIMIT ${candidateLimit} ) ) SELECT id, (SELECT count(*)::int FROM initial_candidates) AS initial_count FROM candidates From 045e2c806bdf41d70f84e8513ad9b4d44050a486 Mon Sep 17 00:00:00 2001 From: Vikhyath Mondreti Date: Wed, 16 Sep 2026 16:38:28 -0700 Subject: [PATCH 2/3] fix(search): keep benchmark corpus defaults within valid bounds --- .../__integration__/search-latency.integration.ts | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts b/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts index eb0e16b9eee..440f5cf70c0 100644 --- a/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts @@ -62,15 +62,17 @@ vi.hoisted(() => { const externalFetch = globalThis.fetch const enabled = process.env.KNOWLEDGE_SEARCH_PERFORMANCE_TEST === 'true' +const batchSize = 1000 +const MIN_CHUNK_COUNT = 5000 const chunkCount = Number(process.env.KNOWLEDGE_SEARCH_PERFORMANCE_CHUNKS ?? 20_000) const unrelatedChunkCount = Number( - process.env.KNOWLEDGE_SEARCH_PERFORMANCE_UNRELATED_CHUNKS ?? chunkCount / 2 + process.env.KNOWLEDGE_SEARCH_PERFORMANCE_UNRELATED_CHUNKS ?? + Math.max(MIN_CHUNK_COUNT, Math.ceil(chunkCount / (2 * batchSize)) * batchSize) ) const evictSharedBuffers = process.env.KNOWLEDGE_SEARCH_PERFORMANCE_EVICT_BUFFERS === 'true' const dimensions = 1536 const candidateDimensions = 512 const chunksPerDocument = 4 -const batchSize = 1000 const logger = createLogger('SearchLatencyIntegration') const fixtureSchema = z.object({ aliceId: z.uuid(), @@ -497,7 +499,10 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu if ( [chunkCount, unrelatedChunkCount].some( (count) => - !Number.isInteger(count) || count < 5_000 || count > 200_000 || count % batchSize !== 0 + !Number.isInteger(count) || + count < MIN_CHUNK_COUNT || + count > 200_000 || + count % batchSize !== 0 ) ) throw new Error( From 44cc97c906c6b316d49677071539239068eb7d13 Mon Sep 17 00:00:00 2001 From: Vikhyath Mondreti Date: Wed, 16 Sep 2026 17:00:22 -0700 Subject: [PATCH 3/3] fix(search): retain ANN recall settings and verify distinct queries --- .../__integration__/search-latency.integration.ts | 7 ++++--- apps/sim/lib/knowledge/search/queries.ts | 7 ++++--- 2 files changed, 8 insertions(+), 6 deletions(-) diff --git a/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts b/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts index 440f5cf70c0..0c89994ee4f 100644 --- a/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/search-latency.integration.ts @@ -452,7 +452,8 @@ async function sample( query.query.includes('WITH visible_search_documents') || (query.query.includes('from "embedding_search"') && query.query.includes('order by')) ) { - await tx.unsafe('SET LOCAL hnsw.ef_search = 100') + await tx.unsafe('SET LOCAL hnsw.max_scan_tuples = 1000') + await tx.unsafe('SET LOCAL hnsw.ef_search = 1000') await tx.unsafe('SET LOCAL hnsw.scan_mem_multiplier = 2') } return tx.unsafe( @@ -1185,7 +1186,7 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu it('runs two independent Assistant searches concurrently', async () => { diagnosticLog?.mockClear() const start = performance.now() - const results = await Promise.all([search(), search(ids.aliceId, 'Engineering operations')]) + const results = await Promise.all([search(), search(ids.aliceId, 'Topic 11 deployment')]) const completed = diagnosticLog!.mock.calls .filter(([message]) => message === 'Knowledge search completed') .map(([, metadata]) => diagnosticSchema.parse(metadata)) @@ -1277,7 +1278,7 @@ describe.skipIf(!enabled)('Assistant search latency on a realistic indexed corpu const started = performance.now() const results = await Promise.all([ search(ids.aliceId, 'Orion deployment', {}, true, 20), - search(ids.aliceId, 'Engineering operations', {}, true, 20), + search(ids.aliceId, 'Topic 11 deployment', {}, true, 20), ]) const diagnostics = diagnosticLog!.mock.calls .filter(([message]) => message === 'Knowledge search completed') diff --git a/apps/sim/lib/knowledge/search/queries.ts b/apps/sim/lib/knowledge/search/queries.ts index 5cca4482d2c..84ace644039 100644 --- a/apps/sim/lib/knowledge/search/queries.ts +++ b/apps/sim/lib/knowledge/search/queries.ts @@ -49,8 +49,9 @@ const logger = createLogger('KnowledgeSearchQueries') const UNDEFINED_OBJECT_SQLSTATE = '42704' /** Tuples a relaxed-order scan may visit before giving up on filling the limit. */ const HNSW_MAX_SCAN_TUPLES = '20000' -/** Iterative scans expand a modest initial neighborhood when permission filters reject neighbors. */ -const CANDIDATE_HNSW_EF_SEARCH = '100' +/** Stop a permission-starved graph walk early enough to scan the filtered projection instead. */ +const CANDIDATE_HNSW_MAX_SCAN_TUPLES = '1000' +const CANDIDATE_HNSW_EF_SEARCH = '1000' const CANDIDATE_HNSW_SCAN_MEM_MULTIPLIER = '2' const MIN_VECTOR_RERANK_CANDIDATES = 400 const MAX_VECTOR_RERANK_CANDIDATES = 1600 @@ -85,7 +86,7 @@ async function withVectorScanSettings( await measureSearchStage('vector.settings', () => tx.execute( ranking === 'candidate' - ? sql`SELECT set_config('hnsw.iterative_scan', 'relaxed_order', true), set_config('hnsw.max_scan_tuples', ${HNSW_MAX_SCAN_TUPLES}, true), set_config('hnsw.ef_search', ${CANDIDATE_HNSW_EF_SEARCH}, true), set_config('hnsw.scan_mem_multiplier', ${CANDIDATE_HNSW_SCAN_MEM_MULTIPLIER}, true)` + ? sql`SELECT set_config('hnsw.iterative_scan', 'relaxed_order', true), set_config('hnsw.max_scan_tuples', ${CANDIDATE_HNSW_MAX_SCAN_TUPLES}, true), set_config('hnsw.ef_search', ${CANDIDATE_HNSW_EF_SEARCH}, true), set_config('hnsw.scan_mem_multiplier', ${CANDIDATE_HNSW_SCAN_MEM_MULTIPLIER}, true)` : sql`SELECT set_config('hnsw.iterative_scan', 'relaxed_order', true), set_config('hnsw.max_scan_tuples', ${HNSW_MAX_SCAN_TUPLES}, true)` ) )