Skip to content

Commit 3c2bb77

Browse files
committed
improvement(knowledge): resolve connector state once per search, not per candidate
Deletion, archival, a pending access rewrite, the organization's integration approval and the access mode are facts about a connector, so a search resolves them once and filters candidates by the resulting ids. A connector whose reader access is proven live per request keeps its per-row lookup, which is where that proof lives. On an organization index, the candidate predicate's connector lookup drops from one per document examined to one per document of a live-proof source: 24,571 evaluations to 610 for a member reading ~25k documents, and the access check from ~390ms to ~160ms.
1 parent 9d879b8 commit 3c2bb77

5 files changed

Lines changed: 267 additions & 31 deletions

File tree

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
import { db } from '@sim/db'
2+
import { knowledgeConnector } from '@sim/db/schema'
3+
import { and, eq, inArray, isNull, sql } from 'drizzle-orm'
4+
import type { KnowledgeConnectorEligibility } from '@/lib/knowledge/access/predicate'
5+
import { searchIntegrationAccessCondition } from '@/lib/knowledge/search/integration-policy'
6+
7+
/**
8+
* The connectors a search may read from, grouped by access mode, with the ones whose reader access
9+
* must be proven live marked.
10+
*
11+
* Deletion, archival, a pending access rewrite and the organization's integration approval are
12+
* facts about a connector. Resolving them once per query — there are tens of connectors against
13+
* hundreds of thousands of documents — leaves each candidate its own columns to check.
14+
*/
15+
export async function resolveConnectorEligibility(
16+
knowledgeBaseIds: readonly string[]
17+
): Promise<KnowledgeConnectorEligibility> {
18+
const eligibility: {
19+
workspace: string[]
20+
admin: string[]
21+
members: string[]
22+
liveProofRequired: string[]
23+
} = { workspace: [], admin: [], members: [], liveProofRequired: [] }
24+
if (knowledgeBaseIds.length === 0) return eligibility
25+
const rows = await db
26+
.select({
27+
id: knowledgeConnector.id,
28+
accessMode: knowledgeConnector.accessMode,
29+
connectorType: knowledgeConnector.connectorType,
30+
/** A GitHub connector is gated only where it names the immutable repository behind a grant. */
31+
githubRepository: sql<boolean>`${knowledgeConnector.sourceConfig}::jsonb ? 'githubRepositoryId'`,
32+
})
33+
.from(knowledgeConnector)
34+
.where(
35+
and(
36+
inArray(knowledgeConnector.knowledgeBaseId, [...knowledgeBaseIds]),
37+
isNull(knowledgeConnector.deletedAt),
38+
isNull(knowledgeConnector.archivedAt),
39+
eq(knowledgeConnector.accessRewritePending, false),
40+
searchIntegrationAccessCondition()
41+
)
42+
)
43+
for (const row of rows) {
44+
if (row.accessMode === 'workspace') eligibility.workspace.push(row.id)
45+
else if (row.accessMode === 'admin') eligibility.admin.push(row.id)
46+
else if (row.accessMode === 'members') eligibility.members.push(row.id)
47+
else continue
48+
const live =
49+
(row.connectorType === 'github' && row.githubRepository) ||
50+
(row.connectorType === 'confluence' && row.accessMode === 'admin')
51+
if (live) eligibility.liveProofRequired.push(row.id)
52+
}
53+
return eligibility
54+
}

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

Lines changed: 73 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
* @vitest-environment node
33
*/
44
import { readFile } from 'node:fs/promises'
5+
import type { SQL } from 'drizzle-orm'
56
import type postgres from 'postgres'
67
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
78
import { createEnterpriseSearchMigrationFixture } from '@/lib/knowledge/__integration__/migration-fixture'
@@ -28,7 +29,11 @@ const { mergeMirroredAcls, hideUnlistedDocuments } = await import(
2829
'@/lib/knowledge/connectors/mirrored-acls'
2930
)
3031
const { PgDialect } = await import('drizzle-orm/pg-core')
31-
const { knowledgeAccessCondition } = await import('@/lib/knowledge/access/predicate')
32+
const {
33+
knowledgeAccessCondition,
34+
knowledgeCandidateAccessConditionForConnectors,
35+
knowledgeMetadataCandidateAccessCondition,
36+
} = await import('@/lib/knowledge/access/predicate')
3237
const { confluencePageAcl } = await import('@/lib/knowledge/access/confluence-permissions')
3338

