Skip to content

Commit cd20265

Browse files
sunnylqmclaude
andauthored
feat(upload): send large files to GCS as parallel parts (#94)
* feat(upload): send large files to GCS as parallel parts A single form POST crawls over a long, lossy path (mainland China to GCS through a proxy). The CLI now announces the file size (chunked: true, size); a server that supports it answers with a chunked plan, and the file is posted as byte ranges, 4 at a time, each with its own policy and two retries, then POST /upload/complete composes them. Servers that do not know the fields keep the single POST. GCS hands out the same url as primary and backup, so the TCP probe (up to 1 s, and GCS is unreachable from mainland China anyway) is skipped when there is nothing to switch to. progress and form-data are loaded through their default export when one exists, like node-fetch; package-optimization.test no longer mocks them (it never uploads), which leaked into every other test of the process. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * test(upload): pin the base url instead of relying on RNU_API getBaseUrl is memoized per process; on CI another test file resolved it first, so the chunked upload test saw the default endpoint. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
1 parent f67ce04 commit cd20265

3 files changed

Lines changed: 312 additions & 31 deletions

File tree

‎src/api.ts‎

Lines changed: 136 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -307,19 +307,27 @@ function loadNodeFetch(): typeof import('node-fetch').default {
307307
return (mod.default ?? mod) as typeof import('node-fetch').default;
308308
}
309309

310+
/** A byte range of the file (end exclusive) and its retry budget. */
311+
interface UploadSlice {
312+
start: number;
313+
end: number;
314+
retries: number;
315+
}
316+
310317
/**
311-
* Send the file with a size-scaled deadline and one retry on a transient
312-
* network error. `buildRequest` is invoked per attempt with a fresh file
313-
* stream (a consumed stream/form cannot be replayed).
318+
* Send the file (or one slice of it) with a size-scaled deadline and retries
319+
* on a transient network error. `buildRequest` is invoked per attempt with a
320+
* fresh file stream (a consumed stream/form cannot be replayed).
314321
*/
315322
async function sendUpload(
316323
fn: string,
317324
realUrl: string,
318325
fileSize: number,
319326
bar: ProgressBar,
320327
buildRequest: (fileStream: fs.ReadStream) => NodeFetchRequestInit,
328+
slice: UploadSlice = { start: 0, end: fileSize, retries: UPLOAD_MAX_RETRIES },
321329
): Promise<NodeFetchResponse> {
322-
const timeoutMs = uploadTimeoutMs(fileSize);
330+
const timeoutMs = uploadTimeoutMs(slice.end - slice.start);
323331
const nodeFetch = loadNodeFetch();
324332
// HTTP(S)_PROXY / NO_PROXY, like every other request of the CLI
325333
const agent = proxyAgentFor(realUrl);
@@ -330,8 +338,14 @@ async function sendUpload(
330338
timedOut = true;
331339
controller.abort();
332340
}, timeoutMs);
333-
const fileStream = fs.createReadStream(fn);
341+
const fileStream = fs.createReadStream(fn, {
342+
start: slice.start,
343+
// an empty file has no last byte; end = start reads nothing
344+
end: Math.max(slice.start, slice.end - 1),
345+
});
346+
let sent = 0;
334347
fileStream.on('data', (data) => {
348+
sent += data.length;
335349
bar.tick(data.length);
336350
});
337351
try {
@@ -343,11 +357,11 @@ async function sendUpload(
343357
} catch (rawError) {
344358
fileStream.destroy();
345359
const error = timedOut ? new UploadTimeoutError(timeoutMs) : rawError;
346-
if (attempt < UPLOAD_MAX_RETRIES && isTransientUploadError(error)) {
360+
if (attempt < slice.retries && isTransientUploadError(error)) {
347361
const reason = error instanceof Error ? error.message : String(error);
348362
console.warn(`\nUpload interrupted (${reason}), retrying...`);
349-
// restart the bar from zero for the second pass
350-
bar.curr = 0;
363+
// take this attempt's bytes back off the bar before the next pass
364+
bar.curr = Math.max(0, bar.curr - sent);
351365
continue;
352366
}
353367
throw error;
@@ -357,20 +371,91 @@ async function sendUpload(
357371
}
358372
}
359373

374+
/** Parts in flight at once; each is its own TCP connection. */
375+
const CHUNK_CONCURRENCY = 4;
376+
/** A part is small, so it can afford more retries than the whole file. */
377+
const CHUNK_MAX_RETRIES = 2;
378+
379+
interface ChunkedInstruction {
380+
key: string;
381+
partSize: number;
382+
parts: Record<string, string>[];
383+
}
384+
385+
/** A part the store answered with an HTTP error: retrying will not help. */
386+
class ChunkRejectedError extends Error {}
387+
388+
/**
389+
* Post the parts of a chunked upload, CHUNK_CONCURRENCY at a time. The first
390+
* failure stops scheduling new parts; the ones in flight finish, then it is
391+
* thrown.
392+
*/
393+
async function sendChunks(
394+
fn: string,
395+
realUrl: string,
396+
fileSize: number,
397+
bar: ProgressBar,
398+
chunked: ChunkedInstruction,
399+
postForm: (
400+
fields: Record<string, string>,
401+
byteCount: number,
402+
) => (fileStream: fs.ReadStream) => NodeFetchRequestInit,
403+
): Promise<void> {
404+
let next = 0;
405+
let failure: unknown;
406+
const worker = async () => {
407+
while (failure === undefined && next < chunked.parts.length) {
408+
const index = next++;
409+
const start = index * chunked.partSize;
410+
const end = Math.min(fileSize, start + chunked.partSize);
411+
try {
412+
const res = await sendUpload(
413+
fn,
414+
realUrl,
415+
fileSize,
416+
bar,
417+
postForm(chunked.parts[index], end - start),
418+
{ start, end, retries: CHUNK_MAX_RETRIES },
419+
);
420+
if (res.status > 299) {
421+
throw new ChunkRejectedError(
422+
`${res.status}: ${res.statusText || 'Upload failed'} (part ${index + 1}/${chunked.parts.length})`,
423+
);
424+
}
425+
} catch (error) {
426+
failure ??= error;
427+
}
428+
}
429+
};
430+
await Promise.all(
431+
Array.from(
432+
{ length: Math.min(CHUNK_CONCURRENCY, chunked.parts.length) },
433+
worker,
434+
),
435+
);
436+
if (failure !== undefined) {
437+
throw failure;
438+
}
439+
}
440+
360441
export async function uploadFile(
361442
fn: string,
362443
key?: string,
363444
appId?: string | number,
364445
) {
446+
const fileSize = fs.statSync(fn).size;
365447
// appId 用于服务端路由:绑定了自托管节点(rnu-node)的应用,
366-
// 上传指令会指向节点或其对象存储
448+
// 上传指令会指向节点或其对象存储。chunked + size 让支持的服务端(GCS)
449+
// 改发分片并行上传的指令;不认识这两个字段的服务端照旧整文件上传
367450
const resp = await post('/upload', {
368451
ext: path.extname(fn),
369452
...(appId ? { appId: Number(appId) } : {}),
453+
...(key ? {} : { chunked: true, size: fileSize }),
370454
});
371455
const { url, backupUrl, formData, maxSize } = resp;
372456
let realUrl = url;
373-
if (backupUrl) {
457+
// GCS hands out the same url twice: nothing to switch to, no probe needed
458+
if (backupUrl && backupUrl !== url) {
374459
if (global.USE_ACC_OSS) {
375460
realUrl = backupUrl;
376461
} else if (!resolveProxy(url)) {
@@ -387,7 +472,6 @@ export async function uploadFile(
387472
// console.log({realUrl});
388473
}
389474

390-
const fileSize = fs.statSync(fn).size;
391475
if (maxSize && fileSize > filesizeParser(maxSize)) {
392476
const readableFileSize = `${(fileSize / 1048576).toFixed(1)}m`;
393477
throw new Error(
@@ -400,7 +484,9 @@ export async function uploadFile(
400484
}
401485

402486
// progress/form-data are only needed here; keep them off the startup path
403-
const ProgressBarImpl = require('progress') as typeof import('progress');
487+
const progressModule = require('progress');
488+
const ProgressBarImpl = (progressModule.default ??
489+
progressModule) as typeof import('progress');
404490
const bar = new ProgressBarImpl(' Uploading [:bar] :percent :etas', {
405491
complete: '=',
406492
incomplete: ' ',
@@ -443,23 +529,49 @@ export async function uploadFile(
443529
return { hash: resp.key };
444530
}
445531

446-
const FormData = require('form-data') as typeof import('form-data');
447-
let res: NodeFetchResponse;
448-
try {
449-
res = await sendUpload(fn, realUrl, fileSize, bar, (fileStream) => {
532+
const formDataModule = require('form-data');
533+
const FormData = (formDataModule.default ??
534+
formDataModule) as typeof import('form-data');
535+
const postForm =
536+
(fields: Record<string, string>, byteCount: number) =>
537+
(fileStream: fs.ReadStream): NodeFetchRequestInit => {
450538
const form = new FormData();
451-
for (const [k, v] of Object.entries(formData)) {
539+
for (const [k, v] of Object.entries(fields)) {
452540
form.append(k, v);
453541
}
454-
if (key) {
455-
form.append('key', key);
456-
}
457-
// With every part's length known node-fetch sends Content-Length instead
458-
// of a chunked body: what object stores expect, and the only framing
459-
// that survives http-proxy-agent's rewrite of the buffered request head.
460-
form.append('file', fileStream, { knownLength: fileSize });
542+
// With every part's length known node-fetch sends Content-Length
543+
// instead of a chunked body: what object stores expect, and the only
544+
// framing that survives http-proxy-agent's rewrite of the buffered
545+
// request head.
546+
form.append('file', fileStream, { knownLength: byteCount });
461547
return { method: 'POST', body: form };
548+
};
549+
550+
if (resp.chunked && !key) {
551+
try {
552+
await sendChunks(fn, realUrl, fileSize, bar, resp.chunked, postForm);
553+
} catch (error) {
554+
if (error instanceof ChunkRejectedError) {
555+
throw createRequestError(error.message, realUrl);
556+
}
557+
return rethrowUploadError(error);
558+
}
559+
await post('/upload/complete', {
560+
key: resp.chunked.key,
561+
parts: resp.chunked.parts.length,
462562
});
563+
return { hash: resp.chunked.key as string };
564+
}
565+
566+
let res: NodeFetchResponse;
567+
try {
568+
res = await sendUpload(
569+
fn,
570+
realUrl,
571+
fileSize,
572+
bar,
573+
postForm(key ? { ...formData, key } : formData, fileSize),
574+
);
463575
} catch (error) {
464576
return rethrowUploadError(error);
465577
}

‎tests/package-optimization.test.ts‎

Lines changed: 0 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -2,13 +2,6 @@ import { describe, expect, mock, spyOn, test } from 'bun:test';
22

33
// Mock modules before any imports
44
mock.module('filesize-parser', () => ({ default: () => 0 }));
5-
mock.module('form-data', () => ({ default: class {} }));
6-
mock.module('node-fetch', () => ({ default: () => {} }));
7-
mock.module('progress', () => ({
8-
default: class {
9-
tick() {}
10-
},
11-
}));
125
mock.module('tty-table', () => {
136
const mockTable = () => ({ render: () => '' });
147
return { default: mockTable };

0 commit comments

Comments
 (0)