From 61a8f27ff032ac01450383b2ac6b1b094fe891e4 Mon Sep 17 00:00:00 2001 From: Thor Haaland Date: Sun, 4 Oct 2026 10:45:51 +0000 Subject: [PATCH 1/3] fix(upload): refuse a stored upload that lost bytes on the way in The api counts the bytes it forwards from each multipart file; the file server counts the bytes it reads, stats the stored object and reports its size. Any mismatch deletes the object and fails that file, so a caller never receives a reference to a truncated input. Seen on KS-7 and AX41: the agent CLI package stored tail-truncated while the upload reported success. --- service/src/file-server.ts | 27 +++++-- service/src/service/router.ts | 37 ++++----- service/src/service/upload-forward.test.ts | 87 ++++++++++++++++++++++ service/src/service/upload-forward.ts | 72 ++++++++++++++++++ service/src/types/files.ts | 5 ++ 5 files changed, 202 insertions(+), 26 deletions(-) create mode 100644 service/src/service/upload-forward.test.ts create mode 100644 service/src/service/upload-forward.ts 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..5b309ee9 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 { forwardUploadToFileServer } from './upload-forward'; import { prepareSandboxJobSecurity } from '../sandbox-egress'; import logger from '../logger'; @@ -456,20 +457,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 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); @@ -696,18 +694,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 => 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..828f31fa --- /dev/null +++ b/service/src/service/upload-forward.test.ts @@ -0,0 +1,87 @@ +import { afterEach, describe, expect, test } from 'bun:test'; +import http from 'http'; +import type { AddressInfo } from 'net'; +import { Readable } from 'stream'; +import { UploadIncompleteError, 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 server = http.createServer((req, res) => { + if (req.method === 'DELETE') { + deletes.push(req.url ?? ''); + res.end('{}'); + return; + } + 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[]; 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, 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([]); + }); + + 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); + }); +}); diff --git a/service/src/service/upload-forward.ts b/service/src/service/upload-forward.ts new file mode 100644 index 00000000..6c171612 --- /dev/null +++ b/service/src/service/upload-forward.ts @@ -0,0 +1,72 @@ +import axios from 'axios'; +import { Transform } from 'stream'; +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'; + } +} + +/** + * Streams one multipart file to the file server and proves the stored object + * holds every byte the api read from the upload. 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 { + let forwardedBytes = 0; + const counter = new Transform({ + transform(chunk: Buffer, _encoding, callback) { + forwardedBytes += chunk.length; + callback(null, chunk); + }, + }); + file.once('error', (error) => counter.destroy(error)); + file.pipe(counter); + const response = await axios.put(url, counter, { + headers: internalServiceHeaders(headers), + 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 }; +} 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 { From 49c542d191c688a0f7c8b09d2599e9629119dad0 Mon Sep 17 00:00:00 2001 From: Thor Haaland Date: Sun, 4 Oct 2026 10:54:57 +0000 Subject: [PATCH 2/3] fix(upload): send each file as one Buffer with a Content-Length Root cause (tcpdump on KS-7, Agent/MCP): streaming the busboy file into axios.put under Bun ended a well-formed chunked request early, dropping chunks written under backpressure; the file server stored exactly what it was sent. The api now stages the file (already capped at the plan size by busboy) and PUTs it as a single fixed-length body. The stored-size checks stay as the guard. --- service/src/service/upload-forward.test.ts | 11 +++++-- service/src/service/upload-forward.ts | 38 ++++++++++++---------- 2 files changed, 30 insertions(+), 19 deletions(-) diff --git a/service/src/service/upload-forward.test.ts b/service/src/service/upload-forward.test.ts index 828f31fa..0a6820df 100644 --- a/service/src/service/upload-forward.test.ts +++ b/service/src/service/upload-forward.test.ts @@ -7,12 +7,17 @@ import { UploadIncompleteError, forwardUploadToFileServer } from './upload-forwa /** 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', () => { @@ -20,10 +25,10 @@ function fileServer(storedSize: (received: number) => number) { res.end(JSON.stringify({ filename: 'cli.tgz', fileId: 'f1', size: storedSize(received) })); }); }); - return new Promise<{ url: string; deletes: string[]; close: () => void }>((resolve) => { + 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, close: () => server.close() }); + resolve({ url: `http://127.0.0.1:${port}/sessions/s1/objects/f1`, deletes, framing, close: () => server.close() }); }); }); } @@ -51,6 +56,8 @@ describe('forwardUploadToFileServer', () => { }); 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 () => { diff --git a/service/src/service/upload-forward.ts b/service/src/service/upload-forward.ts index 6c171612..38915d85 100644 --- a/service/src/service/upload-forward.ts +++ b/service/src/service/upload-forward.ts @@ -1,5 +1,4 @@ import axios from 'axios'; -import { Transform } from 'stream'; import type { Readable } from 'stream'; import type * as t from '../types'; import { internalServiceHeaders } from '../internal-service-auth'; @@ -17,11 +16,18 @@ export class UploadIncompleteError extends Error { } /** - * Streams one multipart file to the file server and proves the stored object - * holds every byte the api read from the upload. 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. + * 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, @@ -36,17 +42,15 @@ export async function forwardUploadToFileServer({ maxBytes: number; signal: AbortSignal; }): Promise { - let forwardedBytes = 0; - const counter = new Transform({ - transform(chunk: Buffer, _encoding, callback) { - forwardedBytes += chunk.length; - callback(null, chunk); - }, - }); - file.once('error', (error) => counter.destroy(error)); - file.pipe(counter); - const response = await axios.put(url, counter, { - headers: internalServiceHeaders(headers), + 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); + const forwardedBytes = body.length; + const response = await axios.put(url, body, { + headers: internalServiceHeaders({ ...headers, 'Content-Length': String(forwardedBytes) }), maxBodyLength: maxBytes, maxContentLength: maxBytes, signal, From 5c8c3fe205e8eb6658d868306849a03c1dea2c97 Mon Sep 17 00:00:00 2001 From: Thor Haaland Date: Sun, 4 Oct 2026 11:04:12 +0000 Subject: [PATCH 3/3] fix(upload): stage one file at a time per upload request Review P1: staging every file at once removed busboy's backpressure, so a batch against a slow file server held every file in the api. Forwards now run one at a time per request (createForwardQueue); a file is read only on its turn and busboy waits on the rest, so the api holds at most one staged file per request. The staged chunks are released once joined. --- service/src/service/router.ts | 14 +++-- service/src/service/upload-forward.test.ts | 67 +++++++++++++++++++++- service/src/service/upload-forward.ts | 20 +++++++ 3 files changed, 95 insertions(+), 6 deletions(-) diff --git a/service/src/service/router.ts b/service/src/service/router.ts index 5b309ee9..2728e0a0 100644 --- a/service/src/service/router.ts +++ b/service/src/service/router.ts @@ -25,7 +25,7 @@ import { captureTraceCarrier, withSpan } from '../telemetry'; import { Jobs, Languages } from '../enum'; import { FileRefAuthorizationError, authorizeRequestedFiles } from './file-authorization'; import { createUploadSessionRegistrar } from './upload-session'; -import { forwardUploadToFileServer } from './upload-forward'; +import { createForwardQueue, forwardUploadToFileServer } from './upload-forward'; import { prepareSandboxJobSecurity } from '../sandbox-egress'; import logger from '../logger'; @@ -378,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') { @@ -457,13 +459,13 @@ 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 forwardUploadToFileServer({ + return enqueueForward(() => forwardUploadToFileServer({ file, url: `${env.FILE_SERVER_URL}/sessions/${session_id}/objects/${fileId}`, headers: putHeaders, maxBytes: planFileSize, signal: abortController.signal, - }); + })); }) .then(result => { clearTimeout(uploadTimeout); @@ -585,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') { @@ -694,13 +698,13 @@ 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 => forwardUploadToFileServer({ + const forwardFile = (): Promise => enqueueForward(() => forwardUploadToFileServer({ file, url: `${env.FILE_SERVER_URL}/sessions/${session_id}/objects/${fileId}`, headers: putHeaders, maxBytes: planFileSize, signal: abortController.signal, - }).then(result => { + })).then(result => { clearTimeout(uploadTimeout); resolve({ status: 'success', filename: result.filename, fileId: result.fileId }); }, resolveUploadFailure); diff --git a/service/src/service/upload-forward.test.ts b/service/src/service/upload-forward.test.ts index 0a6820df..a33e8756 100644 --- a/service/src/service/upload-forward.test.ts +++ b/service/src/service/upload-forward.test.ts @@ -1,8 +1,10 @@ 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 { UploadIncompleteError, forwardUploadToFileServer } from './upload-forward'; +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) { @@ -92,3 +94,66 @@ describe('forwardUploadToFileServer', () => { ).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 index 38915d85..e6457a7e 100644 --- a/service/src/service/upload-forward.ts +++ b/service/src/service/upload-forward.ts @@ -48,6 +48,8 @@ export async function forwardUploadToFileServer({ * 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) }), @@ -74,3 +76,21 @@ export async function forwardUploadToFileServer({ } 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; + }; +}