Skip to content

Commit 70f3dfb

Browse files
committed
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.
1 parent 624793f commit 70f3dfb

57 files changed

Lines changed: 31771 additions & 1374 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)