Skip to content

Commit a2c22f2

Browse files
authored
fix(knowledge): run the projection source/ACL backfill in the background instead of the deploy (#8058)
* fix(knowledge): run the projection source/ACL backfill in the background instead of the deploy Script migration 0021 filled the ranking projections' connector_id and acl columns synchronously, 500 rows per page under a 60s statement timeout. On embedding_search every filled row is re-inserted into each HNSW index, so a page's cost is index maintenance rather than the plan: one page ran past the timeout and the migration, and the deploy, failed. The migration now installs the triggers and builds the indexes only, both idempotent, and the backfill runs from a Trigger.dev task in keyset pages of unfilled rows, paced with a pause between pages and chained across bounded runs. Search does not depend on it: an unfilled row passes the on-row candidate predicate and is decided at hydration under the full document predicate, exactly as every candidate was before the columns existed. Claude-Session: https://claude.ai/code/session_01XU6c7pKRpa5CMoMHKDdqxX * 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 * fix(knowledge): keep the outbox event pending until a deployment without a worker fills the projections Without a Trigger.dev worker the outbox handler no longer detaches the backfill and completes the event; it runs one bounded slice per outbox run and yields with continueOutboxHandler until both projections are filled, so a restart loses at most one slice and the event is never marked done ahead of the work. The index builds run on one reserved connection, so the session-scoped lock timeout covers every build and its reset. Claude-Session: https://claude.ai/code/session_01XU6c7pKRpa5CMoMHKDdqxX * fix(knowledge): start the projection backfill the way the table backfill is started The outbox event and its slice-and-defer handler are gone; nothing else in the repo starts background work that way. The fill is started from the app side as the table backfill is: tasks.trigger on the Trigger.dev worker when one is configured, runDetached in-process otherwise, with the manual script as the operator path after a deploy. The registered migration is now 0022_projection_source_acl_backfill, which supersedes 0021 so a database that already recorded the synchronous shape still gains the unfilled indexes. The page statement timeout is exported so a caller bounding a run can leave it as headroom. Claude-Session: https://claude.ai/code/session_01XU6c7pKRpa5CMoMHKDdqxX * fix(knowledge): refuse a projection backfill page size that is not a positive integer The page size is interpolated into the page statement and a page of nothing would report the projection filled, so a payload that asks for either is refused before the first page rather than quietly reshaped. Claude-Session: https://claude.ai/code/session_01XU6c7pKRpa5CMoMHKDdqxX * fix(knowledge): check the projection backfill budget after the pause between pages A pause that crossed the budget still let the next iteration open a page; the deadline is now checked after the pause, so a bounded run stops before its next page rather than after it. Claude-Session: https://claude.ai/code/session_01XU6c7pKRpa5CMoMHKDdqxX
1 parent 9573ecc commit a2c22f2

16 files changed

Lines changed: 719 additions & 64 deletions
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
import { task, tasks } from '@trigger.dev/sdk'
2+
import { resolveTriggerRegion } from '@/lib/core/async-jobs/region'
3+
import {
4+
PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID,
5+
type ProjectionSourceAclBackfillPayload,
6+
runProjectionSourceAclBackfill,
7+
} from '@/lib/knowledge/search/projection-source-acl-backfill'
8+
9+
/** One run's share of the backfill, inside the worker's run ceiling with room to end its page. */
10+
const RUN_BUDGET_MS = 60 * 60 * 1000
11+
12+
/**
13+
* Trigger.dev wrapper around `runProjectionSourceAclBackfill`. A run fills unset rows for up to
14+
* {@link RUN_BUDGET_MS}, then triggers its continuation from the cursor it reached, so the whole
15+
* projection is filled across as many bounded runs as it takes. Retry-safe: every run writes only
16+
* rows still unset, so a retried or restarted run repeats no write. The queue admits one run at a
17+
* time, so two starts never fill the same pages against each other.
18+
*/
19+
export const projectionSourceAclBackfillTask = task({
20+
id: PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID,
21+
machine: 'small-1x',
22+
retry: { maxAttempts: 3 },
23+
queue: {
24+
name: PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID,
25+
concurrencyLimit: 1,
26+
},
27+
run: async (payload: ProjectionSourceAclBackfillPayload) => {
28+
const cursor = await runProjectionSourceAclBackfill(payload, { budgetMs: RUN_BUDGET_MS })
29+
if (!cursor) return
30+
const continuation: ProjectionSourceAclBackfillPayload = { ...payload, cursor }
31+
await tasks.trigger(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, continuation, {
32+
region: await resolveTriggerRegion(),
33+
})
34+
},
35+
})

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

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -622,6 +622,24 @@ describe.runIf(Boolean(databaseUrl))('knowledge ACLs in PostgreSQL', () => {
622622
)
623623
expect(await admits(perRow, 'members-other')).toBe(false)
624624
expect(await onRowAdmits('members-other')).toBe(false)
625+
/**
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.
629+
*/
630+
await connection.unsafe(
631+
"INSERT INTO document(id, connector_id, acl, acl_verified_at) VALUES ('unfilled-other', 'members', ARRAY['s:slack:-:bob'], statement_timestamp())"
632+
)
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')`)
636+
await connection.unsafe(
637+
"DELETE FROM embedding_search WHERE id IN ('admin-current-chunk', 'members-current-chunk', 'deleted-connector-chunk')"
638+
)
639+
expect(await onRowAdmits('unfilled-other')).toBe(false)
640+
expect(await onRowAdmits('admin-current')).toBe(true)
641+
expect(await onRowAdmits('members-current')).toBe(true)
642+
expect(await onRowAdmits('deleted-connector')).toBe(false)
625643
/**
626644
* Candidate ranking defers the live source proof, as the per-row candidate predicate does: a
627645
* caller holds those grants only after authorization, so applying the clause during ranking

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

Lines changed: 44 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,13 +17,56 @@ vi.unmock('@sim/db/schema')
1717
process.env.DATABASE_URL ??= 'postgresql://user:pass@localhost:5432/test'
1818

1919
const { PgDialect } = await import('drizzle-orm/pg-core')
20-
const { knowledgeAccessCondition } = await import('@/lib/knowledge/access/predicate')
20+
const { embeddingSearch } = await import('@sim/db/schema')
21+
const { knowledgeAccessCondition, projectionCandidateAccessCondition } = await import(
22+
'@/lib/knowledge/access/predicate'
23+
)
2124
const { SYSTEM_ACCESS_SCOPE } = await import('@/lib/knowledge/access/types')
2225

2326
function render(condition: ReturnType<typeof knowledgeAccessCondition>) {
2427
return new PgDialect().sqlToQuery(condition)
2528
}
2629

30+
describe('projectionCandidateAccessCondition', () => {
31+
const plan = {
32+
connectors: { workspace: ['ws-src'], admin: [], members: [], liveProofRequired: [] },
33+
observers: { confirmed: [], observed: [] },
34+
memberSources: [],
35+
}
36+
37+
it('decides a filled row on its mirrored columns and an unfilled row on its document', () => {
38+
const { sql, params } = render(
39+
projectionCandidateAccessCondition(
40+
embeddingSearch,
41+
{ kind: 'user', userId: 'user-1', tokens: ['ws', 'u:alice'] },
42+
plan
43+
)
44+
)
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 ('
47+
)
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)
55+
})
56+
57+
it('still denies everything for an empty token set', () => {
58+
expect(
59+
render(
60+
projectionCandidateAccessCondition(
61+
embeddingSearch,
62+
{ kind: 'user', userId: 'user-1', tokens: [] },
63+
plan
64+
)
65+
).sql
66+
).toBe('false')
67+
})
68+
})
69+
2770
describe('knowledgeAccessCondition', () => {
2871
it('overlaps the ACL with the tokens as a literal array of scalar binds', () => {
2972
const { sql, params } = render(

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

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -312,9 +312,19 @@ export function knowledgeCandidateAccessConditionForConnectors(
312312
* observation's freshness, and requirement clauses live on the document. Both are refused there,
313313
* under the full predicate, before content is returned — this predicate only decides what is worth
314314
* ranking.
315+
*
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.
315321
*/
316322
export function projectionCandidateAccessCondition(
317-
projection: { connectorId: AnyPgColumn | SQL; acl: AnyPgColumn | SQL },
323+
projection: {
324+
connectorId: AnyPgColumn | SQL
325+
acl: AnyPgColumn | SQL
326+
documentId: AnyPgColumn | SQL
327+
},
318328
scope: KnowledgeAccessScope | SystemAccessScope,
319329
plan: SearchAccessPlan
320330
): SQL {
@@ -330,8 +340,13 @@ export function projectionCandidateAccessCondition(
330340
...plan.connectors.admin,
331341
...plan.connectors.members,
332342
]
333-
return sql`(${projection.acl} && ${tokens}
334-
AND (${projection.connectorId} IS NULL OR ${inSources(mirrored)}))`
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}
349+
AND (${projection.connectorId} IS NULL OR ${inSources(mirrored)})))`
335350
}
336351

337352
/**
Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { beforeEach, describe, expect, it, vi } from 'vitest'
5+
6+
const { mockBackfill, mockEnd, mockPostgres, mockTasksTrigger } = vi.hoisted(() => ({
7+
mockBackfill: vi.fn(),
8+
mockEnd: vi.fn(async () => undefined),
9+
mockPostgres: vi.fn(),
10+
mockTasksTrigger: vi.fn(async () => ({ id: 'run-1' })),
11+
}))
12+
13+
vi.mock('@sim/db', () => ({ resolveDbUrl: () => 'postgres://localhost:5432/sim' }))
14+
vi.mock('@sim/db/script-migrations/0021_embedding_search_connector', () => ({
15+
PROJECTION_SOURCE_ACL_TABLES: ['embedding_search', 'embedding_keyword_tin'],
16+
backfillProjectionSourceAcl: mockBackfill,
17+
}))
18+
vi.mock('postgres', () => ({ default: mockPostgres }))
19+
vi.mock('@trigger.dev/sdk', () => ({ tasks: { trigger: mockTasksTrigger } }))
20+
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: async () => 'us-east-1' }))
21+
vi.mock('@/lib/core/utils/background', () => ({
22+
runDetached: (_label: string, work: () => Promise<unknown>) => {
23+
void work()
24+
},
25+
}))
26+
27+
import {
28+
enqueueProjectionSourceAclBackfill,
29+
runProjectionSourceAclBackfill,
30+
} from '@/lib/knowledge/search/projection-source-acl-backfill'
31+
32+
const connection = { end: mockEnd }
33+
34+
describe('runProjectionSourceAclBackfill', () => {
35+
beforeEach(() => {
36+
vi.clearAllMocks()
37+
mockPostgres.mockReturnValue(connection)
38+
mockBackfill.mockImplementation(async (_sql, projection, options) => ({
39+
projection,
40+
scanned: 0,
41+
written: 0,
42+
afterId: options?.afterId ?? '',
43+
done: true,
44+
}))
45+
})
46+
47+
it('fills both projections in order on its own connection and closes it', async () => {
48+
await expect(runProjectionSourceAclBackfill({ pageSize: 50, pauseMs: 10 })).resolves.toBeNull()
49+
expect(mockBackfill.mock.calls.map(([, projection]) => projection)).toEqual([
50+
'embedding_search',
51+
'embedding_keyword_tin',
52+
])
53+
for (const [sql, , options] of mockBackfill.mock.calls) {
54+
expect(sql).toBe(connection)
55+
expect(options).toMatchObject({ afterId: undefined, pageSize: 50, pauseMs: 10 })
56+
}
57+
expect(mockEnd).toHaveBeenCalledTimes(1)
58+
})
59+
60+
it('resumes after the cursor in its projection and from the start of the next', async () => {
61+
await runProjectionSourceAclBackfill({
62+
cursor: { projection: 'embedding_keyword_tin', afterId: 'chunk-9' },
63+
})
64+
expect(mockBackfill).toHaveBeenCalledTimes(1)
65+
expect(mockBackfill.mock.calls[0][1]).toBe('embedding_keyword_tin')
66+
expect(mockBackfill.mock.calls[0][2]).toMatchObject({ afterId: 'chunk-9' })
67+
})
68+
69+
it('returns where a budgeted run stopped so the next run can carry on', async () => {
70+
mockBackfill.mockResolvedValueOnce({
71+
projection: 'embedding_search',
72+
scanned: 100,
73+
written: 100,
74+
afterId: 'chunk-100',
75+
done: false,
76+
})
77+
await expect(runProjectionSourceAclBackfill({}, { budgetMs: 1000 })).resolves.toEqual({
78+
projection: 'embedding_search',
79+
afterId: 'chunk-100',
80+
})
81+
expect(mockBackfill).toHaveBeenCalledTimes(1)
82+
expect(mockBackfill.mock.calls[0][2].budgetMs).toBeLessThanOrEqual(1000)
83+
expect(mockEnd).toHaveBeenCalledTimes(1)
84+
})
85+
86+
it('closes the connection when a page fails', async () => {
87+
mockBackfill.mockRejectedValueOnce(new Error('canceling statement due to statement timeout'))
88+
await expect(runProjectionSourceAclBackfill({})).rejects.toThrow('statement timeout')
89+
expect(mockEnd).toHaveBeenCalledTimes(1)
90+
})
91+
})
92+
93+
describe('enqueueProjectionSourceAclBackfill', () => {
94+
beforeEach(() => {
95+
vi.clearAllMocks()
96+
mockPostgres.mockReturnValue(connection)
97+
mockBackfill.mockResolvedValue({
98+
projection: 'embedding_search',
99+
scanned: 0,
100+
written: 0,
101+
afterId: '',
102+
done: true,
103+
})
104+
})
105+
106+
it('hands the backfill to the Trigger.dev worker when one is configured', async () => {
107+
await expect(enqueueProjectionSourceAclBackfill({ pageSize: 25 }, true)).resolves.toEqual({
108+
runId: 'run-1',
109+
})
110+
expect(mockTasksTrigger).toHaveBeenCalledWith(
111+
'projection-source-acl-backfill',
112+
{ pageSize: 25 },
113+
{ region: 'us-east-1' }
114+
)
115+
expect(mockBackfill).not.toHaveBeenCalled()
116+
})
117+
118+
it('fills the projections detached in this process without one', async () => {
119+
await expect(enqueueProjectionSourceAclBackfill({}, false)).resolves.toBeNull()
120+
expect(mockTasksTrigger).not.toHaveBeenCalled()
121+
await vi.waitFor(() => expect(mockBackfill).toHaveBeenCalledTimes(2))
122+
})
123+
})
Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,98 @@
1+
import { resolveDbUrl } from '@sim/db'
2+
import {
3+
backfillProjectionSourceAcl,
4+
PROJECTION_SOURCE_ACL_TABLES,
5+
type ProjectionSourceAclTable,
6+
} from '@sim/db/script-migrations/0021_embedding_search_connector'
7+
import { createLogger } from '@sim/logger'
8+
import postgres from 'postgres'
9+
import { resolveTriggerRegion } from '@/lib/core/async-jobs/region'
10+
import { env } from '@/lib/core/config/env'
11+
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
12+
import { runDetached } from '@/lib/core/utils/background'
13+
14+
const logger = createLogger('ProjectionSourceAclBackfill')
15+
16+
export const PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID = 'projection-source-acl-backfill'
17+
18+
/** Where a run stopped, so the next one carries on from there instead of rescanning. */
19+
export interface ProjectionSourceAclBackfillCursor {
20+
projection: ProjectionSourceAclTable
21+
afterId: string
22+
}
23+
24+
export interface ProjectionSourceAclBackfillPayload {
25+
/** The first projection's first page when absent. */
26+
cursor?: ProjectionSourceAclBackfillCursor
27+
pageSize?: number
28+
pauseMs?: number
29+
}
30+
31+
export interface ProjectionSourceAclBackfillRunOptions {
32+
/** Stop once this much time has passed and return where to resume; unbounded otherwise. */
33+
budgetMs?: number
34+
}
35+
36+
/**
37+
* Fills the ranking projections' source and ACL columns from the cursor onwards, one projection
38+
* after the other, on a connection of its own: the page statement binds the keyset cursor as a
39+
* scalar and needs no array parameter, so the pool's options would serve, but a run this long
40+
* should not hold one of the worker's pooled connections. Returns the cursor to continue from when
41+
* the budget ran out, `null` once both projections are filled.
42+
*/
43+
export async function runProjectionSourceAclBackfill(
44+
payload: ProjectionSourceAclBackfillPayload,
45+
options: ProjectionSourceAclBackfillRunOptions = {}
46+
): Promise<ProjectionSourceAclBackfillCursor | null> {
47+
const url = resolveDbUrl('DATABASE_URL', process.env.SIM_DB_ROLE?.trim() || 'web')
48+
if (!url) throw new Error('DATABASE_URL is required to backfill the projection source and ACL')
49+
const sql = postgres(url, { max: 1, max_lifetime: null, onnotice: () => undefined })
50+
const startedAt = Date.now()
51+
try {
52+
const start = payload.cursor
53+
? PROJECTION_SOURCE_ACL_TABLES.indexOf(payload.cursor.projection)
54+
: 0
55+
if (start < 0) throw new Error(`Unknown projection ${payload.cursor?.projection}`)
56+
for (const projection of PROJECTION_SOURCE_ACL_TABLES.slice(start)) {
57+
const budgetMs =
58+
options.budgetMs === undefined
59+
? undefined
60+
: Math.max(0, options.budgetMs - (Date.now() - startedAt))
61+
const progress = await backfillProjectionSourceAcl(sql, projection, {
62+
afterId: payload.cursor?.projection === projection ? payload.cursor.afterId : undefined,
63+
pageSize: payload.pageSize,
64+
pauseMs: payload.pauseMs,
65+
budgetMs,
66+
})
67+
if (!progress.done) return { projection, afterId: progress.afterId }
68+
}
69+
logger.info('Projection source and ACL backfill complete', {
70+
elapsedMs: Date.now() - startedAt,
71+
})
72+
return null
73+
} finally {
74+
await sql.end()
75+
}
76+
}
77+
78+
/**
79+
* Starts the backfill the way the table backfill is started: on the deployment's Trigger.dev
80+
* worker when one is configured, where bounded runs chain until both projections are filled, and
81+
* detached in this process otherwise. Safe to call again at any time — a run only fills rows still
82+
* unset.
83+
*/
84+
export async function enqueueProjectionSourceAclBackfill(
85+
payload: ProjectionSourceAclBackfillPayload = {},
86+
useTrigger = Boolean(isTriggerDevEnabled && env.TRIGGER_SECRET_KEY)
87+
): Promise<{ runId: string } | null> {
88+
if (useTrigger) {
89+
const { tasks } = await import('@trigger.dev/sdk')
90+
const handle = await tasks.trigger(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, payload, {
91+
region: await resolveTriggerRegion(),
92+
})
93+
logger.info('Projection source and ACL backfill enqueued', { runId: handle.id })
94+
return { runId: handle.id }
95+
}
96+
runDetached(PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID, () => runProjectionSourceAclBackfill(payload))
97+
return null
98+
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1527,6 +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 decided on its document instead. */
1531+
expect(statement).toContain('IS NULL AND EXISTS (')
1532+
expect(statement).toContain('ranked_tin_chunks.document_id')
15301533
})
15311534

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

0 commit comments

Comments
 (0)