-
Notifications
You must be signed in to change notification settings - Fork 4
feat(upload): send large files to GCS as parallel parts #94
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -307,19 +307,27 @@ 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, | ||||||||||||||||||||||||||||||||
| realUrl: string, | ||||||||||||||||||||||||||||||||
| fileSize: number, | ||||||||||||||||||||||||||||||||
| bar: ProgressBar, | ||||||||||||||||||||||||||||||||
| buildRequest: (fileStream: fs.ReadStream) => NodeFetchRequestInit, | ||||||||||||||||||||||||||||||||
| slice: UploadSlice = { start: 0, end: fileSize, retries: UPLOAD_MAX_RETRIES }, | ||||||||||||||||||||||||||||||||
| ): Promise<NodeFetchResponse> { | ||||||||||||||||||||||||||||||||
| 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<string, string>[]; | ||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||
| /** 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<string, string>, | ||||||||||||||||||||||||||||||||
| byteCount: number, | ||||||||||||||||||||||||||||||||
| ) => (fileStream: fs.ReadStream) => NodeFetchRequestInit, | ||||||||||||||||||||||||||||||||
| ): Promise<void> { | ||||||||||||||||||||||||||||||||
| 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<string, string>, 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); | ||||||||||||||||||||||||||||||||
|
Comment on lines
+550
to
+552
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🗄️ Data Integrity & Integration | 🟡 Minor | ⚡ Quick win 🔎 Supported by static analysis🏁 Script executed: sed -n '370,440p' src/api.ts
sed -n '440,580p' src/api.ts
rg -n 'chunked|partSize|upload/complete' tests src | head -90Repository: reactnativecn/react-native-update-cli Length of output: 8867 🏁 Script executed: #!/bin/bash
printf '%s\n' '--- src/api.ts sendUpload and related helpers ---'
sed -n '250,380p' src/api.ts
printf '%s\n' '--- tests/upload-chunked.test.ts ---'
cat -n tests/upload-chunked.test.ts
printf '%s\n' '--- completion/chunk contract references ---'
rg -n -S 'upload/complete|parts:|partSize|chunked' --glob '!node_modules/**' --glob '!dist/**' --glob '!build/**' .
printf '%s\n' '--- file-size and empty-file test references ---'
rg -n -S 'fileSize|empty|size: 0|statSync' tests src README.md docs 2>/dev/null | head -120Repository: reactnativecn/react-native-update-cli Length of output: 22148 Validate the chunked plan before sending parts.
Reject plans with a non-positive, non-integer 🛡️ Suggested fix if (resp.chunked && !key) {
+ const { partSize, parts } = resp.chunked;
+ if (
+ !Number.isSafeInteger(partSize) ||
+ partSize <= 0 ||
+ !Array.isArray(parts) ||
+ parts.length !== Math.max(1, Math.ceil(fileSize / partSize))
+ ) {
+ throw createRequestError('Invalid chunked upload instruction', realUrl);
+ }
try {📝 Committable suggestion
Suggested change
🤖 Prompt for AI Agents |
||||||||||||||||||||||||||||||||
| } 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); | ||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win
🔎 Supported by static analysis
🏁 Script executed:
Repository: reactnativecn/react-native-update-cli
Length of output: 4992
🏁 Script executed:
Repository: reactnativecn/react-native-update-cli
Length of output: 25619
🏁 Script executed:
Repository: reactnativecn/react-native-update-cli
Length of output: 18132
🌐 Web query:
official Google Cloud Storage documentation retryable HTTP status codes uploads 429 5xx💡 Result:
Retry transient HTTP responses for each part and roll back only uploaded bytes.
sendUploadretries only network exceptions. Therefore,sendChunkstreats HTTP 429 and 5xx responses as unretryable. Handle HTTP 429 and 5xx responses withinsendUpload. The retry condition must include 429. Subtractsent, not the full slice length, becausesendUploadtracks the bytes that reached the progress bar.Keep other 4xx responses non-retryable.
🔧 Suggested fix
🤖 Prompt for AI Agents