Skip to content

Commit 8a64dd1

Browse files
committed
fix(knowledge): decide unfilled projection rows on their document and start the backfill from the outbox
An unfilled row no longer passes the on-row candidate predicate outright: it is decided on its document under the resolved candidate predicate, the join per candidate every row paid before the columns existed, so the bounded candidate pools and the exact slice hold only rows that hydration will keep. A partial index on the unfilled rows keeps that branch, and the backfill's keyset pages, an index probe. The migration also leaves one outbox event whose handler starts the backfill task, so it runs after every deploy without an operator. Claude-Session: https://claude.ai/code/session_01XU6c7pKRpa5CMoMHKDdqxX
1 parent 1dddb8b commit 8a64dd1

11 files changed

Lines changed: 167 additions & 57 deletions

‎apps/sim/lib/core/outbox/processor.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import { slackSearchOutboxHandlers } from '@/lib/knowledge/application/slack-sea
1818
import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error'
1919
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
2020
import { recoverKnowledgeDocumentProcessing } from '@/lib/knowledge/documents/processing-recovery'
21+
import { projectionSourceAclBackfillOutboxHandlers } from '@/lib/knowledge/search/projection-source-acl-backfill'
2122
import { inboxCleanupOutboxHandlers } from '@/lib/mothership/inbox/cleanup-outbox'
2223
import { organizationResourceCleanupOutboxHandlers } from '@/lib/organizations/resource-cleanup'
2324
import { workspaceFileLiveDocOutboxHandlers } from '@/lib/uploads/contexts/workspace/workspace-file-live-doc-outbox'
@@ -42,6 +43,7 @@ const handlers = {
4243
...invitationMigrationOutboxHandlers,
4344
...directGrantOutboxHandlers,
4445
...knowledgeDocumentProcessingOutboxHandlers,
46+
...projectionSourceAclBackfillOutboxHandlers,
4547
...organizationResourceCleanupOutboxHandlers,
4648
...inboxCleanupOutboxHandlers,
4749
...permissionAccessRequestOutboxHandlers,

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

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -623,19 +623,23 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
623623
expect(await admits(perRow, 'members-other')).toBe(false)
624624
expect(await onRowAdmits('members-other')).toBe(false)
625625
/**
626-
* A chunk the backfill has not reached carries no source or ACL yet. It is ranked, and decided
627-
* at hydration under the full predicate, exactly as every candidate was before the columns
628-
* existed — search must not depend on the backfill, and must not lose a document to it.
626+
* A chunk the backfill has not reached carries no source or ACL yet and is decided on its
627+
* document, as every candidate was before the columns existed: search must not depend on the
628+
* backfill, must not lose a document to it, and must not rank an unreadable one because of it.
629629
*/
630630
await connection.unsafe(
631631
"INSERT INTO document(id, connector_id, acl, acl_verified_at) VALUES ('unfilled-other', 'members', ARRAY['s:slack:-:bob'], statement_timestamp())"
632632
)
633+
await connection.unsafe(`INSERT INTO embedding_search(id, document_id) VALUES
634+
('unfilled-other-chunk', 'unfilled-other'), ('unfilled-admin-chunk', 'admin-current'),
635+
('unfilled-members-chunk', 'members-current'), ('unfilled-gone-chunk', 'deleted-connector')`)
633636
await connection.unsafe(
634-
"INSERT INTO embedding_search(id, document_id) VALUES ('unfilled-other-chunk', 'unfilled-other'), ('admin-current-unfilled', 'admin-current')"
637+
"DELETE FROM embedding_search WHERE id IN ('admin-current-chunk', 'members-current-chunk', 'deleted-connector-chunk')"
635638
)
636-
expect(await onRowAdmits('unfilled-other')).toBe(true)
637-
expect(await admits(perRow, 'unfilled-other')).toBe(false)
639+
expect(await onRowAdmits('unfilled-other')).toBe(false)
638640
expect(await onRowAdmits('admin-current')).toBe(true)
641+
expect(await onRowAdmits('members-current')).toBe(true)
642+
expect(await onRowAdmits('deleted-connector')).toBe(false)
639643
/**
640644
* Candidate ranking defers the live source proof, as the per-row candidate predicate does: a
641645
* caller holds those grants only after authorization, so applying the clause during ranking

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

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -34,19 +34,24 @@ describe('projectionCandidateAccessCondition', () => {
3434
memberSources: [],
3535
}
3636

37-
it('admits a row the backfill has not filled, and decides a filled row on its mirrored columns', () => {
37+
it('decides a filled row on its mirrored columns and an unfilled row on its document', () => {
3838
const { sql, params } = render(
3939
projectionCandidateAccessCondition(
4040
embeddingSearch,
4141
{ kind: 'user', userId: 'user-1', tokens: ['ws', 'u:alice'] },
4242
plan
4343
)
4444
)
45-
expect(sql).toBe(
46-
'("embedding_search"."acl" IS NULL OR ("embedding_search"."acl" && ARRAY[$1, $2]::text[]\n' +
47-
' AND ("embedding_search"."connector_id" IS NULL OR "embedding_search"."connector_id" = ANY(ARRAY[$3]::text[]))))'
45+
expect(sql).toContain(
46+
'("embedding_search"."acl" IS NULL AND EXISTS (\n SELECT 1 FROM "document"\n WHERE "document"."id" = "embedding_search"."document_id"\n AND ('
4847
)
49-
expect(params).toEqual(['ws', 'u:alice', 'ws-src'])
48+
expect(sql).toContain('"document"."acl" && ARRAY[$1, $2]::text[]')
49+
expect(sql).toMatch(
50+
/OR \("embedding_search"\."acl" && ARRAY\[\$\d+, \$\d+\]::text\[\]\n {4}AND \("embedding_search"\."connector_id" IS NULL OR "embedding_search"\."connector_id" = ANY\(ARRAY\[\$\d+\]::text\[\]\)\)\)\)$/
51+
)
52+
expect(params.slice(0, 2)).toEqual(['ws', 'u:alice'])
53+
expect(params.slice(-3)).toEqual(['ws', 'u:alice', 'ws-src'])
54+
for (const param of params) expect(Array.isArray(param)).toBe(false)
5055
})
5156

5257
it('still denies everything for an empty token set', () => {

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

Lines changed: 16 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -313,13 +313,18 @@ export function knowledgeCandidateAccessConditionForConnectors(
313313
* under the full predicate, before content is returned — this predicate only decides what is worth
314314
* ranking.
315315
*
316-
* A row the backfill has not reached yet carries no ACL (`acl IS NULL`) and passes: it is decided at
317-
* hydration under the full document predicate, exactly as every candidate was before the columns
318-
* existed. The backfill runs in the background, so search never waits on it and never loses a
319-
* document to it.
316+
* A row the backfill has not reached yet carries no ACL (`acl IS NULL`) and is decided on its
317+
* document instead, under {@link knowledgeCandidateAccessConditionForConnectors} — the join per
318+
* candidate that every row paid before the columns existed. The backfill runs in the background,
319+
* so search never waits on it, never loses a document to it, and never ranks an unreadable one
320+
* into a bounded candidate pool because of it.
320321
*/
321322
export function projectionCandidateAccessCondition(
322-
projection: { connectorId: AnyPgColumn | SQL; acl: AnyPgColumn | SQL },
323+
projection: {
324+
connectorId: AnyPgColumn | SQL
325+
acl: AnyPgColumn | SQL
326+
documentId: AnyPgColumn | SQL
327+
},
323328
scope: KnowledgeAccessScope | SystemAccessScope,
324329
plan: SearchAccessPlan
325330
): SQL {
@@ -335,7 +340,12 @@ export function projectionCandidateAccessCondition(
335340
...plan.connectors.admin,
336341
...plan.connectors.members,
337342
]
338-
return sql`(${projection.acl} IS NULL OR (${projection.acl} && ${tokens}
343+
const unfilled = sql`(${projection.acl} IS NULL AND EXISTS (
344+
SELECT 1 FROM ${document}
345+
WHERE ${document.id} = ${projection.documentId}
346+
AND ${knowledgeCandidateAccessConditionForConnectors(scope, plan)}
347+
))`
348+
return sql`(${unfilled} OR (${projection.acl} && ${tokens}
339349
AND (${projection.connectorId} IS NULL OR ${inSources(mirrored)})))`
340350
}
341351

‎apps/sim/lib/knowledge/search/projection-source-acl-backfill.test.ts‎

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,14 +13,21 @@ const { mockBackfill, mockEnd, mockPostgres, mockTasksTrigger } = vi.hoisted(()
1313
vi.mock('@sim/db', () => ({ resolveDbUrl: () => 'postgres://localhost:5432/sim' }))
1414
vi.mock('@sim/db/script-migrations/0021_embedding_search_connector', () => ({
1515
PROJECTION_SOURCE_ACL_TABLES: ['embedding_search', 'embedding_keyword_tin'],
16+
PROJECTION_SOURCE_ACL_BACKFILL_OUTBOX_EVENT: 'knowledge.projection.source_acl.backfill',
1617
backfillProjectionSourceAcl: mockBackfill,
1718
}))
1819
vi.mock('postgres', () => ({ default: mockPostgres }))
1920
vi.mock('@trigger.dev/sdk', () => ({ tasks: { trigger: mockTasksTrigger } }))
2021
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: async () => 'us-east-1' }))
22+
vi.mock('@/lib/core/utils/background', () => ({
23+
runDetached: (_label: string, work: () => Promise<unknown>) => {
24+
void work()
25+
},
26+
}))
2127

2228
import {
2329
enqueueProjectionSourceAclBackfill,
30+
projectionSourceAclBackfillOutboxHandlers,
2431
runProjectionSourceAclBackfill,
2532
} from '@/lib/knowledge/search/projection-source-acl-backfill'
2633

@@ -110,9 +117,15 @@ describe('enqueueProjectionSourceAclBackfill', () => {
110117
expect(mockBackfill).not.toHaveBeenCalled()
111118
})
112119

113-
it('fills the projections inline without one', async () => {
120+
it('fills the projections detached in this process without one', async () => {
114121
await expect(enqueueProjectionSourceAclBackfill({}, false)).resolves.toBeNull()
115122
expect(mockTasksTrigger).not.toHaveBeenCalled()
116-
expect(mockBackfill).toHaveBeenCalledTimes(2)
123+
await vi.waitFor(() => expect(mockBackfill).toHaveBeenCalledTimes(2))
124+
})
125+
126+
it('starts from the outbox event the migration leaves behind', () => {
127+
expect(Object.keys(projectionSourceAclBackfillOutboxHandlers)).toEqual([
128+
'knowledge.projection.source_acl.backfill',
129+
])
117130
})
118131
})

‎apps/sim/lib/knowledge/search/projection-source-acl-backfill.ts‎

Lines changed: 17 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { resolveDbUrl } from '@sim/db'
22
import {
33
backfillProjectionSourceAcl,
4+
PROJECTION_SOURCE_ACL_BACKFILL_OUTBOX_EVENT,
45
PROJECTION_SOURCE_ACL_TABLES,
56
type ProjectionSourceAclTable,
67
} from '@sim/db/script-migrations/0021_embedding_search_connector'
@@ -9,6 +10,8 @@ import postgres from 'postgres'
910
import { resolveTriggerRegion } from '@/lib/core/async-jobs/region'
1011
import { env } from '@/lib/core/config/env'
1112
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
13+
import type { OutboxHandlerRegistry } from '@/lib/core/outbox/service'
14+
import { runDetached } from '@/lib/core/utils/background'
1215

1316
const logger = createLogger('ProjectionSourceAclBackfill')
1417

@@ -76,8 +79,9 @@ export async function runProjectionSourceAclBackfill(
7679

7780
/**
7881
* Starts the backfill the way the table backfill is started: on the deployment's Trigger.dev
79-
* worker when there is one, where bounded runs chain until both projections are filled, and inline
80-
* otherwise. Safe to call again at any time — a run only fills rows still unset.
82+
* worker when there is one, where bounded runs chain until both projections are filled, and
83+
* detached in this process otherwise. Safe to call again at any time — a run only fills rows still
84+
* unset.
8185
*/
8286
export async function enqueueProjectionSourceAclBackfill(
8387
payload: ProjectionSourceAclBackfillPayload = {},
@@ -91,6 +95,16 @@ export async function enqueueProjectionSourceAclBackfill(
9195
logger.info('Projection source and ACL backfill enqueued', { runId: handle.id })
9296
return { runId: handle.id }
9397
}
94-
await runProjectionSourceAclBackfill(payload)
98+
runDetached(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, () => runProjectionSourceAclBackfill(payload))
9599
return null
96100
}
101+
102+
/**
103+
* The event script migration `0021_embedding_search_connector` leaves behind, so the backfill it
104+
* hands off starts once the app that ships this handler is up, on every deployment.
105+
*/
106+
export const projectionSourceAclBackfillOutboxHandlers: OutboxHandlerRegistry = {
107+
[PROJECTION_SOURCE_ACL_BACKFILL_OUTBOX_EVENT]: async () => {
108+
await enqueueProjectionSourceAclBackfill()
109+
},
110+
}

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1527,8 +1527,9 @@ describe('permitted-document planner', () => {
15271527
/** The ranked CTE carries the mirrored source and ACL the predicate tests. */
15281528
expect(statement).toContain('AS connector_id')
15291529
expect(statement).toContain('ranked_tin_chunks.acl')
1530-
/** A row the backfill has not filled (`acl IS NULL`) is ranked and decided at hydration. */
1531-
expect(statement).toContain('" IS NULL OR ("')
1530+
/** A row the backfill has not filled (`acl IS NULL`) is decided on its document instead. */
1531+
expect(statement).toContain('IS NULL AND EXISTS (')
1532+
expect(statement).toContain('ranked_tin_chunks.document_id')
15321533
})
15331534

15341535
it('widens the window for a broad resolved scope whose first page came back short', async () => {

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

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1819,7 +1819,11 @@ export async function executeKeywordSearch(params: KeywordSearchParams): Promise
18191819
const onRowKeywordVisibility = (excludedSources: readonly string[]) =>
18201820
and(
18211821
projectionCandidateAccessCondition(
1822-
{ connectorId: sql`ranked_tin_chunks.connector_id`, acl: sql`ranked_tin_chunks.acl` },
1822+
{
1823+
connectorId: sql`ranked_tin_chunks.connector_id`,
1824+
acl: sql`ranked_tin_chunks.acl`,
1825+
documentId: sql`ranked_tin_chunks.document_id`,
1826+
},
18231827
access,
18241828
accessPlan!
18251829
),

‎apps/sim/scripts/backfill-projection-source-acl.ts‎

Lines changed: 23 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -2,30 +2,40 @@
22

33
/**
44
* Starts the projection source and ACL backfill that script migration
5-
* `0021_embedding_search_connector` leaves to the background: on a deployment with Trigger.dev it
6-
* enqueues the `projection-source-acl-backfill` task, which chains bounded runs until both ranking
7-
* projections are filled; without one it fills them here, paced the same way. Safe to run again at
8-
* any time — a run only fills rows still unset.
5+
* `0021_embedding_search_connector` leaves to the background, by hand. The migration's outbox
6+
* event already starts it after a deploy; this is for starting it again — on a deployment with
7+
* Trigger.dev it enqueues the `projection-source-acl-backfill` task, which chains bounded runs
8+
* until both ranking projections are filled; without one it fills them here, paced the same way.
9+
* Safe to run at any time — a run only fills rows still unset.
910
*
1011
* Usage:
1112
* bun apps/sim/scripts/backfill-projection-source-acl.ts
1213
*/
1314

1415
import { createLogger } from '@sim/logger'
1516
import { toError } from '@sim/utils/errors'
16-
import { enqueueProjectionSourceAclBackfill } from '@/lib/knowledge/search/projection-source-acl-backfill'
17+
import { env } from '@/lib/core/config/env'
18+
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
19+
import {
20+
enqueueProjectionSourceAclBackfill,
21+
runProjectionSourceAclBackfill,
22+
} from '@/lib/knowledge/search/projection-source-acl-backfill'
1723

1824
const logger = createLogger('BackfillProjectionSourceAcl')
1925

26+
async function main(): Promise<void> {
27+
if (isTriggerDevEnabled && env.TRIGGER_SECRET_KEY) {
28+
const handle = await enqueueProjectionSourceAclBackfill({}, true)
29+
logger.info('Backfill enqueued on the Trigger.dev worker', handle ?? {})
30+
return
31+
}
32+
await runProjectionSourceAclBackfill({})
33+
logger.info('Backfill complete')
34+
}
35+
2036
if (import.meta.main) {
21-
enqueueProjectionSourceAclBackfill().then(
22-
(handle) => {
23-
logger.info(
24-
handle ? 'Backfill enqueued on the Trigger.dev worker' : 'Backfill complete',
25-
handle ?? {}
26-
)
27-
process.exit(0)
28-
},
37+
main().then(
38+
() => process.exit(0),
2939
(error) => {
3040
logger.error('Backfill failed', toError(error))
3141
process.exit(1)

‎packages/db/script-migrations/0021_embedding_search_connector.postgres.test.ts‎

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,10 @@ describe.runIf(Boolean(databaseUrl))('projection source and ACL backfill in Post
4040
await sql`CREATE TABLE document (
4141
id text PRIMARY KEY, connector_id text, acl text[] NOT NULL DEFAULT '{ws}'
4242
)`
43+
await sql`CREATE TABLE outbox_event (
44+
id text PRIMARY KEY, event_type text NOT NULL, payload json NOT NULL,
45+
status text NOT NULL DEFAULT 'pending'
46+
)`
4347
for (const projection of ['embedding_search', 'embedding_keyword_tin']) {
4448
await sql`CREATE TABLE ${sql(projection)} (
4549
id text PRIMARY KEY, document_id text NOT NULL, enabled boolean NOT NULL DEFAULT true,
@@ -61,17 +65,27 @@ describe.runIf(Boolean(databaseUrl))('projection source and ACL backfill in Post
6165
await sql`ALTER TABLE embedding_keyword_tin DISABLE TRIGGER embedding_keyword_tin_source_acl_set`
6266
})
6367

64-
it('installs its triggers and indexes again without failing, so a cut-short deploy completes', async () => {
68+
it('installs its triggers, indexes and start event again without failing, so a cut-short deploy completes', async () => {
6569
await expect(embeddingSearchConnectorMigration.up(sql)).resolves.toBeUndefined()
6670
const indexes = await sql<{ indexname: string }[]>`
6771
SELECT indexname FROM pg_indexes WHERE schemaname = ${schemaName} ORDER BY indexname`
6872
expect(indexes.map((row) => row.indexname)).toEqual(
6973
expect.arrayContaining([
7074
'embedding_keyword_tin_acl_gin_idx',
75+
'embedding_keyword_tin_acl_unfilled_idx',
7176
'embedding_search_acl_gin_idx',
77+
'embedding_search_acl_unfilled_idx',
7278
'embedding_search_source_idx',
7379
])
7480
)
81+
/** Two runs of the migration leave the app one event, so the backfill is started once. */
82+
expect(await sql`SELECT id, event_type, status FROM outbox_event`).toEqual([
83+
{
84+
id: 'projection-source-acl-backfill:0021',
85+
event_type: 'knowledge.projection.source_acl.backfill',
86+
status: 'pending',
87+
},
88+
])
7589
})
7690

7791
it('fills only the rows still unset, in pages, and leaves a chunk that changed documents to its trigger', async () => {

0 commit comments

Comments
 (0)