Skip to content

Commit ef5bbe1

Browse files
committed
improvement(knowledge): decide mirrored-permission freshness from the document row
Every ACL writer now records `document.acl_valid_until`: admin mode a fixed window after its verification evidence, members mode the earliest expiry among the observations its ACL was materialized from, and NULL where no source evidence applies. Behind the `knowledge-acl-expiry` flag the read predicate compares that column instead of re-deriving freshness per candidate through the connector's members and observations, which costs a lookup per document examined. The stored expiry is never later than the evidence behind it, so the flagged predicate can hide a document the per-candidate proof would still admit and can never admit one it would refuse. Script migration 0020 backfills existing rows.
1 parent d5df6f6 commit ef5bbe1

21 files changed

Lines changed: 28403 additions & 22 deletions

‎apps/sim/lib/core/config/env.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -630,6 +630,7 @@ export const env = createEnv({
630630
CREDENTIAL_GROUPS: z.boolean().optional(), // Enable enterprise Credential Groups globally
631631
KNOWLEDGE_MEMBER_ACCESS: z.boolean().optional(), // Enable per-member knowledge connectors and hybrid-by-default retrieval globally
632632
KNOWLEDGE_TIN_KEYWORD: z.boolean().optional(), // Rank large-scope keyword retrieval through the Tin text index where it exists
633+
KNOWLEDGE_ACL_EXPIRY: z.boolean().optional(), // Read mirrored-permission freshness from document.acl_valid_until instead of per-candidate evidence
633634

634635
// Organizations - for self-hosted deployments
635636
ORGANIZATIONS_ENABLED: z.boolean().optional(), // Enable organizations on self-hosted (bypasses plan requirements)

‎apps/sim/lib/core/config/feature-flags.test.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -177,6 +177,19 @@ describe('isFeatureEnabled', () => {
177177
})
178178
})
179179

