Skip to content

Commit 2d66fda

Browse files
committed
improvement(knowledge): rank each readable source in its own vector index
pgvector post-filters, so a walk over every source spends its scan budget on the sources a caller cannot read and returns few of their true neighbours: measured recall for a member reading half an organization index ranged from 0.00 to 1.00, averaging 0.80 across nine queries, with two queries returning nothing of the exact page. Readability follows sources, so each one is ranked on its own terms. A member of a source reads essentially all of it and is ranked by walking that source's index, built after a sync grows it past the threshold and dropped with the connector. Every other source is sliced — mirrored permissions give a caller their own mail, their own files — and those slices are ranked exactly in one statement, which is cheaper than a walk and exact by construction. The lists merge by distance. Recall over the same nine queries rises to 0.95 with no query below 0.75, at 398ms of database time against 95ms for the walk it replaces. Retrieval never waits on an index existing: a source without one is ranked exactly.
1 parent 4c76932 commit 2d66fda

15 files changed

Lines changed: 28631 additions & 7 deletions

‎apps/sim/lib/knowledge/connectors/sync-engine.ts‎

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import {
66
knowledgeConnectorSyncLog,
77
} from '@sim/db/schema'
88
import { createLogger } from '@sim/logger'
9-
import { toError } from '@sim/utils/errors'
9+
import { getErrorMessage, toError } from '@sim/utils/errors'
1010
import { generateId } from '@sim/utils/id'
1111
import { randomInt } from '@sim/utils/random'
1212
import { and, asc, eq, exists, gt, inArray, isNotNull, isNull, or, sql } from 'drizzle-orm'
@@ -85,6 +85,7 @@ import {
8585
} from '@/lib/knowledge/connectors/sync-primitives'
8686
import { hardDeleteDocuments } from '@/lib/knowledge/documents/service'
8787
import { getRetryAfterMs, isRateLimitError } from '@/lib/knowledge/documents/utils'
88+
import { ensureSourceVectorIndex } from '@/lib/knowledge/search/source-vector-indexes'
8889
import { CONNECTOR_REGISTRY } from '@/connectors/registry.server'
8990
import type {
9091
ConnectorAuthConfig,
@@ -391,7 +392,7 @@ export async function completeSuccessfulSync(
391392
const completionNotice =
392393
[directoryNotice, listingNotice, contentNotice].filter(Boolean).join('\n') || null
393394
try {
394-
return await db.transaction(async (tx) => {
395+
const completed = await db.transaction(async (tx) => {
395396
const [lockedKnowledgeBase] = await tx
396397
.select({ id: knowledgeBase.id })
397398
.from(knowledgeBase)
@@ -493,6 +494,18 @@ export async function completeSuccessfulSync(
493494

494495
return true
495496
})
497+
/**
498+
* A source that has grown past the threshold gets its own vector index, so a member who reads
499+
* it whole is ranked through a walk of their own documents. Retrieval ranks exactly without
500+
* it, so a failure here is logged and left for the next sync.
501+
*/
502+
await ensureSourceVectorIndex(connectorId).catch((error: unknown) => {
503+
logger.warn('Could not ensure the source vector index', {
504+
connectorId,
505+
error: getErrorMessage(error),
506+
})
507+
})
508+
return completed
496509
} catch (error) {
497510
if (error instanceof SyncCompletionOwnershipLost) return false
498511
throw error

‎apps/sim/lib/knowledge/orchestration/connectors.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ import {
99
knowledgeConnectorMember,
1010
} from '@sim/db/schema'
1111
import { createLogger } from '@sim/logger'
12-
import { getPostgresErrorCode } from '@sim/utils/errors'
12+
import { getErrorMessage, getPostgresErrorCode } from '@sim/utils/errors'
1313
import { generateId } from '@sim/utils/id'
1414
import { and, eq, isNull, sql } from 'drizzle-orm'
1515
import { encryptApiKey } from '@/lib/api-key/crypto'
@@ -60,6 +60,7 @@ import {
6060
type KnowledgeOperationContext,
6161
type KnowledgeOrchestrationResult,
6262
} from '@/lib/knowledge/orchestration/shared'
63+
import { dropSourceVectorIndex } from '@/lib/knowledge/search/source-vector-indexes'
6364
import { createTagDefinition } from '@/lib/knowledge/tags/service'
6465
import { captureServerEvent } from '@/lib/posthog/server'
6566
import { searchSourceIdentity } from '@/lib/sim-search/source-identity'
@@ -1411,6 +1412,14 @@ export async function performDeleteKnowledgeConnector(
14111412
})
14121413
}
14131414

1415+
/** The source is gone, so its vector index is too; ranking falls back to the exact path. */
1416+
await dropSourceVectorIndex(connectorId).catch((error: unknown) => {
1417+
logger.warn('Could not drop the source vector index', {
1418+
connectorId,
1419+
error: getErrorMessage(error),
1420+
})
1421+
})
1422+
14141423
return {
14151424
success: true,
14161425
documentsDeleted: deleteDocuments ? docCount : 0,

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

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,9 @@ export type SearchStage =
2929
| 'defaults'
3030
| 'retrieval'
3131
| 'connector_eligibility'
32+
| 'vector.source_plan'
33+
| 'vector.source_exact'
34+
| 'vector.source_walk'
3235
| 'permitted_documents'
3336
| 'result_provenance'
3437
| 'reranking'
@@ -82,7 +85,7 @@ export interface SearchDiagnosticMetadata {
8285
searchMode?: 'hybrid' | 'vector'
8386
boostRecency?: boolean
8487
embeddingDimensions?: number
85-
vectorRanking?: 'exact' | 'exact-candidates' | 'candidate-rerank'
88+
vectorRanking?: 'exact' | 'exact-candidates' | 'candidate-rerank' | 'per-source'
8689
vectorCandidateStorage?: 'stored-halfvec'
8790
/**
8891
* Whether the bounded traversal filled its candidate limit. `underfilled` means visibility
@@ -102,6 +105,8 @@ export interface SearchDiagnosticMetadata {
102105
permittedDocuments?: 'bounded' | 'unbounded'
103106
/** Documents in a bounded permitted set. */
104107
permittedDocumentCount?: number
108+
vectorSourcesSliced?: number
109+
vectorSourcesWalked?: number
105110
/**
106111
* Which index ranked an unbounded keyword leg: `tin` ranks by BM25 and checks access on the top
107112
* of that ranking; `gin` ranks every match. Absent when the leg ranked inside a bounded set.

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

Lines changed: 58 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1432,19 +1432,27 @@ describe('permitted-document planner', () => {
14321432

14331433
let probeRows: Array<{ id: string | null; connectorId: string | null; saturated: boolean }>
14341434
let exactRows: Array<{ id: string }>
1435-
let traversedRows: Array<{ id: string }>
1435+
let traversedRows: Array<{ id: string; distance?: number }>
14361436
let rerankRows: Array<ReturnType<typeof hit>>
1437+
let memberedSources: Array<{ connectorId: string }>
1438+
let indexedSourceRows: Array<{ name: string; connectorId: string }>
1439+
let sourceExactRows: Array<{ id: string; distance: number }>
14371440

14381441
beforeEach(() => {
14391442
resetDbChainMock()
14401443
probeRows = []
14411444
exactRows = []
14421445
traversedRows = []
14431446
rerankRows = []
1447+
sourceExactRows = []
1448+
memberedSources = []
1449+
indexedSourceRows = []
14441450
dbChainMockFns.execute.mockImplementation(async (query) => {
14451451
const statement = render(query).sql
1452+
if (statement.includes('pg_index')) return indexedSourceRows
14461453
if (statement.includes('AS visible')) return traversedRows
14471454
if (statement.includes('WITH scored_search_candidates')) return rerankRows
1455+
if (statement.includes('WITH readable_documents')) return sourceExactRows
14481456
if (isExactRanking(statement)) return exactRows
14491457
if (isProbeStatement(statement)) return probeRows
14501458
return []
@@ -1485,6 +1493,55 @@ describe('permitted-document planner', () => {
14851493
expect(sqls.some(isExactRanking)).toBe(false)
14861494
})
14871495

1496+
it('walks a source the caller is a member of and ranks every other source exactly', async () => {
1497+
const eligibility = {
1498+
workspace: [],
1499+
admin: ['sliced-src'],
1500+
members: ['member-src'],
1501+
liveProofRequired: [],
1502+
}
1503+
/** Membership decides the walk; the sliced source contributes enumerated documents. */
1504+
memberedSources = [{ connectorId: 'member-src' }]
1505+
indexedSourceRows = [{ name: 'idx', connectorId: 'member-src' }]
1506+
queueTableRows(schemaMock.knowledgeConnectorMember, memberedSources)
1507+
sourceExactRows = [{ id: 'sliced-hit', distance: 0.05 }]
1508+
traversedRows = [{ id: 'walked-hit', distance: 0.2 }]
1509+
rerankRows = [hit('sliced-hit', 'sliced-src'), hit('walked-hit', 'member-src')]
1510+
queueTableRows(schemaMock.embedding, rerankRows)
1511+
await handleVectorOnlySearch({
1512+
...params,
1513+
permitted: { kind: 'unbounded' },
1514+
connectorEligibility: eligibility,
1515+
})
1516+
const walks = statements().filter((query) => query.sql.includes('AS visible'))
1517+
expect(walks).toHaveLength(1)
1518+
expect(
1519+
walks[0].params.some((param) => JSON.stringify(param).includes('"right":"member-src"'))
1520+
).toBe(true)
1521+
/** The sliced sources resolve their documents inside one statement, not through the app. */
1522+
const exact = statements().filter((query) => query.sql.includes('WITH readable_documents'))
1523+
expect(exact).toHaveLength(1)
1524+
expect(JSON.stringify(exact[0])).toContain('sliced-src')
1525+
})
1526+
1527+
it('ranks every source exactly when the caller is a member of none', async () => {
1528+
memberedSources = []
1529+
sourceExactRows = [{ id: 'sliced-hit', distance: 0.05 }]
1530+
rerankRows = [hit('sliced-hit', 'sliced-src')]
1531+
queueTableRows(schemaMock.embedding, rerankRows)
1532+
await handleVectorOnlySearch({
1533+
...params,
1534+
permitted: { kind: 'unbounded' },
1535+
connectorEligibility: {
1536+
workspace: [],
1537+
admin: ['sliced-src'],
1538+
members: [],
1539+
liveProofRequired: [],
1540+
},
1541+
})
1542+
expect(statements().filter((query) => query.sql.includes('AS visible'))).toHaveLength(0)
1543+
})
1544+
14881545
it('confines keyword matching to the bounded permitted set', async () => {
14891546
await executeKeywordSearch({
14901547
...params,

0 commit comments

Comments
 (0)