Skip to content

Commit a28b053

Browse files
fix(knowledge): make document uploads durable
1 parent 455e24f commit a28b053

12 files changed

Lines changed: 568 additions & 130 deletions

File tree

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

Lines changed: 13 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@ const {
1717
mockUploadDocument,
1818
mockReadFormData,
1919
mockReadFile,
20-
mockUploadWorkspaceFile,
2120
mockPlatformUploaded,
2221
mockCapture,
2322
mockIsPayloadSizeLimitError,
@@ -26,7 +25,6 @@ const {
2625
mockUploadDocument: vi.fn(),
2726
mockReadFormData: vi.fn(),
2827
mockReadFile: vi.fn(),
29-
mockUploadWorkspaceFile: vi.fn(),
3028
mockPlatformUploaded: vi.fn(),
3129
mockCapture: vi.fn(),
3230
mockIsPayloadSizeLimitError: vi.fn(),
@@ -58,10 +56,6 @@ vi.mock('@/lib/core/utils/stream-limits', () => ({
5856
readFileToBufferWithLimit: mockReadFile,
5957
}))
6058

61-
vi.mock('@/lib/uploads/contexts/workspace', () => ({
62-
uploadWorkspaceFile: mockUploadWorkspaceFile,
63-
}))
64-
6559
vi.mock('@/lib/core/telemetry', () => ({
6660
PlatformEvents: { knowledgeBaseDocumentsUploaded: mockPlatformUploaded },
6761
}))
@@ -102,13 +96,11 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
10296
knowledgeBaseId: 'kb-1',
10397
knowledgeBaseName: 'Support docs',
10498
workspaceId: WORKSPACE_ID,
105-
storageActorUserId: 'user-1',
10699
})
107100
const formData = new FormData()
108101
formData.set('file', new File(['hello'], 'support.txt', { type: 'text/plain' }))
109102
mockReadFormData.mockResolvedValue(formData)
110103
mockReadFile.mockResolvedValue(Buffer.from('hello'))
111-
mockUploadWorkspaceFile.mockResolvedValue({ url: 's3://workspace/support.txt' })
112104
mockUploadDocument.mockResolvedValue({
113105
created: true,
114106
document: {
@@ -141,21 +133,14 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
141133
input: { knowledgeBaseId: 'kb-1', assertedWorkspaceId: WORKSPACE_ID },
142134
request,
143135
})
144-
expect(mockUploadWorkspaceFile).toHaveBeenCalledWith(
145-
WORKSPACE_ID,
146-
'user-1',
147-
Buffer.from('hello'),
148-
'support.txt',
149-
'text/plain'
150-
)
151136
expect(mockUploadDocument).toHaveBeenCalledWith({
152137
principal: PRINCIPAL,
153138
input: {
154139
knowledgeBaseId: 'kb-1',
155140
assertedWorkspaceId: WORKSPACE_ID,
156-
document: {
141+
file: {
142+
buffer: Buffer.from('hello'),
157143
filename: 'support.txt',
158-
fileUrl: 's3://workspace/support.txt',
159144
fileSize: 5,
160145
mimeType: 'text/plain',
161146
},
@@ -194,7 +179,6 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
194179
error: { code: 'USAGE_LIMIT_EXCEEDED', message: 'Upgrade required' },
195180
})
196181
expect(mockReadFormData).not.toHaveBeenCalled()
197-
expect(mockUploadWorkspaceFile).not.toHaveBeenCalled()
198182
expect(mockUploadDocument).not.toHaveBeenCalled()
199183
})
200184

@@ -214,7 +198,7 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
214198
expect(mockCapture).not.toHaveBeenCalled()
215199
})
216200