3439
/** Explicit opt-in; every table and index belongs to an isolated disposable schema. */
@@ -113,6 +118,20 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
113118
return rows.length > 0
114119
}
115120

121+
/** Runs any predicate over one document, so two shapes can be compared row by row. */
122+
async function admits(condition: SQL, documentId: string): Promise<boolean> {
123+
const query = new PgDialect().sqlToQuery(condition)
124+
const values = query.params.map((value: unknown) => {
125+
if (typeof value === 'string' || typeof value === 'number') return value
126+
throw new Error('The access predicate must bind scalar strings and numbers')
127+
})
128+
const rows = await connection.unsafe(
129+
`SELECT document.id FROM document WHERE ${query.sql} AND document.id = $${values.length + 1}`,
130+
[...values, documentId]
131+
)
132+
return rows.length > 0
133+
}
134+
116135
async function putDocument(id: string, acl: string[], requirements: string[][] = []) {
117136
await connection.unsafe(
118137
`INSERT INTO document(id, connector_id, acl, acl_requirements, acl_verified_at)
@@ -577,6 +596,59 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
577596
expect(await stored(['ws'], 'upload')).toBe(true)
578597
})
579598

599+
it('admits the same documents whether connector state is proven per row or resolved per query', async () => {
600+
const scope = { kind: 'user' as const, userId: 'reader', tokens: [alice, 'u:alice@corp.com'] }
601+
await connection.unsafe(
602+
`INSERT INTO knowledge_connector(id, access_mode, deleted_at, archived_at) VALUES
603+
('gone', 'admin', statement_timestamp(), NULL), ('shelved', 'admin', NULL, statement_timestamp())`
604+
)
605+
await connection.unsafe(
606+
"INSERT INTO knowledge_connector(id, access_mode) VALUES ('ws-mode', 'workspace')"
607+
)
608+
await connection.unsafe(
609+
`INSERT INTO knowledge_connector_member(id, workspace_id, connector_id, subject_token, status)
610+
VALUES ('m-alice', 'workspace', 'members', $1, 'active')`,
611+
[alice]
612+
)
613+
const cases: Array<[string, string, string[]]> = [
614+
['admin-current', 'admin', ['u:alice@corp.com']],
615+
['members-current', 'members', [alice]],
616+
['workspace-doc', 'ws-mode', ['ws']],
617+
['deleted-connector', 'gone', ['u:alice@corp.com']],
618+
['archived-connector', 'shelved', ['u:alice@corp.com']],
619+
]
620+
for (const [id, connectorId, acl] of cases) {
621+
await connection.unsafe(
622+
`INSERT INTO document(id, connector_id, acl, acl_verified_at, acl_valid_until)
623+
VALUES ($1, $2, string_to_array($3, E'\n'), statement_timestamp(), statement_timestamp() + interval '1 hour')`,
624+
[id, connectorId, acl.join('\n')]
625+
)
626+
}
627+
await connection.unsafe(
628+
`INSERT INTO knowledge_document_observation VALUES ('members-current', 'm-alice', statement_timestamp())`
629+
)
630+
await connection.unsafe("INSERT INTO document(id) VALUES ('upload-doc')")
631+
/** What `resolveConnectorEligibility` returns for this base: the connectors it admits, by mode. */
632+
const eligibility = {
633+
workspace: ['ws-mode'],
634+
admin: ['admin'],
635+
members: ['members'],
636+
liveProofRequired: [],
637+
}
638+
const perRow = knowledgeMetadataCandidateAccessCondition(scope)
639+
const perQuery = knowledgeCandidateAccessConditionForConnectors(scope, eligibility)
640+
for (const id of [...cases.map(([documentId]) => documentId), 'upload-doc']) {
641+
expect([id, await admits(perQuery, id)]).toEqual([id, await admits(perRow, id)])
642+
}
643+
/** A connector left out of the resolution is refused, however current its documents are. */
644+
expect(
645+
await admits(
646+
knowledgeCandidateAccessConditionForConnectors(scope, { ...eligibility, admin: [] }),
647+
'admin-current'
648+
)
649+
).toBe(false)
650+
})
651+
580652
it('does not let one member refresh another member’s stale observation', async () => {
581653
await connection.unsafe(
582654
"INSERT INTO document(id, connector_id, acl) VALUES ('shared', 'members', string_to_array($1, ','))",

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

Lines changed: 101 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -176,9 +176,7 @@ function githubInstallationAccessCondition(scope: KnowledgeAccessScope): SQL {
176176
export function knowledgeAccessCondition(scope: KnowledgeAccessScope | SystemAccessScope): SQL {
177177
return storedKnowledgeAccessCondition(
178178
scope,
179-
scope.kind === 'system'
180-
? sql`true`
181-
: sql`(${githubInstallationAccessCondition(scope)} AND ${confluenceSiteAccessCondition(scope)})`
179+
scope.kind === 'system' ? sql`true` : liveSourceAccessCondition(scope)
182180
)
183181
}
184182

@@ -193,6 +191,103 @@ export function knowledgeMetadataCandidateAccessCondition(
193191
return storedKnowledgeAccessCondition(scope, sql`true`)
194192
}
195193

194+
/**
195+
* Mirrored permissions are current while the row says so. Every writer of an ACL records when its
196+
* evidence expires, so a candidate costs one comparison instead of a lookup per document into its
197+
* connector's members and observations. The stored expiry is never later than the evidence it was
198+
* written from: a members-mode document expires with the earliest observation its ACL names, so
199+
* this can hide a document the evidence would still admit, and can never admit one it would
200+
* refuse.
201+
*
202+
* A row without an expiry is one no current writer has touched — a backfill has not reached it, or
203+
* an older app version wrote its ACL during a deploy — so it falls back to proving freshness from
204+
* the evidence itself. That branch is dead once every row carries an expiry.
205+
*/
206+
function isCurrentFrom(evidence: SQL): SQL {
207+
return sql`(
208+
${document.aclValidUntil} > statement_timestamp()
209+
OR (${document.aclValidUntil} IS NULL AND ${evidence})
210+
)`
211+
}
212+
213+
/** Every requirement clause must reach the caller, which preserves source permission intersections. */
214+
function aclRequirementsSatisfied(tokens: SQL): SQL {
215+
return sql`NOT EXISTS (
216+
SELECT 1 FROM jsonb_array_elements(${document.aclRequirements}) AS required_clause(tokens)
217+
WHERE NOT (required_clause.tokens ?| ${tokens})
218+
)`
219+
}
220+
221+
/** The live source proofs this request carries, as the connector-scoped clause both shapes apply. */
222+
function liveSourceAccessCondition(scope: KnowledgeAccessScope): SQL {
223+
return sql`(${githubInstallationAccessCondition(scope)} AND ${confluenceSiteAccessCondition(scope)})`
224+
}
225+
226+
/**
227+
* The connectors a search may read from, resolved once per query: their ids grouped by the shape
228+
* their documents' ACLs take, and separately those whose reader access is proven live per request.
229+
*/
230+
export interface KnowledgeConnectorEligibility {
231+
/** Documents carry the workspace ACL. */
232+
workspace: readonly string[]
233+
/** Documents carry mirrored source permissions verified as a whole. */
234+
admin: readonly string[]
235+
/** Documents carry the subject tokens of the members who observe them. */
236+
members: readonly string[]
237+
/** Of the above, those that additionally require this request's live source proof. */
238+
liveProofRequired: readonly string[]
239+
}
240+
241+
/**
242+
* The candidate predicate with connector state resolved ahead of the query instead of per row.
243+
*
244+
* Deletion, archival, a pending access rewrite, the organization's integration approval and the
245+
* access mode are facts about a connector, not a document, so checking them once per query leaves
246+
* each candidate an id comparison plus its own columns. A connector that still needs this
247+
* request's live source proof keeps the per-row lookup, which is where that proof is expressed.
248+
*
249+
* It narrows exactly as {@link knowledgeMetadataCandidateAccessCondition} does: the eligible ids
250+
* are the connectors that predicate's `EXISTS` would admit, and every document-level clause is
251+
* carried over unchanged.
252+
*/
253+
export function knowledgeCandidateAccessConditionForConnectors(
254+
scope: KnowledgeAccessScope | SystemAccessScope,
255+
eligibility: KnowledgeConnectorEligibility
256+
): SQL {
257+
if (scope.kind === 'system') return documentConnectorIsActive()
258+
if (scope.tokens.length === 0) return sql`false`
259+
const tokens = textArrayLiteral(scope.tokens)
260+
const cutoff = aclFreshnessCutoff()
261+
const liveProof = new Set(eligibility.liveProofRequired)
262+
const inConnectors = (ids: readonly string[]): SQL =>
263+
ids.length === 0
264+
? sql`false`
265+
: sql`${document.connectorId} = ANY(${textArrayLiteral([...ids])})`
266+
const mirrored = (ids: readonly string[], evidence: SQL): SQL => {
267+
const direct = ids.filter((id) => !liveProof.has(id))
268+
const gated = ids.filter((id) => liveProof.has(id))
269+
const current = sql`${document.acl} <> ARRAY['ws']::text[] AND ${isCurrentFrom(evidence)}`
270+
return sql`(
271+
(${inConnectors(direct)} AND ${current})
272+
OR (${inConnectors(gated)} AND ${current} AND EXISTS (
273+
SELECT 1 FROM ${knowledgeConnector}
274+
WHERE ${knowledgeConnector.id} = ${document.connectorId}
275+
AND ${liveSourceAccessCondition(scope)}
276+
))
277+
)`
278+
}
279+
return sql`(
280+
${aclOverlap(tokens)}
281+
AND ${aclRequirementsSatisfied(tokens)}
282+
AND (
283+
((${document.connectorId} IS NULL OR ${inConnectors(eligibility.workspace)})
284+
AND ${document.acl} = ARRAY['ws']::text[])
285+
OR ${mirrored(eligibility.admin, sql`${document.aclVerifiedAt} > ${cutoff}`)}
286+
OR ${mirrored(eligibility.members, memberObservationCondition(tokens, cutoff))}
287+
)
288+
)`
289+
}
290+
196291
/**
197292
* A members-mode document is readable while one of the caller's active member identities on its
198293
* connector still observes it, freshly. Correlated on the document so each check is a lookup on
@@ -221,29 +316,9 @@ function storedKnowledgeAccessCondition(
221316
if (scope.tokens.length === 0) return sql`false`
222317
const tokens = textArrayLiteral(scope.tokens)
223318
const cutoff = aclFreshnessCutoff()
224-
/**
225-
* Mirrored permissions are current while the row says so. Every writer of an ACL records when
226-
* its evidence expires, so a candidate costs one comparison instead of a lookup per document
227-
* into its connector's members and observations. The stored expiry is never later than the
228-
* evidence it was written from: a members-mode document expires with the earliest observation
229-
* its ACL names, so this can hide a document the evidence would still admit, and can never
230-
* admit one it would refuse.
231-
*
232-
* A row without an expiry is one no current writer has touched — the rollout's backfill has not
233-
* reached it, or an older app version wrote its ACL during a deploy — so it falls back to
234-
* proving freshness from the evidence itself. That branch is dead once every row carries an
235-
* expiry.
236-
*/
237-
const isCurrent = (evidence: SQL): SQL => sql`(
238-
${document.aclValidUntil} > statement_timestamp()
239-
OR (${document.aclValidUntil} IS NULL AND ${evidence})
240-
)`
241319
return sql`(
242320
${aclOverlap(tokens)}
243-
AND NOT EXISTS (
244-
SELECT 1 FROM jsonb_array_elements(${document.aclRequirements}) AS required_clause(tokens)
245-
WHERE NOT (required_clause.tokens ?| ${tokens})
246-
)
321+
AND ${aclRequirementsSatisfied(tokens)}
247322
AND (
248323
(${document.connectorId} IS NULL AND ${document.acl} = ARRAY['ws']::text[])
249324
OR EXISTS (
@@ -257,8 +332,8 @@ function storedKnowledgeAccessCondition(
257332
AND (
258333
(${knowledgeConnector.accessMode} = 'workspace' AND ${document.acl} = ARRAY['ws']::text[])
259334
OR (${document.acl} <> ARRAY['ws']::text[] AND (
260-
(${knowledgeConnector.accessMode} = 'admin' AND ${isCurrent(sql`${document.aclVerifiedAt} > ${cutoff}`)})
261-
OR (${knowledgeConnector.accessMode} = 'members' AND ${isCurrent(memberObservationCondition(tokens, cutoff))})
335+
(${knowledgeConnector.accessMode} = 'admin' AND ${isCurrentFrom(sql`${document.aclVerifiedAt} > ${cutoff}`)})
336+
OR (${knowledgeConnector.accessMode} = 'members' AND ${isCurrentFrom(memberObservationCondition(tokens, cutoff))})
262337
))
263338
)
264339
)

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ export type SearchStage =
2828
| 'access_scope'
2929
| 'defaults'
3030
| 'retrieval'
31+
| 'connector_eligibility'
3132
| 'permitted_documents'
3233
| 'result_provenance'
3334
| 'reranking'

0 commit comments

Comments
 (0)