diff --git a/src/api.ts b/src/api.ts index 74c8a34..8bba1da 100644 --- a/src/api.ts +++ b/src/api.ts @@ -307,10 +307,17 @@ function loadNodeFetch(): typeof import('node-fetch').default { return (mod.default ?? mod) as typeof import('node-fetch').default; } +/** A byte range of the file (end exclusive) and its retry budget. */ +interface UploadSlice { + start: number; + end: number; + retries: number; +} + /** - * Send the file with a size-scaled deadline and one retry on a transient - * network error. `buildRequest` is invoked per attempt with a fresh file - * stream (a consumed stream/form cannot be replayed). + * Send the file (or one slice of it) with a size-scaled deadline and retries + * on a transient network error. `buildRequest` is invoked per attempt with a + * fresh file stream (a consumed stream/form cannot be replayed). */ async function sendUpload( fn: string, @@ -318,8 +325,9 @@ async function sendUpload( fileSize: number, bar: ProgressBar, buildRequest: (fileStream: fs.ReadStream) => NodeFetchRequestInit, + slice: UploadSlice = { start: 0, end: fileSize, retries: UPLOAD_MAX_RETRIES }, ): Promise { - const timeoutMs = uploadTimeoutMs(fileSize); + const timeoutMs = uploadTimeoutMs(slice.end - slice.start); const nodeFetch = loadNodeFetch(); // HTTP(S)_PROXY / NO_PROXY, like every other request of the CLI const agent = proxyAgentFor(realUrl); @@ -330,8 +338,14 @@ async function sendUpload( timedOut = true; controller.abort(); }, timeoutMs); - const fileStream = fs.createReadStream(fn); + const fileStream = fs.createReadStream(fn, { + start: slice.start, + // an empty file has no last byte; end = start reads nothing + end: Math.max(slice.start, slice.end - 1), + }); + let sent = 0; fileStream.on('data', (data) => { + sent += data.length; bar.tick(data.length); }); try { @@ -343,11 +357,11 @@ async function sendUpload( } catch (rawError) { fileStream.destroy(); const error = timedOut ? new UploadTimeoutError(timeoutMs) : rawError; - if (attempt < UPLOAD_MAX_RETRIES && isTransientUploadError(error)) { + if (attempt < slice.retries && isTransientUploadError(error)) { const reason = error instanceof Error ? error.message : String(error); console.warn(`\nUpload interrupted (${reason}), retrying...`); - // restart the bar from zero for the second pass - bar.curr = 0; + // take this attempt's bytes back off the bar before the next pass + bar.curr = Math.max(0, bar.curr - sent); continue; } throw error; @@ -357,20 +371,91 @@ async function sendUpload( } } +/** Parts in flight at once; each is its own TCP connection. */ +const CHUNK_CONCURRENCY = 4; +/** A part is small, so it can afford more retries than the whole file. */ +const CHUNK_MAX_RETRIES = 2; + +interface ChunkedInstruction { + key: string; + partSize: number; + parts: Record[]; +} + +/** A part the store answered with an HTTP error: retrying will not help. */ +class ChunkRejectedError extends Error {} + +/** + * Post the parts of a chunked upload, CHUNK_CONCURRENCY at a time. The first + * failure stops scheduling new parts; the ones in flight finish, then it is + * thrown. + */ +async function sendChunks( + fn: string, + realUrl: string, + fileSize: number, + bar: ProgressBar, + chunked: ChunkedInstruction, + postForm: ( + fields: Record, + byteCount: number, + ) => (fileStream: fs.ReadStream) => NodeFetchRequestInit, +): Promise { + let next = 0; + let failure: unknown; + const worker = async () => { + while (failure === undefined && next < chunked.parts.length) { + const index = next++; + const start = index * chunked.partSize; + const end = Math.min(fileSize, start + chunked.partSize); + try { + const res = await sendUpload( + fn, + realUrl, + fileSize, + bar, + postForm(chunked.parts[index], end - start), + { start, end, retries: CHUNK_MAX_RETRIES }, + ); + if (res.status > 299) { + throw new ChunkRejectedError( + `${res.status}: ${res.statusText || 'Upload failed'} (part ${index + 1}/${chunked.parts.length})`, + ); + } + } catch (error) { + failure ??= error; + } + } + }; + await Promise.all( + Array.from( + { length: Math.min(CHUNK_CONCURRENCY, chunked.parts.length) }, + worker, + ), + ); + if (failure !== undefined) { + throw failure; + } +} + export async function uploadFile( fn: string, key?: string, appId?: string | number, ) { + const fileSize = fs.statSync(fn).size; // appId 用于服务端路由:绑定了自托管节点(rnu-node)的应用, - // 上传指令会指向节点或其对象存储 + // 上传指令会指向节点或其对象存储。chunked + size 让支持的服务端(GCS) + // 改发分片并行上传的指令;不认识这两个字段的服务端照旧整文件上传 const resp = await post('/upload', { ext: path.extname(fn), ...(appId ? { appId: Number(appId) } : {}), + ...(key ? {} : { chunked: true, size: fileSize }), }); const { url, backupUrl, formData, maxSize } = resp; let realUrl = url; - if (backupUrl) { + // GCS hands out the same url twice: nothing to switch to, no probe needed + if (backupUrl && backupUrl !== url) { if (global.USE_ACC_OSS) { realUrl = backupUrl; } else if (!resolveProxy(url)) { @@ -387,7 +472,6 @@ export async function uploadFile( // console.log({realUrl}); } - const fileSize = fs.statSync(fn).size; if (maxSize && fileSize > filesizeParser(maxSize)) { const readableFileSize = `${(fileSize / 1048576).toFixed(1)}m`; throw new Error( @@ -400,7 +484,9 @@ export async function uploadFile( } // progress/form-data are only needed here; keep them off the startup path - const ProgressBarImpl = require('progress') as typeof import('progress'); + const progressModule = require('progress'); + const ProgressBarImpl = (progressModule.default ?? + progressModule) as typeof import('progress'); const bar = new ProgressBarImpl(' Uploading [:bar] :percent :etas', { complete: '=', incomplete: ' ', @@ -443,23 +529,49 @@ export async function uploadFile( return { hash: resp.key }; } - const FormData = require('form-data') as typeof import('form-data'); - let res: NodeFetchResponse; - try { - res = await sendUpload(fn, realUrl, fileSize, bar, (fileStream) => { + const formDataModule = require('form-data'); + const FormData = (formDataModule.default ?? + formDataModule) as typeof import('form-data'); + const postForm = + (fields: Record, byteCount: number) => + (fileStream: fs.ReadStream): NodeFetchRequestInit => { const form = new FormData(); - for (const [k, v] of Object.entries(formData)) { + for (const [k, v] of Object.entries(fields)) { form.append(k, v); } - if (key) { - form.append('key', key); - } - // With every part's length known node-fetch sends Content-Length instead - // of a chunked body: what object stores expect, and the only framing - // that survives http-proxy-agent's rewrite of the buffered request head. - form.append('file', fileStream, { knownLength: fileSize }); + // With every part's length known node-fetch sends Content-Length + // instead of a chunked body: what object stores expect, and the only + // framing that survives http-proxy-agent's rewrite of the buffered + // request head. + form.append('file', fileStream, { knownLength: byteCount }); return { method: 'POST', body: form }; + }; + + if (resp.chunked && !key) { + try { + await sendChunks(fn, realUrl, fileSize, bar, resp.chunked, postForm); + } catch (error) { + if (error instanceof ChunkRejectedError) { + throw createRequestError(error.message, realUrl); + } + return rethrowUploadError(error); + } + await post('/upload/complete', { + key: resp.chunked.key, + parts: resp.chunked.parts.length, }); + return { hash: resp.chunked.key as string }; + } + + let res: NodeFetchResponse; + try { + res = await sendUpload( + fn, + realUrl, + fileSize, + bar, + postForm(key ? { ...formData, key } : formData, fileSize), + ); } catch (error) { return rethrowUploadError(error); } diff --git a/tests/package-optimization.test.ts b/tests/package-optimization.test.ts index 0e44918..c7b4d2d 100644 --- a/tests/package-optimization.test.ts +++ b/tests/package-optimization.test.ts @@ -2,13 +2,6 @@ import { describe, expect, mock, spyOn, test } from 'bun:test'; // Mock modules before any imports mock.module('filesize-parser', () => ({ default: () => 0 })); -mock.module('form-data', () => ({ default: class {} })); -mock.module('node-fetch', () => ({ default: () => {} })); -mock.module('progress', () => ({ - default: class { - tick() {} - }, -})); mock.module('tty-table', () => { const mockTable = () => ({ render: () => '' }); return { default: mockTable }; diff --git a/tests/upload-chunked.test.ts b/tests/upload-chunked.test.ts new file mode 100644 index 0000000..de94731 --- /dev/null +++ b/tests/upload-chunked.test.ts @@ -0,0 +1,176 @@ +import { afterEach, describe, expect, mock, spyOn, test } from 'bun:test'; +import fs from 'fs'; +import os from 'os'; +import path from 'path'; +import { PassThrough } from 'stream'; + +type FetchCall = { url: string; body: Buffer }; +const fetchCalls: FetchCall[] = []; +let fetchImpl: (call: FetchCall) => Promise<{ status: number }> = async () => ({ + status: 204, +}); + +// Read the streamed multipart body the way the socket would. +async function drain(body: NodeJS.ReadableStream): Promise { + const chunks: Buffer[] = []; + // form-data is an old-style stream: pipe it into a modern one + const pass = new PassThrough(); + body.pipe(pass); + for await (const chunk of pass) { + chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); + } + return Buffer.concat(chunks); +} + +mock.module('node-fetch', () => ({ + default: async (url: string, init: { body: NodeJS.ReadableStream }) => { + const call = { url, body: await drain(init.body) }; + fetchCalls.push(call); + const response = await fetchImpl(call); + return { statusText: '', ...response }; + }, +})); +const { uploadFile } = await import('../src/api'); +const runtime = await import('../src/utils/runtime'); +const httpHelper = await import('../src/utils/http-helper'); + +const tmp = path.join(os.tmpdir(), `rnu-chunked-${process.pid}.ppk`); +// distinct bytes per MiB, so a misplaced range is visible +const content = Buffer.concat( + Array.from({ length: 5 }, (_, index) => Buffer.alloc(1 << 20, 65 + index)), +).subarray(0, (5 << 20) - 123); + +describe('chunked upload', () => { + let spies: { mockRestore: () => void }[] = []; + const env = { ...process.env }; + + afterEach(() => { + for (const spy of spies) spy.mockRestore(); + spies = []; + fetchCalls.length = 0; + fetchImpl = async () => ({ status: 204 }); + process.env = { ...env }; + fs.rmSync(tmp, { force: true }); + }); + + function setup(instruction: Record) { + fs.writeFileSync(tmp, content); + for (const name of [ + 'HTTPS_PROXY', + 'https_proxy', + 'HTTP_PROXY', + 'http_proxy', + ]) { + delete process.env[name]; + } + const apiCalls: { url: string; body: any }[] = []; + spies.push( + spyOn(runtime, 'runtimeFetch').mockImplementation( + async (url: string, init?: any) => { + apiCalls.push({ + url, + body: init?.body ? JSON.parse(init.body) : undefined, + }); + const payload = url.endsWith('/upload/complete') + ? { key: 'k' } + : instruction; + return { + status: 200, + statusText: 'OK', + text: async () => JSON.stringify(payload), + }; + }, + ), + spyOn(runtime, 'measureTcpLatency').mockResolvedValue(10), + // the base url is memoized per process: another test file may have + // resolved it already, so pin it instead of setting RNU_API + spyOn(httpHelper, 'getBaseUrl').mockResolvedValue('https://api.test'), + spyOn(console, 'warn').mockImplementation(() => {}), + ); + return apiCalls; + } + + const gcs = 'https://storage.googleapis.com/cresc-storage'; + const partSize = 2 << 20; + const chunkedInstruction = { + url: gcs, + backupUrl: gcs, + formData: { key: 'k' }, + chunked: { + key: 'k', + partSize, + parts: [1, 2, 3].map((n) => ({ key: `k.part${n}`, policy: `p${n}` })), + }, + }; + + test('announces the size, posts every byte range once, then completes', async () => { + const apiCalls = setup(chunkedInstruction); + const result = await uploadFile(tmp, undefined, 9); + + expect(apiCalls[0].url).toBe('https://api.test/upload'); + expect(apiCalls[0].body).toEqual({ + ext: '.ppk', + appId: 9, + chunked: true, + size: content.length, + }); + expect(fetchCalls).toHaveLength(3); + for (const [index, name] of ['k.part1', 'k.part2', 'k.part3'].entries()) { + const call = fetchCalls.find((c) => + c.body.includes(`\r\n\r\n${name}\r\n`), + ); + expect(call).toBeDefined(); + const slice = content.subarray(index * partSize, (index + 1) * partSize); + expect(call!.url).toBe(gcs); + expect(call!.body.includes(slice)).toBe(true); + expect(call!.body.length).toBeLessThan(slice.length + 2048); + } + expect(apiCalls[1]).toEqual({ + url: 'https://api.test/upload/complete', + body: { key: 'k', parts: 3 }, + }); + expect(result).toEqual({ hash: 'k' }); + }); + + test('retries a part after a reset and still completes', async () => { + const apiCalls = setup(chunkedInstruction); + let failed = false; + fetchImpl = async (call) => { + if (!failed && call.body.includes('\r\n\r\nk.part2\r\n')) { + failed = true; + throw Object.assign(new Error('socket hang up'), { + code: 'ECONNRESET', + }); + } + return { status: 204 }; + }; + await uploadFile(tmp); + expect(fetchCalls).toHaveLength(4); + expect(apiCalls.at(-1)?.url).toBe('https://api.test/upload/complete'); + }); + + test('a rejected part fails the upload without completing it', async () => { + const apiCalls = setup(chunkedInstruction); + fetchImpl = async (call) => ({ + status: call.body.includes('\r\n\r\nk.part1\r\n') ? 403 : 204, + }); + await expect(uploadFile(tmp)).rejects.toThrow('403'); + expect(apiCalls.some((c) => c.url.endsWith('/upload/complete'))).toBe( + false, + ); + }); + + test('a server without chunking gets the single POST and no probe', async () => { + const apiCalls = setup({ + url: gcs, + backupUrl: gcs, + formData: { key: 'k' }, + }); + const result = await uploadFile(tmp); + expect(fetchCalls).toHaveLength(1); + expect(fetchCalls[0].body.includes(content)).toBe(true); + expect(runtime.measureTcpLatency).not.toHaveBeenCalled(); + expect(apiCalls).toHaveLength(1); + expect(result).toEqual({ hash: 'k' }); + }); +});