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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 22 additions & 5 deletions service/src/file-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -232,7 +232,7 @@ async function uploadFile(
mimetype: string,
existingFileId?: string,
readOnly = false,
): Promise<t.UploadResult> {
): Promise<t.StoredUploadResult> {
const fileId = existingFileId ?? nanoid();
const fileExtension = path.extname(filename);
const objectName = `${session_id}/${fileId}${fileExtension}`;
Expand All @@ -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,
};
}

Expand Down
41 changes: 20 additions & 21 deletions service/src/service/router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand Down Expand Up @@ -377,6 +378,8 @@ router.post('/upload', uploadLimiter, async (req: t.AuthenticatedRequest, res: R
});

const uploadPromises: Promise<t.UploadResult>[] = [];
/* 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') {
Expand Down Expand Up @@ -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<t.UploadResult>(
`${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);
Expand Down Expand Up @@ -587,6 +587,8 @@ router.post('/upload/batch', uploadLimiter, async (req: t.AuthenticatedRequest,
});

const uploadPromises: Promise<t.BatchUploadFileResult>[] = [];
/* 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') {
Expand Down Expand Up @@ -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<void> => axios.put<t.UploadResult>(
`${env.FILE_SERVER_URL}/sessions/${session_id}/objects/${fileId}`,
const forwardFile = (): Promise<void> => 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)
Expand Down
159 changes: 159 additions & 0 deletions service/src/service/upload-forward.test.ts
Original file line number Diff line number Diff line change
@@ -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<unknown>[] = [];
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<void>((resolve) => fileServerStub.listen(0, '127.0.0.1', resolve));
await new Promise<void>((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');
});
});
96 changes: 96 additions & 0 deletions service/src/service/upload-forward.ts
Original file line number Diff line number Diff line change
@@ -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<string, string>;
maxBytes: number;
signal: AbortSignal;
}): Promise<t.UploadResult> {
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<t.StoredUploadResult>(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(): <T>(forward: () => Promise<T>) => Promise<T> {
let tail: Promise<unknown> = Promise.resolve();
return <T>(forward: () => Promise<T>): Promise<T> => {
const run = tail.then(forward);
tail = run.catch(() => undefined);
return run;
};
}
Loading
Loading