diff --git a/service/src/file-server.ts b/service/src/file-server.ts index b4e698a8..4e0604de 100644 --- a/service/src/file-server.ts +++ b/service/src/file-server.ts @@ -4,7 +4,7 @@ import IORedis from 'ioredis'; import express from 'express'; import { Client } from 'minio'; import { nanoid } from 'nanoid'; -import { PassThrough } from 'stream'; +import { PassThrough, Transform } from 'stream'; import { pipeline } from 'stream/promises'; import type { BucketItem, BucketItemStat, ClientOptions } from 'minio'; import type { Readable } from 'stream'; @@ -232,7 +232,7 @@ async function uploadFile( mimetype: string, existingFileId?: string, readOnly = false, -): Promise { +): Promise { const fileId = existingFileId ?? nanoid(); const fileExtension = path.extname(filename); const objectName = `${session_id}/${fileId}${fileExtension}`; @@ -253,23 +253,40 @@ async function uploadFile( metaData['X-Amz-Meta-Read-Only'] = 'true'; } - /* Note: this returns UploadedObjectInfo */ const sessionKey = await redisClient.get(`session:${session_id}`); const peeked = await peekStreamForEmpty(fileStream); + let receivedBytes = 0; if (peeked.empty) { /* Empty file: explicit single PUT with size=0 — multipart fails with * "You must specify at least one part" on zero-byte streams. */ await minioClient.putObject(bucketName, objectName, Buffer.alloc(0), 0, metaData); } else { - await minioClient.putObject(bucketName, objectName, peeked.body, undefined, metaData); + const counter = new Transform({ + transform(chunk: Buffer, _encoding, callback) { + receivedBytes += chunk.length; + callback(null, chunk); + }, + }); + peeked.body.once('error', (error) => counter.destroy(error)); + await minioClient.putObject(bucketName, objectName, peeked.body.pipe(counter), undefined, metaData); + } + /* The stored size is the uploader's proof that nothing was lost on the way + * in. An object that does not hold every byte this server read is removed + * rather than left behind as a truncated input. */ + const { size } = await minioClient.statObject(bucketName, objectName); + if (size !== receivedBytes) { + logger.error('Stored object size differs from the bytes received', { session_id, fileId, receivedBytes, storedBytes: size }); + await minioClient.removeObject(bucketName, objectName); + throw new Error('Stored object is incomplete'); } - logger.info(`[${INSTANCE_ID}] File ID: ${fileId} | Filename: ${filename} | Session key: ${sessionKey}`); + logger.info(`[${INSTANCE_ID}] File ID: ${fileId} | Filename: ${filename} | Session key: ${sessionKey} | Bytes: ${size}`); await redisClient.set(`upload:${sessionKey}${session_id}${fileId}`, 'true', 'EX', env.SESSION_CACHE_TTL); fileUploads.inc(); return { filename, fileId, + size, }; } diff --git a/service/src/service/router.ts b/service/src/service/router.ts index b64d4a67..2728e0a0 100644 --- a/service/src/service/router.ts +++ b/service/src/service/router.ts @@ -25,6 +25,7 @@ import { captureTraceCarrier, withSpan } from '../telemetry'; import { Jobs, Languages } from '../enum'; import { FileRefAuthorizationError, authorizeRequestedFiles } from './file-authorization'; import { createUploadSessionRegistrar } from './upload-session'; +import { createForwardQueue, forwardUploadToFileServer } from './upload-forward'; import { prepareSandboxJobSecurity } from '../sandbox-egress'; import logger from '../logger'; @@ -377,6 +378,8 @@ router.post('/upload', uploadLimiter, async (req: t.AuthenticatedRequest, res: R }); const uploadPromises: Promise[] = []; + /* One staged file in flight per request; busboy waits on the rest. */ + const enqueueForward = createForwardQueue(); bb.on('field', (fieldname: string, val: string) => { if (fieldname === 'kind') { @@ -456,20 +459,17 @@ router.post('/upload', uploadLimiter, async (req: t.AuthenticatedRequest, res: R connection.set(`session:${session_id}`, sessionKey, 'EX', env.SESSION_CACHE_TTL) .then(() => { logger.info(`[${INSTANCE_ID}] Upload: Session ID: ${session_id} | User ID: ${userId} | Session key: ${sessionKey}`); - return axios.put( - `${env.FILE_SERVER_URL}/sessions/${session_id}/objects/${fileId}`, + return enqueueForward(() => forwardUploadToFileServer({ file, - { - headers: internalServiceHeaders(putHeaders), - maxBodyLength: planFileSize, - maxContentLength: planFileSize, - signal: abortController.signal, - }, - ); + url: `${env.FILE_SERVER_URL}/sessions/${session_id}/objects/${fileId}`, + headers: putHeaders, + maxBytes: planFileSize, + signal: abortController.signal, + })); }) - .then(response => { + .then(result => { clearTimeout(uploadTimeout); - resolve(response.data); + resolve(result); }) .catch(error => { clearTimeout(uploadTimeout); @@ -587,6 +587,8 @@ router.post('/upload/batch', uploadLimiter, async (req: t.AuthenticatedRequest, }); const uploadPromises: Promise[] = []; + /* One staged file in flight per request; busboy waits on the rest. */ + const enqueueForward = createForwardQueue(); bb.on('field', (fieldname: string, val: string) => { if (fieldname === 'kind') { @@ -696,18 +698,15 @@ router.post('/upload/batch', uploadLimiter, async (req: t.AuthenticatedRequest, logger.error(`[${INSTANCE_ID}] Batch upload file failed: ${filename} | Session: ${session_id}`, { error: message }); resolve({ status: 'error', filename, error: message }); }; - const forwardFile = (): Promise => axios.put( - `${env.FILE_SERVER_URL}/sessions/${session_id}/objects/${fileId}`, + const forwardFile = (): Promise => enqueueForward(() => forwardUploadToFileServer({ file, - { - headers: internalServiceHeaders(putHeaders), - maxBodyLength: planFileSize, - maxContentLength: planFileSize, - signal: abortController.signal, - }, - ).then(response => { + url: `${env.FILE_SERVER_URL}/sessions/${session_id}/objects/${fileId}`, + headers: putHeaders, + maxBytes: planFileSize, + signal: abortController.signal, + })).then(result => { clearTimeout(uploadTimeout); - resolve({ status: 'success', filename: response.data.filename, fileId: response.data.fileId }); + resolve({ status: 'success', filename: result.filename, fileId: result.fileId }); }, resolveUploadFailure); void ensureSessionRegistered(sessionKey) diff --git a/service/src/service/upload-forward.test.ts b/service/src/service/upload-forward.test.ts new file mode 100644 index 00000000..a33e8756 --- /dev/null +++ b/service/src/service/upload-forward.test.ts @@ -0,0 +1,159 @@ +import { afterEach, describe, expect, test } from 'bun:test'; +import busboy from 'busboy'; +import http from 'http'; +import type { AddressInfo } from 'net'; +import { Readable } from 'stream'; +import { setImmediate as yieldTurn } from 'timers/promises'; +import { UploadIncompleteError, createForwardQueue, forwardUploadToFileServer } from './upload-forward'; + +/** A file server stand-in that reads the PUT body and reports a stored size. */ +function fileServer(storedSize: (received: number) => number) { + const deletes: string[] = []; + const framing: Array<{ contentLength?: string; transferEncoding?: string }> = []; + const server = http.createServer((req, res) => { + if (req.method === 'DELETE') { + deletes.push(req.url ?? ''); + res.end('{}'); + return; + } + framing.push({ + contentLength: req.headers['content-length'], + transferEncoding: req.headers['transfer-encoding'], + }); + let received = 0; + req.on('data', (chunk: Buffer) => (received += chunk.length)); + req.on('end', () => { + res.setHeader('content-type', 'application/json'); + res.end(JSON.stringify({ filename: 'cli.tgz', fileId: 'f1', size: storedSize(received) })); + }); + }); + return new Promise<{ url: string; deletes: string[]; framing: typeof framing; close: () => void }>((resolve) => { + server.listen(0, '127.0.0.1', () => { + const { port } = server.address() as AddressInfo; + resolve({ url: `http://127.0.0.1:${port}/sessions/s1/objects/f1`, deletes, framing, close: () => server.close() }); + }); + }); +} + +function upload(bytes: number): Readable { + const body = Buffer.alloc(bytes, 7); + return Readable.from((function* () { + for (let offset = 0; offset < body.length; offset += 65536) yield body.subarray(offset, offset + 65536); + })()); +} + +let close: (() => void) | undefined; +afterEach(() => close?.()); + +describe('forwardUploadToFileServer', () => { + test('returns the stored file when the file server holds every forwarded byte', async () => { + const server = await fileServer((received) => received); + close = server.close; + const result = await forwardUploadToFileServer({ + file: upload(262600), + url: server.url, + headers: { 'Content-Type': 'application/gzip', 'X-Original-Filename': 'cli.tgz' }, + maxBytes: 1 << 20, + signal: new AbortController().signal, + }); + expect(result).toEqual({ filename: 'cli.tgz', fileId: 'f1' }); + expect(server.deletes).toEqual([]); + /* The whole file goes out as one fixed-length body, never chunked. */ + expect(server.framing).toEqual([{ contentLength: '262600', transferEncoding: undefined }]); + }); + + test('refuses and deletes a stored object shorter than what was forwarded', async () => { + const server = await fileServer(() => 225423); + close = server.close; + const forwarding = forwardUploadToFileServer({ + file: upload(262600), + url: server.url, + headers: { 'Content-Type': 'application/gzip', 'X-Original-Filename': 'cli.tgz' }, + maxBytes: 1 << 20, + signal: new AbortController().signal, + }); + await expect(forwarding).rejects.toBeInstanceOf(UploadIncompleteError); + await forwarding.catch((error: UploadIncompleteError) => { + expect(error.forwardedBytes).toBe(262600); + expect(error.storedBytes).toBe(225423); + }); + expect(server.deletes).toEqual(['/sessions/s1/objects/f1']); + }); + + test('refuses a file server answer that carries no stored size', async () => { + const server = await fileServer(() => undefined as unknown as number); + close = server.close; + await expect( + forwardUploadToFileServer({ + file: upload(1024), + url: server.url, + headers: { 'Content-Type': 'application/gzip', 'X-Original-Filename': 'cli.tgz' }, + maxBytes: 1 << 20, + signal: new AbortController().signal, + }), + ).rejects.toBeInstanceOf(UploadIncompleteError); + }); +}); + +describe('createForwardQueue', () => { + test('keeps one file in flight per multipart upload, however many files it carries', async () => { + /* A file server that answers a few event-loop turns late and records how + * many PUTs are open at once. */ + let open = 0; + let maxOpen = 0; + const fileServerStub = http.createServer((req, res) => { + open += 1; + maxOpen = Math.max(maxOpen, open); + let received = 0; + req.on('data', (chunk: Buffer) => (received += chunk.length)); + req.on('end', async () => { + for (let turn = 0; turn < 5; turn++) await yieldTurn(); + open -= 1; + res.setHeader('content-type', 'application/json'); + res.end(JSON.stringify({ filename: 'f', fileId: req.url, size: received })); + }); + }); + /* The api side, wired as router.ts wires /upload/batch. */ + const api = http.createServer((req, res) => { + const enqueueForward = createForwardQueue(); + const forwards: Promise[] = []; + const bb = busboy({ headers: req.headers }); + let n = 0; + bb.on('file', (_field, file) => { + const { port } = fileServerStub.address() as AddressInfo; + const url = `http://127.0.0.1:${port}/sessions/s/objects/${n++}`; + forwards.push(enqueueForward(() => forwardUploadToFileServer({ + file, + url, + headers: { 'Content-Type': 'application/octet-stream', 'X-Original-Filename': 'f' }, + maxBytes: 1 << 22, + signal: new AbortController().signal, + }))); + }); + bb.on('finish', async () => res.end(JSON.stringify((await Promise.allSettled(forwards)).map(r => r.status)))); + req.pipe(bb); + }); + await new Promise((resolve) => fileServerStub.listen(0, '127.0.0.1', resolve)); + await new Promise((resolve) => api.listen(0, '127.0.0.1', resolve)); + close = () => { + fileServerStub.close(); + api.close(); + }; + + const form = new FormData(); + for (let i = 0; i < 12; i++) form.append('file', new Blob([Buffer.alloc(256 * 1024, i)]), `f${i}`); + const { port } = api.address() as AddressInfo; + const statuses = await (await fetch(`http://127.0.0.1:${port}/`, { method: 'POST', body: form })).json(); + + expect(statuses).toEqual(Array(12).fill('fulfilled')); + expect(maxOpen).toBe(1); + }); + + test('runs the forwards after a failed one', async () => { + const enqueue = createForwardQueue(); + const failed = enqueue(() => Promise.reject(new Error('lost'))); + const next = enqueue(() => Promise.resolve('stored')); + await expect(failed).rejects.toThrow(); + expect(await next).toBe('stored'); + }); +}); diff --git a/service/src/service/upload-forward.ts b/service/src/service/upload-forward.ts new file mode 100644 index 00000000..e6457a7e --- /dev/null +++ b/service/src/service/upload-forward.ts @@ -0,0 +1,96 @@ +import axios from 'axios'; +import type { Readable } from 'stream'; +import type * as t from '../types'; +import { internalServiceHeaders } from '../internal-service-auth'; +import logger from '../logger'; + +/** The file server stored fewer (or more) bytes than the api forwarded. */ +export class UploadIncompleteError extends Error { + constructor( + readonly forwardedBytes: number, + readonly storedBytes: number | undefined, + ) { + super('Upload incomplete'); + this.name = 'UploadIncompleteError'; + } +} + +/** + * Sends one multipart file to the file server and proves the stored object + * holds every byte the api read from the upload. + * + * The file is staged in memory (busboy already caps it at the plan's file + * size) and sent as one Buffer with an explicit Content-Length. Streaming the + * busboy file straight into axios.put lost the tail under Bun: the chunked + * request was ended cleanly while chunks written under backpressure were + * dropped, and the file server stored a well-formed but short body. + * + * The file server reports the stored object's size. A different size (or + * none) means bytes were lost on the way, so the object is deleted and the + * upload fails instead of handing the caller a reference to a truncated file. + */ +export async function forwardUploadToFileServer({ + file, + url, + headers, + maxBytes, + signal, +}: { + file: Readable; + url: string; + headers: Record; + maxBytes: number; + signal: AbortSignal; +}): Promise { + const chunks: Buffer[] = []; + for await (const chunk of file) chunks.push(chunk as Buffer); + /* busboy ends a file early at the plan's size limit; the route aborts the + * signal for that case and reports the limit, so nothing partial is sent. */ + signal.throwIfAborted(); + const body = Buffer.concat(chunks); + /* Only the joined copy stays alive while the request is in flight. */ + chunks.length = 0; + const forwardedBytes = body.length; + const response = await axios.put(url, body, { + headers: internalServiceHeaders({ ...headers, 'Content-Length': String(forwardedBytes) }), + maxBodyLength: maxBytes, + maxContentLength: maxBytes, + signal, + }); + const storedBytes = response.data.size; + if (storedBytes !== forwardedBytes) { + logger.error('Stored upload is incomplete', { + fileId: response.data.fileId, + forwardedBytes, + storedBytes, + }); + await axios + .delete(url, { headers: internalServiceHeaders() }) + .catch((error: unknown) => { + logger.error('Failed to delete incomplete upload', { + fileId: response.data.fileId, + status: axios.isAxiosError(error) ? error.response?.status : undefined, + }); + }); + throw new UploadIncompleteError(forwardedBytes, storedBytes); + } + return { filename: response.data.filename, fileId: response.data.fileId }; +} + +/** + * Runs one upload's forwards one at a time, in the order they are queued. + * + * A file is read only when its turn comes, and busboy does not reach the next + * part until the current file stream is drained. So the request body is held + * back while a file is in flight, and the api holds at most one staged file + * per upload request, whatever the number of files in it. A failed forward + * does not stop the ones after it. + */ +export function createForwardQueue(): (forward: () => Promise) => Promise { + let tail: Promise = Promise.resolve(); + return (forward: () => Promise): Promise => { + const run = tail.then(forward); + tail = run.catch(() => undefined); + return run; + }; +} diff --git a/service/src/types/files.ts b/service/src/types/files.ts index 3124b89f..46c620a7 100644 --- a/service/src/types/files.ts +++ b/service/src/types/files.ts @@ -4,6 +4,11 @@ export interface UploadResult { fileId: string; } +/** The file server's answer to an object PUT: the stored object's size in bytes. */ +export interface StoredUploadResult extends UploadResult { + size: number; +} + export type SimpleObject = string; export interface SummaryObject {