180+
describe('knowledge-acl-expiry flag', () => {
181+
it('is a global switch', async () => {
182+
expect(await isFeatureEnabled('knowledge-acl-expiry')).toBe(false)
183+
envRef.KNOWLEDGE_ACL_EXPIRY = true
184+
expect(await isFeatureEnabled('knowledge-acl-expiry')).toBe(true)
185+
})
186+
187+
it('follows an AppConfig global rule', async () => {
188+
withAppConfig({ 'knowledge-acl-expiry': { enabled: true } })
189+
expect(await isFeatureEnabled('knowledge-acl-expiry')).toBe(true)
190+
})
191+
})
192+
180193
describe('knowledge-member-access flag', () => {
181194
it('uses a global fallback switch off AppConfig', async () => {
182195
expect(await isFeatureEnabled('knowledge-member-access')).toBe(false)

‎apps/sim/lib/core/config/feature-flags.ts‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,15 @@ const FEATURE_FLAGS = {
9393
'KNOWLEDGE_MEMBER_ACCESS.',
9494
fallback: 'KNOWLEDGE_MEMBER_ACCESS',
9595
},
96+
'knowledge-acl-expiry': {
97+
description:
98+
'Read mirrored-permission freshness from document.acl_valid_until, written by every ACL ' +
99+
'writer, instead of re-deriving it per candidate from verification and observation ' +
100+
'timestamps. Strictly no wider than the per-candidate proof: a stored expiry is never ' +
101+
'later than the evidence behind it. Requires the 0020 backfill. Off-AppConfig falls back ' +
102+
'to KNOWLEDGE_ACL_EXPIRY.',
103+
fallback: 'KNOWLEDGE_ACL_EXPIRY',
104+
},
96105
'knowledge-tin-keyword': {
97106
description:
98107
'Rank keyword retrieval for members whose permitted set is too large to enumerate through ' +

‎apps/sim/lib/knowledge/__integration__/migration-fixture.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -56,7 +56,7 @@ export async function createEnterpriseSearchMigrationFixture(databaseUrl: string
5656
CREATE TABLE document (
5757
id text PRIMARY KEY, external_id text, connector_id text, knowledge_base_id text,
5858
tag1 text, tag2 text, tag3 text, tag4 text, tag5 text, tag6 text, tag7 text,
59-
acl text[] NOT NULL DEFAULT '{ws}', storage_key text,
59+
acl text[] NOT NULL DEFAULT '{ws}', storage_key text, acl_valid_until timestamp,
6060
user_excluded boolean NOT NULL DEFAULT false, archived_at timestamp
6161
);
6262
CREATE INDEX doc_connector_id_idx ON document(connector_id);
Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,2 +1,14 @@
1+
import { type SQL, sql } from 'drizzle-orm'
2+
13
/** Source authorization evidence expires even when its sync scheduler has stopped. */
24
export const SOURCE_ACL_MAX_AGE_MS = 24 * 60 * 60 * 1000
5+
6+
/** How long evidence recorded at `recordedAt` stays current, as a SQL expression. */
7+
export function aclValidUntilFrom(recordedAt: SQL): SQL {
8+
return sql`(${recordedAt}) + (${SOURCE_ACL_MAX_AGE_MS} * interval '1 millisecond')`
9+
}
10+
11+
/** Evidence older than this is no longer current. */
12+
export function aclFreshnessCutoff(): SQL {
13+
return sql`statement_timestamp() - (${SOURCE_ACL_MAX_AGE_MS} * interval '1 millisecond')`
14+
}

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

Lines changed: 84 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ vi.mock('@/connectors/registry.server', () => ({ CONNECTOR_REGISTRY: {} }))
2121
const { drizzle } = await import('drizzle-orm/postgres-js')
2222
const schema = await import('@sim/db/schema')
2323
const { persistDocumentAcls } = await import('@/lib/knowledge/connectors/sync-persistence')
24+
const { materializeDocumentAcls } = await import('@/lib/knowledge/connectors/member-observations')
2425
const { mergeMirroredAcls, hideUnlistedDocuments } = await import(
2526
'@/lib/knowledge/connectors/mirrored-acls'
2627
)
@@ -87,7 +88,8 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
8788
join = false,
8889
githubInstallationGrants?: GitHubInstallationReadGrant[],
8990
userId = 'reader',
90-
confluenceSiteGrants?: ConfluenceSiteReadGrant[]
91+
confluenceSiteGrants?: ConfluenceSiteReadGrant[],
92+
storedAclFreshness = false
9193
): Promise<boolean> {
9294
const query = new PgDialect().sqlToQuery(
9395
knowledgeAccessCondition({
@@ -96,6 +98,7 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
9698
tokens,
9799
githubInstallationGrants,
98100
confluenceSiteGrants,
101+
storedAclFreshness,
99102
})
100103
)
101104
const values = query.params.map((value: unknown) => {
@@ -495,6 +498,75 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
495498
expect(await readable(['ws'], 'upload')).toBe(true)
496499
})
497500

501+
/**
502+
* The stored expiry every ACL writer records must decide mirrored freshness exactly as the
503+
* per-candidate proof does, except where it is deliberately stricter: a members-mode ACL expires
504+
* with the earliest observation it names, so one stale observer hides the document from readers
505+
* whose own observation is still current, until the sweep drops it and re-materializes.
506+
*/
507+
it('decides mirrored freshness from the row, never admitting what per-candidate evidence refuses', async () => {
508+
const proof = (tokens: string[], id: string) => readable(tokens, id)
509+
const stored = (tokens: string[], id: string) =>
510+
readable(tokens, id, false, undefined, 'reader', undefined, true)
511+
512+
await putDocument('admin-fresh', ['u:alice@corp.com'])
513+
await connection.unsafe(
514+
"UPDATE document SET acl_valid_until = acl_verified_at + interval '24 hours' WHERE id = 'admin-fresh'"
515+
)
516+
await putDocument('admin-stale', ['u:alice@corp.com'])
517+
await connection.unsafe(`UPDATE document SET acl_verified_at = statement_timestamp() - interval '25 hours',
518+
acl_valid_until = statement_timestamp() - interval '1 hour' WHERE id = 'admin-stale'`)
519+
await putDocument('admin-unbackfilled', ['u:alice@corp.com'])
520+
await connection.unsafe(
521+
"UPDATE document SET acl_valid_until = NULL WHERE id = 'admin-unbackfilled'"
522+
)
523+
524+
await connection.unsafe(
525+
`INSERT INTO knowledge_connector_member(id, workspace_id, connector_id, subject_token, status) VALUES
526+
('alice', 'workspace', 'members', $1, 'active'), ('bob', 'workspace', 'members', $2, 'active')`,
527+
[alice, bob]
528+
)
529+
const putMembers = async (id: string, seen: string[]) => {
530+
await connection.unsafe(
531+
"INSERT INTO document(id, connector_id, acl) VALUES ($1, 'members', string_to_array($2, ','))",
532+
[id, [alice, bob].join(',')]
533+
)
534+
await connection.unsafe(
535+
`INSERT INTO knowledge_document_observation VALUES
536+
($1, 'alice', statement_timestamp() - ($2)::interval),
537+
($1, 'bob', statement_timestamp() - ($3)::interval)`,
538+
[id, seen[0], seen[1]]
539+
)
540+
/** The real writer decides both the ACL and its expiry, so this cannot drift from it. */
541+
await materializeDocumentAcls('members', [id], drizzle(connection, { schema }))
542+
}
543+
await putMembers('members-fresh', ['1 hour', '2 hours'])
544+
await putMembers('members-stale', ['25 hours', '26 hours'])
545+
await putMembers('members-mixed', ['1 hour', '26 hours'])
546+
547+
for (const [id, tokens, expected] of [
548+
['admin-fresh', ['u:alice@corp.com'], true],
549+
['admin-stale', ['u:alice@corp.com'], false],
550+
['members-fresh', [alice], true],
551+
['members-stale', [alice], false],
552+
] as const) {
553+
expect(await proof([...tokens], id)).toBe(expected)
554+
expect(await stored([...tokens], id)).toBe(expected)
555+
}
556+
/** Stricter, never wider: the fresh observer waits for the sweep. */
557+
expect(await proof([alice], 'members-mixed')).toBe(true)
558+
expect(await stored([alice], 'members-mixed')).toBe(false)
559+
expect(await proof([bob], 'members-mixed')).toBe(false)
560+
expect(await stored([bob], 'members-mixed')).toBe(false)
561+
/** A row the backfill has not reached yet is unreadable rather than assumed current. */
562+
expect(await proof(['u:alice@corp.com'], 'admin-unbackfilled')).toBe(true)
563+
expect(await stored(['u:alice@corp.com'], 'admin-unbackfilled')).toBe(false)
564+
/** Uploads and workspace-mode documents carry no expiry and are unaffected. */
565+
await connection.unsafe("INSERT INTO document(id) VALUES ('upload')")
566+
expect(await proof(['ws'], 'upload')).toBe(true)
567+
expect(await stored(['ws'], 'upload')).toBe(true)
568+
})
569+
498570
it('does not let one member refresh another member’s stale observation', async () => {
499571
await connection.unsafe(
500572
"INSERT INTO document(id, connector_id, acl) VALUES ('shared', 'members', string_to_array($1, ','))",
@@ -545,6 +617,17 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
545617
"SELECT jsonb_typeof(acl_requirements) AS shape, acl_requirements FROM document WHERE id = 'persisted'"
546618
)
547619
expect(stored).toEqual({ shape: 'array', acl_requirements: [[space], [page]] })
620+
/** The writer records when its evidence expires, so a reader never re-derives it. */
621+
const [expiry] = await connection.unsafe(
622+
`SELECT acl_valid_until = acl_verified_at + interval '24 hours' AS matches_evidence
623+
FROM document WHERE id = 'persisted'`
624+
)
625+
expect(expiry).toEqual({ matches_evidence: true })
626+
await persistDocumentAcls('admin', new Map([['page', { acl: [], requirements: [] }]]), executor)
627+
const [revoked] = await connection.unsafe(
628+
"SELECT acl_valid_until FROM document WHERE id = 'persisted'"
629+
)
630+
expect(revoked).toEqual({ acl_valid_until: null })
548631
})
549632

550633
it.each([

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

Lines changed: 18 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ import {
1212
} from '@sim/db/schema'
1313
import { type SQL, sql } from 'drizzle-orm'
1414
import { EXTERNAL_GROUP_STALE_AFTER_MS } from '@/lib/knowledge/access/external-groups'
15-
import { SOURCE_ACL_MAX_AGE_MS } from '@/lib/knowledge/access/freshness'
15+
import { aclFreshnessCutoff } from '@/lib/knowledge/access/freshness'
1616
import { confluenceReaderGroupCondition } from '@/lib/knowledge/access/group-membership'
1717
import type { KnowledgeAccessScope, SystemAccessScope } from '@/lib/knowledge/access/types'
1818
import { documentConnectorIsActive } from '@/lib/knowledge/documents/connector-lifecycle'
@@ -220,7 +220,21 @@ function storedKnowledgeAccessCondition(
220220
if (scope.kind === 'system') return documentConnectorIsActive()
221221
if (scope.tokens.length === 0) return sql`false`
222222
const tokens = textArrayLiteral(scope.tokens)
223-
const cutoff = sql`statement_timestamp() - (${SOURCE_ACL_MAX_AGE_MS} * interval '1 millisecond')`
223+
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 is checked with one comparison instead of a lookup per
227+
* document into its connector's members and observations. The stored expiry is never later than
228+
* the evidence it was written from: a members-mode document expires with the earliest
229+
* observation its ACL names, so this can hide a document the per-candidate proof would still
230+
* admit, and can never admit one it would refuse.
231+
*/
232+
const mirroredIsCurrent = (mode: 'admin' | 'members'): SQL =>
233+
scope.storedAclFreshness
234+
? sql`${knowledgeConnector.accessMode} = ${mode} AND ${document.aclValidUntil} > statement_timestamp()`
235+
: mode === 'admin'
236+
? sql`${knowledgeConnector.accessMode} = 'admin' AND ${document.aclVerifiedAt} > ${cutoff}`
237+
: sql`${knowledgeConnector.accessMode} = 'members' AND ${memberObservationCondition(tokens, cutoff)}`
224238
return sql`(
225239
${aclOverlap(tokens)}
226240
AND NOT EXISTS (
@@ -240,8 +254,8 @@ function storedKnowledgeAccessCondition(
240254
AND (
241255
(${knowledgeConnector.accessMode} = 'workspace' AND ${document.acl} = ARRAY['ws']::text[])
242256
OR (${document.acl} <> ARRAY['ws']::text[] AND (
243-
(${knowledgeConnector.accessMode} = 'admin' AND ${document.aclVerifiedAt} > ${cutoff})
244-
OR (${knowledgeConnector.accessMode} = 'members' AND ${memberObservationCondition(tokens, cutoff)})
257+
(${mirroredIsCurrent('admin')})
258+
OR (${mirroredIsCurrent('members')})
245259
))
246260
)
247261
)

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

Lines changed: 34 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import {
1616
import { createLogger } from '@sim/logger'
1717
import { getErrorMessage } from '@sim/utils/errors'
1818
import { and, eq, gte, inArray, isNull, or, sql } from 'drizzle-orm'
19+
import { isFeatureEnabled } from '@/lib/core/config/feature-flags'
1920
import { OrchestrationError } from '@/lib/core/orchestration/types'
2021
import { type ResourceScope, resourceScopeFromOwner } from '@/lib/core/resource-scope'
2122
import { resourceScopeCondition } from '@/lib/core/resource-scope.server'
@@ -355,25 +356,46 @@ async function resolveKnowledgeIdentity(
355356
if (subject?.kind !== 'sim_user') {
356357
if (context.organizationId)
357358
throw new OrchestrationError('forbidden', 'Organization search requires a user subject')
358-
return { access: WORKSPACE_ACCESS_SCOPE, githubReaders: [], confluenceReaders: [] }
359+
const storedAclFreshness = await storedAclFreshnessEnabled()
360+
return {
361+
/** The shared constant while the rollout has not reached this request. */
362+
access: storedAclFreshness.storedAclFreshness
363+
? { ...WORKSPACE_ACCESS_SCOPE, ...storedAclFreshness }
364+
: WORKSPACE_ACCESS_SCOPE,
365+
githubReaders: [],
366+
confluenceReaders: [],
367+
}
359368
}
360369
return resolveUserKnowledgeIdentity(subject.userId, context)
361370
}
362371

363372
async function resolveUserKnowledgeIdentity(userId: string, context: KnowledgeAccessScopeContext) {
364373
resourceScopeFromOwner(context)
365-
const {
366-
tokens,
367-
githubReaders = [],
368-
confluenceReaders = [],
369-
} = await loadUserAccess(userId, context)
374+
const [{ tokens, githubReaders = [], confluenceReaders = [] }, storedAclFreshness] =
375+
await Promise.all([loadUserAccess(userId, context), storedAclFreshnessEnabled()])
370376
return {
371-
access: { kind: 'user' as const, userId, tokens },
377+
access: { kind: 'user' as const, userId, tokens, ...storedAclFreshness },
372378
githubReaders,
373379
confluenceReaders,
374380
}
375381
}
376382

383+
/**
384+
* Whether this request reads mirrored-permission freshness from the document row. Resolved once
385+
* per scope so the predicate stays synchronous; a configuration read that fails keeps the
386+
* per-candidate proof, which is correct and only slower.
387+
*/
388+
async function storedAclFreshnessEnabled(): Promise<{ storedAclFreshness?: true }> {
389+
try {
390+
return (await isFeatureEnabled('knowledge-acl-expiry')) ? { storedAclFreshness: true } : {}
391+
} catch (error) {
392+
logger.warn('Feature flag read failed; deriving ACL freshness per candidate', {
393+
error: getErrorMessage(error),
394+
})
395+
return {}
396+
}
397+
}
398+
377399
/**
378400
* The scope of a person identified only by user id — the shape session-backed
379401
* routes outside the application layer have in hand. Never call this with a
@@ -384,7 +406,11 @@ export async function resolveUserKnowledgeAccessScope(
384406
userId: string,
385407
workspaceId: string | undefined
386408
): Promise<KnowledgeAccessScope> {
387-
return { kind: 'user', userId, tokens: (await loadUserAccess(userId, { workspaceId })).tokens }
409+
const [access, storedAclFreshness] = await Promise.all([
410+
loadUserAccess(userId, { workspaceId }),
411+
storedAclFreshnessEnabled(),
412+
])
413+
return { kind: 'user', userId, tokens: access.tokens, ...storedAclFreshness }
388414
}
389415

390416
/** Memoises {@link resolveKnowledgeAccessScope} for one operation; a failed lookup is retried on the next call. */

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

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,8 @@ export const WORKSPACE_ACCESS_TOKENS = [PUBLIC_ACCESS_TOKEN, WORKSPACE_ACCESS_TO
3030
export interface WorkspaceAccessScope {
3131
kind: 'workspace'
3232
tokens: typeof WORKSPACE_ACCESS_TOKENS
33+
/** See {@link UserAccessScope.storedAclFreshness}. */
34+
storedAclFreshness?: true
3335
}
3436

3537
export interface UserAccessScope {
@@ -45,6 +47,14 @@ export interface UserAccessScope {
4547
githubInstallationGrants?: readonly GitHubInstallationReadGrant[]
4648
/** Current reader access to the immutable Confluence site behind a central crawl. */
4749
confluenceSiteGrants?: readonly ConfluenceSiteReadGrant[]
50+
/**
51+
* Whether mirrored-permission freshness is read from `document.acl_valid_until` instead of
52+
* being re-derived per candidate from observation and verification timestamps. Resolved once
53+
* per scope from the `knowledge-acl-expiry` rollout flag, so a predicate stays synchronous and
54+
* every surface that carries a scope agrees on one answer for the whole request. Absent while
55+
* the flag is off, so a scope is unchanged until the rollout reaches it.
56+
*/
57+
storedAclFreshness?: true
4858
}
4959

5060
export interface ConfluenceSiteReadGrant {

‎apps/sim/lib/knowledge/connectors/member-observations.ts‎

Lines changed: 27 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import {
2121
sql,
2222
} from 'drizzle-orm'
2323
import type { DbOrTx } from '@/lib/db/types'
24+
import { aclValidUntilFrom } from '@/lib/knowledge/access/freshness'
2425
import { textArrayLiteral } from '@/lib/knowledge/access/predicate'
2526
import {
2627
MEMBER_OBSERVATION_STALE_AFTER_HOURS,
@@ -66,6 +67,25 @@ function observedAcl() {
6667
), '{}'::text[])`
6768
}
6869

70+
/**
71+
* When a members-mode document's ACL stops being current evidence: the earliest expiry among the
72+
* active observations it was materialized from, so the row fails closed as soon as any observer
73+
* it names goes stale. The staleness sweep then drops that observation and re-materializes, which
74+
* restores the document for the observers that are still fresh.
75+
*/
76+
function observedAclValidUntil() {
77+
return sql<Date | null>`(
78+
SELECT min(${aclValidUntilFrom(
79+
sql`GREATEST(${knowledgeDocumentObservation.lastSeenAt}, ${knowledgeConnectorMember.memberSyncedThrough})`
80+
)})
81+
FROM ${knowledgeDocumentObservation}
82+
JOIN ${knowledgeConnectorMember}
83+
ON ${knowledgeConnectorMember.id} = ${knowledgeDocumentObservation.memberId}
84+
AND ${knowledgeConnectorMember.status} = 'active'
85+
WHERE ${knowledgeDocumentObservation.documentId} = ${document.id}
86+
)`
87+
}
88+
6989
function observationQuery() {
7090
return db
7191
.select({ one: sql`1` })
@@ -340,12 +360,17 @@ export async function materializeDocumentAcls(
340360
const batch = ids.slice(offset, offset + MATERIALIZE_BATCH_SIZE)
341361
const rows = await executor
342362
.update(document)
343-
.set({ acl: observedAcl(), aclRequirements: [], aclVerifiedAt: null })
363+
.set({
364+
acl: observedAcl(),
365+
aclRequirements: [],
366+
aclVerifiedAt: null,
367+
aclValidUntil: observedAclValidUntil(),
368+
})
344369
.where(
345370
and(
346371
inArray(document.id, batch),
347372
eq(document.connectorId, connectorId),
348-
sql`(${document.acl} IS DISTINCT FROM ${observedAcl()} OR ${document.aclRequirements} <> '[]'::jsonb OR ${document.aclVerifiedAt} IS NOT NULL)`
373+
sql`(${document.acl} IS DISTINCT FROM ${observedAcl()} OR ${document.aclRequirements} <> '[]'::jsonb OR ${document.aclVerifiedAt} IS NOT NULL OR ${document.aclValidUntil} IS DISTINCT FROM ${observedAclValidUntil()})`
349374
)
350375
)
351376
.returning({ id: document.id })

0 commit comments

Comments
 (0)