Skip to content

Commit b958c90

Browse files
committed
fix(knowledge): refill after a denied gated source, reserve the index build's session, lock the backfill's documents
- once a live proof shows the caller does not hold a gated source, rebuild the candidate pages without it, so its candidates no longer hold the slots of sources the caller can read - run the source index build on one reserved connection, so its memory setting and its reset apply to the session that builds - share-lock a page's documents before copying their source, so a detachment in flight waits for the page and then fans its own change out through the trigger
1 parent 0f01d94 commit b958c90

5 files changed

Lines changed: 149 additions & 34 deletions

File tree

‎apps/sim/lib/knowledge/search/queries.test.ts‎

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1660,6 +1660,48 @@ describe('permitted-document planner', () => {
16601660
expect(getForConnectors).not.toHaveBeenCalled()
16611661
})
16621662

1663+
it('rebuilds the pages without a gated source the caller turns out not to hold', async () => {
1664+
queueTableRows(schemaMock.knowledgeConnector, [
1665+
{
1666+
id: 'gated-src',
1667+
accessMode: 'admin',
1668+
connectorType: 'confluence',
1669+
githubRepository: false,
1670+
},
1671+
])
1672+
/**
1673+
* The first pool is filled by the gated source alone; only a pool built without it — the
1674+
* exclusion carries the source id into the statement — reaches the accessible candidate.
1675+
*/
1676+
dbChainMockFns.execute.mockImplementation(async (query) => {
1677+
const statement = render(query).sql
1678+
/** The exclusion is the only clause that negates a connector membership. */
1679+
const rebuilt = JSON.stringify(query).includes('OR NOT (')
1680+
if (statement.includes('WITH scored_search_candidates'))
1681+
return rebuilt ? [hit('b', 'other-src')] : [hit('a', 'gated-src')]
1682+
if (isExactRanking(statement)) return [{ id: rebuilt ? 'b' : 'a' }]
1683+
if (isProbeStatement(statement))
1684+
return [{ id: 'doc-a', connectorId: 'gated-src', saturated: false }]
1685+
return []
1686+
})
1687+
queueTableRows(schemaMock.embedding, [])
1688+
queueTableRows(schemaMock.embedding, [hit('b', 'other-src')])
1689+
/** No grants come back, so the gated source is denied. */
1690+
const getForConnectors = vi.fn<KnowledgeAccessProvider['getForConnectors']>(async () => reader)
1691+
const result = await retrieveKnowledgeSearch({
1692+
...liveSearch,
1693+
searchMode: 'vector',
1694+
access: reader,
1695+
accessProvider: { ...provider, getForConnectors },
1696+
})
1697+
expect(getForConnectors).toHaveBeenCalledOnce()
1698+
expect(result.rows.map((row) => row.id)).toEqual(['b'])
1699+
const reranks = statements().filter((query) => query.sql.includes('scored_search_candidates'))
1700+
expect(reranks).toHaveLength(2)
1701+
expect(JSON.stringify(reranks[0])).not.toContain('OR NOT (')
1702+
expect(JSON.stringify(reranks[1])).toContain('OR NOT (')
1703+
})
1704+
16631705
it('asks a live source for its grants once, when a candidate of its own is read', async () => {
16641706
queueTableRows(schemaMock.knowledgeConnector, [
16651707
{

‎apps/sim/lib/knowledge/search/queries.ts‎

Lines changed: 71 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -596,7 +596,8 @@ const SEARCH_READ_CANDIDATE_FIELDS = {
596596
*/
597597
export interface LiveSourceAccess {
598598
gates: (connectorId: string) => boolean
599-
resolve: () => Promise<KnowledgeAccessScope>
599+
/** The caller's scope with its grants, and the gated sources those grants do not cover. */
600+
resolve: () => Promise<{ access: KnowledgeAccessScope; denied: ReadonlySet<string> }>
600601
}
601602

602603
/** Binds a search's gated sources to one memoized resolution of the caller's grants. */
@@ -608,7 +609,7 @@ export function liveSourceAccessFor(
608609
): LiveSourceAccess | undefined {
609610
const gated = new Set(plan?.connectors.liveProofRequired ?? [])
610611
if (!accessProvider || gated.size === 0) return undefined
611-
let pending: Promise<KnowledgeAccessScope> | undefined
612+
let pending: Promise<{ access: KnowledgeAccessScope; denied: ReadonlySet<string> }> | undefined
612613
/** The provider authorizes connectors in bounded pages, so a wide scope resolves page by page. */
613614
const pages: string[][] = []
614615
for (const id of gated) {
@@ -624,7 +625,6 @@ export function liveSourceAccessFor(
624625
accessProvider.getForConnectors(page, signal)
625626
)
626627
const [first] = scopes
627-
if (scopes.length === 1 || first.kind !== 'user') return first
628628
const githubInstallationGrants: GitHubInstallationReadGrant[] = []
629629
const confluenceSiteGrants: ConfluenceSiteReadGrant[] = []
630630
for (const scope of scopes) {
@@ -633,7 +633,16 @@ export function liveSourceAccessFor(
633633
githubInstallationGrants.push(...scope.githubInstallationGrants)
634634
if (scope.confluenceSiteGrants) confluenceSiteGrants.push(...scope.confluenceSiteGrants)
635635
}
636-
return { ...first, githubInstallationGrants, confluenceSiteGrants }
636+
const granted = new Set([
637+
...githubInstallationGrants.map((grant) => grant.connectorId),
638+
...confluenceSiteGrants.map((grant) => grant.connectorId),
639+
])
640+
const denied = new Set([...gated].filter((id) => !granted.has(id)))
641+
const merged =
642+
first.kind === 'user'
643+
? { ...first, githubInstallationGrants, confluenceSiteGrants }
644+
: first
645+
return { access: merged, denied }
637646
})
638647
return pending
639648
},
@@ -655,7 +664,11 @@ async function selectAuthorizedSearchResults(input: {
655664
signal?: AbortSignal
656665
budget?: SearchBudget
657666
topK: number
658-
selectPage: (limit: number, offset: number) => Promise<SearchReadCandidatePage>
667+
selectPage: (
668+
limit: number,
669+
offset: number,
670+
excludedSources: readonly string[]
671+
) => Promise<SearchReadCandidatePage>
659672
compareResults?: (a: SearchResult, b: SearchResult) => number
660673
hydrate: (ids: string[], access: KnowledgeAccessScope) => Promise<SearchResult[]>
661674
liveSourceAccess?: LiveSourceAccess
@@ -664,6 +677,8 @@ async function selectAuthorizedSearchResults(input: {
664677
const pageSize = Math.min(AUTHORIZED_SEARCH_PAGE_SIZE, Math.max(input.topK, 20))
665678
const results = new Map<string, SearchResult>()
666679
const considered = new Set<string>()
680+
/** Gated sources the caller turned out not to hold: left out of every page once that is known. */
681+
let excluded: ReadonlySet<string> = new Set()
667682
let scanned = 0
668683
let offset = 0
669684
try {
@@ -675,7 +690,7 @@ async function selectAuthorizedSearchResults(input: {
675690
input.signal?.throwIfAborted()
676691
input.budget?.remaining()
677692
const page = await measureSearchStage(`${input.leg}.candidates`, () =>
678-
input.selectPage(pageSize, offset)
693+
input.selectPage(pageSize, offset, [...excluded])
679694
)
680695
if (!page.candidates.length) break
681696
scanned += page.candidates.length
@@ -691,11 +706,25 @@ async function selectAuthorizedSearchResults(input: {
691706
* its own reaches this page, and then once for the whole search: a scope that ranks none
692707
* of them — most scopes — never asks, and one that ranks many asks once.
693708
*/
694-
const access = candidates.some(
695-
(candidate) => candidate.connectorId && input.liveSourceAccess?.gates(candidate.connectorId)
696-
)
697-
? await input.liveSourceAccess!.resolve()
698-
: input.access
709+
const proof = input.liveSourceAccess
710+
const gatedOnPage =
711+
proof !== undefined &&
712+
candidates.some((candidate) => candidate.connectorId && proof.gates(candidate.connectorId))
713+
let access = input.access
714+
let refill = false
715+
if (gatedOnPage && proof) {
716+
const resolved = await proof.resolve()
717+
access = resolved.access
718+
/**
719+
* A denied source's candidates cannot hydrate, yet they took the slots of sources the
720+
* caller does hold. Once the denial is known the pages are rebuilt without that source,
721+
* from the start; the ids already seen are not read twice.
722+
*/
723+
if (resolved.denied.size > excluded.size) {
724+
excluded = resolved.denied
725+
refill = true
726+
}
727+
}
699728
const hydrated = await measureSearchStage(`${input.leg}.hydration`, () =>
700729
input.hydrate(
701730
candidates.map((candidate) => candidate.id),
@@ -713,6 +742,10 @@ async function selectAuthorizedSearchResults(input: {
713742
results.clear()
714743
for (const row of ranked) results.set(row.id, row)
715744
}
745+
if (refill) {
746+
offset = 0
747+
continue
748+
}
716749
/** A short page is the end of the candidates, whether or not they were reordered. */
717750
if (page.candidates.length < pageSize) break
718751
}
@@ -723,6 +756,13 @@ async function selectAuthorizedSearchResults(input: {
723756
return [...results.values()]
724757
}
725758

759+
/** Keeps the candidates of sources the caller turned out not to hold out of a page. */
760+
function excludeSearchSources(sourceIds: readonly string[]): SQL | undefined {
761+
return sourceIds.length
762+
? sql`(${document.connectorId} IS NULL OR NOT (${inArray(document.connectorId, [...sourceIds])}))`
763+
: undefined
764+
}
765+
726766
/**
727767
* Loads the content of candidates that survived ranking, under the read predicate.
728768
*
@@ -799,7 +839,7 @@ export async function handleTagOnlySearch(params: SearchParams): Promise<SearchR
799839
signal: params.signal,
800840
budget: params.budget,
801841
topK,
802-
selectPage: async (limit, offset) => {
842+
selectPage: async (limit, offset, excludedSources) => {
803843
const candidates = await runSearchQuery(params.budget, 'tags.sql', (executor) =>
804844
executor
805845
.select(SEARCH_READ_CANDIDATE_FIELDS)
@@ -812,7 +852,8 @@ export async function handleTagOnlySearch(params: SearchParams): Promise<SearchR
812852
access,
813853
params.filters,
814854
candidateAccessCondition(access, params.accessPlan)
815-
)
855+
),
856+
excludeSearchSources(excludedSources)
816857
)
817858
)
818859
.orderBy(embedding.id)
@@ -1277,7 +1318,7 @@ async function selectVectorResults(params: SearchParams): Promise<SearchResult[]
12771318
* refill reuses the pool it already has. Excluding another source is the only thing that
12781319
* changes which candidates belong in it, and that resets the offset to zero anyway.
12791320
*/
1280-
let candidatePool: Array<{ id: string }> | undefined
1321+
let candidatePool: { excludedKey: string; ids: Array<{ id: string }> } | undefined
12811322
return selectAuthorizedSearchResults({
12821323
leg: 'vector',
12831324
access: params.access,
@@ -1287,9 +1328,11 @@ async function selectVectorResults(params: SearchParams): Promise<SearchResult[]
12871328
budget: params.budget,
12881329
topK: params.topK,
12891330
compareResults: (a, b) => a.distance - b.distance,
1290-
selectPage: async (limit, offset) => {
1331+
selectPage: async (limit, offset, excludedSources) => {
1332+
const excludedKey = [...excludedSources].sort().join(',')
12911333
const visibility = [
12921334
...getVisibilityConditions(params.access, params.filters, candidateAccess),
1335+
excludeSearchSources(excludedSources),
12931336
]
12941337
const candidateDocumentVisibility = [
12951338
...candidateDocumentConditions(
@@ -1298,6 +1341,7 @@ async function selectVectorResults(params: SearchParams): Promise<SearchResult[]
12981341
params.filters,
12991342
candidateAccess
13001343
),
1344+
excludeSearchSources(excludedSources),
13011345
]
13021346
/** Explicit document IDs are already a bounded scope, and retain exhaustive ordering. */
13031347
const exactPage = async () => {
@@ -1315,7 +1359,7 @@ async function selectVectorResults(params: SearchParams): Promise<SearchResult[]
13151359
return { candidates, nextOffset: offset + candidates.length }
13161360
}
13171361
if (params.filters?.documentIds?.length) return exactPage()
1318-
if (!candidatePool) {
1362+
if (candidatePool?.excludedKey !== excludedKey) {
13191363
annotateSearchDiagnostics({
13201364
vectorRanking: 'candidate-rerank',
13211365
vectorCandidateStorage: 'stored-halfvec',
@@ -1421,13 +1465,13 @@ async function selectVectorResults(params: SearchParams): Promise<SearchResult[]
14211465
}
14221466
}
14231467
}
1424-
candidatePool = selected
1468+
candidatePool = { excludedKey, ids: selected }
14251469
annotateSearchDiagnostics({
14261470
vectorCandidateCount: selected.length,
14271471
vectorCandidateScan: selected.length < candidateLimit ? 'underfilled' : 'planned',
14281472
})
14291473
}
1430-
const identities = candidatePool
1474+
const identities = candidatePool.ids
14311475
if (!identities.length) return { candidates: [], nextOffset: offset }
14321476
/** Score each bounded candidate once; sorting the materialized scalar cannot invoke HNSW again. */
14331477
const page = await runSearchQuery(params.budget, 'vector.rerank', (executor) =>
@@ -1555,14 +1599,15 @@ export async function executeKeywordSearch(params: KeywordSearchParams): Promise
15551599
? { keywordRanking: tinQuery ? 'tin' : 'gin' }
15561600
: {}),
15571601
})
1558-
const documentConditions = () =>
1602+
const documentConditions = (excludedSources: readonly string[]) =>
15591603
and(
15601604
...candidateDocumentConditions(
15611605
knowledgeBaseIds,
15621606
access,
15631607
params.filters,
15641608
candidateAccessCondition(access, params.accessPlan)
1565-
)
1609+
),
1610+
excludeSearchSources(excludedSources)
15661611
)
15671612
/**
15681613
* One page from the top of Tin's ranking. The window of ranked chunks widens while too few of
@@ -1572,7 +1617,8 @@ export async function executeKeywordSearch(params: KeywordSearchParams): Promise
15721617
const selectTinPage = async (
15731618
scopedQuery: SQL,
15741619
limit: number,
1575-
offset: number
1620+
offset: number,
1621+
excludedSources: readonly string[]
15761622
): Promise<SearchReadCandidatePage | null> => {
15771623
for (const window of TIN_KEYWORD_WINDOWS) {
15781624
if (window < offset + limit) continue
@@ -1590,7 +1636,7 @@ export async function executeKeywordSearch(params: KeywordSearchParams): Promise
15901636
SELECT ${document.id} AS id FROM ${document}
15911637
WHERE ${and(
15921638
sql`${document.id} = ANY (ARRAY(SELECT document_id FROM ranked_tin_chunks))`,
1593-
documentConditions()
1639+
documentConditions(excludedSources)
15941640
)}
15951641
), page AS (
15961642
SELECT ranked_tin_chunks.id, ${document.id} AS "documentId",
@@ -1635,7 +1681,7 @@ export async function executeKeywordSearch(params: KeywordSearchParams): Promise
16351681
signal: params.signal,
16361682
budget: params.budget,
16371683
topK,
1638-
selectPage: async (limit, offset) => {
1684+
selectPage: async (limit, offset, excludedSources) => {
16391685
/**
16401686
* A bounded permitted set confines matching to the chunks the caller may read, so a term
16411687
* common across the index is ranked only where it can surface. The visibility CTE below
@@ -1647,7 +1693,7 @@ export async function executeKeywordSearch(params: KeywordSearchParams): Promise
16471693
: undefined
16481694
if (permittedIds?.length === 0) return { candidates: [], nextOffset: offset }
16491695
if (tinScope && !permittedIds) {
1650-
const tinPage = await selectTinPage(tinScope, limit, offset)
1696+
const tinPage = await selectTinPage(tinScope, limit, offset, excludedSources)
16511697
if (tinPage) return tinPage
16521698
annotateSearchDiagnostics({ keywordRanking: 'gin' })
16531699
}
@@ -1694,7 +1740,7 @@ export async function executeKeywordSearch(params: KeywordSearchParams): Promise
16941740
SELECT ${document.id} AS id FROM ${document}
16951741
WHERE ${and(
16961742
sql`${document.id} = ANY (ARRAY(SELECT document_id FROM matched_keyword_chunks))`,
1697-
documentConditions()
1743+
documentConditions(excludedSources)
16981744
)}
16991745
), ranked_keyword_candidates AS MATERIALIZED (
17001746
SELECT matched_keyword_chunks.id, matched_keyword_chunks.document_id,

‎apps/sim/lib/knowledge/search/source-vector-indexes.test.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
/**
22
* @vitest-environment node
33
*/
4+
import { db } from '@sim/db'
45
import { dbChainMockFns, resetDbChainMock } from '@sim/testing'
56
import { beforeEach, describe, expect, it } from 'vitest'
67
import {
@@ -27,6 +28,18 @@ describe('source vector indexes', () => {
2728
if (text.includes('sampled')) return [{ column: 'vector_512' }]
2829
return []
2930
})
31+
/** The build reserves one connection; its statements are recorded with the rest. */
32+
Object.assign(db, {
33+
$client: {
34+
reserve: async () => ({
35+
unsafe: async (text: string) => {
36+
statements.push(text)
37+
return []
38+
},
39+
release: () => undefined,
40+
}),
41+
},
42+
})
3043
/** Warms the catalog cache with this case's state, so a build decision is not a stale read. */
3144
expect((await indexedVectorSources()).size).toBe(indexed.length)
3245
statements = []

‎apps/sim/lib/knowledge/search/source-vector-indexes.ts‎

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -103,14 +103,18 @@ export async function ensureSourceVectorIndex(connectorId: string): Promise<bool
103103
if (!column) return false
104104
const name = indexName(connectorId)
105105
const startedAt = Date.now()
106+
/**
107+
* The memory setting and the build must share a session, and `CONCURRENTLY` forbids a
108+
* transaction, so one connection is reserved from the pool for the whole build and its setting
109+
* is reset before the connection goes back.
110+
*/
111+
const session = await db.$client.reserve()
106112
try {
107113
await dropInvalidIndex(name)
108-
await db.execute(sql`SET maintenance_work_mem = '2GB'`)
109-
await db.execute(
110-
sql.raw(`CREATE INDEX CONCURRENTLY "${name}" ON embedding_search
114+
await session.unsafe("SET maintenance_work_mem = '2GB'")
115+
await session.unsafe(`CREATE INDEX CONCURRENTLY "${name}" ON embedding_search
111116
USING hnsw (${column} halfvec_cosine_ops) WITH (m = 16, ef_construction = 64)
112117
WHERE connector_id = '${connectorId}' AND enabled`)
113-
)
114118
logger.info('Built a source vector index', {
115119
connectorId,
116120
documents,
@@ -123,7 +127,8 @@ export async function ensureSourceVectorIndex(connectorId: string): Promise<bool
123127
await dropInvalidIndex(name)
124128
return false
125129
} finally {
126-
await db.execute(sql`RESET maintenance_work_mem`)
130+
await session.unsafe('RESET maintenance_work_mem').catch(() => undefined)
131+
session.release()
127132
}
128133
}
129134

0 commit comments

Comments
 (0)