Skip to content

Commit 41488fd

Browse files
authored
fix(knowledge): give document processing per-tenant queue lanes (#7937)
Document processing ran through one Trigger.dev queue with a global concurrency limit and no concurrency key, so a single connector backfill could hold every slot and leave every other tenant's uploads queued behind it for hours. Dispatch now names a lane — interactive for work a person is waiting on, backfill for connector-driven ingestion — and keys each run by the entity that owns the knowledge base, so the limit applies per tenant copy of the queue rather than fleet-wide. - Key on the owning entity, not the workspace: an organization-scoped knowledge base carries a null workspace id and would otherwise share one bucket with every other such tenant. Owned scopes reuse resourceScopeKey so tenant identity keeps one spelling. - Stamp the lane on the payload so quota and capacity continuations resume in the lane they were admitted against instead of promoting themselves out of the backfill ceiling. - Parse an absent or unrecognized lane as backfill rather than throwing, so payloads written before the lanes existed and payloads stamped by a newer version mid-rollout cannot burn a run's retry budget. - Keep interactive work on the pre-existing queue name; a queue the running worker has not registered parks its runs in PENDING_VERSION, and the app deploys separately from the worker. - Stop absorbing a rejected backfill chunk in-process: those documents are reported failed and reclaimed by the stuck-document sweep instead of outlasting the connector lease they run under.
1 parent 279898b commit 41488fd

25 files changed

Lines changed: 640 additions & 68 deletions

‎apps/sim/app/api/v1/knowledge/[id]/documents/route.test.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -245,7 +245,8 @@ describe('v1 knowledge document upload route', () => {
245245
'kb-1',
246246
{},
247247
'req-1',
248-
SYSTEM_BILLING_ATTRIBUTION
248+
SYSTEM_BILLING_ATTRIBUTION,
249+
'interactive'
249250
)
250251
})
251252
})

‎apps/sim/background/knowledge-processing.test.ts‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,18 +7,24 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
77
const {
88
mockAssertBillingAttributionSnapshot,
99
mockProcessDocumentAsync,
10+
mockQueue,
1011
mockResolveTriggerRegion,
1112
mockTask,
1213
mockTrigger,
1314
} = vi.hoisted(() => ({
1415
mockAssertBillingAttributionSnapshot: vi.fn(),
1516
mockProcessDocumentAsync: vi.fn(),
1617
mockResolveTriggerRegion: vi.fn(),
18+
mockQueue: vi.fn((config) => config),
1719
mockTask: vi.fn((config) => config),
1820
mockTrigger: vi.fn(),
1921
}))
2022

