Skip to content

Commit 7481589

Browse files
authored
feat(knowledge): project document ACL and chunk changes asynchronously (#8199)
* feat(knowledge): project document ACL and chunk changes asynchronously Document ACL/source changes and chunk writes mark their document in knowledge_projection_dirty (one upsert with a generation bump) in the writer's transaction. A knowledge projector converges every search projection per document in short pages bounded by chunk rows, several documents at once, and removes a mark only on the generation it read. Writers that declare sim.projection_mode = 'async' skip the synchronous embedding and document fan-out triggers. The processing commit and every connector-lease ACL page declare it while the knowledge-async-projection flag is on; every other writer, and every release before this one, keeps writing projection rows itself. Marks are written in both modes so a projector pass can never settle over a concurrent synchronous write. Migrations 0021-0023 keep their trigger body; 0024 alone installs the marking. Search decides a marked document's rows on the document itself, so a pending projection never admits a revoked grant. A document that moved sources joins the new source's per-source ranking after its pass. The projector runs as a Trigger.dev task, requested after writes and swept every minute while there is work, or in-process without Trigger.dev. A pass opens one worker per mark up to KB_CONFIG_PROJECTION_CONCURRENCY. The separate source/ACL backfill is folded into it behind knowledge-projection-fill. * fix(knowledge): keep projection guards out of historical migrations Restores 0016, 0019 and 0021 to their staging bodies; 0024 alone re-creates the projection triggers with the mode guard, in the same transaction that installs the marks. The fill starts only inside the pass budget and reports what it marked as remaining. Passes dispatch to Trigger.dev by the rule document processing uses, and the inline coalescer clears its running flag in the same step it reads the owed flag. Bumps the Helm chart for the projection cron job. * test(knowledge): count projection requests per connector ACL case
1 parent 624793f commit 7481589

54 files changed

Lines changed: 31926 additions & 921 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.github/workflows/test-build.yml‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -282,6 +282,8 @@ jobs:
282282
lib/knowledge/__integration__/kb-block-search.integration.ts
283283
lib/knowledge/__integration__/gitlab-workspace.integration.ts
284284
lib/knowledge/__integration__/unfilled-projection-source.integration.ts
285+
lib/knowledge/__integration__/knowledge-projection.integration.ts
286+
lib/knowledge/__integration__/async-projection-processing.integration.ts
285287
lib/knowledge/__integration__/purged-detach-reservation.integration.ts
286288
lib/core/outbox/service.integration.ts
287289
lib/knowledge/__integration__/connector-upload.integration.ts

‎apps/docs/content/docs/platform/self-hosting/background-jobs.mdx‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,7 @@ Point cron at an **internal** address where possible (the in-cluster Service, or
4747
| Time pause/resume | `/api/resume/poll` | `*/1 * * * *` | Workflows paused on a timer |
4848
| Outbox processing | `/api/webhooks/outbox/process` | `*/1 * * * *` | Transactional-outbox retries for billing, membership, enterprise issuance, and workflow-deployment side effects |
4949
| Workspace file search dispatch | `/api/cron/workspace-file-search-dispatch` | `*/1 * * * *` | Dispatches indexing work for workspace file search |
50+
| Knowledge projection | `/api/cron/knowledge-projection` | `*/1 * * * *` | Brings knowledge base search up to date with document, permission, and chunk changes |
5051
| Connector sync | `/api/knowledge/connectors/sync` | `*/5 * * * *` | Knowledge base connector syncs |
5152
| Connector member sync | `/api/knowledge/connectors/member-sync` | `*/5 * * * *` | Per-member access sync for permission-aware connectors |
5253
| Connector directory sync | `/api/knowledge/connectors/directory-sync` | `*/5 * * * *` | Refreshes the directory groups administrator-mode connectors mirror, so a membership change takes effect without waiting for a content sync |
Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,86 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { createMockRequest } from '@sim/testing'
5+
import { beforeEach, describe, expect, it, vi } from 'vitest'
6+
7+
const mocks = vi.hoisted(() => ({
8+
enqueueSweep: vi.fn(),
9+
verifyCronAuth: vi.fn(),
10+
}))
11+
12+
vi.mock('@/lib/auth/internal', () => ({ verifyCronAuth: mocks.verifyCronAuth }))
13+
vi.mock('@/lib/knowledge/projection/enqueue', () => ({
14+
enqueueKnowledgeProjectionSweep: mocks.enqueueSweep,
15+
}))
16+
17+
import { GET } from '@/app/api/cron/knowledge-projection/route'
18+
19+
function request() {
20+
return createMockRequest(
21+
'GET',
22+
undefined,
23+
{},
24+
'http://localhost:3000/api/cron/knowledge-projection'
25+
)
26+
}
27+
28+
describe('knowledge projection sweep route', () => {
29+
beforeEach(() => {
30+
vi.clearAllMocks()
31+
mocks.verifyCronAuth.mockReturnValue(null)
32+
})
33+
34+
it('returns as soon as Trigger.dev accepts the pass', async () => {
35+
mocks.enqueueSweep.mockResolvedValue({
36+
triggered: true,
37+
backend: 'trigger-dev',
38+
jobId: 'run-1',
39+
})
40+
41+
const response = await GET(request())
42+
43+
expect(response.status).toBe(202)
44+
await expect(response.json()).resolves.toEqual({
45+
success: true,
46+
triggered: true,
47+
backend: 'trigger-dev',
48+
jobId: 'run-1',
49+
})
50+
})
51+
52+
it('answers 200 without a pass when the projector has nothing to do', async () => {
53+
mocks.enqueueSweep.mockResolvedValue({ triggered: false, backend: null, jobId: null })
54+
55+
const response = await GET(request())
56+
57+
expect(response.status).toBe(200)
58+
await expect(response.json()).resolves.toEqual({
59+
success: true,
60+
triggered: false,
61+
backend: null,
62+
jobId: null,
63+
})
64+
})
65+
66+
it('returns the cron auth refusal without enqueueing', async () => {
67+
mocks.verifyCronAuth.mockReturnValue(new Response(null, { status: 401 }))
68+
69+
const response = await GET(request())
70+
71+
expect(response.status).toBe(401)
72+
expect(mocks.enqueueSweep).not.toHaveBeenCalled()
73+
})
74+
75+
it('fails closed when Trigger.dev does not accept the pass', async () => {
76+
mocks.enqueueSweep.mockRejectedValue(new Error('trigger unavailable'))
77+
78+
const response = await GET(request())
79+
80+
expect(response.status).toBe(500)
81+
await expect(response.json()).resolves.toEqual({
82+
success: false,
83+
error: 'Sweep enqueue failed',
84+
})
85+
})
86+
})
Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
import { createLogger } from '@sim/logger'
2+
import { getErrorMessage } from '@sim/utils/errors'
3+
import { type NextRequest, NextResponse } from 'next/server'
4+
import { verifyCronAuth } from '@/lib/auth/internal'
5+
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
6+
import { enqueueKnowledgeProjectionSweep } from '@/lib/knowledge/projection/enqueue'
7+
8+
const logger = createLogger('KnowledgeProjectionSweepRoute')
9+
10+
export const dynamic = 'force-dynamic'
11+
export const maxDuration = 60
12+
13+
/**
14+
* The knowledge projector's periodic sweep: enqueues one pass per window while there is work, and
15+
* returns once Trigger.dev accepts it. Writers ask for passes as they commit; this converges
16+
* whatever those requests missed.
17+
*/
18+
export const GET = withRouteHandler(async (request: NextRequest) => {
19+
const authError = verifyCronAuth(request, 'Knowledge projection sweep')
20+
if (authError) return authError
21+
22+
try {
23+
const result = await enqueueKnowledgeProjectionSweep()
24+
return NextResponse.json({ success: true, ...result }, { status: result.triggered ? 202 : 200 })
25+
} catch (error) {
26+
logger.error('Knowledge projection sweep enqueue failed', { error: getErrorMessage(error) })
27+
return NextResponse.json({ success: false, error: 'Sweep enqueue failed' }, { status: 500 })
28+
}
29+
})
Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
1+
import { task } from '@trigger.dev/sdk'
2+
import {
3+
type BackgroundRetryPolicy,
4+
backgroundRetryAttemptCeiling,
5+
getBackgroundRetryDecision,
6+
} from '@/lib/core/errors/background-retry'
7+
import {
8+
KNOWLEDGE_PROJECTION_PASS_BUDGET_MS,
9+
KNOWLEDGE_PROJECTION_TASK_ID,
10+
requestKnowledgeProjection,
11+
} from '@/lib/knowledge/projection/enqueue'
12+
import { runKnowledgeProjectionPass } from '@/lib/knowledge/projection/run'
13+
14+
/**
15+
* A pass gives a single document up on a lock or statement timeout without failing, so a failed
16+
* pass lost its connection or its database. Those back off for minutes; the sweep starts a fresh
17+
* pass every minute regardless, so a few attempts are enough.
18+
*/
19+
export const KNOWLEDGE_PROJECTION_RETRY_POLICY: BackgroundRetryPolicy = {
20+
maxAttempts: 2,
21+
database: { maxAttempts: 3, baseDelayMs: 60 * 1000, maxDelayMs: 5 * 60 * 1000 },
22+
}
23+
24+
/**
25+
* Runs one knowledge projector pass. One pass runs at a time and projects several documents at
26+
* once itself; the prompt requests and the sweep collapse into whichever pass is queued. A pass
27+
* that ran out of budget with marks left asks for the next one. Retry-safe: a pass writes only rows
28+
* that differ from their source and removes a mark only on the generation it read.
29+
*/
30+
export const knowledgeProjectionTask = task({
31+
id: KNOWLEDGE_PROJECTION_TASK_ID,
32+
machine: 'small-1x',
33+
maxDuration: 15 * 60,
34+
retry: { maxAttempts: backgroundRetryAttemptCeiling(KNOWLEDGE_PROJECTION_RETRY_POLICY) },
35+
queue: { name: KNOWLEDGE_PROJECTION_TASK_ID, concurrencyLimit: 1 },
36+
catchError: async ({ error, ctx }) =>
37+
getBackgroundRetryDecision(error, ctx.attempt.number, KNOWLEDGE_PROJECTION_RETRY_POLICY),
38+
run: async () => {
39+
const result = await runKnowledgeProjectionPass({
40+
budgetMs: KNOWLEDGE_PROJECTION_PASS_BUDGET_MS,
41+
})
42+
if (result.remaining) await requestKnowledgeProjection()
43+
return result
44+
},
45+
})

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

Lines changed: 0 additions & 44 deletions
This file was deleted.

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -464,6 +464,7 @@ export const env = createEnv({
464464
KB_CONFIG_CONCURRENCY_LIMIT: z.number().optional().default(20), // Per-tenant concurrent document-processing runs in the interactive lane
465465
KB_CONFIG_BACKFILL_CONCURRENCY_LIMIT: z.number().optional().default(20), // Per-tenant concurrent document-processing runs in the connector-backfill lane
466466
KB_CONFIG_EMBEDDING_CONCURRENCY: z.number().optional().default(8), // Concurrent embedding API requests within one embed call
467+
KB_CONFIG_PROJECTION_CONCURRENCY: z.number().optional().default(8), // Most documents one knowledge projector pass projects at once, each on its own connection
467468
/** Deployment operating budgets shared by every caller using the same provider credential. */
468469
KB_CONFIG_EMBEDDING_REQUESTS_PER_MINUTE: z.number().positive().optional().default(600),
469470
KB_CONFIG_EMBEDDING_TOKENS_PER_MINUTE: z.number().positive().optional().default(600000),
@@ -634,6 +635,8 @@ export const env = createEnv({
634635
CREDENTIAL_GROUPS: z.boolean().optional(), // Enable enterprise Credential Groups globally
635636
KNOWLEDGE_MEMBER_ACCESS: z.boolean().optional(), // Enable per-member knowledge connectors and hybrid-by-default retrieval globally
636637
KNOWLEDGE_TIN_KEYWORD: z.boolean().optional(), // Rank large-scope keyword retrieval through the Tin text index where it exists
638+
KNOWLEDGE_ASYNC_PROJECTION: z.boolean().optional(), // Knowledge writers leave search projection rows to the background projector
639+
KNOWLEDGE_PROJECTION_FILL: z.boolean().optional(), // The knowledge projector fills projection rows written before they carried a source and ACL
637640

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

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

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,8 @@ const { mockFetch, mockIsPlatformAdmin, envRef } = vi.hoisted(() => ({
1010
mockIsPlatformAdmin: vi.fn(),
1111
envRef: {
1212
APPCONFIG_APPLICATION: 'sim-staging' as string | undefined,
13+
KNOWLEDGE_PROJECTION_FILL: undefined as boolean | undefined,
14+
KNOWLEDGE_ASYNC_PROJECTION: undefined as boolean | undefined,
1315
APPCONFIG_ENVIRONMENT: 'staging' as string | undefined,
1416
TABLES_V2_API: undefined as boolean | undefined,
1517
TABLE_ROW_TTL: undefined as boolean | undefined,
@@ -148,6 +150,8 @@ describe('isFeatureEnabled', () => {
148150
envRef.CREDENTIAL_GROUPS = undefined
149151
envRef.KNOWLEDGE_MEMBER_ACCESS = undefined
150152
envRef.KNOWLEDGE_TIN_KEYWORD = undefined
153+
envRef.KNOWLEDGE_ASYNC_PROJECTION = undefined
154+
envRef.KNOWLEDGE_PROJECTION_FILL = undefined
151155
envRef.SLACK_SEARCH_SHARED_APP = undefined
152156
})
153157

@@ -197,6 +201,32 @@ describe('isFeatureEnabled', () => {
197201
})
198202
})
199203

204+
describe('knowledge-async-projection flag', () => {
205+
it('is a global switch', async () => {
206+
expect(await isFeatureEnabled('knowledge-async-projection')).toBe(false)
207+
envRef.KNOWLEDGE_ASYNC_PROJECTION = true
208+
expect(await isFeatureEnabled('knowledge-async-projection')).toBe(true)
209+
})
210+
211+
it('follows an AppConfig global rule', async () => {
212+
withAppConfig({ 'knowledge-async-projection': { enabled: true } })
213+
expect(await isFeatureEnabled('knowledge-async-projection')).toBe(true)
214+
})
215+
})
216+
217+
describe('knowledge-projection-fill flag', () => {
218+
it('is a global switch', async () => {
219+
expect(await isFeatureEnabled('knowledge-projection-fill')).toBe(false)
220+
envRef.KNOWLEDGE_PROJECTION_FILL = true
221+
expect(await isFeatureEnabled('knowledge-projection-fill')).toBe(true)
222+
})
223+
224+
it('follows an AppConfig global rule', async () => {
225+
withAppConfig({ 'knowledge-projection-fill': { enabled: true } })
226+
expect(await isFeatureEnabled('knowledge-projection-fill')).toBe(true)
227+
})
228+
})
229+
200230
describe('knowledge-member-access flag', () => {
201231
it('uses a global fallback switch off AppConfig', async () => {
202232
expect(await isFeatureEnabled('knowledge-member-access')).toBe(false)

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,23 @@ const FEATURE_FLAGS = {
105105
'invalid. Off-AppConfig falls back to KNOWLEDGE_TIN_KEYWORD.',
106106
fallback: 'KNOWLEDGE_TIN_KEYWORD',
107107
},
108+
'knowledge-async-projection': {
109+
description:
110+
'Knowledge writers (document processing and connector ACL writes) leave search projection ' +
111+
'rows to the background knowledge projector instead of rewriting them in their own ' +
112+
'transaction. Global on/off only; turn it on only once no release older than the ' +
113+
'projector serves search. Off-AppConfig falls back to KNOWLEDGE_ASYNC_PROJECTION.',
114+
fallback: 'KNOWLEDGE_ASYNC_PROJECTION',
115+
},
116+
'knowledge-projection-fill': {
117+
description:
118+
'The knowledge projector also fills search projection rows written before they carried ' +
119+
"their document's source and ACL, marking at most 100 documents at once so fresh writes " +
120+
'never wait behind much of it. Global on/off only; off pauses the fill, and search keeps ' +
121+
'deciding unfilled rows on their document. Off-AppConfig falls back to ' +
122+
'KNOWLEDGE_PROJECTION_FILL.',
123+
fallback: 'KNOWLEDGE_PROJECTION_FILL',
124+
},
108125
} satisfies Record<string, FeatureFlagDefinition>
109126

110127
/**

0 commit comments

Comments
 (0)