Skip to content

Commit 4601d47

Browse files
committed
improvement(knowledge): resolve the caller's member identities with the rest of the plan
A members-mode document is readable while one of the caller's active members observes it, freshly, and which members those are is a fact about the caller. Resolving them with the connectors — one query, one plan — turns each candidate's check into a lookup on the observation key instead of a join to the member behind it, and lets the vector planner read the sources the caller belongs to from the same resolution rather than asking again.
1 parent d7d57fc commit 4601d47

5 files changed

Lines changed: 170 additions & 79 deletions

File tree

‎apps/sim/lib/knowledge/access/connector-eligibility.ts‎

Lines changed: 69 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,14 @@
11
import { db } from '@sim/db'
2-
import { knowledgeConnector } from '@sim/db/schema'
2+
import { knowledgeConnector, knowledgeConnectorMember } from '@sim/db/schema'
33
import { and, eq, inArray, isNull, sql } from 'drizzle-orm'
4-
import type { KnowledgeConnectorEligibility } from '@/lib/knowledge/access/predicate'
4+
import { SOURCE_ACL_MAX_AGE_MS } from '@/lib/knowledge/access/freshness'
5+
import {
6+
type KnowledgeConnectorEligibility,
7+
type KnowledgeMemberObservers,
8+
type SearchAccessPlan,
9+
textArrayLiteral,
10+
} from '@/lib/knowledge/access/predicate'
11+
import type { KnowledgeAccessScope } from '@/lib/knowledge/access/types'
512
import { searchIntegrationAccessCondition } from '@/lib/knowledge/search/integration-policy'
613

