diff --git a/apps/sim/app/api/v1/knowledge/[id]/documents/route.test.ts b/apps/sim/app/api/v1/knowledge/[id]/documents/route.test.ts index 56d10d72095..39d4ab1c816 100644 --- a/apps/sim/app/api/v1/knowledge/[id]/documents/route.test.ts +++ b/apps/sim/app/api/v1/knowledge/[id]/documents/route.test.ts @@ -245,7 +245,8 @@ describe('v1 knowledge document upload route', () => { 'kb-1', {}, 'req-1', - SYSTEM_BILLING_ATTRIBUTION + SYSTEM_BILLING_ATTRIBUTION, + 'interactive' ) }) }) diff --git a/apps/sim/background/knowledge-processing.test.ts b/apps/sim/background/knowledge-processing.test.ts index 64f5f7c556f..2623409ec68 100644 --- a/apps/sim/background/knowledge-processing.test.ts +++ b/apps/sim/background/knowledge-processing.test.ts @@ -7,6 +7,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const { mockAssertBillingAttributionSnapshot, mockProcessDocumentAsync, + mockQueue, mockResolveTriggerRegion, mockTask, mockTrigger, @@ -14,11 +15,16 @@ const { 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, diff --git a/apps/sim/background/knowledge-processing.ts b/apps/sim/background/knowledge-processing.ts index 7981c351208..8ca8526ae26 100644 --- a/apps/sim/background/knowledge-processing.ts +++ b/apps/sim/background/knowledge-processing.ts @@ -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, @@ -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, @@ -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), @@ -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), }) diff --git a/apps/sim/lib/core/config/env.ts b/apps/sim/lib/core/config/env.ts index 2fcfd17badb..780a569bcf6 100644 --- a/apps/sim/lib/core/config/env.ts +++ b/apps/sim/lib/core/config/env.ts @@ -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), diff --git a/apps/sim/lib/embeddings/client.ts b/apps/sim/lib/embeddings/client.ts index 846ba759f0d..aea273979a7 100644 --- a/apps/sim/lib/embeddings/client.ts +++ b/apps/sim/lib/embeddings/client.ts @@ -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 diff --git a/apps/sim/lib/knowledge/__integration__/embedding-processing-recovery.integration.ts b/apps/sim/lib/knowledge/__integration__/embedding-processing-recovery.integration.ts index 141d16eeaa6..e23caca3fde 100644 --- a/apps/sim/lib/knowledge/__integration__/embedding-processing-recovery.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/embedding-processing-recovery.integration.ts @@ -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({ diff --git a/apps/sim/lib/knowledge/__integration__/ocr-input-failures.integration.ts b/apps/sim/lib/knowledge/__integration__/ocr-input-failures.integration.ts index 477fb7d5253..3dd20af7932 100644 --- a/apps/sim/lib/knowledge/__integration__/ocr-input-failures.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/ocr-input-failures.integration.ts @@ -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({ diff --git a/apps/sim/lib/knowledge/__integration__/provider-processing-recovery.integration.ts b/apps/sim/lib/knowledge/__integration__/provider-processing-recovery.integration.ts index 98831b8765b..c653e69eac0 100644 --- a/apps/sim/lib/knowledge/__integration__/provider-processing-recovery.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/provider-processing-recovery.integration.ts @@ -292,7 +292,8 @@ describe('provider throttling resumes the shared indexing pipeline', () => { ids.knowledgeBaseId, {}, requestId, - billing + billing, + 'interactive' ) if (holdParentHandoff) { await Promise.race([ diff --git a/apps/sim/lib/knowledge/connectors/sync-primitives.test.ts b/apps/sim/lib/knowledge/connectors/sync-primitives.test.ts index 6ae03568626..4c8d403b02c 100644 --- a/apps/sim/lib/knowledge/connectors/sync-primitives.test.ts +++ b/apps/sim/lib/knowledge/connectors/sync-primitives.test.ts @@ -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 }) @@ -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 }) diff --git a/apps/sim/lib/knowledge/connectors/sync-primitives.ts b/apps/sim/lib/knowledge/connectors/sync-primitives.ts index 719144ae011..5f20e765fca 100644 --- a/apps/sim/lib/knowledge/connectors/sync-primitives.ts +++ b/apps/sim/lib/knowledge/connectors/sync-primitives.ts @@ -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' @@ -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 { @@ -1012,6 +1003,7 @@ export async function processDocOps(input: ProcessDocOpsInput): Promise {}, generateId(), billingAttribution, + 'backfill', { connectorId, stillHeld: input.lease.stillHeld } ) result.processingDispatch.accepted += dispatch.accepted @@ -1510,6 +1502,7 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom {}, generateId(), billingAttribution, + 'backfill', { connectorId, stillHeld: input.lease.stillHeld } ) result.processingDispatch.accepted += dispatch.accepted diff --git a/apps/sim/lib/knowledge/documents/document-processing-source.test.ts b/apps/sim/lib/knowledge/documents/document-processing-source.test.ts index dd0941ad827..1be1439f8f8 100644 --- a/apps/sim/lib/knowledge/documents/document-processing-source.test.ts +++ b/apps/sim/lib/knowledge/documents/document-processing-source.test.ts @@ -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: [] }) @@ -1648,6 +1649,7 @@ describe('in-process quota continuation dispatch', () => { {}, 'request-1', BILLING_ATTRIBUTION, + 'interactive', undefined, context ) @@ -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: [] }) @@ -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() @@ -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: [] }) @@ -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: [] }) diff --git a/apps/sim/lib/knowledge/documents/processing-continuation-dispatch.test.ts b/apps/sim/lib/knowledge/documents/processing-continuation-dispatch.test.ts new file mode 100644 index 00000000000..83abbbad122 --- /dev/null +++ b/apps/sim/lib/knowledge/documents/processing-continuation-dispatch.test.ts @@ -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', + }) + }) +}) diff --git a/apps/sim/lib/knowledge/documents/processing-continuation-dispatch.ts b/apps/sim/lib/knowledge/documents/processing-continuation-dispatch.ts index 99a228cf5d5..74335ad8d5c 100644 --- a/apps/sim/lib/knowledge/documents/processing-continuation-dispatch.ts +++ b/apps/sim/lib/knowledge/documents/processing-continuation-dispatch.ts @@ -5,6 +5,7 @@ import { resolveTriggerRegion } from '@/lib/core/async-jobs/region' import { env } from '@/lib/core/config/env' import { isTriggerDevEnabled } from '@/lib/core/config/env-flags' import { isInsideTriggerRun } from '@/lib/core/config/trigger-runtime' +import { documentProcessingQueueOptions } from '@/lib/knowledge/documents/processing-lane' import type { DocumentProcessingPayload } from '@/lib/knowledge/documents/processing-payload' export interface DocumentProcessingContinuation { @@ -30,6 +31,13 @@ export async function dispatchDocumentProcessingContinuation( delay: deferredUntil, idempotencyKey, tags: [`knowledgeBaseId:${payload.knowledgeBaseId}`, `documentId:${payload.documentId}`], + /** + * A continuation stays in the lane its original dispatch was admitted + * against. Re-deriving the lane here would let a backfill document that + * deferred on quota resume as interactive work and escape the tenant's + * bulk ceiling — the retry path would become the way around the limit. + */ + ...documentProcessingQueueOptions(payload), region, }) return diff --git a/apps/sim/lib/knowledge/documents/processing-dispatch.ts b/apps/sim/lib/knowledge/documents/processing-dispatch.ts index 4f3fb6f235e..6b67d8ce81d 100644 --- a/apps/sim/lib/knowledge/documents/processing-dispatch.ts +++ b/apps/sim/lib/knowledge/documents/processing-dispatch.ts @@ -77,7 +77,8 @@ export async function dispatchDocumentProcessing({ knowledgeBaseId, processingOptions, requestId, - billingAttribution + billingAttribution, + 'interactive' ) failedDocumentIds = dispatch.failedDocumentIds failureMessage = 'Document processing dispatch was not accepted' diff --git a/apps/sim/lib/knowledge/documents/processing-lane.test.ts b/apps/sim/lib/knowledge/documents/processing-lane.test.ts new file mode 100644 index 00000000000..24178163dd4 --- /dev/null +++ b/apps/sim/lib/knowledge/documents/processing-lane.test.ts @@ -0,0 +1,132 @@ +/** + * @vitest-environment node + */ +import { describe, expect, it } from 'vitest' +import { + BACKFILL_PROCESSING_QUEUE_NAME, + documentProcessingQueueOptions, + INTERACTIVE_PROCESSING_QUEUE_NAME, +} from '@/lib/knowledge/documents/processing-lane' +import type { DocumentProcessingPayload } from '@/lib/knowledge/documents/processing-payload' +import { resolveDocumentProcessingLane } from '@/lib/knowledge/documents/processing-payload' + +const BILLING_ATTRIBUTION = { + actorUserId: 'user-1', + workspaceId: 'workspace-1', + organizationId: null, + billedAccountUserId: 'owner-1', + billingEntity: { type: 'user' as const, id: 'owner-1' }, + billingPeriod: { start: '2026-09-01T00:00:00.000Z', end: '2026-10-01T00:00:00.000Z' }, + payerSubscription: null, +} + +const BASE = { + knowledgeBaseId: 'knowledge-base-1', + documentId: 'document-1', + processingLane: 'interactive', + docData: { filename: 'a.txt', fileUrl: '/a.txt', fileSize: 1, mimeType: 'text/plain' }, + processingOptions: {}, + requestId: 'request-1', +} satisfies Partial + +function workspacePayload( + overrides: Partial = {} +): DocumentProcessingPayload { + return { + ...BASE, + billingScope: 'workspace', + actorUserId: 'user-1', + workspaceId: 'workspace-1', + billingAttribution: BILLING_ATTRIBUTION, + ...overrides, + } +} + +function organizationPayload(organizationId: string): DocumentProcessingPayload { + return { + ...BASE, + processingLane: 'backfill', + billingScope: 'organization', + actorUserId: 'user-1', + workspaceId: null, + organizationId, + billingAttribution: { ...BILLING_ATTRIBUTION, workspaceId: null, organizationId }, + } +} + +describe('resolveDocumentProcessingLane', () => { + it('keeps an explicit interactive stamp', () => { + expect(resolveDocumentProcessingLane('interactive')).toBe('interactive') + }) + + /** + * Payloads written before the lanes existed carry no lane, and a rolling + * deploy can hand this version one stamped by a newer. Neither may throw: + * failing the run would burn its retry budget for the length of a rollout. + */ + it.each([undefined, null, '', 'backfill', 'express', 42, {}, 'Interactive', ' interactive'])( + 'reads %p as backfill rather than throwing or widening', + (value) => { + expect(resolveDocumentProcessingLane(value)).toBe('backfill') + } + ) +}) + +describe('documentProcessingQueueOptions', () => { + it('routes interactive work to the interactive queue keyed by workspace', () => { + expect(documentProcessingQueueOptions(workspacePayload())).toEqual({ + queue: INTERACTIVE_PROCESSING_QUEUE_NAME, + concurrencyKey: 'workspace:workspace-1', + }) + }) + + it('routes backfill to its own queue under the same key', () => { + expect( + documentProcessingQueueOptions(workspacePayload({ processingLane: 'backfill' })) + ).toEqual({ + queue: BACKFILL_PROCESSING_QUEUE_NAME, + concurrencyKey: 'workspace:workspace-1', + }) + }) + + /** + * The incident this split exists for: organization-scoped knowledge bases + * carry `workspaceId: null`, so a workspace-only key would collapse every one + * of them onto a single shared queue and starve the rest of the fleet again. + */ + it('keys an organization-owned knowledge base on its organization', () => { + expect(documentProcessingQueueOptions(organizationPayload('org-1')).concurrencyKey).toBe( + 'organization:org-1' + ) + }) + + it('keys a knowledge base with no workspace or organization on its actor', () => { + const payload: DocumentProcessingPayload = { + ...BASE, + billingScope: 'non-workspace', + actorUserId: 'user-9', + workspaceId: null, + } + expect(documentProcessingQueueOptions(payload).concurrencyKey).toBe('user:user-9') + }) + + it('puts two tenants in the same lane on separate keys', () => { + const first = documentProcessingQueueOptions(organizationPayload('org-1')) + const second = documentProcessingQueueOptions(organizationPayload('org-2')) + expect(first.queue).toBe(second.queue) + expect(first.concurrencyKey).not.toBe(second.concurrencyKey) + }) + + it('keeps an organization and a workspace sharing an id on separate keys', () => { + expect( + documentProcessingQueueOptions(organizationPayload('shared-id')).concurrencyKey + ).not.toBe( + documentProcessingQueueOptions( + workspacePayload({ + workspaceId: 'shared-id', + billingAttribution: { ...BILLING_ATTRIBUTION, workspaceId: 'shared-id' }, + }) + ).concurrencyKey + ) + }) +}) diff --git a/apps/sim/lib/knowledge/documents/processing-lane.ts b/apps/sim/lib/knowledge/documents/processing-lane.ts new file mode 100644 index 00000000000..9a64f51113a --- /dev/null +++ b/apps/sim/lib/knowledge/documents/processing-lane.ts @@ -0,0 +1,57 @@ +import { resourceScopeFromOwner, resourceScopeKey } from '@/lib/core/resource-scope' +import type { + DocumentProcessingBillingContext, + DocumentProcessingPayload, +} from '@/lib/knowledge/documents/processing-payload' + +/** + * Queue backing the interactive lane. + * + * Deliberately the name the single shared queue already used in every deployed + * environment. A queue named at trigger time that the running Trigger.dev + * version has not registered leaves its run in `PENDING_VERSION` until the + * worker deploy catches up, and the app and worker deploy separately. Keeping + * person-facing work on the pre-existing name means only the backfill lane can + * be caught by that window, and stranded backfill work is exactly what the + * stuck-document sweep already recovers. + */ +export const INTERACTIVE_PROCESSING_QUEUE_NAME = 'document-processing-queue' + +export const BACKFILL_PROCESSING_QUEUE_NAME = 'document-processing-backfill-queue' + +/** + * The fairness key a lane's concurrency limit is applied per copy of. + * + * Keyed on the entity that owns the knowledge base, which is the entity whose + * sync can produce unbounded work — not on the workspace alone. An + * organization-scoped knowledge base carries `workspaceId: null`, so keying on + * the workspace would collapse every such tenant onto one shared bucket and + * reproduce the starvation the lanes exist to prevent. + * + * Owned scopes defer to {@link resourceScopeKey} so tenant identity has one + * spelling across the codebase. A knowledge base with neither owner has no + * `ResourceScope`, and falls back to the actor that created it. + */ +function documentProcessingTenantKey(context: DocumentProcessingBillingContext): string { + return context.billingScope === 'non-workspace' + ? `user:${context.actorUserId}` + : resourceScopeKey(resourceScopeFromOwner(context)) +} + +/** + * Trigger.dev options placing one payload in its lane's per-tenant queue copy. + * Both must be set together: the queue name alone is shared by every tenant, + * and the key alone would split whichever queue the payload happened to land in. + */ +export function documentProcessingQueueOptions(payload: DocumentProcessingPayload): { + queue: string + concurrencyKey: string +} { + return { + queue: + payload.processingLane === 'interactive' + ? INTERACTIVE_PROCESSING_QUEUE_NAME + : BACKFILL_PROCESSING_QUEUE_NAME, + concurrencyKey: documentProcessingTenantKey(payload), + } +} diff --git a/apps/sim/lib/knowledge/documents/processing-outbox-event.ts b/apps/sim/lib/knowledge/documents/processing-outbox-event.ts index ae893680044..ea8d2071ae1 100644 --- a/apps/sim/lib/knowledge/documents/processing-outbox-event.ts +++ b/apps/sim/lib/knowledge/documents/processing-outbox-event.ts @@ -1,6 +1,7 @@ import type { db } from '@sim/db' import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' import { enqueueOutboxEvent } from '@/lib/core/outbox/service' +import type { DocumentProcessingLane } from '@/lib/knowledge/documents/processing-payload' import type { ProcessingOptions } from '@/lib/knowledge/documents/service' export const KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT = 'knowledge.document.processing.dispatch' @@ -10,6 +11,11 @@ export interface KnowledgeDocumentProcessingOutboxPayload { documentId: string processingOptions: ProcessingOptions billingAttribution: BillingAttributionSnapshot + /** + * Required so the relay never has to infer a lane. Rows enqueued before the + * lanes existed carry none and are replayed conservatively as bulk. + */ + processingLane: DocumentProcessingLane } /** Enqueues durable processing in the same transaction that creates the document. */ diff --git a/apps/sim/lib/knowledge/documents/processing-outbox-handler.test.ts b/apps/sim/lib/knowledge/documents/processing-outbox-handler.test.ts index 4ea7dc6b006..22df7231552 100644 --- a/apps/sim/lib/knowledge/documents/processing-outbox-handler.test.ts +++ b/apps/sim/lib/knowledge/documents/processing-outbox-handler.test.ts @@ -56,6 +56,7 @@ const PAYLOAD = { documentId: 'document-1', processingOptions: { recipe: 'default', lang: 'en' }, billingAttribution: BILLING_ATTRIBUTION, + processingLane: 'interactive', } function createContext(eventId = 'outbox-event-1'): OutboxEventContext { @@ -120,7 +121,7 @@ describe('knowledge document processing outbox handler', () => { const context = { ...createContext(), deadlineAt: Date.now() + 550_000 } await handler()(PAYLOAD, context) expect(handler().timeoutMs).toBe(550_000) - expect(mocks.processDocumentsWithQueue.mock.calls[0][6]).toEqual({ + expect(mocks.processDocumentsWithQueue.mock.calls[0][7]).toEqual({ signal: context.signal, deadlineAt: context.deadlineAt, }) @@ -148,6 +149,7 @@ describe('knowledge document processing outbox handler', () => { { recipe: 'default', lang: 'en' }, 'outbox-event-stable', BILLING_ATTRIBUTION, + 'interactive', undefined, { signal: expect.any(AbortSignal), deadlineAt: undefined } ) @@ -215,6 +217,7 @@ describe('knowledge document processing outbox handler', () => { { recipe: 'default', lang: 'en' }, 'outbox-event-retry', BILLING_ATTRIBUTION, + 'interactive', undefined, { signal: expect.any(AbortSignal), deadlineAt: undefined } ) diff --git a/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts b/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts index eef0ceb2a25..112b3cef991 100644 --- a/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts +++ b/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts @@ -35,6 +35,7 @@ import { } from '@/lib/knowledge/documents/processing-outbox-event' import { assertDocumentProcessingPayload, + resolveDocumentProcessingLane, shouldRefundDocumentProcessingPredecessor, } from '@/lib/knowledge/documents/processing-payload' import { scheduleDocumentProcessingProviderContinuation } from '@/lib/knowledge/documents/processing-provider-continuation' @@ -100,6 +101,7 @@ function parsePayload(payload: unknown): KnowledgeDocumentProcessingOutboxPayloa documentId: requireNonEmptyString(record.documentId, 'documentId'), processingOptions: parseProcessingOptions(record.processingOptions), billingAttribution: assertBillingAttributionSnapshot(record.billingAttribution), + processingLane: resolveDocumentProcessingLane(record.processingLane), } } @@ -144,6 +146,7 @@ const processKnowledgeDocument: OutboxHandler = async (rawPayload, cont payload.processingOptions, context.eventId, payload.billingAttribution, + payload.processingLane, undefined, { signal: context.signal, deadlineAt: context.deadlineAt } ) diff --git a/apps/sim/lib/knowledge/documents/processing-payload.ts b/apps/sim/lib/knowledge/documents/processing-payload.ts index 83ccbe39878..8cd6a4a797b 100644 --- a/apps/sim/lib/knowledge/documents/processing-payload.ts +++ b/apps/sim/lib/knowledge/documents/processing-payload.ts @@ -3,10 +3,46 @@ import { assertBillingAttributionSnapshot, type BillingAttributionSnapshot, } from '@/lib/billing/core/billing-attribution' +/** + * Which shared processing queue a document's indexing pass is admitted through. + * + * `interactive` is work a person is waiting on — a direct upload, an upload + * session, a retry click, including a retry of a connector-owned document. + * `backfill` is connector-driven bulk ingestion, which is throughput-shaped + * rather than latency-shaped and is always reclaimable by the next sync's + * stuck-document sweep. + * + * The lanes exist because both used to share one queue: a single connector sync + * holding every slot left every other tenant's uploads waiting behind it. + * + * Unrelated to the `bulk` flag in `@/lib/embeddings/client`, which separates + * document indexing from query-time embedding. A person's upload is + * `interactive` here and `bulk` there, in the same request. + */ +export type DocumentProcessingLane = 'interactive' | 'backfill' + +/** + * Reads a lane off an untrusted payload, where only an explicit `interactive` + * stamp earns the interactive lane. + * + * Never throws, and never widens. Payloads written before the lanes existed + * carry no lane, and a rolling deploy can hand this version a payload stamped + * by a newer one; rejecting either would fail the run into its retry budget for + * the length of a rollout. Everything unrecognized reads as backfill, so an + * unlabeled pass cannot claim capacity it was not admitted against. + * + * Only payload parsers need this. A dispatch site names its lane as a required + * argument, so omitting one there is a compile error, not a silent downgrade. + */ +export function resolveDocumentProcessingLane(value: unknown): DocumentProcessingLane { + return value === 'interactive' ? 'interactive' : 'backfill' +} export interface DocumentProcessingPayloadBase { knowledgeBaseId: string documentId: string + /** Required so a new dispatch site cannot silently inherit interactive capacity. */ + processingLane: DocumentProcessingLane docData: { filename: string fileUrl: string @@ -371,6 +407,7 @@ export function assertDocumentProcessingPayload(value: unknown): DocumentProcess return { knowledgeBaseId: value.knowledgeBaseId, documentId: value.documentId, + processingLane: resolveDocumentProcessingLane(value.processingLane), docData: { filename: docData.filename, fileUrl: docData.fileUrl, diff --git a/apps/sim/lib/knowledge/documents/processing-queue.test.ts b/apps/sim/lib/knowledge/documents/processing-queue.test.ts index b4fdac5d664..64debf56f43 100644 --- a/apps/sim/lib/knowledge/documents/processing-queue.test.ts +++ b/apps/sim/lib/knowledge/documents/processing-queue.test.ts @@ -106,7 +106,8 @@ describe('processDocumentsWithQueue billing attribution', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) const jobs = mockBatchTrigger.mock.calls[0][1] @@ -119,6 +120,7 @@ describe('processDocumentsWithQueue billing attribution', () => { expect(structuredClone(jobs[0].payload)).toEqual({ knowledgeBaseId: 'knowledge-base-1', documentId: 'document-1', + processingLane: 'interactive', docData: { filename: 'document.txt', fileUrl: 'https://example.com/document.txt', @@ -183,7 +185,14 @@ describe('processDocumentsWithQueue billing attribution', () => { ]) await expect( - processDocumentsWithQueue([DOCUMENT], 'knowledge-base-1', {}, 'request-1', undefined) + processDocumentsWithQueue( + [DOCUMENT], + 'knowledge-base-1', + {}, + 'request-1', + undefined, + 'interactive' + ) ).rejects.toThrow('Workspace document processing requires a billing attribution snapshot') expect(mockBatchTrigger).not.toHaveBeenCalled() const withdrawal = dbChainMockFns.set.mock.calls.find( @@ -204,7 +213,8 @@ describe('processDocumentsWithQueue billing attribution', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) ).rejects.toThrow('Document processing workspace does not match billing attribution') expect(mockBatchTrigger).not.toHaveBeenCalled() @@ -216,7 +226,14 @@ describe('processDocumentsWithQueue billing attribution', () => { ]) await expect( - processDocumentsWithQueue([DOCUMENT], 'knowledge-base-1', {}, 'request-1', undefined) + processDocumentsWithQueue( + [DOCUMENT], + 'knowledge-base-1', + {}, + 'request-1', + undefined, + 'interactive' + ) ).rejects.toThrow('Document processing requires a workspace or organization owner') expect(mockBatchTrigger).not.toHaveBeenCalled() }) @@ -325,12 +342,114 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(mockBatchTrigger).toHaveBeenCalledTimes(1) }) + it('places each job in its lane queue under the owning tenant concurrency key', async () => { + markInsideTriggerRun() + + await processDocumentsWithQueue( + [DOCUMENT], + 'knowledge-base-1', + {}, + 'request-1', + BILLING_ATTRIBUTION, + 'interactive' + ) + + const items = mockBatchTrigger.mock.calls[0][1] + expect(items[0].options).toMatchObject({ + queue: 'document-processing-queue', + concurrencyKey: 'workspace:workspace-1', + }) + }) + + /** + * The starvation this split exists to prevent: connector backfill must not be + * able to occupy the queue a person's upload is admitted through. + */ + it('sends connector backfill to a different queue than interactive work', async () => { + markInsideTriggerRun() + + await processDocumentsWithQueue( + [DOCUMENT], + 'knowledge-base-1', + {}, + 'request-1', + BILLING_ATTRIBUTION, + 'backfill' + ) + + const items = mockBatchTrigger.mock.calls[0][1] + expect(items[0].options.queue).toBe('document-processing-backfill-queue') + expect(items[0].options.concurrencyKey).toBe('workspace:workspace-1') + }) + + it('keys an organization-owned knowledge base on its organization, not a null workspace', async () => { + markInsideTriggerRun() + dbChainMockFns.limit.mockResolvedValue([ + { userId: 'knowledge-owner', workspaceId: null, organizationId: 'org-1' }, + ]) + + await processDocumentsWithQueue( + [DOCUMENT], + 'knowledge-base-1', + {}, + 'request-1', + { + ...BILLING_ATTRIBUTION, + workspaceId: null, + organizationId: 'org-1', + billingEntity: { type: 'organization', id: 'org-1' }, + } satisfies BillingAttributionSnapshot, + 'backfill' + ) + + const items = mockBatchTrigger.mock.calls[0][1] + expect(items[0].options.concurrencyKey).toBe('organization:org-1') + }) + + /** + * A rejected backfill chunk means the tenant's own queue is already saturated. + * Absorbing those documents into the dispatching worker would turn queue + * backpressure into unbounded in-process fan-out on the one machine still + * holding the connector lease, so they are reported failed and left for the + * next sync's stuck-document sweep. The interactive lane still falls back, + * because nothing else would ever come back for a person's upload. + */ + it('leaves a rejected backfill chunk for the sweep instead of processing it in-process', async () => { + markInsideTriggerRun() + const documents = Array.from({ length: 1001 }, (_, index) => ({ + ...DOCUMENT, + documentId: `document-${index}`, + })) + dbChainMockFns.returning.mockResolvedValueOnce(documents.map((doc) => ({ id: doc.documentId }))) + mockBatchTrigger + .mockResolvedValueOnce({ batchId: 'batch-1' }) + .mockRejectedValueOnce(new Error('queue is full')) + + const result = await processDocumentsWithQueue( + documents, + 'knowledge-base-1', + {}, + 'request-1', + BILLING_ATTRIBUTION, + 'backfill' + ) + + expect(mockBatchTrigger).toHaveBeenCalledTimes(2) + expect(result).toMatchObject({ + requested: 1001, + accepted: 1000, + failed: 1, + failedDocumentIds: ['document-1000'], + }) + }) + it('returns acceptance separately from eventual child completion', async () => { markInsideTriggerRun() @@ -339,7 +458,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(result).toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) @@ -361,7 +481,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(result).toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) @@ -381,7 +502,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(result).toEqual({ @@ -504,7 +626,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(result).toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) @@ -580,7 +703,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(result).toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) @@ -625,7 +749,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(result).toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) @@ -692,7 +817,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) ).rejects.toThrow('document processing dispatches failed') @@ -728,7 +854,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(result).toEqual({ @@ -749,7 +876,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(result).toEqual({ requested: 2, accepted: 2, failed: 0, failedDocumentIds: [] }) @@ -761,7 +889,14 @@ describe('processDocumentsWithQueue dispatch backend', () => { it('returns an empty dispatch summary without resolving billing context', async () => { await expect( - processDocumentsWithQueue([], 'missing-knowledge-base', {}, 'request-1', undefined) + processDocumentsWithQueue( + [], + 'missing-knowledge-base', + {}, + 'request-1', + undefined, + 'interactive' + ) ).resolves.toEqual({ requested: 0, accepted: 0, failed: 0, failedDocumentIds: [] }) expect(mockBatchTrigger).not.toHaveBeenCalled() @@ -777,7 +912,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect(mockBatchTrigger).toHaveBeenCalledTimes(1) @@ -792,7 +928,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) ).resolves.toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) @@ -834,7 +971,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) ).rejects.toThrow('document processing dispatches failed') @@ -876,7 +1014,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) ).resolves.toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) @@ -958,7 +1097,8 @@ describe('processDocumentsWithQueue dispatch backend', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) ).resolves.toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) @@ -999,7 +1139,8 @@ describe('processDocumentsWithQueue attempt refund', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) ).rejects.toThrow('trigger.dev region unavailable') @@ -1048,7 +1189,8 @@ describe('processDocumentsWithQueue attempt refund', () => { 'knowledge-base-1', {}, 'request-1', - BILLING_ATTRIBUTION + BILLING_ATTRIBUTION, + 'interactive' ) expect( @@ -1092,6 +1234,7 @@ describe('processDocumentsWithQueue under a connector sync lease', () => { {}, 'request-1', BILLING_ATTRIBUTION, + 'backfill', lease ) ).rejects.toBeInstanceOf(SyncLockLostException) @@ -1115,6 +1258,7 @@ describe('processDocumentsWithQueue under a connector sync lease', () => { {}, 'request-1', BILLING_ATTRIBUTION, + 'backfill', lease ) ).resolves.toEqual({ requested: 1, accepted: 1, failed: 0, failedDocumentIds: [] }) diff --git a/apps/sim/lib/knowledge/documents/processing-recovery.ts b/apps/sim/lib/knowledge/documents/processing-recovery.ts index 8d558070ff1..c9aeb3206ed 100644 --- a/apps/sim/lib/knowledge/documents/processing-recovery.ts +++ b/apps/sim/lib/knowledge/documents/processing-recovery.ts @@ -170,6 +170,8 @@ async function recoverStoredDocuments( mimeType: doc.mimeType, }, processingOptions: {}, + /** This sweep only selects connector-owned documents. */ + processingLane: 'backfill', requestId: token, processingQueueToken: token, processingQueuedAt: queuedAt.toISOString(), diff --git a/apps/sim/lib/knowledge/documents/service.ts b/apps/sim/lib/knowledge/documents/service.ts index 1aac135fb6e..16531d6fa01 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -103,6 +103,7 @@ import { recordUndispatchedDocumentFailure, } from '@/lib/knowledge/documents/processing-claim' import type { DocumentProcessingContinuation } from '@/lib/knowledge/documents/processing-continuation-dispatch' +import { documentProcessingQueueOptions } from '@/lib/knowledge/documents/processing-lane' import { enqueueKnowledgeDocumentProcessing } from '@/lib/knowledge/documents/processing-outbox-event' import { assertDocumentProcessingBillingContext, @@ -110,6 +111,7 @@ import { createOrganizationDocumentProcessingBillingContext, createWorkspaceDocumentProcessingBillingContext, type DocumentProcessingBillingContext, + type DocumentProcessingLane, type DocumentProcessingPayload, hasDocumentProcessingBillingScope, } from '@/lib/knowledge/documents/processing-payload' @@ -720,12 +722,14 @@ function buildJobPayload( processingQueueToken: string, processingQueuedAt: Date, chargedAtDispatch: boolean, - billingContext: DocumentProcessingBillingContext + billingContext: DocumentProcessingBillingContext, + lane: DocumentProcessingLane ): DocumentProcessingPayload { return createDocumentProcessingPayload( { knowledgeBaseId, documentId: doc.documentId, + processingLane: lane, docData: { filename: doc.filename, fileUrl: doc.fileUrl, @@ -1070,6 +1074,11 @@ async function bestEffortWithdrawDocumentsQueued( * pass. A successful Trigger.dev hand-off is only an accepted child run, not a * claim about its eventual processing outcome. A connector sync passes its * lease, and the queue write then lands only while the run still holds it. + * + * `lane` decides which per-tenant queue the work is admitted through, and is + * required rather than inferred so a new caller has to state whether someone is + * waiting on the result. It is stamped onto every payload so continuations of + * this pass resume in the same lane. */ export async function processDocumentsWithQueue( createdDocuments: DocumentData[], @@ -1077,6 +1086,7 @@ export async function processDocumentsWithQueue( processingOptions: ProcessingOptions, requestId: string, billingAttribution: BillingAttributionSnapshot | undefined, + lane: DocumentProcessingLane, lease?: ProcessingDispatchLease, executionContext?: DocumentProcessingExecutionContext ): Promise { @@ -1150,7 +1160,8 @@ export async function processDocumentsWithQueue( requestId, generation.processingQueuedAt, generation.chargedAtDispatch, - billingContext + billingContext, + lane ) }) @@ -1239,6 +1250,7 @@ async function dispatchViaBatchTrigger( `knowledgeBaseId:${payload.knowledgeBaseId}`, `documentId:${payload.documentId}`, ], + ...documentProcessingQueueOptions(payload), region, }, })) @@ -1260,12 +1272,27 @@ async function dispatchViaBatchTrigger( * Only a total dispatch failure raises, so a chunk failing alone would leave its * documents at `pending` with nothing recording why. Processing them here is * slower than the queue but does not drop the work. + * + * Backfill is excluded deliberately. Such a chunk is rejected precisely when + * the tenant is already saturated, and the worker absorbing it is the one + * holding the connector lease: the in-flight count stays bounded, but the + * total work does not, so a rejected chunk can outlast the lease it is + * running under. Leaving those documents undispatched reports them failed, + * which is what the next sync's stuck-document sweep reclaims. Interactive + * work has no such sweep behind it, so it still falls back. */ - if (undispatched.length > 0) { + const fallback = undispatched.filter((payload) => payload.processingLane !== 'backfill') + const swept = undispatched.length - fallback.length + if (swept > 0) { + logger.warn( + `[${requestId}] Leaving ${swept} backfill documents for the stuck-document sweep after failed enqueue` + ) + } + if (fallback.length > 0) { logger.warn( - `[${requestId}] Processing ${undispatched.length} documents in-process after failed enqueue` + `[${requestId}] Processing ${fallback.length} documents in-process after failed enqueue` ) - const directlyDispatchedIds = await dispatchInProcess(undispatched, requestId, executionContext) + const directlyDispatchedIds = await dispatchInProcess(fallback, requestId, executionContext) for (const documentId of directlyDispatchedIds) dispatchedIds.add(documentId) } @@ -3102,6 +3129,12 @@ export async function createSingleDocument( documentId, processingOptions: options.processing.processingOptions, billingAttribution: options.processing.billingAttribution, + /** + * Every caller of this function is a person adding a document — an + * upload session, a workspace-file attach, a direct create. Connector + * ingestion builds its documents through the sync engine instead. + */ + processingLane: 'interactive', }) } @@ -3544,7 +3577,9 @@ export async function retryDocumentProcessing( knowledgeBaseId, {}, requestId, - billingAttribution + billingAttribution, + /** One document, retried by hand: a person is waiting on this one. */ + 'interactive' ) if (dispatch.failed > 0 || dispatch.accepted !== 1) { throw new Error(`Document processing dispatch was not accepted for ${documentId}`) diff --git a/apps/sim/lib/knowledge/documents/storage-billing.test.ts b/apps/sim/lib/knowledge/documents/storage-billing.test.ts index 07677ae676e..3a327f7d25d 100644 --- a/apps/sim/lib/knowledge/documents/storage-billing.test.ts +++ b/apps/sim/lib/knowledge/documents/storage-billing.test.ts @@ -236,6 +236,7 @@ describe('knowledge document storage attribution', () => { expect(mockEnqueueKnowledgeDocumentProcessing).toHaveBeenCalledWith(dbChainMock.db, { knowledgeBaseId: 'knowledge-base-1', documentId: 'document-1', + processingLane: 'interactive', processingOptions: { lang: 'en' }, billingAttribution, }) diff --git a/apps/sim/lib/knowledge/documents/types.ts b/apps/sim/lib/knowledge/documents/types.ts index 3b0bd225f99..4bf7435f690 100644 --- a/apps/sim/lib/knowledge/documents/types.ts +++ b/apps/sim/lib/knowledge/documents/types.ts @@ -27,11 +27,11 @@ export const MAX_PROCESSING_ATTEMPTS = 5 * `STALE_PROCESSING_MINUTES` bounds a run that has already begun, derived from * the task's own duration and retry budget. Queue *wait* is a different * quantity: it is backlog / concurrency, not run duration. - * `document-processing-queue` has a global concurrency shared by every - * workspace, so a corpus large enough to approach - * `CONNECTOR_SYNC_MAX_DURATION_SECONDS` enqueues thousands of documents that - * drain in waves of that width — at roughly a minute of occupancy each, a few - * hours, and longer while other workspaces hold slots. + * Backfill drains through a per-tenant copy of the backfill queue, so a corpus + * large enough to approach `CONNECTOR_SYNC_MAX_DURATION_SECONDS` enqueues + * thousands of documents that drain in waves of that concurrency — at roughly a + * minute of occupancy each, a few hours. Another tenant's corpus no longer + * extends that wait, but the shared environment concurrency limit still can. * * Four hours is chosen against three bounds that are all constants in this * repository rather than any one deployment's corpus: it is well above that