21-
vi.mock('@trigger.dev/sdk', () => ({ task: mockTask, tasks: { trigger: mockTrigger } }))
23+
vi.mock('@trigger.dev/sdk', () => ({
24+
queue: mockQueue,
25+
task: mockTask,
26+
tasks: { trigger: mockTrigger },
27+
}))
2228
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: mockResolveTriggerRegion }))
2329
vi.mock('@/lib/billing/core/billing-attribution', () => ({
2430
assertBillingAttributionSnapshot: mockAssertBillingAttributionSnapshot,

‎apps/sim/background/knowledge-processing.ts‎

Lines changed: 39 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import { createLogger } from '@sim/logger'
2-
import { task } from '@trigger.dev/sdk'
2+
import { queue, task } from '@trigger.dev/sdk'
33
import { env, envNumber } from '@/lib/core/config/env'
44
import {
55
BYOK_EMBEDDING_CREDENTIAL_REJECTION_MESSAGE,
@@ -13,6 +13,10 @@ import {
1313
isPermanentDocumentProcessingError,
1414
isUsageLimitDocumentProcessingError,
1515
} from '@/lib/knowledge/documents/document-processing-error'
16+
import {
17+
BACKFILL_PROCESSING_QUEUE_NAME,
18+
INTERACTIVE_PROCESSING_QUEUE_NAME,
19+
} from '@/lib/knowledge/documents/processing-lane'
1620
import {
1721
assertDocumentProcessingBillingContext,
1822
assertDocumentProcessingPayload,
@@ -205,6 +209,34 @@ export async function runDocumentProcessing(
205209
}
206210
}
207211

212+
/**
213+
* Both lanes are keyed by tenant at dispatch, so `concurrencyLimit` is the
214+
* ceiling one tenant may hold in that lane, not a ceiling for the fleet. The
215+
* shared bound is the Trigger.dev environment concurrency limit, which is where
216+
* a global ceiling belongs; observed peak there is ~97 across every task.
217+
*
218+
* Both default to the limit the single shared queue carried, which is what
219+
* keeps this split from ever draining slower than the queue it replaces: the
220+
* busiest case it has to beat is one tenant alone, and one tenant alone still
221+
* gets the same slots it used to get for backfill plus a separate allowance for
222+
* work someone is waiting on. Any second tenant is pure gain, because under the
223+
* shared queue it got whatever the first one left.
224+
*
225+
* Splitting the two into separate variables is for operating them, not for
226+
* sizing them: backfill is the one to lower when the environment ceiling is the
227+
* binding constraint, and lowering it must not slow down a person's upload.
228+
*/
229+
export const interactiveProcessingQueue = queue({
230+
name: INTERACTIVE_PROCESSING_QUEUE_NAME,
231+
concurrencyLimit: envNumber(env.KB_CONFIG_CONCURRENCY_LIMIT, 20),
232+
})
233+
234+
/** Referenced by no dispatch site: named per trigger, declared here so the deploy registers it. */
235+
export const backfillProcessingQueue = queue({
236+
name: BACKFILL_PROCESSING_QUEUE_NAME,
237+
concurrencyLimit: envNumber(env.KB_CONFIG_BACKFILL_CONCURRENCY_LIMIT, 20),
238+
})
239+
208240
export const processDocument = task({
209241
id: 'knowledge-process-document',
210242
maxDuration: envNumber(env.KB_CONFIG_MAX_DURATION, 600),
@@ -227,10 +259,12 @@ export const processDocument = task({
227259
*/
228260
outOfMemory: { machine: 'large-2x' },
229261
},
230-
queue: {
231-
concurrencyLimit: envNumber(env.KB_CONFIG_CONCURRENCY_LIMIT, 20),
232-
name: 'document-processing-queue',
233-
},
262+
/**
263+
* The lane every dispatch names explicitly. Declared here as well so a
264+
* trigger that somehow omits the option still lands on a registered queue
265+
* rather than waiting in `PENDING_VERSION` for one that does not exist.
266+
*/
267+
queue: interactiveProcessingQueue,
234268
run: (payload: DocumentProcessingPayload, { ctx }) =>
235269
runDocumentProcessing(payload, ctx.attempt.number),
236270
})

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -459,7 +459,8 @@ export const env = createEnv({
459459
KB_CONFIG_RETRY_FACTOR: z.number().optional().default(2), // Retry backoff factor
460460
KB_CONFIG_MIN_TIMEOUT: z.number().optional().default(1000), // Min timeout in ms
461461
KB_CONFIG_MAX_TIMEOUT: z.number().optional().default(10000), // Max timeout in ms
462-
KB_CONFIG_CONCURRENCY_LIMIT: z.number().optional().default(20), // Concurrent document-processing task runs (Trigger.dev queue depth)
462+
KB_CONFIG_CONCURRENCY_LIMIT: z.number().optional().default(20), // Per-tenant concurrent document-processing runs in the interactive lane
463+
KB_CONFIG_BACKFILL_CONCURRENCY_LIMIT: z.number().optional().default(20), // Per-tenant concurrent document-processing runs in the connector-backfill lane
463464
KB_CONFIG_EMBEDDING_CONCURRENCY: z.number().optional().default(8), // Concurrent embedding API requests within one embed call
464465
/** Deployment operating budgets shared by every caller using the same provider credential. */
465466
KB_CONFIG_EMBEDDING_REQUESTS_PER_MINUTE: z.number().positive().optional().default(600),

‎apps/sim/lib/embeddings/client.ts‎

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -70,11 +70,20 @@ const logger = createLogger('EmbeddingClient')
7070
* Embedding requests issued concurrently within a single embed call.
7171
*
7272
* A provider's rate limit is per API key, so this multiplies with however many
73-
* documents are being processed at once: the document-processing queue admits
74-
* {@link env.KB_CONFIG_CONCURRENCY_LIMIT} task runs, each reaching here. It was
75-
* previously read from that same variable, so one knob set both factors and the
76-
* product reached four figures of in-flight requests against one key — enough to
77-
* hold a provider at its limit indefinitely, which no retry policy can absorb.
73+
* documents are being processed at once. That document count is no longer a
74+
* single number: the processing queues admit
75+
* {@link env.KB_CONFIG_CONCURRENCY_LIMIT} interactive and
76+
* {@link env.KB_CONFIG_BACKFILL_CONCURRENCY_LIMIT} backfill runs *per tenant*,
77+
* bounded in aggregate by the Trigger.dev environment concurrency limit, and
78+
* each run reaches here. The product is held down instead by the durable
79+
* per-credential token bucket in `waitForProviderAdmission`, which every one of
80+
* those runs shares. This factor was previously read from the same variable as
81+
* the queue depth, so one knob set both and the product reached four figures of
82+
* in-flight requests against one key — enough to hold a provider at its limit
83+
* indefinitely, which no retry policy can absorb.
84+
*
85+
* The `bulk` parameter below is a different axis: it marks document indexing as
86+
* opposed to query-time embedding, and is true for an interactive upload too.
7887
*/
7988
const DEFAULT_CONCURRENT_BATCHES = 8
8089
const MAX_ALLOWED_CONCURRENT_BATCHES = 16

‎apps/sim/lib/knowledge/__integration__/embedding-processing-recovery.integration.ts‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -174,7 +174,14 @@ describe('embedding progress survives a processing slice', () => {
174174
})
175175
})
176176
expect(
177-
await processDocumentsWithQueue([file], ids.knowledgeBaseId, {}, requestId, billing)
177+
await processDocumentsWithQueue(
178+
[file],
179+
ids.knowledgeBaseId,
180+
{},
181+
requestId,
182+
billing,
183+
'interactive'
184+
)
178185
).toMatchObject({ accepted: 1, failed: 0 })
179186
const [deferred] = await db.select().from(document).where(eq(document.id, file.documentId))
180187
expect(deferred).toMatchObject({

‎apps/sim/lib/knowledge/__integration__/ocr-input-failures.integration.ts‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -168,7 +168,14 @@ describe('OCR input failures stop without partial indexing or futile retries', (
168168
workspaceId: ids.workspaceId,
169169
})
170170
expect(
171-
await processDocumentsWithQueue([file], ids.knowledgeBaseId, {}, generateId(), billing)
171+
await processDocumentsWithQueue(
172+
[file],
173+
ids.knowledgeBaseId,
174+
{},
175+
generateId(),
176+
billing,
177+
'interactive'
178+
)
172179
).toMatchObject({ accepted: 1, failed: 0 })
173180
const [failed] = await db.select().from(document).where(eq(document.id, file.documentId))
174181
expect(failed).toMatchObject({

‎apps/sim/lib/knowledge/__integration__/provider-processing-recovery.integration.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -292,7 +292,8 @@ describe('provider throttling resumes the shared indexing pipeline', () => {
292292
ids.knowledgeBaseId,
293293
{},
294294
requestId,
295-
billing
295+
billing,
296+
'interactive'
296297
)
297298
if (holdParentHandoff) {
298299
await Promise.race([

‎apps/sim/lib/knowledge/connectors/sync-primitives.test.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -383,6 +383,7 @@ describe('processDocOps dispatch buffering', () => {
383383
{},
384384
expect.any(String),
385385
input.billingAttribution,
386+
'backfill',
386387
{ connectorId: 'connector', stillHeld: input.lease.stillHeld }
387388
)
388389
expect(input.state.result).toMatchObject({ docsAdded: 1, docsUpdated: 1 })
@@ -405,6 +406,7 @@ describe('processDocOps dispatch buffering', () => {
405406
{},
406407
expect.any(String),
407408
input.billingAttribution,
409+
'backfill',
408410
{ connectorId: 'connector', stillHeld: input.lease.stillHeld }
409411
)
410412
expect(input.state.result.processingDispatch).toEqual({ requested: 3, accepted: 0, failed: 0 })

‎apps/sim/lib/knowledge/connectors/sync-primitives.ts‎

Lines changed: 5 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,6 @@ import { toError } from '@sim/utils/errors'
1111
import { generateId } from '@sim/utils/id'
1212
import { and, asc, desc, eq, inArray, isNotNull, isNull, lt, ne, sql } from 'drizzle-orm'
1313
import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution'
14-
import { env, envNumber } from '@/lib/core/config/env'
1514
import { ProviderCapacityDeferredError } from '@/lib/core/rate-limiter/provider-capacity-error'
1615
import { withDatabaseReadRetry } from '@/lib/db/read-retry'
1716
import type { ConnectorAccessMode } from '@/lib/knowledge/connectors/access-modes'
@@ -81,20 +80,12 @@ export const CONNECTOR_SYNC_MAX_SOURCE_PAYLOAD_BYTES = 256 * 1024 * 1024
8180
const PROCESSING_DISPATCH_BATCH_SIZE = 25
8281

8382
/**
84-
* Bounds each sync's contribution to the shared processing queue. Oldest eligible
85-
* documents drain first; the remaining backlog stays eligible for subsequent syncs.
83+
* Bounds each sync's contribution to this tenant's bulk processing queue. Oldest
84+
* eligible documents drain first; the remaining backlog stays eligible for
85+
* subsequent syncs.
8686
*/
8787
export const STUCK_RETRY_MAX_CANDIDATES_PER_SYNC = 200
8888

89-
/**
90-
* Concurrent `knowledge-process-document` runs, shared by every workspace.
91-
*
92-
* Read from the same env var the task itself is configured with rather than
93-
* restated, so the drain estimate below cannot describe a queue depth the
94-
* deployment does not actually run.
95-
*/
96-
const PROCESSING_QUEUE_CONCURRENCY = envNumber(env.KB_CONFIG_CONCURRENCY_LIMIT, 20)
97-
9889
export class ConnectorSyncCapacityError extends Error {}
9990

10091
export function sourcePageFitsSyncWorkingSet(rowsAlreadyLoaded: number, pageRows: number): boolean {
@@ -1012,6 +1003,7 @@ export async function processDocOps(input: ProcessDocOpsInput): Promise<boolean>
10121003
{},
10131004
generateId(),
10141005
billingAttribution,
1006+
'backfill',
10151007
{ connectorId, stillHeld: input.lease.stillHeld }
10161008
)
10171009
result.processingDispatch.accepted += dispatch.accepted
@@ -1510,6 +1502,7 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
15101502
{},
15111503
generateId(),
15121504
billingAttribution,
1505+
'backfill',
15131506
{ connectorId, stillHeld: input.lease.stillHeld }
15141507
)
15151508
result.processingDispatch.accepted += dispatch.accepted

0 commit comments

Comments
 (0)