714
/**
@@ -12,7 +19,7 @@ import { searchIntegrationAccessCondition } from '@/lib/knowledge/search/integra
1219
* facts about a connector. Resolving them once per query — there are tens of connectors against
1320
* hundreds of thousands of documents — leaves each candidate its own columns to check.
1421
*/
15-
export async function resolveConnectorEligibility(
22+
async function resolveConnectorEligibility(
1623
knowledgeBaseIds: readonly string[]
1724
): Promise<KnowledgeConnectorEligibility> {
1825
const eligibility: {
@@ -52,3 +59,62 @@ export async function resolveConnectorEligibility(
5259
}
5360
return eligibility
5461
}
62+
63+
/**
64+
* The caller's own member identities on these connectors, split by whether the member's change
65+
* feed is itself current.
66+
*
67+
* A members-mode document is readable while one of the caller's active members observes it,
68+
* freshly — and which members those are is a fact about the caller, not about any document. A
69+
* member whose feed drained recently confirms every observation it holds, so its observations need
70+
* no age check at all; the rest are checked against the age of the observation itself. Resolved
71+
* once, the per-document check becomes one lookup on the observation key, with no join to the
72+
* member behind it.
73+
*/
74+
async function resolveMemberObservers(
75+
access: KnowledgeAccessScope,
76+
connectorIds: readonly string[]
77+
): Promise<{ observers: KnowledgeMemberObservers; memberSources: string[] }> {
78+
if (access.kind !== 'user' || connectorIds.length === 0 || access.tokens.length === 0) {
79+
return { observers: { confirmed: [], observed: [] }, memberSources: [] }
80+
}
81+
const rows = await db
82+
.select({
83+
id: knowledgeConnectorMember.id,
84+
connectorId: knowledgeConnectorMember.connectorId,
85+
syncedThrough: knowledgeConnectorMember.memberSyncedThrough,
86+
})
87+
.from(knowledgeConnectorMember)
88+
.where(
89+
and(
90+
inArray(knowledgeConnectorMember.connectorId, [...connectorIds]),
91+
eq(knowledgeConnectorMember.status, 'active'),
92+
sql`${knowledgeConnectorMember.subjectToken} = ANY(${textArrayLiteral([...access.tokens])})`
93+
)
94+
)
95+
const cutoff = Date.now() - SOURCE_ACL_MAX_AGE_MS
96+
const confirmed: string[] = []
97+
const observed: string[] = []
98+
const memberSources = new Set<string>()
99+
for (const row of rows) {
100+
if (row.syncedThrough !== null && row.syncedThrough.getTime() > cutoff) confirmed.push(row.id)
101+
else observed.push(row.id)
102+
memberSources.add(row.connectorId)
103+
}
104+
return { observers: { confirmed, observed }, memberSources: [...memberSources] }
105+
}
106+
107+
/**
108+
* Everything a search needs to know about its sources and the caller's standing in them, resolved
109+
* once: which connectors it may read, the caller's member identities there, and the sources they
110+
* are a member of. Each is a fact about a connector or a caller, so deriving them per candidate
111+
* document is what made retrieval cost grow with the size of what someone may read.
112+
*/
113+
export async function resolveSearchAccessPlan(
114+
knowledgeBaseIds: readonly string[],
115+
access: KnowledgeAccessScope
116+
): Promise<SearchAccessPlan> {
117+
const connectors = await resolveConnectorEligibility(knowledgeBaseIds)
118+
const { observers, memberSources } = await resolveMemberObservers(access, connectors.members)
119+
return { connectors, observers, memberSources }
120+
}

‎apps/sim/lib/knowledge/access/predicate.postgres.test.ts‎

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -553,8 +553,11 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
553553
members: ['members'],
554554
liveProofRequired: [],
555555
}
556+
/** What `resolveSearchAccessPlan` resolves for this caller: their member identity, confirmed. */
557+
const observers = { confirmed: ['m-alice'], observed: [] }
556558
const perRow = knowledgeMetadataCandidateAccessCondition(scope)
557-
const perQuery = knowledgeCandidateAccessConditionForConnectors(scope, eligibility)
559+
const plan = { connectors: eligibility, observers, memberSources: ['members'] }
560+
const perQuery = knowledgeCandidateAccessConditionForConnectors(scope, plan)
558561
for (const id of [...cases.map(([documentId]) => documentId), 'upload-doc']) {
559562
expect([id, await admits(perQuery, id)]).toEqual([id, await admits(perRow, id)])
560563
}
@@ -563,7 +566,7 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
563566
* caller holds those grants only after authorization, so applying the clause during ranking
564567
* would drop every candidate of a gated source before it could be proven.
565568
*/
566-
const gated = { ...eligibility, liveProofRequired: ['admin'] }
569+
const gated = { ...plan, connectors: { ...eligibility, liveProofRequired: ['admin'] } }
567570
expect(
568571
await admits(knowledgeCandidateAccessConditionForConnectors(scope, gated), 'admin-current')
569572
).toBe(true)
@@ -576,7 +579,10 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
576579
/** A connector left out of the resolution is refused, however current its documents are. */
577580
expect(
578581
await admits(
579-
knowledgeCandidateAccessConditionForConnectors(scope, { ...eligibility, admin: [] }),
582+
knowledgeCandidateAccessConditionForConnectors(scope, {
583+
...plan,
584+
connectors: { ...eligibility, admin: [] },
585+
}),
580586
'admin-current'
581587
)
582588
).toBe(false)

‎apps/sim/lib/knowledge/access/predicate.ts‎

Lines changed: 45 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -208,6 +208,25 @@ export function liveSourceAccessCondition(scope: KnowledgeAccessScope): SQL {
208208
* The connectors a search may read from, resolved once per query: their ids grouped by the shape
209209
* their documents' ACLs take, and separately those whose reader access is proven live per request.
210210
*/
211+
/**
212+
* The caller's active member identities on the connectors a search reads, by what makes their
213+
* observations current: `confirmed` members drained their change feed inside the freshness window,
214+
* so every observation they hold stands; `observed` members are trusted only where the observation
215+
* itself is recent.
216+
*/
217+
export interface KnowledgeMemberObservers {
218+
confirmed: readonly string[]
219+
observed: readonly string[]
220+
}
221+
222+
/** What a search resolves once about its sources and the caller's standing in them. */
223+
export interface SearchAccessPlan {
224+
connectors: KnowledgeConnectorEligibility
225+
observers: KnowledgeMemberObservers
226+
/** Connectors the caller is an active member of, whose documents they read broadly. */
227+
memberSources: readonly string[]
228+
}
229+
211230
export interface KnowledgeConnectorEligibility {
212231
/** Documents carry the workspace ACL. */
213232
workspace: readonly string[]
@@ -238,9 +257,10 @@ export interface KnowledgeConnectorEligibility {
238257
*/
239258
export function knowledgeCandidateAccessConditionForConnectors(
240259
scope: KnowledgeAccessScope | SystemAccessScope,
241-
eligibility: KnowledgeConnectorEligibility,
260+
plan: SearchAccessPlan,
242261
liveSourceAccess: SQL = sql`true`
243262
): SQL {
263+
const eligibility = plan.connectors
244264
if (scope.kind === 'system') return documentConnectorIsActive()
245265
if (scope.tokens.length === 0) return sql`false`
246266
const tokens = textArrayLiteral(scope.tokens)
@@ -270,7 +290,7 @@ export function knowledgeCandidateAccessConditionForConnectors(
270290
((${document.connectorId} IS NULL OR ${inConnectors(eligibility.workspace)})
271291
AND ${document.acl} = ARRAY['ws']::text[])
272292
OR ${mirrored(eligibility.admin, sql`${document.aclVerifiedAt} > ${cutoff}`)}
273-
OR ${mirrored(eligibility.members, memberObservationCondition(tokens, cutoff))}
293+
OR ${mirrored(eligibility.members, resolvedObservationCondition(plan.observers, cutoff))}
274294
)
275295
)`
276296
}
@@ -282,6 +302,29 @@ export function knowledgeCandidateAccessConditionForConnectors(
282302
* inside the access predicate's `OR`, PostgreSQL instead hashes every observation in the table
283303
* once per statement, a fixed cost paid by every query that carries the predicate.
284304
*/
305+
/**
306+
* The same membership, resolved ahead of the query: each candidate costs one lookup on the
307+
* observation key instead of a join to the member behind it. Equivalent by construction — the ids
308+
* are the members that join would have matched, and each one's freshness rule is carried over.
309+
*/
310+
function resolvedObservationCondition(observers: KnowledgeMemberObservers, cutoff: SQL): SQL {
311+
if (observers.confirmed.length === 0 && observers.observed.length === 0) return sql`false`
312+
const byMember = (ids: readonly string[]): SQL =>
313+
sql`${knowledgeDocumentObservation.memberId} = ANY(${textArrayLiteral([...ids])})`
314+
const current =
315+
observers.confirmed.length === 0
316+
? sql`${byMember(observers.observed)} AND ${knowledgeDocumentObservation.lastSeenAt} > ${cutoff}`
317+
: observers.observed.length === 0
318+
? byMember(observers.confirmed)
319+
: sql`(${byMember(observers.confirmed)}
320+
OR (${byMember(observers.observed)} AND ${knowledgeDocumentObservation.lastSeenAt} > ${cutoff}))`
321+
return sql`EXISTS (
322+
SELECT 1 FROM ${knowledgeDocumentObservation}
323+
WHERE ${knowledgeDocumentObservation.documentId} = ${document.id}
324+
AND ${current}
325+
)`
326+
}
327+
285328
function memberObservationCondition(tokens: SQL, cutoff: SQL): SQL {
286329
return sql`EXISTS (
287330
SELECT 1 FROM ${knowledgeDocumentObservation}

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

Lines changed: 7 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1434,7 +1434,6 @@ describe('permitted-document planner', () => {
14341434
let exactRows: Array<{ id: string }>
14351435
let traversedRows: Array<{ id: string; distance?: number }>
14361436
let rerankRows: Array<ReturnType<typeof hit>>
1437-
let memberedSources: Array<{ connectorId: string }>
14381437
let indexedSourceRows: Array<{ name: string; connectorId: string }>
14391438
let sourceExactRows: Array<{ id: string; distance: number }>
14401439

@@ -1445,7 +1444,6 @@ describe('permitted-document planner', () => {
14451444
traversedRows = []
14461445
rerankRows = []
14471446
sourceExactRows = []
1448-
memberedSources = []
14491447
indexedSourceRows = []
14501448
dbChainMockFns.execute.mockImplementation(async (query) => {
14511449
const statement = render(query).sql
@@ -1501,17 +1499,17 @@ describe('permitted-document planner', () => {
15011499
liveProofRequired: [],
15021500
}
15031501
/** Membership decides the walk; the sliced source contributes enumerated documents. */
1504-
memberedSources = [{ connectorId: 'member-src' }]
1502+
const memberSources = ['member-src']
1503+
const observers = { confirmed: ['m-1'], observed: [] }
15051504
indexedSourceRows = [{ name: 'idx', connectorId: 'member-src' }]
1506-
queueTableRows(schemaMock.knowledgeConnectorMember, memberedSources)
15071505
sourceExactRows = [{ id: 'sliced-hit', distance: 0.05 }]
15081506
traversedRows = [{ id: 'walked-hit', distance: 0.2 }]
15091507
rerankRows = [hit('sliced-hit', 'sliced-src'), hit('walked-hit', 'member-src')]
15101508
queueTableRows(schemaMock.embedding, rerankRows)
15111509
await handleVectorOnlySearch({
15121510
...params,
15131511
permitted: { kind: 'unbounded' },
1514-
connectorEligibility: eligibility,
1512+
accessPlan: { connectors: eligibility, observers, memberSources },
15151513
})
15161514
const walks = statements().filter((query) => query.sql.includes('AS visible'))
15171515
expect(walks).toHaveLength(1)
@@ -1525,18 +1523,16 @@ describe('permitted-document planner', () => {
15251523
})
15261524

15271525
it('ranks every source exactly when the caller is a member of none', async () => {
1528-
memberedSources = []
15291526
sourceExactRows = [{ id: 'sliced-hit', distance: 0.05 }]
15301527
rerankRows = [hit('sliced-hit', 'sliced-src')]
15311528
queueTableRows(schemaMock.embedding, rerankRows)
15321529
await handleVectorOnlySearch({
15331530
...params,
15341531
permitted: { kind: 'unbounded' },
1535-
connectorEligibility: {
1536-
workspace: [],
1537-
admin: ['sliced-src'],
1538-
members: [],
1539-
liveProofRequired: [],
1532+
accessPlan: {
1533+
connectors: { workspace: [], admin: ['sliced-src'], members: [], liveProofRequired: [] },
1534+
observers: { confirmed: [], observed: [] },
1535+
memberSources: [],
15401536
},
15411537
})
15421538
expect(statements().filter((query) => query.sql.includes('AS visible'))).toHaveLength(0)

0 commit comments

Comments
 (0)