217-
it('preserves the malformed multipart envelope without transferring storage', async () => {
201+
it('preserves the malformed multipart envelope without entering the upload operation', async () => {
218202
mockReadFormData.mockRejectedValueOnce(new Error('multipart boundary missing'))
219203

220204
const response = await POST(buildRequest(), { params: Promise.resolve({ id: 'kb-1' }) })
@@ -223,12 +207,11 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
223207
expect(await response.json()).toEqual({
224208
error: { code: 'BAD_REQUEST', message: 'Request body must be valid multipart form data' },
225209
})
226-
expect(mockUploadWorkspaceFile).not.toHaveBeenCalled()
227210
expect(mockUploadDocument).not.toHaveBeenCalled()
228211
expect(mockPlatformUploaded).not.toHaveBeenCalled()
229212
})
230213

231-
it('preserves bounded multipart rejection and stops before storage transfer', async () => {
214+
it('preserves bounded multipart rejection and stops before the upload operation', async () => {
232215
const error = new Error('knowledge document upload body exceeds maximum size')
233216
mockReadFormData.mockRejectedValueOnce(error)
234217
mockIsPayloadSizeLimitError.mockImplementation((candidate: unknown) => candidate === error)
@@ -239,11 +222,10 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
239222
expect(await response.json()).toEqual({
240223
error: { code: 'PAYLOAD_TOO_LARGE', message: error.message },
241224
})
242-
expect(mockUploadWorkspaceFile).not.toHaveBeenCalled()
243225
expect(mockUploadDocument).not.toHaveBeenCalled()
244226
})
245227

246-
it('requires a file form field before storage transfer', async () => {
228+
it('requires a file form field before the upload operation', async () => {
247229
mockReadFormData.mockResolvedValueOnce(new FormData())
248230

249231
const response = await POST(buildRequest(), { params: Promise.resolve({ id: 'kb-1' }) })
@@ -252,7 +234,7 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
252234
expect(await response.json()).toEqual({
253235
error: { code: 'BAD_REQUEST', message: 'file form field is required' },
254236
})
255-
expect(mockUploadWorkspaceFile).not.toHaveBeenCalled()
237+
expect(mockUploadDocument).not.toHaveBeenCalled()
256238
})
257239

258240
it('preserves the exact file-size rejection before reading file bytes', async () => {
@@ -269,7 +251,7 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
269251
error: { code: 'PAYLOAD_TOO_LARGE', message: 'File size exceeds 100MB limit (100.00MB)' },
270252
})
271253
expect(mockReadFile).not.toHaveBeenCalled()
272-
expect(mockUploadWorkspaceFile).not.toHaveBeenCalled()
254+
expect(mockUploadDocument).not.toHaveBeenCalled()
273255
})
274256

275257
it('preserves unsupported file-type validation before reading file bytes', async () => {
@@ -286,24 +268,24 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
286268
error: { code: 'UNSUPPORTED_MEDIA_TYPE', message: expectedMessage },
287269
})
288270
expect(mockReadFile).not.toHaveBeenCalled()
289-
expect(mockUploadWorkspaceFile).not.toHaveBeenCalled()
271+
expect(mockUploadDocument).not.toHaveBeenCalled()
290272
})
291273

292-
it('does not register or emit effects when storage transfer fails', async () => {
293-
mockUploadWorkspaceFile.mockRejectedValueOnce(new Error('storage unavailable'))
274+
it('does not emit effects when the upload operation fails', async () => {
275+
mockUploadDocument.mockRejectedValueOnce(new Error('storage unavailable'))
294276

295277
const response = await POST(buildRequest(), { params: Promise.resolve({ id: 'kb-1' }) })
296278

297279
expect(response.status).toBe(500)
298280
expect(await response.json()).toEqual({
299281
error: { code: 'INTERNAL_ERROR', message: 'Internal server error' },
300282
})
301-
expect(mockUploadDocument).not.toHaveBeenCalled()
283+
expect(mockUploadDocument).toHaveBeenCalledOnce()
302284
expect(mockPlatformUploaded).not.toHaveBeenCalled()
303285
expect(mockCapture).not.toHaveBeenCalled()
304286
})
305287

306-
it('preserves application authorization errors after storage transfer', async () => {
288+
it('preserves final application authorization errors', async () => {
307289
mockUploadDocument.mockRejectedValueOnce(
308290
new OrchestrationError('forbidden', 'Insufficient workspace permissions')
309291
)
@@ -314,7 +296,7 @@ describe('POST /api/v2/knowledge/[id]/documents', () => {
314296
expect(await response.json()).toEqual({
315297
error: { code: 'FORBIDDEN', message: 'Insufficient workspace permissions' },
316298
})
317-
expect(mockUploadWorkspaceFile).toHaveBeenCalledOnce()
299+
expect(mockUploadDocument).toHaveBeenCalledOnce()
318300
expect(mockPlatformUploaded).not.toHaveBeenCalled()
319301
expect(mockCapture).not.toHaveBeenCalled()
320302
})

apps/sim/app/api/v2/knowledge/[id]/documents/route.ts

Lines changed: 3 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,6 @@ import {
2626
import { knowledgeOperations } from '@/lib/knowledge/application/operations'
2727
import { KnowledgeDocumentUnsupportedMediaTypeError } from '@/lib/knowledge/application/upload-sessions'
2828
import { captureServerEvent } from '@/lib/posthog/server'
29-
import { uploadWorkspaceFile } from '@/lib/uploads/contexts/workspace'
3029
import { MAX_KNOWLEDGE_DOCUMENT_FILE_SIZE } from '@/lib/uploads/shared/types'
3130
import { validateFileType } from '@/lib/uploads/utils/validation'
3231
import { serializeDate } from '@/app/api/v1/knowledge/utils'
@@ -138,20 +137,12 @@ export const POST = defineV2BodyLifecycleRoute({
138137
})
139138
return { file: rawFile, buffer, contentType }
140139
},
141-
transfer: ({ admission, body }) =>
142-
uploadWorkspaceFile(
143-
admission.workspaceId,
144-
admission.storageActorUserId,
145-
body.buffer,
146-
body.file.name,
147-
body.contentType
148-
),
149-
mapInput: ({ parsed, body, transfer }) => ({
140+
mapInput: ({ parsed, body }) => ({
150141
knowledgeBaseId: parsed.params.id,
151142
assertedWorkspaceId: parsed.query.workspaceId,
152-
document: {
143+
file: {
144+
buffer: body.buffer,
153145
filename: body.file.name,
154-
fileUrl: transfer.url,
155146
fileSize: body.file.size,
156147
mimeType: body.contentType,
157148
},

apps/sim/app/api/webhooks/outbox/process/route.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import { billingOutboxHandlers } from '@/lib/billing/webhooks/outbox-handlers'
99
import { processOutboxEvents } from '@/lib/core/outbox/service'
1010
import { generateRequestId } from '@/lib/core/utils/request'
1111
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
12+
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
1213
import { workflowDeploymentOutboxHandlers } from '@/lib/workflows/deployment-outbox'
1314
import { invitationMigrationOutboxHandlers } from '@/lib/workspaces/admin-move'
1415
import { reapStaleBackgroundWork } from '@/ee/workspace-forking/lib/background-work/store'
@@ -23,6 +24,7 @@ const handlers = {
2324
...membershipBillingOutboxHandlers,
2425
...enterpriseIssuanceOutboxHandlers,
2526
...invitationMigrationOutboxHandlers,
27+
...knowledgeDocumentProcessingOutboxHandlers,
2628
...workflowDeploymentOutboxHandlers,
2729
} as const
2830

apps/sim/lib/api/server/routes/v2-body-lifecycle-route.test.ts

Lines changed: 17 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ const contract = defineRouteContract({
4343

4444
class StageRejection extends Error {}
4545

46-
type RejectableStage = 'admission' | 'body' | 'transfer' | 'application' | 'presenter' | 'effects'
46+
type RejectableStage = 'admission' | 'body' | 'application' | 'presenter' | 'effects'
4747

4848
let rejectedStage: RejectableStage | null = null
4949

@@ -80,11 +80,10 @@ function buildHandler() {
8080
rejectAt('body')
8181
return { bytes: Buffer.from('body') }
8282
},
83-
async transfer({ admission }) {
84-
rejectAt('transfer')
85-
return { url: `stored://${admission.canonicalWorkspaceId}` }
86-
},
87-
mapInput: ({ parsed, transfer }) => ({ id: parsed.params.id, url: transfer.url }),
83+
mapInput: ({ parsed, admission }) => ({
84+
id: parsed.params.id,
85+
url: `stored://${admission.canonicalWorkspaceId}`,
86+
}),
8887
useCase: {
8988
operation,
9089
async execute({ input }) {
@@ -161,7 +160,6 @@ describe('defineV2BodyLifecycleRoute', () => {
161160
errorPolicy: { render: () => null },
162161
admission: { mapInput: () => ({}), useCase },
163162
readBody: async () => Buffer.alloc(0),
164-
transfer: async () => ({ url: 'stored://item-1' }),
165163
mapInput: () => ({}),
166164
useCase,
167165
present: () => ({ data: { id: 'item-1' } }),
@@ -183,7 +181,6 @@ describe('defineV2BodyLifecycleRoute', () => {
183181
'contract',
184182
'admission',
185183
'body',
186-
'transfer',
187184
'application',
188185
'presenter',
189186
'effects',
@@ -260,24 +257,20 @@ describe('defineV2BodyLifecycleRoute', () => {
260257
])
261258
})
262259

263-
it.each<RejectableStage>([
264-
'admission',
265-
'body',
266-
'transfer',
267-
'application',
268-
'presenter',
269-
'effects',
270-
])('renders typed %s rejection without entering later phases', async (stage) => {
271-
rejectedStage = stage
260+
it.each<RejectableStage>(['admission', 'body', 'application', 'presenter', 'effects'])(
261+
'renders typed %s rejection without entering later phases',
262+
async (stage) => {
263+
rejectedStage = stage
272264

273-
const response = await buildHandler()(buildRequest(), context())
265+
const response = await buildHandler()(buildRequest(), context())
274266

275-
expect(response.status).toBe(409)
276-
expect(await response.json()).toEqual({
277-
error: { code: 'CONFLICT', message: `${stage} rejected` },
278-
})
279-
expect(mocks.order.at(-1)).toBe(stage)
280-
})
267+
expect(response.status).toBe(409)
268+
expect(await response.json()).toEqual({
269+
error: { code: 'CONFLICT', message: `${stage} rejected` },
270+
})
271+
expect(mocks.order.at(-1)).toBe(stage)
272+
}
273+
)
281274

282275
it.each(['authentication', 'rollout', 'rate_limit'] as const)(
283276
'maps %s infrastructure failures to service unavailable',

apps/sim/lib/api/server/routes/v2-body-lifecycle-route.ts

Lines changed: 7 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -41,18 +41,13 @@ interface V2BodyLifecycleContext<C extends JsonApiRouteContract, A> {
4141
admission: A
4242
}
4343

44-
interface V2BodyLifecycleTransferContext<C extends JsonApiRouteContract, A, B>
44+
interface V2BodyLifecycleInputContext<C extends JsonApiRouteContract, A, B>
4545
extends V2BodyLifecycleContext<C, A> {
4646
body: B
4747
}
4848

49-
interface V2BodyLifecycleInputContext<C extends JsonApiRouteContract, A, B, T>
50-
extends V2BodyLifecycleTransferContext<C, A, B> {
51-
transfer: T
52-
}
53-
54-
interface V2BodyLifecycleSuccessContext<C extends JsonApiRouteContract, A, B, T, I, R>
55-
extends V2BodyLifecycleInputContext<C, A, B, T> {
49+
interface V2BodyLifecycleSuccessContext<C extends JsonApiRouteContract, A, B, I, R>
50+
extends V2BodyLifecycleInputContext<C, A, B> {
5651
input: I
5752
result: R
5853
}
@@ -63,7 +58,6 @@ interface V2BodyLifecycleRouteOptions<
6358
AI,
6459
A,
6560
B,
66-
T,
6761
I,
6862
R,
6963
> {
@@ -75,12 +69,11 @@ interface V2BodyLifecycleRouteOptions<
7569
parseOptions?: Omit<ParseRequestOptions, 'validationErrorResponse'>
7670
admission: V2BodyLifecycleAdmission<O, C, AI, A>
7771
readBody(context: V2BodyLifecycleContext<C, A>): Promise<B>
78-
transfer(context: V2BodyLifecycleTransferContext<C, A, B>): Promise<T>
79-
mapInput(context: V2BodyLifecycleInputContext<C, A, B, T>): I
72+
mapInput(context: V2BodyLifecycleInputContext<C, A, B>): I
8073
useCase: OperationUseCase<NoInfer<O>, I, R>
8174
present(result: R): ContractJsonResponse<C> | Promise<ContractJsonResponse<C>>
8275
onSuccess?(
83-
context: V2BodyLifecycleSuccessContext<C, A, B, T, NoInfer<I>, NoInfer<R>>
76+
context: V2BodyLifecycleSuccessContext<C, A, B, NoInfer<I>, NoInfer<R>>
8477
): void | Promise<void>
8578
}
8679

@@ -95,10 +88,9 @@ export function defineV2BodyLifecycleRoute<
9588
AI,
9689
A,
9790
B,
98-
T,
9991
I,
10092
R,
101-
>(options: V2BodyLifecycleRouteOptions<C, O, AI, A, B, T, I, R>): JsonNextRouteHandler {
93+
>(options: V2BodyLifecycleRouteOptions<C, O, AI, A, B, I, R>): JsonNextRouteHandler {
10294
if (options.contract.body) {
10395
throw new Error(
10496
`${options.contract.method} ${options.contract.path} must omit its body schema so admission precedes body reads`
@@ -153,8 +145,7 @@ export function defineV2BodyLifecycleRoute<
153145
})
154146
const lifecycleContext = { request, principal, parsed: parsed.data, admission }
155147
const body = await options.readBody(lifecycleContext)
156-
const transfer = await options.transfer({ ...lifecycleContext, body })
157-
const inputContext = { ...lifecycleContext, body, transfer }
148+
const inputContext = { ...lifecycleContext, body }
158149
const input = options.mapInput(inputContext)
159150
const result = await options.useCase.execute({ principal, input, request })
160151
const responseBody = options.contract.response.schema.parse(await options.present(result))

0 commit comments

Comments
 (0)