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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ describe('API-key KB block fan-out', () => {
it.each([false, true])(
'completes 18 concurrent KB searches with access checks intact (tag filter: %s)',
async (withTags) => {
/** The projection-fill memo outlives an iteration; each one must read it once, like a cold process. */
/** A cold process must also skip the global readiness probe for ordinary KBs. */
forgetProjectionFilled()
const previousDebug = db.$client.options.debug
const statements: string[] = []
Expand Down Expand Up @@ -135,20 +135,20 @@ describe('API-key KB block fan-out', () => {
statements.filter((query) => query.includes(fragment))
/**
* Every statement runs under the leg's deadline: the candidate search applies it with the
* scan settings in one statement, and the probe, the exact ranking and hydration each
* open with one of their own. The projection-fill read is shared by the searches that
* miss its memo together, so it appears once.
* scan settings in one statement, and the probe, exact ranking, document-backed page
* and hydration each open with one of their own.
*/
expect(matching('statement_timeout')).toHaveLength(bases.length * 4 + 1)
expect(matching('statement_timeout')).toHaveLength(bases.length * 5)
expect(matching('IS NOT NULL AS unfilled')).toHaveLength(0)
/**
* A scope this small leaves the bounded traversal short of its candidate limit, so every
* search probes once and rescues once — never a widening retry loop.
*/
expect(matching('hnsw.iterative_scan')).toHaveLength(bases.length)
expect(matching('AS visible')).toHaveLength(bases.length)
expect(matching(') + 0 LIMIT')).toHaveLength(bases.length)
/** The walk carries each candidate's identities, so a filled projection reads no page. */
expect(matching('"embedding_search"."id" = ANY(')).toHaveLength(0)
/** Ordinary KBs read page identities from documents without requiring a filled projection. */
expect(matching('"embedding_search"."id" = ANY(')).toHaveLength(bases.length)
/** The probe enumerates visible documents and reports saturation; it never ranks them. */
expect(
statements.filter(
Expand Down
203 changes: 120 additions & 83 deletions apps/sim/lib/knowledge/search/queries.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -441,6 +441,26 @@ describe('workspace-scoped vector retrieval', () => {
vi.useRealTimers()
})

it.each([false, undefined])(
'retrieves ordinary KB results without global projection readiness (searchIndexOnly=%s)',
async (searchIndexOnly) => {
queueTableRows(schemaMock.embedding, [...ranked].reverse())
const result = await retrieveKnowledgeSearch({
...params,
searchIndexOnly,
searchMode: 'vector',
query: 'What is the capital of France?',
})
expect(result.retrieval).toEqual({ status: 'complete', timedOutLegs: [] })
expect(result.rows.map((row) => row.id)).toEqual(['near', 'far'])
expect(statements().filter((query) => query.sql.includes('AS unfilled'))).toHaveLength(0)
expect(statements().some((query) => isPageStatement(query.sql))).toBe(true)
const walk = statements().find((query) => isWalk(query.sql))!
expect(JSON.stringify(walk)).toContain('required_clause')
expect(JSON.stringify(walk)).toContain(String(schemaMock.document.acl))
}
)

it.each([handleVectorOnlySearch, handleTagAndVectorSearch])(
'does not acquire a connection or start SQL after the KB retrieval deadline',
async (search) => {
Expand Down Expand Up @@ -752,7 +772,7 @@ describe('workspace-scoped vector retrieval', () => {
statements()
.filter((query) => query.sql.includes('statement_timeout'))
.map((query) => query.params[0])
).toEqual(['100', '100', '40', '20'])
).toEqual(['100', '40', '20'])
})

it('applies the scan settings in the deadline statement rather than one of their own', async () => {
Expand Down Expand Up @@ -808,49 +828,53 @@ describe('workspace-scoped vector retrieval', () => {
)
})

it('reports incomplete retrieval for 18 expired pool waiters without starting their SQL later', async () => {
vi.useFakeTimers()
const release: Array<() => void> = []
const transactions: Array<Promise<unknown>> = []
vi.spyOn(db, 'transaction').mockImplementation((callback) => {
const transaction = new Promise<void>((resolve) => release.push(resolve)).then(() =>
callback(db as never)
it.each([false, true])(
'expires pool waiters without starting SQL later (searchIndexOnly=%s)',
async (searchIndexOnly) => {
vi.useFakeTimers()
const release: Array<() => void> = []
const transactions: Array<Promise<unknown>> = []
vi.spyOn(db, 'transaction').mockImplementation((callback) => {
const transaction = new Promise<void>((resolve) => release.push(resolve)).then(() =>
callback(db as never)
)
transactions.push(transaction)
return transaction as ReturnType<typeof db.transaction>
})
const pending = Promise.all(
Array.from({ length: 18 }, (_, index) =>
retrieveKnowledgeSearch({
...params,
knowledgeBaseIds: [`kb-${index}`],
searchIndexOnly,
query: 'fixture policy',
searchMode: 'vector',
vectorBudgetMs: 50,
})
)
)
transactions.push(transaction)
return transaction as ReturnType<typeof db.transaction>
})
const pending = Promise.all(
Array.from({ length: 18 }, (_, index) =>
retrieveKnowledgeSearch({
...params,
knowledgeBaseIds: [`kb-${index}`],
query: 'fixture policy',
searchMode: 'vector',
vectorBudgetMs: 50,
await vi.advanceTimersByTimeAsync(60)
const results = await pending
expect(results).toHaveLength(18)
for (const result of results) {
expect(result).toEqual({
rows: [],
retrieval: { status: 'partial', timedOutLegs: ['vector'] },
})
)
)
await vi.advanceTimersByTimeAsync(60)
const results = await pending
expect(results).toHaveLength(18)
for (const result of results) {
expect(result).toEqual({
rows: [],
retrieval: { status: 'partial', timedOutLegs: ['vector'] },
})
}
for (const resume of release) resume()
const settled = await Promise.allSettled(transactions)
/** The searches that miss the projection-fill memo together share one read; each search's own read is refused at its deadline before it starts. */
expect(settled).toHaveLength(1)
for (const transaction of settled) {
expect(transaction.status).toBe('rejected')
if (transaction.status === 'rejected')
expect(transaction.reason).toBeInstanceOf(SearchDeadlineError)
}
for (const resume of release) resume()
const settled = await Promise.allSettled(transactions)
/** Search indexes share the readiness read; ordinary KBs acquire their own candidate reads. */
expect(settled).toHaveLength(searchIndexOnly ? 1 : 18)
for (const transaction of settled) {
expect(transaction.status).toBe('rejected')
if (transaction.status === 'rejected')
expect(transaction.reason).toBeInstanceOf(SearchDeadlineError)
}
expect(dbChainMockFns.select).not.toHaveBeenCalled()
expect(dbChainMockFns.execute).not.toHaveBeenCalled()
}
expect(dbChainMockFns.select).not.toHaveBeenCalled()
expect(dbChainMockFns.execute).not.toHaveBeenCalled()
})
)
})

describe('workspace search filters before ranking', () => {
Expand Down Expand Up @@ -946,6 +970,7 @@ describe('hydration follows ranked candidates', () => {
}
const params: SearchParams = {
knowledgeBaseIds: ['org-index'],
searchIndexOnly: true,
topK: 1,
access: identity,
accessProvider: provider,
Expand Down Expand Up @@ -1299,6 +1324,7 @@ describe('permitted-document planner', () => {
}
const params: SearchParams = {
knowledgeBaseIds: ['org-index'],
searchIndexOnly: true,
topK: 1,
access: reader,
accessProvider: provider,
Expand Down Expand Up @@ -2281,6 +2307,7 @@ describe('permitted-document planner', () => {

const liveSearch = {
knowledgeBaseIds: ['org-index'],
searchIndexOnly: true,
topK: 1,
searchMode: 'hybrid' as const,
query: 'release',
Expand Down Expand Up @@ -2367,48 +2394,57 @@ describe('permitted-document planner', () => {
expect(statements().some((query) => isPageStatement(query.sql))).toBe(false)
})

it('excludes a denied source through its documents while the projection is unfilled', async () => {
queueTableRows(schemaMock.knowledgeConnector, [
{
id: 'gated-src',
accessMode: 'admin',
connectorType: 'confluence',
githubRepository: false,
},
])
dbChainMockFns.execute.mockImplementation(async (query) => {
const statement = render(query).sql
/** The fill has not reached every row, so a denied source cannot be read off the row. */
if (statement.includes('AS unfilled')) return [{ unfilled: true }]
/** The mock renders nested fragments as parameters, so the marker is found in the whole query. */
const rebuilt = JSON.stringify(query).includes('/* excluded sources */')
if (isPageStatement(statement))
return JSON.stringify(render(query).params).includes('"b"')
? [hit('b', 'other-src')]
: [hit('a', 'gated-src')]
if (isWalk(statement))
return Array.from({ length: 400 }, (_, i) => ({
id: i === 0 ? (rebuilt ? 'b' : 'a') : `w-${i}`,
distance: 0.1,
}))
return []
})
queueTableRows(schemaMock.embedding, [])
queueTableRows(schemaMock.embedding, [hit('b', 'other-src')])
const getForConnectors = vi.fn<KnowledgeAccessProvider['getForConnectors']>(async () => reader)
const result = await retrieveKnowledgeSearch({
...liveSearch,
searchMode: 'vector',
access: reader,
accessProvider: { ...provider, getForConnectors },
})
expect(result.rows.map((row) => row.id)).toEqual(['b'])
const walks = statements().filter((query) => isWalk(query.sql))
expect(walks).toHaveLength(2)
expect(JSON.stringify(walks[0])).not.toContain('/* excluded sources */')
expect(JSON.stringify(walks[1])).toContain('NOT EXISTS (SELECT 1 FROM')
expect(JSON.stringify(walks[1])).toContain('/* excluded sources */')
})
it.each([false, true])(
'excludes a denied source through its documents (searchIndexOnly=%s)',
async (searchIndexOnly) => {
queueTableRows(schemaMock.knowledgeConnector, [
{
id: 'gated-src',
accessMode: 'admin',
connectorType: 'confluence',
githubRepository: false,
},
])
dbChainMockFns.execute.mockImplementation(async (query) => {
const statement = render(query).sql
/** The fill has not reached every row, so a denied source cannot be read off the row. */
if (statement.includes('AS unfilled')) return [{ unfilled: true }]
/** The mock renders nested fragments as parameters, so the marker is found in the whole query. */
const rebuilt = JSON.stringify(query).includes('/* excluded sources */')
if (isPageStatement(statement))
return JSON.stringify(render(query).params).includes('"b"')
? [hit('b', 'other-src')]
: [hit('a', 'gated-src')]
if (isWalk(statement))
return Array.from({ length: 400 }, (_, i) => ({
id: i === 0 ? (rebuilt ? 'b' : 'a') : `w-${i}`,
distance: 0.1,
}))
return []
})
queueTableRows(schemaMock.embedding, [])
queueTableRows(schemaMock.embedding, [hit('b', 'other-src')])
const getForConnectors = vi.fn<KnowledgeAccessProvider['getForConnectors']>(
async () => reader
)
const result = await retrieveKnowledgeSearch({
...liveSearch,
searchIndexOnly,
searchMode: 'vector',
access: reader,
accessProvider: { ...provider, getForConnectors },
})
expect(result.rows.map((row) => row.id)).toEqual(['b'])
expect(statements().filter((query) => query.sql.includes('AS unfilled'))).toHaveLength(
searchIndexOnly ? 1 : 0
)
const walks = statements().filter((query) => isWalk(query.sql))
expect(walks).toHaveLength(2)
expect(JSON.stringify(walks[0])).not.toContain('/* excluded sources */')
expect(JSON.stringify(walks[1])).toContain('NOT EXISTS (SELECT 1 FROM')
expect(JSON.stringify(walks[1])).toContain('/* excluded sources */')
}
)

it('hands back the unread slices of a page a denied source made it rebuild', async () => {
queueTableRows(schemaMock.knowledgeConnector, [
Expand Down Expand Up @@ -2491,6 +2527,7 @@ describe('filters on a resolved scope', () => {
}
const params: SearchParams = {
knowledgeBaseIds: ['org-index'],
searchIndexOnly: true,
topK: 1,
access: reader,
accessProvider: provider,
Expand Down
6 changes: 5 additions & 1 deletion apps/sim/lib/knowledge/search/queries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -369,6 +369,8 @@ export interface SearchParams {
permitted?: PermittedDocuments
/** Connector state resolved once per search, so no candidate re-derives it. */
accessPlan?: SearchAccessPlan
/** Every searched base is a Sim Search index; ordinary KBs use document-backed pages. */
searchIndexOnly?: boolean
}

/** All valid tag slot keys */
Expand Down Expand Up @@ -1832,7 +1834,9 @@ async function selectVectorResults(params: SearchParams): Promise<SearchResult[]
const plan = params.access.kind === 'user' ? params.accessPlan : undefined
/** Two remembered facts, read together when neither is remembered. */
const [filled, plannedIndexedSources] = await Promise.all([
isProjectionFilled('embedding_search', 'vector.projection_filled', params.budget),
params.searchIndexOnly === true
? isProjectionFilled('embedding_search', 'vector.projection_filled', params.budget)
: false,
plan?.memberSources.length ? indexedVectorSources(params.budget) : undefined,
])
/**
Expand Down
Loading