Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion apps/sim/app/api/v1/knowledge/[id]/documents/route.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -245,7 +245,8 @@ describe('v1 knowledge document upload route', () => {
'kb-1',
{},
'req-1',
SYSTEM_BILLING_ATTRIBUTION
SYSTEM_BILLING_ATTRIBUTION,
'interactive'
)
})
})
8 changes: 7 additions & 1 deletion apps/sim/background/knowledge-processing.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,18 +7,24 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
const {
mockAssertBillingAttributionSnapshot,
mockProcessDocumentAsync,
mockQueue,
mockResolveTriggerRegion,
mockTask,
mockTrigger,
} = vi.hoisted(() => ({
mockAssertBillingAttributionSnapshot: vi.fn(),
mockProcessDocumentAsync: vi.fn(),
mockResolveTriggerRegion: vi.fn(),
mockQueue: vi.fn((config) => config),
mockTask: vi.fn((config) => config),
mockTrigger: vi.fn(),
}))

vi.mock('@trigger.dev/sdk', () => ({ task: mockTask, tasks: { trigger: mockTrigger } }))
vi.mock('@trigger.dev/sdk', () => ({
queue: mockQueue,
task: mockTask,
tasks: { trigger: mockTrigger },
}))
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: mockResolveTriggerRegion }))
vi.mock('@/lib/billing/core/billing-attribution', () => ({
assertBillingAttributionSnapshot: mockAssertBillingAttributionSnapshot,
Expand Down
44 changes: 39 additions & 5 deletions apps/sim/background/knowledge-processing.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { createLogger } from '@sim/logger'
import { task } from '@trigger.dev/sdk'
import { queue, task } from '@trigger.dev/sdk'
import { env, envNumber } from '@/lib/core/config/env'
import {
BYOK_EMBEDDING_CREDENTIAL_REJECTION_MESSAGE,
Expand All @@ -13,6 +13,10 @@ import {
isPermanentDocumentProcessingError,
isUsageLimitDocumentProcessingError,
} from '@/lib/knowledge/documents/document-processing-error'
import {
BACKFILL_PROCESSING_QUEUE_NAME,
INTERACTIVE_PROCESSING_QUEUE_NAME,
} from '@/lib/knowledge/documents/processing-lane'
import {
assertDocumentProcessingBillingContext,
assertDocumentProcessingPayload,
Expand Down Expand Up @@ -205,6 +209,34 @@ export async function runDocumentProcessing(
}
}

/**
* Both lanes are keyed by tenant at dispatch, so `concurrencyLimit` is the
* ceiling one tenant may hold in that lane, not a ceiling for the fleet. The
* shared bound is the Trigger.dev environment concurrency limit, which is where
* a global ceiling belongs; observed peak there is ~97 across every task.
*
* Both default to the limit the single shared queue carried, which is what
* keeps this split from ever draining slower than the queue it replaces: the
* busiest case it has to beat is one tenant alone, and one tenant alone still
* gets the same slots it used to get for backfill plus a separate allowance for
* work someone is waiting on. Any second tenant is pure gain, because under the
* shared queue it got whatever the first one left.
*
* Splitting the two into separate variables is for operating them, not for
* sizing them: backfill is the one to lower when the environment ceiling is the
* binding constraint, and lowering it must not slow down a person's upload.
*/
export const interactiveProcessingQueue = queue({
name: INTERACTIVE_PROCESSING_QUEUE_NAME,
concurrencyLimit: envNumber(env.KB_CONFIG_CONCURRENCY_LIMIT, 20),
})

/** Referenced by no dispatch site: named per trigger, declared here so the deploy registers it. */
export const backfillProcessingQueue = queue({
name: BACKFILL_PROCESSING_QUEUE_NAME,
concurrencyLimit: envNumber(env.KB_CONFIG_BACKFILL_CONCURRENCY_LIMIT, 20),
})

export const processDocument = task({
id: 'knowledge-process-document',
maxDuration: envNumber(env.KB_CONFIG_MAX_DURATION, 600),
Expand All @@ -227,10 +259,12 @@ export const processDocument = task({
*/
outOfMemory: { machine: 'large-2x' },
},
queue: {
concurrencyLimit: envNumber(env.KB_CONFIG_CONCURRENCY_LIMIT, 20),
name: 'document-processing-queue',
},
/**
* The lane every dispatch names explicitly. Declared here as well so a
* trigger that somehow omits the option still lands on a registered queue
* rather than waiting in `PENDING_VERSION` for one that does not exist.
*/
queue: interactiveProcessingQueue,
run: (payload: DocumentProcessingPayload, { ctx }) =>
runDocumentProcessing(payload, ctx.attempt.number),
})
3 changes: 2 additions & 1 deletion apps/sim/lib/core/config/env.ts
Original file line number Diff line number Diff line change
Expand Up @@ -459,7 +459,8 @@ export const env = createEnv({
KB_CONFIG_RETRY_FACTOR: z.number().optional().default(2), // Retry backoff factor
KB_CONFIG_MIN_TIMEOUT: z.number().optional().default(1000), // Min timeout in ms
KB_CONFIG_MAX_TIMEOUT: z.number().optional().default(10000), // Max timeout in ms
KB_CONFIG_CONCURRENCY_LIMIT: z.number().optional().default(20), // Concurrent document-processing task runs (Trigger.dev queue depth)
KB_CONFIG_CONCURRENCY_LIMIT: z.number().optional().default(20), // Per-tenant concurrent document-processing runs in the interactive lane
KB_CONFIG_BACKFILL_CONCURRENCY_LIMIT: z.number().optional().default(20), // Per-tenant concurrent document-processing runs in the connector-backfill lane
KB_CONFIG_EMBEDDING_CONCURRENCY: z.number().optional().default(8), // Concurrent embedding API requests within one embed call
/** Deployment operating budgets shared by every caller using the same provider credential. */
KB_CONFIG_EMBEDDING_REQUESTS_PER_MINUTE: z.number().positive().optional().default(600),
Expand Down
19 changes: 14 additions & 5 deletions apps/sim/lib/embeddings/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -70,11 +70,20 @@ const logger = createLogger('EmbeddingClient')
* Embedding requests issued concurrently within a single embed call.
*
* A provider's rate limit is per API key, so this multiplies with however many
* documents are being processed at once: the document-processing queue admits
* {@link env.KB_CONFIG_CONCURRENCY_LIMIT} task runs, each reaching here. It was
* previously read from that same variable, so one knob set both factors and the
* product reached four figures of in-flight requests against one key — enough to
* hold a provider at its limit indefinitely, which no retry policy can absorb.
* documents are being processed at once. That document count is no longer a
* single number: the processing queues admit
* {@link env.KB_CONFIG_CONCURRENCY_LIMIT} interactive and
* {@link env.KB_CONFIG_BACKFILL_CONCURRENCY_LIMIT} backfill runs *per tenant*,
* bounded in aggregate by the Trigger.dev environment concurrency limit, and
* each run reaches here. The product is held down instead by the durable
* per-credential token bucket in `waitForProviderAdmission`, which every one of
* those runs shares. This factor was previously read from the same variable as
* the queue depth, so one knob set both and the product reached four figures of
* in-flight requests against one key — enough to hold a provider at its limit
* indefinitely, which no retry policy can absorb.
*
* The `bulk` parameter below is a different axis: it marks document indexing as
* opposed to query-time embedding, and is true for an interactive upload too.
*/
const DEFAULT_CONCURRENT_BATCHES = 8
const MAX_ALLOWED_CONCURRENT_BATCHES = 16
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -174,7 +174,14 @@ describe('embedding progress survives a processing slice', () => {
})
})
expect(
await processDocumentsWithQueue([file], ids.knowledgeBaseId, {}, requestId, billing)
await processDocumentsWithQueue(
[file],
ids.knowledgeBaseId,
{},
requestId,
billing,
'interactive'
)
).toMatchObject({ accepted: 1, failed: 0 })
const [deferred] = await db.select().from(document).where(eq(document.id, file.documentId))
expect(deferred).toMatchObject({
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,14 @@ describe('OCR input failures stop without partial indexing or futile retries', (
workspaceId: ids.workspaceId,
})
expect(
await processDocumentsWithQueue([file], ids.knowledgeBaseId, {}, generateId(), billing)
await processDocumentsWithQueue(
[file],
ids.knowledgeBaseId,
{},
generateId(),
billing,
'interactive'
)
).toMatchObject({ accepted: 1, failed: 0 })
const [failed] = await db.select().from(document).where(eq(document.id, file.documentId))
expect(failed).toMatchObject({
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -292,7 +292,8 @@ describe('provider throttling resumes the shared indexing pipeline', () => {
ids.knowledgeBaseId,
{},
requestId,
billing
billing,
'interactive'
)
if (holdParentHandoff) {
await Promise.race([
Expand Down
2 changes: 2 additions & 0 deletions apps/sim/lib/knowledge/connectors/sync-primitives.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -383,6 +383,7 @@ describe('processDocOps dispatch buffering', () => {
{},
expect.any(String),
input.billingAttribution,
'backfill',
{ connectorId: 'connector', stillHeld: input.lease.stillHeld }
)
expect(input.state.result).toMatchObject({ docsAdded: 1, docsUpdated: 1 })
Expand All @@ -405,6 +406,7 @@ describe('processDocOps dispatch buffering', () => {
{},
expect.any(String),
input.billingAttribution,
'backfill',
{ connectorId: 'connector', stillHeld: input.lease.stillHeld }
)
expect(input.state.result.processingDispatch).toEqual({ requested: 3, accepted: 0, failed: 0 })
Expand Down
17 changes: 5 additions & 12 deletions apps/sim/lib/knowledge/connectors/sync-primitives.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@ import { toError } from '@sim/utils/errors'
import { generateId } from '@sim/utils/id'
import { and, asc, desc, eq, inArray, isNotNull, isNull, lt, ne, sql } from 'drizzle-orm'
import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution'
import { env, envNumber } from '@/lib/core/config/env'
import { ProviderCapacityDeferredError } from '@/lib/core/rate-limiter/provider-capacity-error'
import { withDatabaseReadRetry } from '@/lib/db/read-retry'
import type { ConnectorAccessMode } from '@/lib/knowledge/connectors/access-modes'
Expand Down Expand Up @@ -81,20 +80,12 @@ export const CONNECTOR_SYNC_MAX_SOURCE_PAYLOAD_BYTES = 256 * 1024 * 1024
const PROCESSING_DISPATCH_BATCH_SIZE = 25

/**
* Bounds each sync's contribution to the shared processing queue. Oldest eligible
* documents drain first; the remaining backlog stays eligible for subsequent syncs.
* Bounds each sync's contribution to this tenant's bulk processing queue. Oldest
* eligible documents drain first; the remaining backlog stays eligible for
* subsequent syncs.
*/
export const STUCK_RETRY_MAX_CANDIDATES_PER_SYNC = 200

/**
* Concurrent `knowledge-process-document` runs, shared by every workspace.
*
* Read from the same env var the task itself is configured with rather than
* restated, so the drain estimate below cannot describe a queue depth the
* deployment does not actually run.
*/
const PROCESSING_QUEUE_CONCURRENCY = envNumber(env.KB_CONFIG_CONCURRENCY_LIMIT, 20)

export class ConnectorSyncCapacityError extends Error {}

export function sourcePageFitsSyncWorkingSet(rowsAlreadyLoaded: number, pageRows: number): boolean {
Expand Down Expand Up @@ -1012,6 +1003,7 @@ export async function processDocOps(input: ProcessDocOpsInput): Promise<boolean>
{},
generateId(),
billingAttribution,
'backfill',
{ connectorId, stillHeld: input.lease.stillHeld }
)
result.processingDispatch.accepted += dispatch.accepted
Expand Down Expand Up @@ -1510,6 +1502,7 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
{},
generateId(),
billingAttribution,
'backfill',
{ connectorId, stillHeld: input.lease.stillHeld }
)
result.processingDispatch.accepted += dispatch.accepted
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1599,7 +1599,8 @@ describe('in-process quota continuation dispatch', () => {
'knowledge-base-1',
{},
'request-1',
BILLING_ATTRIBUTION
BILLING_ATTRIBUTION,
'interactive'
)
).resolves.toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] })

Expand Down Expand Up @@ -1648,6 +1649,7 @@ describe('in-process quota continuation dispatch', () => {
{},
'request-1',
BILLING_ATTRIBUTION,
'interactive',
undefined,
context
)
Expand All @@ -1671,7 +1673,8 @@ describe('in-process quota continuation dispatch', () => {
'knowledge-base-1',
{},
'request-1',
BILLING_ATTRIBUTION
BILLING_ATTRIBUTION,
'interactive'
)
).resolves.toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] })

Expand All @@ -1695,7 +1698,8 @@ describe('in-process quota continuation dispatch', () => {
'knowledge-base-1',
{},
'request-1',
BILLING_ATTRIBUTION
BILLING_ATTRIBUTION,
'interactive'
)
).resolves.toMatchObject({ accepted: 1, failed: 0 })
expect(mockTrigger).not.toHaveBeenCalled()
Expand Down Expand Up @@ -1876,7 +1880,8 @@ describe('in-process quota continuation dispatch', () => {
'knowledge-base-1',
{},
'request-1',
BILLING_ATTRIBUTION
BILLING_ATTRIBUTION,
'interactive'
)
).resolves.toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] })

Expand All @@ -1901,7 +1906,8 @@ describe('in-process quota continuation dispatch', () => {
'knowledge-base-1',
{},
'request-1',
BILLING_ATTRIBUTION
BILLING_ATTRIBUTION,
'interactive'
)
).resolves.toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] })

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
/**
* @vitest-environment node
*/
import { beforeEach, describe, expect, it, vi } from 'vitest'

const { mockTrigger, mockResolveTriggerRegion } = vi.hoisted(() => ({
mockTrigger: vi.fn(),
mockResolveTriggerRegion: vi.fn().mockResolvedValue('us-east-1'),
}))

vi.mock('@trigger.dev/sdk', () => ({ tasks: { trigger: mockTrigger } }))
vi.mock('@/lib/core/async-jobs/region', () => ({
resolveTriggerRegion: mockResolveTriggerRegion,
}))

import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution'
import { dispatchDocumentProcessingContinuation } from '@/lib/knowledge/documents/processing-continuation-dispatch'
import type {
DocumentProcessingLane,
DocumentProcessingPayload,
} from '@/lib/knowledge/documents/processing-payload'

const BILLING_ATTRIBUTION = {
actorUserId: 'user-1',
workspaceId: null,
organizationId: 'org-1',
billedAccountUserId: 'owner-1',
billingEntity: { type: 'organization', id: 'org-1' },
billingPeriod: { start: '2026-09-01T00:00:00.000Z', end: '2026-10-01T00:00:00.000Z' },
payerSubscription: null,
} satisfies BillingAttributionSnapshot

function payload(lane: DocumentProcessingLane): DocumentProcessingPayload {
return {
knowledgeBaseId: 'knowledge-base-1',
documentId: 'document-1',
processingLane: lane,
docData: { filename: 'a.txt', fileUrl: '/a.txt', fileSize: 1, mimeType: 'text/plain' },
processingOptions: {},
requestId: 'request-1',
billingScope: 'organization',
actorUserId: 'user-1',
workspaceId: null,
organizationId: 'org-1',
billingAttribution: BILLING_ATTRIBUTION,
}
}

describe('dispatchDocumentProcessingContinuation', () => {
beforeEach(() => {
vi.clearAllMocks()
})

/**
* A deferred backfill document resuming as interactive work would be a way
* around the tenant's backfill ceiling: every quota or capacity deferral
* would promote one more run into the lane a person's upload is waiting in.
*/
it.each([
['backfill', 'document-processing-backfill-queue'],
['interactive', 'document-processing-queue'],
] as const)('resumes a %s document in the same lane', async (lane, expectedQueue) => {
await dispatchDocumentProcessingContinuation(
payload(lane),
new Date('2026-09-17T19:00:00.000Z'),
'continuation-key',
true
)

expect(mockTrigger).toHaveBeenCalledTimes(1)
expect(mockTrigger.mock.calls[0][2]).toMatchObject({
queue: expectedQueue,
concurrencyKey: 'organization:org-1',
})
})
})
Loading
Loading