diff --git a/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.test.ts b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.test.ts new file mode 100644 index 00000000000..152186137e7 --- /dev/null +++ b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.test.ts @@ -0,0 +1,63 @@ +/** + * @vitest-environment node + */ + +import { authMockFns, createMockRequest } from '@sim/testing' +import { beforeEach, describe, expect, it, vi } from 'vitest' + +const mocks = vi.hoisted(() => ({ + list: vi.fn(), + create: vi.fn(), + bulk: vi.fn(), +})) + +vi.mock('@/lib/knowledge/application/chunks', () => ({ + listKnowledgeChunks: { + operation: { id: 'knowledge.chunks.list' }, + execute: mocks.list, + }, + createKnowledgeChunk: { + operation: { id: 'knowledge.chunks.create' }, + execute: mocks.create, + }, + bulkUpdateKnowledgeChunks: { + operation: { id: 'knowledge.chunks.bulk' }, + execute: mocks.bulk, + }, +})) + +vi.mock('@/app/api/knowledge/secret-provenance', () => ({ + finalizeKnowledgePersistedResponse: vi.fn(), + finalizeKnowledgeProvenanceResponse: vi.fn(), + resolveKnowledgeWriteSecretProvenance: vi.fn(), +})) + +import { KnowledgeDocumentNotReadyError } from '@/lib/knowledge/application/chunk-errors' +import { GET } from '@/app/api/knowledge/[id]/documents/[documentId]/chunks/route' + +const params = () => ({ + params: Promise.resolve({ id: 'knowledge-1', documentId: 'document-1' }), +}) + +describe('/api/knowledge/[id]/documents/[documentId]/chunks internal route composition', () => { + beforeEach(() => { + vi.clearAllMocks() + authMockFns.mockGetSession.mockResolvedValue({ + user: { id: 'user-1' }, + session: { id: 'session-1' }, + }) + }) + + it('preserves retry metadata when a document is still processing', async () => { + mocks.list.mockRejectedValueOnce(new KnowledgeDocumentNotReadyError('processing')) + + const response = await GET(createMockRequest('GET'), params()) + + expect(response.status).toBe(400) + await expect(response.json()).resolves.toEqual({ + error: 'Document is not ready for access', + details: 'Document status: processing', + retryAfter: 5, + }) + }) +}) diff --git a/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.ts b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.ts index c7461356308..d805898f287 100644 --- a/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.ts +++ b/apps/sim/app/api/knowledge/[id]/documents/[documentId]/chunks/route.ts @@ -54,7 +54,7 @@ export const GET = defineInternalJsonRoute({ auth: internalKnowledgeSessionOrExecutorAuth, operation: knowledgeOperations.listChunks, rateLimit: internalRateLimits.none({ reason: 'Preserve existing internal chunk-list behavior' }), - errorPolicy: internalKnowledgeErrorPolicies.chunks, + errorPolicy: internalKnowledgeErrorPolicies.chunkList, mapInput: ({ params, query }) => ({ knowledgeBaseId: params.id, documentId: params.documentId, diff --git a/apps/sim/app/api/v2/files/[fileId]/route.test.ts b/apps/sim/app/api/v2/files/[fileId]/route.test.ts index 770f1eafddd..a63dbb5c50a 100644 --- a/apps/sim/app/api/v2/files/[fileId]/route.test.ts +++ b/apps/sim/app/api/v2/files/[fileId]/route.test.ts @@ -141,6 +141,25 @@ describe('v2 single-file routes', () => { }) }) + it('encodes special characters in the extended download filename', async () => { + mocks.download.mockResolvedValueOnce({ + file: fileRecord({ name: "it's (final)* café.pdf" }), + stream: new Blob(['pdf']).stream(), + contentType: 'application/pdf', + contentLength: 3, + }) + + const response = await GET( + new NextRequest(`http://localhost:3000/api/v2/files/${FILE_ID}?workspaceId=${WORKSPACE_ID}`), + context + ) + + expect(response.status).toBe(200) + expect(response.headers.get('Content-Disposition')).toBe( + `attachment; filename="it's (final)* caf_.pdf"; filename*=UTF-8''it%27s%20%28final%29%2A%20caf%C3%A9.pdf` + ) + }) + it('conceals cross-workspace download authorization', async () => { mocks.download.mockRejectedValue(new NoWorkspaceAccessError()) diff --git a/apps/sim/app/api/v2/files/[fileId]/route.ts b/apps/sim/app/api/v2/files/[fileId]/route.ts index 10b1ca6c94d..0a7bbc896bc 100644 --- a/apps/sim/app/api/v2/files/[fileId]/route.ts +++ b/apps/sim/app/api/v2/files/[fileId]/route.ts @@ -14,6 +14,7 @@ import { deleteWorkspaceFileOperation } from '@/lib/workspace-files/application/ import { downloadWorkspaceFileStream } from '@/lib/workspace-files/application/download-workspace-file' import { fileOperations } from '@/lib/workspace-files/application/operations' import { renameWorkspaceFile } from '@/lib/workspace-files/application/rename-workspace-file' +import { encodeFilenameForHeader } from '@/app/api/files/utils' import { toV2File } from '@/app/api/v2/files/utils' export const dynamic = 'force-dynamic' @@ -42,7 +43,7 @@ export const GET = defineV2BinaryRoute({ present: ({ file, stream, contentType, contentLength }) => ({ body: stream, contentType, - contentDisposition: `attachment; filename="${file.name.replace(/[^\w.-]/g, '_')}"; filename*=UTF-8''${encodeURIComponent(file.name)}`, + contentDisposition: `attachment; ${encodeFilenameForHeader(file.name)}`, contentLength, }), }) diff --git a/apps/sim/app/api/v2/knowledge/[id]/documents/route.test.ts b/apps/sim/app/api/v2/knowledge/[id]/documents/route.test.ts index 401d0f63f9b..1b03a060bbe 100644 --- a/apps/sim/app/api/v2/knowledge/[id]/documents/route.test.ts +++ b/apps/sim/app/api/v2/knowledge/[id]/documents/route.test.ts @@ -17,7 +17,6 @@ const { mockUploadDocument, mockReadFormData, mockReadFile, - mockUploadWorkspaceFile, mockPlatformUploaded, mockCapture, mockIsPayloadSizeLimitError, @@ -26,7 +25,6 @@ const { mockUploadDocument: vi.fn(), mockReadFormData: vi.fn(), mockReadFile: vi.fn(), - mockUploadWorkspaceFile: vi.fn(), mockPlatformUploaded: vi.fn(), mockCapture: vi.fn(), mockIsPayloadSizeLimitError: vi.fn(), @@ -58,10 +56,6 @@ vi.mock('@/lib/core/utils/stream-limits', () => ({ readFileToBufferWithLimit: mockReadFile, })) -vi.mock('@/lib/uploads/contexts/workspace', () => ({ - uploadWorkspaceFile: mockUploadWorkspaceFile, -})) - vi.mock('@/lib/core/telemetry', () => ({ PlatformEvents: { knowledgeBaseDocumentsUploaded: mockPlatformUploaded }, })) @@ -102,13 +96,11 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { knowledgeBaseId: 'kb-1', knowledgeBaseName: 'Support docs', workspaceId: WORKSPACE_ID, - storageActorUserId: 'user-1', }) const formData = new FormData() formData.set('file', new File(['hello'], 'support.txt', { type: 'text/plain' })) mockReadFormData.mockResolvedValue(formData) mockReadFile.mockResolvedValue(Buffer.from('hello')) - mockUploadWorkspaceFile.mockResolvedValue({ url: 's3://workspace/support.txt' }) mockUploadDocument.mockResolvedValue({ created: true, document: { @@ -141,21 +133,14 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { input: { knowledgeBaseId: 'kb-1', assertedWorkspaceId: WORKSPACE_ID }, request, }) - expect(mockUploadWorkspaceFile).toHaveBeenCalledWith( - WORKSPACE_ID, - 'user-1', - Buffer.from('hello'), - 'support.txt', - 'text/plain' - ) expect(mockUploadDocument).toHaveBeenCalledWith({ principal: PRINCIPAL, input: { knowledgeBaseId: 'kb-1', assertedWorkspaceId: WORKSPACE_ID, - document: { + file: { + buffer: Buffer.from('hello'), filename: 'support.txt', - fileUrl: 's3://workspace/support.txt', fileSize: 5, mimeType: 'text/plain', }, @@ -194,7 +179,6 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { error: { code: 'USAGE_LIMIT_EXCEEDED', message: 'Upgrade required' }, }) expect(mockReadFormData).not.toHaveBeenCalled() - expect(mockUploadWorkspaceFile).not.toHaveBeenCalled() expect(mockUploadDocument).not.toHaveBeenCalled() }) @@ -214,7 +198,7 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { expect(mockCapture).not.toHaveBeenCalled() }) - it('preserves the malformed multipart envelope without transferring storage', async () => { + it('preserves the malformed multipart envelope without entering the upload operation', async () => { mockReadFormData.mockRejectedValueOnce(new Error('multipart boundary missing')) const response = await POST(buildRequest(), { params: Promise.resolve({ id: 'kb-1' }) }) @@ -223,12 +207,11 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { expect(await response.json()).toEqual({ error: { code: 'BAD_REQUEST', message: 'Request body must be valid multipart form data' }, }) - expect(mockUploadWorkspaceFile).not.toHaveBeenCalled() expect(mockUploadDocument).not.toHaveBeenCalled() expect(mockPlatformUploaded).not.toHaveBeenCalled() }) - it('preserves bounded multipart rejection and stops before storage transfer', async () => { + it('preserves bounded multipart rejection and stops before the upload operation', async () => { const error = new Error('knowledge document upload body exceeds maximum size') mockReadFormData.mockRejectedValueOnce(error) mockIsPayloadSizeLimitError.mockImplementation((candidate: unknown) => candidate === error) @@ -239,11 +222,10 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { expect(await response.json()).toEqual({ error: { code: 'PAYLOAD_TOO_LARGE', message: error.message }, }) - expect(mockUploadWorkspaceFile).not.toHaveBeenCalled() expect(mockUploadDocument).not.toHaveBeenCalled() }) - it('requires a file form field before storage transfer', async () => { + it('requires a file form field before the upload operation', async () => { mockReadFormData.mockResolvedValueOnce(new FormData()) const response = await POST(buildRequest(), { params: Promise.resolve({ id: 'kb-1' }) }) @@ -252,7 +234,7 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { expect(await response.json()).toEqual({ error: { code: 'BAD_REQUEST', message: 'file form field is required' }, }) - expect(mockUploadWorkspaceFile).not.toHaveBeenCalled() + expect(mockUploadDocument).not.toHaveBeenCalled() }) it('preserves the exact file-size rejection before reading file bytes', async () => { @@ -269,7 +251,7 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { error: { code: 'PAYLOAD_TOO_LARGE', message: 'File size exceeds 100MB limit (100.00MB)' }, }) expect(mockReadFile).not.toHaveBeenCalled() - expect(mockUploadWorkspaceFile).not.toHaveBeenCalled() + expect(mockUploadDocument).not.toHaveBeenCalled() }) it('preserves unsupported file-type validation before reading file bytes', async () => { @@ -286,11 +268,11 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { error: { code: 'UNSUPPORTED_MEDIA_TYPE', message: expectedMessage }, }) expect(mockReadFile).not.toHaveBeenCalled() - expect(mockUploadWorkspaceFile).not.toHaveBeenCalled() + expect(mockUploadDocument).not.toHaveBeenCalled() }) - it('does not register or emit effects when storage transfer fails', async () => { - mockUploadWorkspaceFile.mockRejectedValueOnce(new Error('storage unavailable')) + it('does not emit effects when the upload operation fails', async () => { + mockUploadDocument.mockRejectedValueOnce(new Error('storage unavailable')) const response = await POST(buildRequest(), { params: Promise.resolve({ id: 'kb-1' }) }) @@ -298,12 +280,12 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { expect(await response.json()).toEqual({ error: { code: 'INTERNAL_ERROR', message: 'Internal server error' }, }) - expect(mockUploadDocument).not.toHaveBeenCalled() + expect(mockUploadDocument).toHaveBeenCalledOnce() expect(mockPlatformUploaded).not.toHaveBeenCalled() expect(mockCapture).not.toHaveBeenCalled() }) - it('preserves application authorization errors after storage transfer', async () => { + it('preserves final application authorization errors', async () => { mockUploadDocument.mockRejectedValueOnce( new OrchestrationError('forbidden', 'Insufficient workspace permissions') ) @@ -314,7 +296,7 @@ describe('POST /api/v2/knowledge/[id]/documents', () => { expect(await response.json()).toEqual({ error: { code: 'FORBIDDEN', message: 'Insufficient workspace permissions' }, }) - expect(mockUploadWorkspaceFile).toHaveBeenCalledOnce() + expect(mockUploadDocument).toHaveBeenCalledOnce() expect(mockPlatformUploaded).not.toHaveBeenCalled() expect(mockCapture).not.toHaveBeenCalled() }) diff --git a/apps/sim/app/api/v2/knowledge/[id]/documents/route.ts b/apps/sim/app/api/v2/knowledge/[id]/documents/route.ts index 62871ad988f..6ba62a1a941 100644 --- a/apps/sim/app/api/v2/knowledge/[id]/documents/route.ts +++ b/apps/sim/app/api/v2/knowledge/[id]/documents/route.ts @@ -26,7 +26,6 @@ import { import { knowledgeOperations } from '@/lib/knowledge/application/operations' import { KnowledgeDocumentUnsupportedMediaTypeError } from '@/lib/knowledge/application/upload-sessions' import { captureServerEvent } from '@/lib/posthog/server' -import { uploadWorkspaceFile } from '@/lib/uploads/contexts/workspace' import { MAX_KNOWLEDGE_DOCUMENT_FILE_SIZE } from '@/lib/uploads/shared/types' import { validateFileType } from '@/lib/uploads/utils/validation' import { serializeDate } from '@/app/api/v1/knowledge/utils' @@ -138,20 +137,12 @@ export const POST = defineV2BodyLifecycleRoute({ }) return { file: rawFile, buffer, contentType } }, - transfer: ({ admission, body }) => - uploadWorkspaceFile( - admission.workspaceId, - admission.storageActorUserId, - body.buffer, - body.file.name, - body.contentType - ), - mapInput: ({ parsed, body, transfer }) => ({ + mapInput: ({ parsed, body }) => ({ knowledgeBaseId: parsed.params.id, assertedWorkspaceId: parsed.query.workspaceId, - document: { + file: { + buffer: body.buffer, filename: body.file.name, - fileUrl: transfer.url, fileSize: body.file.size, mimeType: body.contentType, }, diff --git a/apps/sim/app/api/webhooks/outbox/process/route.ts b/apps/sim/app/api/webhooks/outbox/process/route.ts index 5a150e66ac7..a9edf8bc231 100644 --- a/apps/sim/app/api/webhooks/outbox/process/route.ts +++ b/apps/sim/app/api/webhooks/outbox/process/route.ts @@ -9,6 +9,7 @@ import { billingOutboxHandlers } from '@/lib/billing/webhooks/outbox-handlers' import { processOutboxEvents } from '@/lib/core/outbox/service' import { generateRequestId } from '@/lib/core/utils/request' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' +import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler' import { workflowDeploymentOutboxHandlers } from '@/lib/workflows/deployment-outbox' import { invitationMigrationOutboxHandlers } from '@/lib/workspaces/admin-move' import { reapStaleBackgroundWork } from '@/ee/workspace-forking/lib/background-work/store' @@ -23,6 +24,7 @@ const handlers = { ...membershipBillingOutboxHandlers, ...enterpriseIssuanceOutboxHandlers, ...invitationMigrationOutboxHandlers, + ...knowledgeDocumentProcessingOutboxHandlers, ...workflowDeploymentOutboxHandlers, } as const diff --git a/apps/sim/lib/api/server/routes/v2-body-lifecycle-route.test.ts b/apps/sim/lib/api/server/routes/v2-body-lifecycle-route.test.ts index f0d940e8abb..6885f3d1a4d 100644 --- a/apps/sim/lib/api/server/routes/v2-body-lifecycle-route.test.ts +++ b/apps/sim/lib/api/server/routes/v2-body-lifecycle-route.test.ts @@ -43,7 +43,7 @@ const contract = defineRouteContract({ class StageRejection extends Error {} -type RejectableStage = 'admission' | 'body' | 'transfer' | 'application' | 'presenter' | 'effects' +type RejectableStage = 'admission' | 'body' | 'application' | 'presenter' | 'effects' let rejectedStage: RejectableStage | null = null @@ -80,11 +80,10 @@ function buildHandler() { rejectAt('body') return { bytes: Buffer.from('body') } }, - async transfer({ admission }) { - rejectAt('transfer') - return { url: `stored://${admission.canonicalWorkspaceId}` } - }, - mapInput: ({ parsed, transfer }) => ({ id: parsed.params.id, url: transfer.url }), + mapInput: ({ parsed, admission }) => ({ + id: parsed.params.id, + url: `stored://${admission.canonicalWorkspaceId}`, + }), useCase: { operation, async execute({ input }) { @@ -161,7 +160,6 @@ describe('defineV2BodyLifecycleRoute', () => { errorPolicy: { render: () => null }, admission: { mapInput: () => ({}), useCase }, readBody: async () => Buffer.alloc(0), - transfer: async () => ({ url: 'stored://item-1' }), mapInput: () => ({}), useCase, present: () => ({ data: { id: 'item-1' } }), @@ -183,7 +181,6 @@ describe('defineV2BodyLifecycleRoute', () => { 'contract', 'admission', 'body', - 'transfer', 'application', 'presenter', 'effects', @@ -260,24 +257,20 @@ describe('defineV2BodyLifecycleRoute', () => { ]) }) - it.each([ - 'admission', - 'body', - 'transfer', - 'application', - 'presenter', - 'effects', - ])('renders typed %s rejection without entering later phases', async (stage) => { - rejectedStage = stage + it.each(['admission', 'body', 'application', 'presenter', 'effects'])( + 'renders typed %s rejection without entering later phases', + async (stage) => { + rejectedStage = stage - const response = await buildHandler()(buildRequest(), context()) + const response = await buildHandler()(buildRequest(), context()) - expect(response.status).toBe(409) - expect(await response.json()).toEqual({ - error: { code: 'CONFLICT', message: `${stage} rejected` }, - }) - expect(mocks.order.at(-1)).toBe(stage) - }) + expect(response.status).toBe(409) + expect(await response.json()).toEqual({ + error: { code: 'CONFLICT', message: `${stage} rejected` }, + }) + expect(mocks.order.at(-1)).toBe(stage) + } + ) it.each(['authentication', 'rollout', 'rate_limit'] as const)( 'maps %s infrastructure failures to service unavailable', diff --git a/apps/sim/lib/api/server/routes/v2-body-lifecycle-route.ts b/apps/sim/lib/api/server/routes/v2-body-lifecycle-route.ts index f5e04983e58..a6c415ec493 100644 --- a/apps/sim/lib/api/server/routes/v2-body-lifecycle-route.ts +++ b/apps/sim/lib/api/server/routes/v2-body-lifecycle-route.ts @@ -41,18 +41,13 @@ interface V2BodyLifecycleContext { admission: A } -interface V2BodyLifecycleTransferContext +interface V2BodyLifecycleInputContext extends V2BodyLifecycleContext { body: B } -interface V2BodyLifecycleInputContext - extends V2BodyLifecycleTransferContext { - transfer: T -} - -interface V2BodyLifecycleSuccessContext - extends V2BodyLifecycleInputContext { +interface V2BodyLifecycleSuccessContext + extends V2BodyLifecycleInputContext { input: I result: R } @@ -63,7 +58,6 @@ interface V2BodyLifecycleRouteOptions< AI, A, B, - T, I, R, > { @@ -75,12 +69,11 @@ interface V2BodyLifecycleRouteOptions< parseOptions?: Omit admission: V2BodyLifecycleAdmission readBody(context: V2BodyLifecycleContext): Promise - transfer(context: V2BodyLifecycleTransferContext): Promise - mapInput(context: V2BodyLifecycleInputContext): I + mapInput(context: V2BodyLifecycleInputContext): I useCase: OperationUseCase, I, R> present(result: R): ContractJsonResponse | Promise> onSuccess?( - context: V2BodyLifecycleSuccessContext, NoInfer> + context: V2BodyLifecycleSuccessContext, NoInfer> ): void | Promise } @@ -95,10 +88,9 @@ export function defineV2BodyLifecycleRoute< AI, A, B, - T, I, R, ->(options: V2BodyLifecycleRouteOptions): JsonNextRouteHandler { +>(options: V2BodyLifecycleRouteOptions): JsonNextRouteHandler { if (options.contract.body) { throw new Error( `${options.contract.method} ${options.contract.path} must omit its body schema so admission precedes body reads` @@ -153,8 +145,7 @@ export function defineV2BodyLifecycleRoute< }) const lifecycleContext = { request, principal, parsed: parsed.data, admission } const body = await options.readBody(lifecycleContext) - const transfer = await options.transfer({ ...lifecycleContext, body }) - const inputContext = { ...lifecycleContext, body, transfer } + const inputContext = { ...lifecycleContext, body } const input = options.mapInput(inputContext) const result = await options.useCase.execute({ principal, input, request }) const responseBody = options.contract.response.schema.parse(await options.present(result)) diff --git a/apps/sim/lib/knowledge/api/route-policies.ts b/apps/sim/lib/knowledge/api/route-policies.ts index 22628d413fe..b9ae9031575 100644 --- a/apps/sim/lib/knowledge/api/route-policies.ts +++ b/apps/sim/lib/knowledge/api/route-policies.ts @@ -2,6 +2,7 @@ import { createInternalResourceConcealmentPolicy, createInternalSessionOrExecutorAuth, createV2ResourceConcealmentPolicy, + extendInternalErrorPolicy, type InternalErrorPolicy, internalErrorResponse, internalOrchestrationErrorPolicy, @@ -11,6 +12,7 @@ import { import { isPayloadSizeLimitError } from '@/lib/core/utils/stream-limits' import { KNOWLEDGE_DELEGATION_AUDIENCE } from '@/lib/knowledge/application/authorization' import { KnowledgeUsageLimitExceededError } from '@/lib/knowledge/application/billing' +import { KnowledgeDocumentNotReadyError } from '@/lib/knowledge/application/chunk-errors' import { KnowledgeSearchProvenanceUnavailableError } from '@/lib/knowledge/application/search' import { KnowledgeDocumentUnsupportedMediaTypeError } from '@/lib/knowledge/application/upload-sessions' import { v2Error } from '@/app/api/v2/lib/response' @@ -82,6 +84,19 @@ export const internalKnowledgeErrorPolicies = { chunks: concealKnowledgeBase( internalKnowledgeErrorPolicy('Failed to process knowledge chunk request') ), + chunkList: concealKnowledgeBase( + extendInternalErrorPolicy( + internalKnowledgeErrorPolicy('Failed to process knowledge chunk request'), + (error) => + error instanceof KnowledgeDocumentNotReadyError + ? internalErrorResponse(400, { + error: 'Document is not ready for access', + details: `Document status: ${error.processingStatus}`, + retryAfter: error.processingStatus === 'processing' ? 5 : null, + }) + : null + ) + ), upsert: concealKnowledgeBase(internalKnowledgeUploadErrorPolicy), search: concealKnowledgeBase(internalKnowledgeSearchErrorPolicy), tags: concealKnowledgeBase( diff --git a/apps/sim/lib/knowledge/application/chunk-errors.ts b/apps/sim/lib/knowledge/application/chunk-errors.ts new file mode 100644 index 00000000000..a56c22bb485 --- /dev/null +++ b/apps/sim/lib/knowledge/application/chunk-errors.ts @@ -0,0 +1,8 @@ +import { OrchestrationError } from '@/lib/core/orchestration/types' + +export class KnowledgeDocumentNotReadyError extends OrchestrationError { + constructor(readonly processingStatus: string) { + super('validation', `Document is not ready for access (status: ${processingStatus})`) + this.name = 'KnowledgeDocumentNotReadyError' + } +} diff --git a/apps/sim/lib/knowledge/application/chunks.test.ts b/apps/sim/lib/knowledge/application/chunks.test.ts new file mode 100644 index 00000000000..74b93781ad3 --- /dev/null +++ b/apps/sim/lib/knowledge/application/chunks.test.ts @@ -0,0 +1,79 @@ +/** + * @vitest-environment node + */ + +import { beforeEach, describe, expect, it, vi } from 'vitest' + +const mocks = vi.hoisted(() => ({ + resolveDocument: vi.fn(), + resolvePermission: vi.fn(), + queryChunks: vi.fn(), +})) + +vi.mock('@sim/platform-authz/workspace', () => ({ + permissionSatisfies: (actual: string | null, required: string) => { + const rank = { read: 1, write: 2, admin: 3 } as const + return ( + actual !== null && rank[actual as keyof typeof rank] >= rank[required as keyof typeof rank] + ) + }, + resolveEffectiveWorkspacePermission: mocks.resolvePermission, +})) + +vi.mock('@/lib/knowledge/application/contexts', () => ({ + resolveCanonicalActiveKnowledgeDocumentContext: mocks.resolveDocument, + resolveActiveKnowledgeChunkContext: vi.fn(), +})) + +vi.mock('@/lib/knowledge/chunks/service', () => ({ + batchChunkOperation: vi.fn(), + createChunk: vi.fn(), + deleteChunk: vi.fn(), + queryChunks: mocks.queryChunks, + updateChunk: vi.fn(), +})) + +vi.mock('@/lib/execution/durable-secret-provenance', () => ({ + createDurableSecretProvenanceRegistry: vi.fn(), +})) + +vi.mock('@/lib/knowledge/model-input-provenance', () => ({ + runWithKnowledgeModelInputProvenance: vi.fn(), +})) + +vi.mock('@/providers/utils', () => ({ calculateCost: vi.fn() })) + +import { KnowledgeDocumentNotReadyError } from '@/lib/knowledge/application/chunk-errors' +import { listKnowledgeChunks } from '@/lib/knowledge/application/chunks' + +describe('knowledge chunk application use cases', () => { + beforeEach(() => { + vi.clearAllMocks() + mocks.resolvePermission.mockResolvedValue('read') + mocks.resolveDocument.mockResolvedValue({ + workspaceId: 'workspace-1', + workspaceOrganizationId: null, + allowPersonalApiKeys: true, + billedAccountUserId: 'billing-owner-1', + knowledgeBaseId: 'knowledge-1', + knowledgeBase: { id: 'knowledge-1' }, + documentId: 'document-1', + document: { id: 'document-1', processingStatus: 'processing' }, + }) + }) + + it('returns a typed transient failure before querying chunks for a processing document', async () => { + const promise = listKnowledgeChunks.execute({ + principal: { kind: 'session', userId: 'user-1', sessionId: 'session-1' }, + input: { knowledgeBaseId: 'knowledge-1', documentId: 'document-1' }, + }) + + await expect(promise).rejects.toBeInstanceOf(KnowledgeDocumentNotReadyError) + await expect(promise).rejects.toMatchObject({ + code: 'validation', + processingStatus: 'processing', + message: 'Document is not ready for access (status: processing)', + }) + expect(mocks.queryChunks).not.toHaveBeenCalled() + }) +}) diff --git a/apps/sim/lib/knowledge/application/chunks.ts b/apps/sim/lib/knowledge/application/chunks.ts index 519ac5b1198..3f81a5d5683 100644 --- a/apps/sim/lib/knowledge/application/chunks.ts +++ b/apps/sim/lib/knowledge/application/chunks.ts @@ -7,6 +7,7 @@ import { } from '@/lib/execution/durable-secret-provenance' import { defineAuthorizedKnowledgeUseCase } from '@/lib/knowledge/application/authorized-knowledge-use-case' import { resolveKnowledgeAttributedUserId } from '@/lib/knowledge/application/billing' +import { KnowledgeDocumentNotReadyError } from '@/lib/knowledge/application/chunk-errors' import { type ActiveKnowledgeDocumentContext, resolveActiveKnowledgeChunkContext, @@ -62,10 +63,7 @@ export interface BulkKnowledgeChunksInput extends KnowledgeDocumentChunkInput { function requireChunkReadable(context: ActiveKnowledgeDocumentContext): void { if (context.document.processingStatus !== 'completed') { - throw new OrchestrationError( - 'validation', - `Document is not ready for access (status: ${context.document.processingStatus})` - ) + throw new KnowledgeDocumentNotReadyError(context.document.processingStatus) } } diff --git a/apps/sim/lib/knowledge/application/documents.test.ts b/apps/sim/lib/knowledge/application/documents.test.ts index 21fd1734e69..8df17daf2ee 100644 --- a/apps/sim/lib/knowledge/application/documents.test.ts +++ b/apps/sim/lib/knowledge/application/documents.test.ts @@ -22,6 +22,9 @@ const mocks = vi.hoisted(() => ({ performBulkUpload: vi.fn(), markTimedOut: vi.fn(), retryProcessing: vi.fn(), + uploadStoredFile: vi.fn(), + generateKnowledgeBaseFileKey: vi.fn(), + recordKnowledgeBaseFileOwnership: vi.fn(), recordAudit: vi.fn(), captureServerEvent: vi.fn(), })) @@ -75,6 +78,18 @@ vi.mock('@/lib/knowledge/orchestration/documents', () => ({ performRetryKnowledgeDocumentProcessing: mocks.retryProcessing, })) +vi.mock('@/lib/uploads', () => ({ + StorageService: { uploadFile: mocks.uploadStoredFile }, +})) + +vi.mock('@/lib/uploads/contexts/knowledge-base/knowledge-base-file-manager', () => ({ + generateKnowledgeBaseFileKey: mocks.generateKnowledgeBaseFileKey, +})) + +vi.mock('@/lib/uploads/server/metadata', () => ({ + recordKnowledgeBaseFileOwnership: mocks.recordKnowledgeBaseFileOwnership, +})) + vi.mock('@/lib/posthog/server', () => ({ captureServerEvent: mocks.captureServerEvent })) import { OrchestrationError } from '@/lib/core/orchestration/types' @@ -107,6 +122,20 @@ const document = { uploadedAt: new Date('2026-01-01T00:00:00Z'), } +const uploadFile = { + buffer: Buffer.alloc(42, 'a'), + filename: 'guide.pdf', + fileSize: 42, + mimeType: 'application/pdf', +} + +const storedDocumentInput = { + filename: uploadFile.filename, + fileUrl: '/api/files/serve/kb%2Fupload-1?context=knowledge-base', + fileSize: uploadFile.fileSize, + mimeType: uploadFile.mimeType, +} + describe('knowledge document application use cases', () => { beforeEach(() => { vi.clearAllMocks() @@ -131,6 +160,15 @@ describe('knowledge document application use cases', () => { workspaceId: 'workspace-1', }) mocks.checkUsage.mockResolvedValue({ isExceeded: false }) + mocks.generateKnowledgeBaseFileKey.mockReturnValue('kb/upload-1') + mocks.uploadStoredFile.mockResolvedValue({ + path: '/api/files/serve/kb%2Fupload-1', + key: 'kb/upload-1', + name: uploadFile.filename, + size: uploadFile.fileSize, + type: uploadFile.mimeType, + }) + mocks.recordKnowledgeBaseFileOwnership.mockResolvedValue(undefined) mocks.createDocument.mockResolvedValue(document) mocks.updateDocument.mockResolvedValue(document) mocks.processQueue.mockResolvedValue(undefined) @@ -254,20 +292,52 @@ describe('knowledge document application use cases', () => { input: { knowledgeBaseId: 'knowledge-1', assertedWorkspaceId: 'workspace-1', - document, + file: uploadFile, source: 'v2', }, }) expect(mocks.resolveSystemBilling).toHaveBeenCalledWith('workspace-1') + expect(mocks.recordKnowledgeBaseFileOwnership).toHaveBeenCalledWith({ + key: 'kb/upload-1', + userId: 'billing-owner-1', + workspaceId: 'workspace-1', + originalName: 'guide.pdf', + contentType: 'application/pdf', + size: 42, + }) + expect(mocks.uploadStoredFile).toHaveBeenCalledWith({ + file: uploadFile.buffer, + fileName: uploadFile.filename, + contentType: uploadFile.mimeType, + context: 'knowledge-base', + customKey: 'kb/upload-1', + preserveKey: true, + persistMetadata: false, + }) + expect(mocks.recordKnowledgeBaseFileOwnership.mock.invocationCallOrder[0]).toBeLessThan( + mocks.uploadStoredFile.mock.invocationCallOrder[0] + ) + expect(mocks.uploadStoredFile.mock.invocationCallOrder[0]).toBeLessThan( + mocks.createDocument.mock.invocationCallOrder[0] + ) expect(mocks.createDocument).toHaveBeenCalledWith( - document, + storedDocumentInput, 'knowledge-1', expect.any(String), 'billing-owner-1', undefined, undefined, - { expectedWorkspaceId: 'workspace-1' } + { + expectedWorkspaceId: 'workspace-1', + processing: { + processingOptions: {}, + billingAttribution: { + actorUserId: 'billing-owner-1', + workspaceId: 'workspace-1', + }, + }, + } ) expect(mocks.recordAudit).toHaveBeenCalledWith( expect.objectContaining({ @@ -291,7 +361,7 @@ describe('knowledge document application use cases', () => { input: { knowledgeBaseId: 'knowledge-1', assertedWorkspaceId: 'workspace-1', - document, + file: uploadFile, usageAdmission: 'pre_admitted', }, }) @@ -300,6 +370,29 @@ describe('knowledge document application use cases', () => { expect(mocks.createDocument).toHaveBeenCalledOnce() }) + it('leaves only a sweepable knowledge-base binding when final authorization fails', async () => { + mocks.resolvePermission.mockResolvedValueOnce('write').mockResolvedValueOnce(null) + + await expect( + uploadKnowledgeDocument.execute({ + principal: { kind: 'session', userId: 'user-1', sessionId: 'session-1' }, + input: { + knowledgeBaseId: 'knowledge-1', + assertedWorkspaceId: 'workspace-1', + file: uploadFile, + usageAdmission: 'pre_admitted', + }, + }) + ).rejects.toMatchObject({ code: 'forbidden' }) + + expect(mocks.recordKnowledgeBaseFileOwnership).toHaveBeenCalledOnce() + expect(mocks.uploadStoredFile).toHaveBeenCalledWith( + expect.objectContaining({ context: 'knowledge-base', persistMetadata: false }) + ) + expect(mocks.createDocument).not.toHaveBeenCalled() + expect(mocks.recordAudit).not.toHaveBeenCalled() + }) + it('conceals a cross-knowledge-base document before deletion and audit', async () => { mocks.resolveDocument.mockRejectedValueOnce( new OrchestrationError('not_found', 'Document not found') @@ -441,7 +534,7 @@ describe('knowledge document application use cases', () => { input: { knowledgeBaseId: 'knowledge-1', assertedWorkspaceId: 'workspace-1', - document, + file: uploadFile, }, }) ).rejects.toBe(failure) @@ -449,6 +542,34 @@ describe('knowledge document application use cases', () => { expect(mocks.recordAudit).not.toHaveBeenCalled() }) + it('commits a durable processing intent instead of dispatching from the request', async () => { + await uploadKnowledgeDocument.execute({ + principal: { kind: 'session', userId: 'user-1', sessionId: 'session-1' }, + input: { + knowledgeBaseId: 'knowledge-1', + assertedWorkspaceId: 'workspace-1', + file: uploadFile, + processingOptions: { recipe: 'default', lang: 'en' }, + }, + }) + + expect(mocks.createDocument).toHaveBeenCalledWith( + storedDocumentInput, + 'knowledge-1', + expect.any(String), + 'user-1', + undefined, + undefined, + expect.objectContaining({ + processing: { + processingOptions: { recipe: 'default', lang: 'en' }, + billingAttribution: { actorUserId: 'user-1', workspaceId: 'workspace-1' }, + }, + }) + ) + expect(mocks.processQueue).not.toHaveBeenCalled() + }) + it('bounds bulk document creation before billing or orchestration', async () => { await expect( createKnowledgeDocuments.execute({ diff --git a/apps/sim/lib/knowledge/application/documents.ts b/apps/sim/lib/knowledge/application/documents.ts index 9be50d5c4d6..30b9445fbd3 100644 --- a/apps/sim/lib/knowledge/application/documents.ts +++ b/apps/sim/lib/knowledge/application/documents.ts @@ -39,7 +39,6 @@ import { bulkDocumentOperationByFilter, createDocumentRecords, createSingleDocument, - type DocumentData, deleteDocument, deleteKnowledgeDocumentInKnowledgeBase, getDocuments, @@ -57,6 +56,9 @@ import { performUploadKnowledgeDocuments, } from '@/lib/knowledge/orchestration/documents' import type { KnowledgeDocumentWriteSecretProvenance } from '@/lib/knowledge/secret-provenance' +import { StorageService } from '@/lib/uploads' +import { generateKnowledgeBaseFileKey } from '@/lib/uploads/contexts/knowledge-base/knowledge-base-file-manager' +import { recordKnowledgeBaseFileOwnership } from '@/lib/uploads/server/metadata' import { MAX_KNOWLEDGE_DOCUMENT_FILE_SIZE } from '@/lib/uploads/shared/types' import { validateFileType } from '@/lib/uploads/utils/validation' @@ -101,7 +103,12 @@ export interface KnowledgeDocumentInput { } export interface UploadKnowledgeDocumentInput extends UploadKnowledgeDocumentAdmissionInput { - document: KnowledgeDocumentInput + file: { + buffer: Buffer + filename: string + fileSize: number + mimeType: string + } processingOptions?: ProcessingOptions startProcessing?: boolean /** Code-defined admission state; HTTP/model payloads must never populate it. */ @@ -241,7 +248,6 @@ export const admitKnowledgeDocumentUpload = defineAuthorizedKnowledgeUseCase({ knowledgeBaseId: context.knowledgeBaseId, knowledgeBaseName: context.knowledgeBase.name, workspaceId: context.workspaceId, - storageActorUserId: resolveKnowledgeAttributedUserId(principal, context), } }, }) @@ -251,16 +257,19 @@ export const uploadKnowledgeDocument = defineAuthorizedKnowledgeUseCase({ resolveContext: ({ input }: { input: UploadKnowledgeDocumentInput }) => resolveActiveKnowledgeBaseContext(input), async execute({ principal, input, context }) { - if (input.document.fileSize < 0 || input.document.fileSize > MAX_KNOWLEDGE_DOCUMENT_FILE_SIZE) { + if (input.file.fileSize < 0 || input.file.fileSize > MAX_KNOWLEDGE_DOCUMENT_FILE_SIZE) { throw new OrchestrationError( 'payload_too_large', 'Knowledge document exceeds the 100MB limit' ) } - const fileTypeError = validateFileType(input.document.filename, input.document.mimeType) + if (input.file.fileSize !== input.file.buffer.byteLength) { + throw new Error('Knowledge document upload size does not match its buffered bytes') + } + const fileTypeError = validateFileType(input.file.filename, input.file.mimeType) if (fileTypeError) throw new OrchestrationError('validation', fileTypeError.message) - const billingAttribution = await resolveKnowledgeBillingAttribution(principal, context) if (input.usageAdmission !== 'pre_admitted') { + const billingAttribution = await resolveKnowledgeBillingAttribution(principal, context) const usage = await checkAttributedUsageLimits(billingAttribution) if (usage.isExceeded) { throw new KnowledgeUsageLimitExceededError( @@ -269,38 +278,71 @@ export const uploadKnowledgeDocument = defineAuthorizedKnowledgeUseCase({ } } const requestId = generateRequestId() - const uploadedBy = resolveKnowledgeAttributedUserId(principal, context) + const storageActorUserId = resolveKnowledgeAttributedUserId(principal, context) + const storageKey = generateKnowledgeBaseFileKey(input.file.filename) + await recordKnowledgeBaseFileOwnership({ + key: storageKey, + userId: storageActorUserId, + workspaceId: context.workspaceId, + originalName: input.file.filename, + contentType: input.file.mimeType, + size: input.file.fileSize, + }) + const storedFile = await StorageService.uploadFile({ + file: input.file.buffer, + fileName: input.file.filename, + contentType: input.file.mimeType, + context: 'knowledge-base', + customKey: storageKey, + preserveKey: true, + persistMetadata: false, + }) + if (storedFile.key !== storageKey || storedFile.size !== input.file.fileSize) { + throw new Error('Knowledge document storage did not preserve the admitted file identity') + } + if (storedFile.path.includes('?')) { + throw new Error('Knowledge document storage returned a path with an unexpected query') + } + + const registrationContext = await resolveActiveKnowledgeBaseContext(input) + await authorizeWorkspaceOperation( + principal, + knowledgeOperations.uploadDocument, + registrationContext, + { + delegation: knowledgeDelegationPolicy, + } + ) + const billingAttribution = await resolveKnowledgeBillingAttribution( + principal, + registrationContext + ) + const uploadedBy = resolveKnowledgeAttributedUserId(principal, registrationContext) + const documentInput: KnowledgeDocumentInput = { + filename: input.file.filename, + fileUrl: `${storedFile.path}?context=knowledge-base`, + fileSize: input.file.fileSize, + mimeType: input.file.mimeType, + } const document = await createSingleDocument( - input.document, - context.knowledgeBaseId, + documentInput, + registrationContext.knowledgeBaseId, requestId, uploadedBy, undefined, undefined, - { expectedWorkspaceId: context.workspaceId } - ) - if (input.startProcessing !== false) { - const processingDocument: DocumentData = { - documentId: document.id, - filename: document.filename, - fileUrl: document.fileUrl, - fileSize: document.fileSize, - mimeType: document.mimeType, + { + expectedWorkspaceId: registrationContext.workspaceId, + ...(input.startProcessing !== false + ? { + processing: { + processingOptions: input.processingOptions ?? {}, + billingAttribution, + }, + } + : {}), } - processDocumentsWithQueue( - [processingDocument], - context.knowledgeBaseId, - input.processingOptions ?? {}, - requestId, - billingAttribution - ).catch((error: unknown) => { - logger.error('Knowledge document processing pipeline failed', { - knowledgeBaseId: context.knowledgeBaseId, - documentId: document.id, - error, - }) - }) - } + ) return { document, created: true as const } }, projectAudit: ({ input, context, result }) => ({ 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 4f5379d421f..1a2f7402632 100644 --- a/apps/sim/lib/knowledge/documents/document-processing-source.test.ts +++ b/apps/sim/lib/knowledge/documents/document-processing-source.test.ts @@ -228,4 +228,29 @@ describe('knowledge document processing source', () => { expect(mockProcessDocument).not.toHaveBeenCalled() expect(mockGenerateEmbeddings).not.toHaveBeenCalled() }) + + it('takes over an existing processing attempt', async () => { + dbChainMockFns.limit + .mockReset() + .mockResolvedValueOnce([{ ...PERSISTED_CONTEXT, processingStatus: 'processing' }]) + .mockResolvedValueOnce([PERSISTED_PROVENANCE_ROW]) + .mockResolvedValueOnce([{ id: 'document-1' }]) + dbChainMockFns.returning.mockReset().mockResolvedValueOnce([]) + + await processDocumentAsync('knowledge-base-1', 'document-1', { + filename: 'stale.pdf', + fileUrl: 'https://example.com/stale.pdf', + fileSize: 1, + mimeType: 'text/plain', + }) + + expect(mockProcessDocument).toHaveBeenCalled() + expect(mockGenerateEmbeddings).not.toHaveBeenCalled() + expect(dbChainMockFns.set).toHaveBeenCalledWith( + expect.objectContaining({ + processingStatus: 'processing', + processingStartedAt: expect.any(Date), + }) + ) + }) }) diff --git a/apps/sim/lib/knowledge/documents/processing-claim.test.ts b/apps/sim/lib/knowledge/documents/processing-claim.test.ts new file mode 100644 index 00000000000..6d898f60c40 --- /dev/null +++ b/apps/sim/lib/knowledge/documents/processing-claim.test.ts @@ -0,0 +1,129 @@ +/** + * @vitest-environment node + */ + +import { dbChainMockFns, resetDbChainMock } from '@sim/testing' +import { beforeEach, describe, expect, it, vi } from 'vitest' +import { + failStaleDocumentProcessingClaim, + KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS, + reclaimStaleDocumentProcessingClaim, +} from '@/lib/knowledge/documents/processing-claim' + +const NOW = new Date('2026-08-11T12:00:00.000Z') + +describe('reclaimStaleDocumentProcessingClaim', () => { + beforeEach(() => { + vi.clearAllMocks() + resetDbChainMock() + }) + + it('leaves an active processing claim untouched', async () => { + const reclaimed = await reclaimStaleDocumentProcessingClaim({ + knowledgeBaseId: 'knowledge-base-1', + documentId: 'document-1', + processingStartedAt: new Date( + NOW.getTime() - KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS + ), + now: NOW, + }) + + expect(reclaimed).toBe(false) + expect(dbChainMockFns.set).not.toHaveBeenCalled() + }) + + it.each([null, new Date(NOW.getTime() - KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS - 1)])( + 'reopens an abandoned processing claim started at %s', + async (processingStartedAt) => { + dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'document-1' }]) + + const reclaimed = await reclaimStaleDocumentProcessingClaim({ + knowledgeBaseId: 'knowledge-base-1', + documentId: 'document-1', + processingStartedAt, + now: NOW, + }) + + expect(reclaimed).toBe(true) + expect(dbChainMockFns.set).toHaveBeenCalledWith({ + processingStatus: 'pending', + processingStartedAt: null, + processingCompletedAt: null, + processingError: null, + }) + } + ) + + it('does not report success when the original claim changed before the update', async () => { + dbChainMockFns.returning.mockResolvedValueOnce([]) + + const reclaimed = await reclaimStaleDocumentProcessingClaim({ + knowledgeBaseId: 'knowledge-base-1', + documentId: 'document-1', + processingStartedAt: new Date( + NOW.getTime() - KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS - 1 + ), + now: NOW, + }) + + expect(reclaimed).toBe(false) + }) +}) + +describe('failStaleDocumentProcessingClaim', () => { + beforeEach(() => { + vi.clearAllMocks() + resetDbChainMock() + }) + + it('rejects an active processing claim', async () => { + await expect( + failStaleDocumentProcessingClaim({ + documentId: 'document-1', + processingStartedAt: new Date( + NOW.getTime() - KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS + ), + now: NOW, + }) + ).rejects.toThrow('Document has not been processing long enough to be considered dead') + + expect(dbChainMockFns.set).not.toHaveBeenCalled() + }) + + it('fails the exact abandoned processing claim', async () => { + const processingStartedAt = new Date( + NOW.getTime() - KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS - 1 + ) + dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'document-1' }]) + + const result = await failStaleDocumentProcessingClaim({ + documentId: 'document-1', + processingStartedAt, + now: NOW, + }) + + expect(result).toEqual({ + success: true, + processingDuration: KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS + 1, + }) + expect(dbChainMockFns.set).toHaveBeenCalledWith({ + processingStatus: 'failed', + processingError: 'Processing timed out. Please retry or re-sync the connector.', + processingCompletedAt: NOW, + }) + }) + + it('does not fail a replacement processing claim', async () => { + dbChainMockFns.returning.mockResolvedValueOnce([]) + + const result = await failStaleDocumentProcessingClaim({ + documentId: 'document-1', + processingStartedAt: new Date( + NOW.getTime() - KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS - 1 + ), + now: NOW, + }) + + expect(result.success).toBe(false) + }) +}) diff --git a/apps/sim/lib/knowledge/documents/processing-claim.ts b/apps/sim/lib/knowledge/documents/processing-claim.ts new file mode 100644 index 00000000000..bfb4a2208c4 --- /dev/null +++ b/apps/sim/lib/knowledge/documents/processing-claim.ts @@ -0,0 +1,96 @@ +import { db } from '@sim/db' +import { document } from '@sim/db/schema' +import { and, eq, isNull } from 'drizzle-orm' + +export const KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS = 10 * 60 * 1000 + +interface ReclaimStaleDocumentProcessingClaimParams { + knowledgeBaseId: string + documentId: string + processingStartedAt: Date | null + now?: Date +} + +interface FailStaleDocumentProcessingClaimParams { + documentId: string + processingStartedAt: Date + now?: Date +} + +/** + * Reopens an abandoned processing attempt using its start time as a compare-and-set token. + * The former worker's timestamp-guarded writes then cannot commit after the claim is reclaimed. + */ +export async function reclaimStaleDocumentProcessingClaim({ + knowledgeBaseId, + documentId, + processingStartedAt, + now = new Date(), +}: ReclaimStaleDocumentProcessingClaimParams): Promise { + if ( + processingStartedAt && + now.getTime() - processingStartedAt.getTime() <= + KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS + ) { + return false + } + + const processingStartedAtGuard = processingStartedAt + ? eq(document.processingStartedAt, processingStartedAt) + : isNull(document.processingStartedAt) + const [reclaimed] = await db + .update(document) + .set({ + processingStatus: 'pending', + processingStartedAt: null, + processingCompletedAt: null, + processingError: null, + }) + .where( + and( + eq(document.id, documentId), + eq(document.knowledgeBaseId, knowledgeBaseId), + eq(document.processingStatus, 'processing'), + processingStartedAtGuard, + eq(document.userExcluded, false), + isNull(document.archivedAt), + isNull(document.deletedAt) + ) + ) + .returning({ id: document.id }) + + return Boolean(reclaimed) +} + +/** Marks only the abandoned processing attempt identified by its start-time token as failed. */ +export async function failStaleDocumentProcessingClaim({ + documentId, + processingStartedAt, + now = new Date(), +}: FailStaleDocumentProcessingClaimParams): Promise<{ + success: boolean + processingDuration: number +}> { + const processingDuration = now.getTime() - processingStartedAt.getTime() + if (processingDuration <= KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS) { + throw new Error('Document has not been processing long enough to be considered dead') + } + + const [failed] = await db + .update(document) + .set({ + processingStatus: 'failed', + processingError: 'Processing timed out. Please retry or re-sync the connector.', + processingCompletedAt: now, + }) + .where( + and( + eq(document.id, documentId), + eq(document.processingStatus, 'processing'), + eq(document.processingStartedAt, processingStartedAt) + ) + ) + .returning({ id: document.id }) + + return { success: Boolean(failed), processingDuration } +} diff --git a/apps/sim/lib/knowledge/documents/processing-outbox-event.ts b/apps/sim/lib/knowledge/documents/processing-outbox-event.ts new file mode 100644 index 00000000000..ae893680044 --- /dev/null +++ b/apps/sim/lib/knowledge/documents/processing-outbox-event.ts @@ -0,0 +1,21 @@ +import type { db } from '@sim/db' +import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' +import { enqueueOutboxEvent } from '@/lib/core/outbox/service' +import type { ProcessingOptions } from '@/lib/knowledge/documents/service' + +export const KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT = 'knowledge.document.processing.dispatch' + +export interface KnowledgeDocumentProcessingOutboxPayload { + knowledgeBaseId: string + documentId: string + processingOptions: ProcessingOptions + billingAttribution: BillingAttributionSnapshot +} + +/** Enqueues durable processing in the same transaction that creates the document. */ +export function enqueueKnowledgeDocumentProcessing( + executor: Pick, + payload: KnowledgeDocumentProcessingOutboxPayload +): Promise { + return enqueueOutboxEvent(executor, KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT, payload) +} diff --git a/apps/sim/lib/knowledge/documents/processing-outbox-handler.test.ts b/apps/sim/lib/knowledge/documents/processing-outbox-handler.test.ts new file mode 100644 index 00000000000..9e095e9243a --- /dev/null +++ b/apps/sim/lib/knowledge/documents/processing-outbox-handler.test.ts @@ -0,0 +1,181 @@ +/** + * @vitest-environment node + */ + +import { beforeEach, describe, expect, it, vi } from 'vitest' + +const mocks = vi.hoisted(() => ({ + getKnowledgeDocument: vi.fn(), + processDocumentsWithQueue: vi.fn(), + reclaimStaleDocumentProcessingClaim: vi.fn(), +})) + +vi.mock('@/lib/knowledge/documents/service', () => ({ + getKnowledgeDocument: mocks.getKnowledgeDocument, + processDocumentsWithQueue: mocks.processDocumentsWithQueue, +})) + +vi.mock('@/lib/knowledge/documents/processing-claim', () => ({ + reclaimStaleDocumentProcessingClaim: mocks.reclaimStaleDocumentProcessingClaim, +})) + +import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' +import type { OutboxEventContext } from '@/lib/core/outbox/service' +import { KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT } from '@/lib/knowledge/documents/processing-outbox-event' +import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler' + +const BILLING_ATTRIBUTION = { + actorUserId: 'user-1', + workspaceId: 'workspace-1', + organizationId: null, + billedAccountUserId: 'owner-1', + billingEntity: { type: 'user', id: 'owner-1' }, + billingPeriod: { + start: '2026-08-01T00:00:00.000Z', + end: '2026-09-01T00:00:00.000Z', + }, + payerSubscription: null, +} satisfies BillingAttributionSnapshot + +const DOCUMENT = { + id: 'document-1', + filename: 'guide.pdf', + fileUrl: '/api/files/serve/kb%2Fguide.pdf?context=knowledge-base', + fileSize: 128, + mimeType: 'application/pdf', + processingStatus: 'pending', +} + +const PAYLOAD = { + knowledgeBaseId: 'knowledge-base-1', + documentId: 'document-1', + processingOptions: { recipe: 'default', lang: 'en' }, + billingAttribution: BILLING_ATTRIBUTION, +} + +function createContext(eventId = 'outbox-event-1'): OutboxEventContext { + return { + eventId, + eventType: KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT, + attempts: 0, + maxAttempts: 10, + signal: new AbortController().signal, + checkpointPayload: vi.fn(), + } +} + +function handler() { + const value = + knowledgeDocumentProcessingOutboxHandlers[KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT] + if (!value) throw new Error('Knowledge processing outbox handler is not registered') + return value +} + +describe('knowledge document processing outbox handler', () => { + beforeEach(() => { + vi.clearAllMocks() + mocks.getKnowledgeDocument.mockResolvedValue(DOCUMENT) + mocks.processDocumentsWithQueue.mockResolvedValue(undefined) + mocks.reclaimStaleDocumentProcessingClaim.mockResolvedValue(false) + }) + + it('dispatches the authoritative document with the stable outbox event id', async () => { + await handler()(PAYLOAD, createContext('outbox-event-stable')) + + expect(mocks.getKnowledgeDocument).toHaveBeenCalledWith('knowledge-base-1', 'document-1') + expect(mocks.processDocumentsWithQueue).toHaveBeenCalledWith( + [ + { + documentId: 'document-1', + filename: 'guide.pdf', + fileUrl: '/api/files/serve/kb%2Fguide.pdf?context=knowledge-base', + fileSize: 128, + mimeType: 'application/pdf', + }, + ], + 'knowledge-base-1', + { recipe: 'default', lang: 'en' }, + 'outbox-event-stable', + BILLING_ATTRIBUTION + ) + }) + + it.each([null, { ...DOCUMENT, processingStatus: 'completed' }])( + 'completes without redispatch when the document is absent or completed', + async (document) => { + mocks.getKnowledgeDocument.mockResolvedValueOnce(document) + + await handler()(PAYLOAD, createContext()) + + expect(mocks.processDocumentsWithQueue).not.toHaveBeenCalled() + } + ) + + it('keeps the event retryable while an earlier processing attempt is active', async () => { + const processingStartedAt = new Date() + mocks.getKnowledgeDocument.mockResolvedValueOnce({ + ...DOCUMENT, + processingStatus: 'processing', + processingStartedAt, + }) + + await expect(handler()(PAYLOAD, createContext())).rejects.toThrow( + 'Knowledge document document-1 is already being processed' + ) + + expect(mocks.reclaimStaleDocumentProcessingClaim).toHaveBeenCalledWith({ + knowledgeBaseId: 'knowledge-base-1', + documentId: 'document-1', + processingStartedAt, + }) + expect(mocks.processDocumentsWithQueue).not.toHaveBeenCalled() + }) + + it('reclaims and redispatches an abandoned processing attempt', async () => { + const processingStartedAt = new Date('2026-08-11T11:00:00.000Z') + mocks.getKnowledgeDocument.mockResolvedValueOnce({ + ...DOCUMENT, + processingStatus: 'processing', + processingStartedAt, + }) + mocks.reclaimStaleDocumentProcessingClaim.mockResolvedValueOnce(true) + + await handler()(PAYLOAD, createContext('outbox-event-retry')) + + expect(mocks.reclaimStaleDocumentProcessingClaim).toHaveBeenCalledWith({ + knowledgeBaseId: 'knowledge-base-1', + documentId: 'document-1', + processingStartedAt, + }) + expect(mocks.processDocumentsWithQueue).toHaveBeenCalledWith( + [ + { + documentId: 'document-1', + filename: 'guide.pdf', + fileUrl: '/api/files/serve/kb%2Fguide.pdf?context=knowledge-base', + fileSize: 128, + mimeType: 'application/pdf', + }, + ], + 'knowledge-base-1', + { recipe: 'default', lang: 'en' }, + 'outbox-event-retry', + BILLING_ATTRIBUTION + ) + }) + + it('propagates dispatch failures so the outbox schedules a retry', async () => { + const failure = new Error('queue unavailable') + mocks.processDocumentsWithQueue.mockRejectedValueOnce(failure) + + await expect(handler()(PAYLOAD, createContext())).rejects.toBe(failure) + }) + + it('fails fast on malformed durable processing options', async () => { + await expect( + handler()({ ...PAYLOAD, processingOptions: { unsupported: true } }, createContext()) + ).rejects.toThrow('unsupported processing options') + + expect(mocks.getKnowledgeDocument).not.toHaveBeenCalled() + }) +}) diff --git a/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts b/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts new file mode 100644 index 00000000000..e148de05d12 --- /dev/null +++ b/apps/sim/lib/knowledge/documents/processing-outbox-handler.ts @@ -0,0 +1,94 @@ +import { assertBillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' +import type { OutboxHandler, OutboxHandlerRegistry } from '@/lib/core/outbox/service' +import { reclaimStaleDocumentProcessingClaim } from '@/lib/knowledge/documents/processing-claim' +import { + KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT, + type KnowledgeDocumentProcessingOutboxPayload, +} from '@/lib/knowledge/documents/processing-outbox-event' +import { + getKnowledgeDocument, + type ProcessingOptions, + processDocumentsWithQueue, +} from '@/lib/knowledge/documents/service' + +function requirePayloadRecord(payload: unknown): Record { + if (!payload || typeof payload !== 'object' || Array.isArray(payload)) { + throw new Error('Knowledge document processing outbox payload must be an object') + } + return payload as Record +} + +function requireNonEmptyString(value: unknown, field: string): string { + if (typeof value !== 'string' || value.trim().length === 0) { + throw new Error(`Knowledge document processing outbox payload is missing ${field}`) + } + return value +} + +function parseProcessingOptions(value: unknown): ProcessingOptions { + const record = requirePayloadRecord(value) + const unsupported = Object.keys(record).filter((key) => key !== 'recipe' && key !== 'lang') + if (unsupported.length > 0) { + throw new Error( + `Knowledge document processing outbox payload has unsupported processing options: ${unsupported.join(', ')}` + ) + } + if (record.recipe !== undefined && typeof record.recipe !== 'string') { + throw new Error('Knowledge document processing outbox recipe must be a string') + } + if (record.lang !== undefined && typeof record.lang !== 'string') { + throw new Error('Knowledge document processing outbox lang must be a string') + } + return { + ...(record.recipe !== undefined ? { recipe: record.recipe } : {}), + ...(record.lang !== undefined ? { lang: record.lang } : {}), + } +} + +function parsePayload(payload: unknown): KnowledgeDocumentProcessingOutboxPayload { + const record = requirePayloadRecord(payload) + return { + knowledgeBaseId: requireNonEmptyString(record.knowledgeBaseId, 'knowledgeBaseId'), + documentId: requireNonEmptyString(record.documentId, 'documentId'), + processingOptions: parseProcessingOptions(record.processingOptions), + billingAttribution: assertBillingAttributionSnapshot(record.billingAttribution), + } +} + +const processKnowledgeDocument: OutboxHandler = async (rawPayload, context) => { + const payload = parsePayload(rawPayload) + context.signal.throwIfAborted() + const document = await getKnowledgeDocument(payload.knowledgeBaseId, payload.documentId) + if (!document || document.processingStatus === 'completed') return + if (document.processingStatus === 'processing') { + const reclaimed = await reclaimStaleDocumentProcessingClaim({ + knowledgeBaseId: payload.knowledgeBaseId, + documentId: document.id, + processingStartedAt: document.processingStartedAt, + }) + if (!reclaimed) { + throw new Error(`Knowledge document ${document.id} is already being processed`) + } + } + + context.signal.throwIfAborted() + await processDocumentsWithQueue( + [ + { + documentId: document.id, + filename: document.filename, + fileUrl: document.fileUrl, + fileSize: document.fileSize, + mimeType: document.mimeType, + }, + ], + payload.knowledgeBaseId, + payload.processingOptions, + context.eventId, + payload.billingAttribution + ) +} + +export const knowledgeDocumentProcessingOutboxHandlers = { + [KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT]: processKnowledgeDocument, +} satisfies OutboxHandlerRegistry diff --git a/apps/sim/lib/knowledge/documents/service.ts b/apps/sim/lib/knowledge/documents/service.ts index 5a896d8c9a2..93125c06eb9 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -64,6 +64,8 @@ import { mergeDurableSecretProvenance, } from '@/lib/execution/durable-secret-provenance' import { processDocument } from '@/lib/knowledge/documents/document-processor' +import { failStaleDocumentProcessingClaim } from '@/lib/knowledge/documents/processing-claim' +import { enqueueKnowledgeDocumentProcessing } from '@/lib/knowledge/documents/processing-outbox-event' import { assertDocumentProcessingBillingContext, createDocumentProcessingPayload, @@ -777,6 +779,7 @@ export async function processDocumentAsync( providedBillingContext?: BillingAttributionSnapshot | DocumentProcessingBillingContext ): Promise { const startTime = Date.now() + const processingStartedAt = new Date() try { logger.info(`[${documentId}] Starting document processing`, { knowledgeBaseId, @@ -859,7 +862,7 @@ export async function processDocumentAsync( .update(document) .set({ processingStatus: 'processing', - processingStartedAt: new Date(), + processingStartedAt, processingCompletedAt: null, processingError: null, }) @@ -934,7 +937,13 @@ export async function processDocumentAsync( usageGate.message ?? 'Usage limit exceeded. Please upgrade your plan to continue.', processingCompletedAt: new Date(), }) - .where(eq(document.id, documentId)) + .where( + and( + eq(document.id, documentId), + eq(document.processingStatus, 'processing'), + eq(document.processingStartedAt, processingStartedAt) + ) + ) return } let billableEmbeddingTokens = 0 @@ -953,6 +962,7 @@ export async function processDocumentAsync( currentSourceFileProvenance ) + let processingCommitted = false await withTimeout( runWithKnowledgeModelInputProvenance( documentSecretContext.registry, @@ -1068,7 +1078,7 @@ export async function processDocumentAsync( updatedAt: now, })) - await db.transaction(async (tx) => { + processingCommitted = await db.transaction(async (tx) => { const activeDocument = await tx .select({ id: document.id }) .from(document) @@ -1076,15 +1086,18 @@ export async function processDocumentAsync( .where( and( eq(document.id, documentId), + eq(document.processingStatus, 'processing'), + eq(document.processingStartedAt, processingStartedAt), isNull(document.archivedAt), isNull(document.deletedAt), isNull(knowledgeBase.deletedAt) ) ) + .for('update', { of: document }) .limit(1) if (activeDocument.length === 0) { - return + return false } if (embeddingRecords.length > 0) { @@ -1130,7 +1143,14 @@ export async function processDocumentAsync( processingCompletedAt: now, processingError: null, }) - .where(eq(document.id, documentId)) + .where( + and( + eq(document.id, documentId), + eq(document.processingStatus, 'processing'), + eq(document.processingStartedAt, processingStartedAt) + ) + ) + return true }) }, { @@ -1143,6 +1163,11 @@ export async function processDocumentAsync( 'Document processing' ) + if (!processingCommitted) { + logger.info(`[${documentId}] Discarded output from an obsolete processing attempt`) + return + } + const processingTime = Date.now() - startTime logger.info(`[${documentId}] Successfully processed document in ${processingTime}ms`) @@ -1204,7 +1229,13 @@ export async function processDocumentAsync( processingError: errorMessage, processingCompletedAt: new Date(), }) - .where(eq(document.id, documentId)) + .where( + and( + eq(document.id, documentId), + eq(document.processingStatus, 'processing'), + eq(document.processingStartedAt, processingStartedAt) + ) + ) throw error } @@ -1860,7 +1891,13 @@ export async function createSingleDocument( uploadedBy: string | null = null, documentId = generateId(), secretProvenance?: KnowledgeDocumentWriteSecretProvenance, - options?: { expectedWorkspaceId?: string } + options?: { + expectedWorkspaceId?: string + processing?: { + processingOptions: ProcessingOptions + billingAttribution: BillingAttributionSnapshot + } + } ): Promise<{ id: string knowledgeBaseId: string @@ -2058,6 +2095,15 @@ export async function createSingleDocument( .set({ updatedAt: now }) .where(eq(knowledgeBase.id, knowledgeBaseId)) + if (options?.processing) { + await enqueueKnowledgeDocumentProcessing(tx, { + knowledgeBaseId, + documentId, + processingOptions: options.processing.processingOptions, + billingAttribution: options.processing.billingAttribution, + }) + } + return storageNotification }) @@ -2311,31 +2357,17 @@ export async function markDocumentAsFailedTimeout( processingStartedAt: Date, requestId: string ): Promise<{ success: boolean; processingDuration: number }> { - const now = new Date() - const processingDuration = now.getTime() - processingStartedAt.getTime() - const DEAD_PROCESS_THRESHOLD_MS = 600 * 1000 // 10 minutes + const result = await failStaleDocumentProcessingClaim({ documentId, processingStartedAt }) - if (processingDuration <= DEAD_PROCESS_THRESHOLD_MS) { - throw new Error('Document has not been processing long enough to be considered dead') + if (result.success) { + logger.info( + `[${requestId}] Marked document ${documentId} as failed due to dead process (processing time: ${Math.round(result.processingDuration / 1000)}s)` + ) + } else { + logger.info(`[${requestId}] Did not time out document ${documentId} because its claim changed`) } - await db - .update(document) - .set({ - processingStatus: 'failed', - processingError: 'Processing timed out. Please retry or re-sync the connector.', - processingCompletedAt: now, - }) - .where(eq(document.id, documentId)) - - logger.info( - `[${requestId}] Marked document ${documentId} as failed due to dead process (processing time: ${Math.round(processingDuration / 1000)}s)` - ) - - return { - success: true, - processingDuration, - } + return result } export async function retryDocumentProcessing( diff --git a/apps/sim/lib/knowledge/documents/storage-billing.test.ts b/apps/sim/lib/knowledge/documents/storage-billing.test.ts index 2f0705e3eca..b89971b9f03 100644 --- a/apps/sim/lib/knowledge/documents/storage-billing.test.ts +++ b/apps/sim/lib/knowledge/documents/storage-billing.test.ts @@ -13,6 +13,7 @@ const { mockMaybeNotifyStorageLimitForBillingContext, mockResolveStorageBillingContext, mockGetFileMetadataByKeys, + mockEnqueueKnowledgeDocumentProcessing, } = vi.hoisted(() => ({ mockApplyStorageUsageDeltasInTx: vi.fn(), mockCheckStorageQuota: vi.fn(), @@ -22,6 +23,7 @@ const { mockMaybeNotifyStorageLimitForBillingContext: vi.fn(), mockResolveStorageBillingContext: vi.fn(), mockGetFileMetadataByKeys: vi.fn(), + mockEnqueueKnowledgeDocumentProcessing: vi.fn(), })) vi.mock('@/lib/billing/storage', () => ({ @@ -39,6 +41,10 @@ vi.mock('@/lib/uploads/server/metadata', () => ({ getFileMetadataByKeys: mockGetFileMetadataByKeys, })) +vi.mock('@/lib/knowledge/documents/processing-outbox-event', () => ({ + enqueueKnowledgeDocumentProcessing: mockEnqueueKnowledgeDocumentProcessing, +})) + import { createDocumentRecords, createSingleDocument, @@ -70,6 +76,7 @@ describe('knowledge document storage attribution', () => { mockApplyStorageUsageDeltasInTx.mockResolvedValue(undefined) mockMaybeNotifyStorageLimitForBillingContext.mockResolvedValue(undefined) mockGetFileMetadataByKeys.mockResolvedValue([]) + mockEnqueueKnowledgeDocumentProcessing.mockResolvedValue('outbox-1') }) it.each(['external-collaborator', 'personal-api-key-user'])( @@ -167,6 +174,87 @@ describe('knowledge document storage attribution', () => { ) }) + it('enqueues processing inside the document storage transaction', async () => { + const billingAttribution = { + actorUserId: 'external-collaborator', + workspaceId: 'workspace-1', + organizationId: null, + billedAccountUserId: 'workspace-owner', + billingEntity: { type: 'user' as const, id: 'workspace-owner' }, + billingPeriod: { + start: '2026-08-01T00:00:00.000Z', + end: '2026-09-01T00:00:00.000Z', + }, + payerSubscription: null, + } + + await createSingleDocument( + { + filename: 'note.txt', + fileUrl: 'data:text/plain;base64,SGVsbG8=', + fileSize: 5, + mimeType: 'text/plain', + }, + 'knowledge-base-1', + 'request-1', + 'external-collaborator', + 'document-1', + undefined, + { + expectedWorkspaceId: 'workspace-1', + processing: { processingOptions: { lang: 'en' }, billingAttribution }, + } + ) + + expect(mockEnqueueKnowledgeDocumentProcessing).toHaveBeenCalledWith(dbChainMock.db, { + knowledgeBaseId: 'knowledge-base-1', + documentId: 'document-1', + processingOptions: { lang: 'en' }, + billingAttribution, + }) + }) + + it('rolls back document creation when durable processing enqueue fails', async () => { + const failure = new Error('outbox unavailable') + mockEnqueueKnowledgeDocumentProcessing.mockRejectedValueOnce(failure) + + await expect( + createSingleDocument( + { + filename: 'note.txt', + fileUrl: 'data:text/plain;base64,SGVsbG8=', + fileSize: 5, + mimeType: 'text/plain', + }, + 'knowledge-base-1', + 'request-1', + 'external-collaborator', + 'document-1', + undefined, + { + expectedWorkspaceId: 'workspace-1', + processing: { + processingOptions: {}, + billingAttribution: { + actorUserId: 'external-collaborator', + workspaceId: 'workspace-1', + organizationId: null, + billedAccountUserId: 'workspace-owner', + billingEntity: { type: 'user', id: 'workspace-owner' }, + billingPeriod: { + start: '2026-08-01T00:00:00.000Z', + end: '2026-09-01T00:00:00.000Z', + }, + payerSubscription: null, + }, + }, + } + ) + ).rejects.toBe(failure) + + expect(mockMaybeNotifyStorageLimitForBillingContext).not.toHaveBeenCalled() + }) + it.each(['kb', 'knowledge-base'])( 'uses server-known %s file metadata size for quota, ledger, and document row', async (keyPrefix) => { diff --git a/apps/sim/lib/knowledge/orchestration/documents.test.ts b/apps/sim/lib/knowledge/orchestration/documents.test.ts index e4a66fec7f7..aef8d2eafca 100644 --- a/apps/sim/lib/knowledge/orchestration/documents.test.ts +++ b/apps/sim/lib/knowledge/orchestration/documents.test.ts @@ -378,7 +378,13 @@ describe('performDeleteKnowledgeDocument', () => { }) describe('document processing state changes', () => { - beforeEach(() => vi.clearAllMocks()) + beforeEach(() => { + vi.clearAllMocks() + mockMarkDocumentAsFailedTimeout.mockResolvedValue({ + success: true, + processingDuration: 600_001, + }) + }) it('refuses to time out a document that is not processing', async () => { const outcome = await performMarkKnowledgeDocumentTimedOut({ @@ -409,6 +415,19 @@ describe('document processing state changes', () => { expect(outcome).toMatchObject({ success: false, errorCode: 'validation' }) }) + it('reports a conflict when the processing claim changes before the timeout write', async () => { + mockMarkDocumentAsFailedTimeout.mockResolvedValue({ + success: false, + processingDuration: 600_001, + }) + + const outcome = await performMarkKnowledgeDocumentTimedOut({ + document: { id: 'doc-1', processingStatus: 'processing', processingStartedAt: new Date() }, + }) + + expect(outcome).toMatchObject({ success: false, errorCode: 'conflict' }) + }) + it('refuses to retry a document that has not failed', async () => { const outcome = await performRetryKnowledgeDocumentProcessing({ knowledgeBaseId: 'kb-1', diff --git a/apps/sim/lib/knowledge/orchestration/documents.ts b/apps/sim/lib/knowledge/orchestration/documents.ts index 72c7eb51cf7..60d1d8758e2 100644 --- a/apps/sim/lib/knowledge/orchestration/documents.ts +++ b/apps/sim/lib/knowledge/orchestration/documents.ts @@ -527,7 +527,17 @@ export async function performMarkKnowledgeDocumentTimedOut( } try { - await markDocumentAsFailedTimeout(document.id, document.processingStartedAt, requestId) + const result = await markDocumentAsFailedTimeout( + document.id, + document.processingStartedAt, + requestId + ) + if (!result.success) { + return fail( + 'Document processing attempt changed before timeout could be recorded', + 'conflict' + ) + } } catch (error) { // The service rejects a document that has not been processing long enough // to be presumed dead; that is a caller-fixable "try again later", not a fault.