From fd82dc07f5355ead990a359d68f4f6d0b79efe01 Mon Sep 17 00:00:00 2001 From: feywind <57276408+feywind@users.noreply.github.com> Date: Wed, 30 Sep 2026 17:27:22 +0000 Subject: [PATCH 1/3] fix(gax): propagate chunk granularity on resume and preserve retry error details --- core/packages/gax/src/resumableUpload.ts | 63 ++++- .../test/showcase-resumable-upload/sample.ts | 49 ++++ .../gax/test/system-test/resumableUpload.ts | 44 +++ .../packages/gax/test/unit/resumableUpload.ts | 267 ++++++++++++++++++ 4 files changed, 410 insertions(+), 13 deletions(-) diff --git a/core/packages/gax/src/resumableUpload.ts b/core/packages/gax/src/resumableUpload.ts index dd46f6516eab..ca957515cd0f 100644 --- a/core/packages/gax/src/resumableUpload.ts +++ b/core/packages/gax/src/resumableUpload.ts @@ -291,6 +291,7 @@ interface TransmitResult { interface QueryOffsetResult { offset: number; + granularity: number | null; finalResponse?: {}; } @@ -556,6 +557,11 @@ export class ResumableUploadSession { const queried = await this.queryOffset(params.resumeUrl); this.uploadUrl_ = params.resumeUrl; this.committedBytes_ = queried.offset; + granularity = queried.granularity; + this.effectiveChunkSize_ = this.computeEffectiveChunkSize( + params.chunkSize, + granularity, + ); this.reportProgress(); sessionUrl = params.resumeUrl; finalResponse = queried.finalResponse; @@ -565,6 +571,10 @@ export class ResumableUploadSession { sessionUrl = started.uploadUrl; granularity = started.granularity; this.uploadUrl_ = sessionUrl; + this.effectiveChunkSize_ = this.computeEffectiveChunkSize( + params.chunkSize, + granularity, + ); if (started.committedOffset !== undefined) { this.committedBytes_ = started.committedOffset; this.reportProgress(); @@ -580,10 +590,6 @@ export class ResumableUploadSession { throw err; } - this.effectiveChunkSize_ = this.computeEffectiveChunkSize( - params.chunkSize, - granularity, - ); if (finalResponse !== undefined) { this.state_ = ResumableUploadState.FINALIZING; this.done_ = true; @@ -730,7 +736,7 @@ export class ResumableUploadSession { const queried = await this.queryOffset(uploadUrl); return { uploadUrl, - granularity, + granularity: granularity ?? queried.granularity, committedOffset: queried.offset, finalResponse: queried.finalResponse, }; @@ -773,13 +779,16 @@ export class ResumableUploadSession { const size = parseHeaderInt( response.headers.get(UPLOAD_SIZE_RECEIVED_HEADER), ); + const granularity = parseHeaderInt( + response.headers.get(UPLOAD_CHUNK_GRANULARITY_HEADER), + ); if (uploadStatus === 'final') { const finalResponse = this.decodeFinalResponse(response); const offset = size ?? (this.params ? this.uploadSizeForDeadline(this.params) : 0) ?? this.committedBytes_; - return {offset, finalResponse}; + return {offset, granularity, finalResponse}; } if (size === null) { throw new Category3Error( @@ -787,7 +796,7 @@ export class ResumableUploadSession { `${UPLOAD_SIZE_RECEIVED_HEADER} header.`, ); } - return {offset: size}; + return {offset: size, granularity}; } private async runTransmission(sessionUrl: string): Promise { @@ -1076,6 +1085,12 @@ export class ResumableUploadSession { // align the local state, and re-enter the transmission phase. this.state_ = ResumableUploadState.RECOVERY; const queried = await this.queryOffset(sessionUrl); + if (queried.granularity !== null) { + this.effectiveChunkSize_ = this.computeEffectiveChunkSize( + this.params?.chunkSize, + queried.granularity, + ); + } const serverOffset = queried.offset; this.committedBytes_ = serverOffset; this.reportProgress(); @@ -1097,13 +1112,20 @@ export class ResumableUploadSession { throw new RestartUploadError(serverOffset); } if (serverOffset === currentOffset) { + // If the recovery query negotiated a different effective chunk size + // than the chunk currently in flight, restart from the server offset + // so the source is re-read with the aligned chunk size. + if (!isFinal && currentChunk.length !== this.effectiveChunkSize_) { + throw new RestartUploadError(serverOffset); + } // Nothing was committed; apply retry limit and backoff before // retrying the same chunk so a persistent Category 2 error does not // spin in a tight 0ms loop. if (sameOffsetAttempts >= retry.maxRetries) { throw createGoogleError( `Exceeded the maximum number of recovery retries (${retry.maxRetries}) ` + - `without forward progress at byte offset ${currentOffset}.`, + `without forward progress at byte offset ${currentOffset}: ` + + `${(err as Error).message}`, Status.DEADLINE_EXCEEDED, ); } @@ -1218,6 +1240,12 @@ export class ResumableUploadSession { } this.state_ = ResumableUploadState.RECOVERY; const queried = await this.queryOffset(sessionUrl); + if (queried.granularity !== null) { + this.effectiveChunkSize_ = this.computeEffectiveChunkSize( + this.params?.chunkSize, + queried.granularity, + ); + } const serverOffset = queried.offset; this.committedBytes_ = serverOffset; this.reportProgress(); @@ -1234,7 +1262,8 @@ export class ResumableUploadSession { if (sameOffsetAttempts >= retry.maxRetries) { throw createGoogleError( `Exceeded the maximum number of recovery retries (${retry.maxRetries}) ` + - `while finalizing at byte offset ${currentOffset}.`, + `while finalizing at byte offset ${currentOffset}: ` + + `${(err as Error).message}`, Status.DEADLINE_EXCEEDED, ); } @@ -1315,7 +1344,8 @@ export class ResumableUploadSession { if (attempt >= retry.maxRetries) { throw createGoogleError( `Exceeded the maximum number of retries (${retry.maxRetries}) ` + - `while sending the resumable upload command "${command}".`, + `while sending the resumable upload command "${command}": ` + + `${err.message}`, Status.DEADLINE_EXCEEDED, ); } @@ -1416,7 +1446,8 @@ export class ResumableUploadSession { if ( err instanceof Category2Error || err instanceof Category3Error || - err instanceof GoogleError + err instanceof GoogleError || + err instanceof TransientError ) { throw err; } @@ -1453,8 +1484,14 @@ export class ResumableUploadSession { ); } if (command === COMMAND_START) { - // The start handler reconciles via the returned session URL. - return; + if (response.headers.get(UPLOAD_URL_HEADER) !== null) { + // The start handler reconciles via the returned session URL. + return; + } + throw new TransientError( + 'The resumable upload start response did not include ' + + `${UPLOAD_STATUS_HEADER} or ${UPLOAD_URL_HEADER} headers.`, + ); } // A missing status header on a starting/transmission/finalizing // response is a recoverable state mismatch. diff --git a/core/packages/gax/test/showcase-resumable-upload/sample.ts b/core/packages/gax/test/showcase-resumable-upload/sample.ts index 68cb81eb47b0..a6092c5bf3d7 100644 --- a/core/packages/gax/test/showcase-resumable-upload/sample.ts +++ b/core/packages/gax/test/showcase-resumable-upload/sample.ts @@ -842,6 +842,55 @@ async function testQueryAndChunkGranularityScenarios( console.log( ' PASSED Part B: chunk_granularity rounded 600-byte chunkSize to 512 bytes and completed 1500-byte upload.', ); + + // --- Part C: chunk_granularity with user-style pause & resume --- + console.log( + ' --- Part C: X-Goog-Test-Scenario: chunk_granularity (pause & resume) ---', + ); + const sessionC1 = await clientB.uploadMedia({ + name: path.basename(payloadB.filePath), + }); + let resumeHandle: {uploadUrl: string; chunkSize: number} | null = null; + await sessionC1.start({ + uploadSource: clientB.getResumableSource(payloadB.filePath), + chunkSize: 128, // < 256 server granularity -> rounds up to 256 + startHeaders: { + 'X-Goog-Test-Scenario': 'chunk_granularity', + }, + onProgress: progress => { + if (progress.bytesUploaded >= 256 && !resumeHandle) { + resumeHandle = { + uploadUrl: progress.uploadUrl, + chunkSize: sessionC1.chunkSize!, + }; + throw new Error('User pause after first 256-byte chunk'); + } + }, + }); + await assert.rejects( + sessionC1.finished(), + /User pause after first 256-byte chunk/, + ); + assert.ok(resumeHandle, 'Expected resumeHandle to be captured'); + const handle = resumeHandle as {uploadUrl: string; chunkSize: number}; + assert.strictEqual( + handle.chunkSize, + 256, + `Expected negotiated handle.chunkSize to be 256, got ${handle.chunkSize}`, + ); + + const sessionC2 = await clientB.uploadMedia({ + name: path.basename(payloadB.filePath), + }); + await sessionC2.start({ + uploadSource: clientB.getResumableSource(payloadB.filePath), + resumeUrl: handle.uploadUrl, + chunkSize: handle.chunkSize, + }); + assert.strictEqual(await getFinishedSize(sessionC2), sizeB); + console.log( + ' PASSED Part C: chunk_granularity pause & resume preserved negotiated 256-byte chunkSize and completed 1500-byte upload.', + ); } finally { payloadB.cleanup(); await clientB.close(); diff --git a/core/packages/gax/test/system-test/resumableUpload.ts b/core/packages/gax/test/system-test/resumableUpload.ts index 0b95eb2a5409..595db0e2439f 100644 --- a/core/packages/gax/test/system-test/resumableUpload.ts +++ b/core/packages/gax/test/system-test/resumableUpload.ts @@ -84,6 +84,7 @@ class MockResumableUploadServer { const active = (extra: {[name: string]: string} = {}) => ({ 'x-goog-upload-status': 'active', 'x-goog-upload-size-received': String(this.received.length), + 'x-goog-upload-chunk-granularity': String(GRANULARITY), ...extra, }); @@ -108,6 +109,9 @@ class MockResumableUploadServer { if (offset !== this.received.length) { return send(416, active(), 'offset mismatch'); } + if (command === 'upload' && body.length % GRANULARITY !== 0) { + return send(400, active(), 'chunk size not aligned to granularity'); + } this.received = Buffer.concat([this.received, body]); if (command === 'upload, finalize') { return send( @@ -289,4 +293,44 @@ describe('resumable upload (system)', () => { assert.ok(server.received.equals(data)); assert.strictEqual(resumedProgress[0].bytesUploaded, committed); }); + + it('resumes with a sub-granularity chunkSize and aligns chunks to server granularity', async () => { + server.received = Buffer.alloc(0); + server.commands = []; + + const session1 = new gax.ResumableUploadSession(context); + let pausedHandle: {uploadUrl: string; chunkSize: number} | null = null; + await session1.start({ + uploadSource: gax.resumableSourceFromFile(file), + chunkSize: GRANULARITY / 2, + onProgress: status => { + if (status.bytesUploaded >= GRANULARITY && !pausedHandle) { + pausedHandle = { + uploadUrl: status.uploadUrl, + chunkSize: session1.chunkSize!, + }; + throw new Error('user pause after first aligned chunk'); + } + }, + }); + await assert.rejects( + session1.finished(), + /user pause after first aligned chunk/, + ); + assert.ok(pausedHandle); + const handle = pausedHandle as {uploadUrl: string; chunkSize: number}; + assert.strictEqual(handle.chunkSize, GRANULARITY); + + const session2 = new gax.ResumableUploadSession(context); + await session2.start({ + uploadSource: gax.resumableSourceFromFile(file), + resumeUrl: handle.uploadUrl, + chunkSize: GRANULARITY / 2, + }); + assert.strictEqual(session2.chunkSize, GRANULARITY); + const response = (await session2.finished()) as {status: string}; + + assert.strictEqual(response.status, 'done'); + assert.ok(server.received.equals(data)); + }); }); diff --git a/core/packages/gax/test/unit/resumableUpload.ts b/core/packages/gax/test/unit/resumableUpload.ts index 6d4bbb6ea72d..d9014d4297d7 100644 --- a/core/packages/gax/test/unit/resumableUpload.ts +++ b/core/packages/gax/test/unit/resumableUpload.ts @@ -1664,4 +1664,271 @@ describe('resumable upload', () => { /did not include a x-goog-upload-size-received header/i, ); }); + + it('applies x-goog-upload-chunk-granularity from a recovery query when resuming', async () => { + const requests: MockRequestOptions[] = []; + const observedChunkSizes: Array = []; + const auth = mockAuth(async opts => { + requests.push(opts); + const command = commandOf(opts); + if (command === 'query') { + return resumableUploadResponse(200, { + 'x-goog-upload-status': 'active', + 'x-goog-upload-size-received': String(GRANULARITY), + 'x-goog-upload-chunk-granularity': String(GRANULARITY), + }); + } + if (command === 'upload') { + if (bodyLength(opts) % GRANULARITY !== 0) { + return resumableUploadResponse(400, { + 'x-goog-upload-status': 'active', + }); + } + return resumableUploadResponse(200, {'x-goog-upload-status': 'active'}); + } + if (command === 'upload, finalize') { + return resumableUploadResponse( + 200, + {'x-goog-upload-status': 'final'}, + JSON.stringify({name: 'complete'}), + ); + } + throw new Error(`Unexpected command: ${command}`); + }); + + const payload = Buffer.alloc(2 * GRANULARITY + 128, 0x5a); + const helper = new gax.ResumableUploadSession(buildContext(auth)); + await helper.start({ + uploadSource: bufferSource(payload).source, + chunkSize: GRANULARITY / 2, + resumeUrl: SESSION_URL, + onProgress: () => { + observedChunkSizes.push(helper.chunkSize); + }, + }); + assert.strictEqual(helper.chunkSize, GRANULARITY); + const response = await helper.finished(); + assert.deepStrictEqual(response, {name: 'complete'}); + assert.ok( + observedChunkSizes.every(size => size === GRANULARITY), + `expected all onProgress calls to observe chunkSize=${GRANULARITY}, got ${JSON.stringify(observedChunkSizes)}`, + ); + assert.deepStrictEqual( + requests.map(r => commandOf(r)), + ['query', 'upload', 'upload, finalize'], + ); + assert.strictEqual(offsetOf(requests[1]), GRANULARITY); + assert.strictEqual(bodyLength(requests[1]), GRANULARITY); + assert.strictEqual(offsetOf(requests[2]), 2 * GRANULARITY); + assert.strictEqual(bodyLength(requests[2]), 128); + }); + + it('resumes across sessions using the negotiated chunkSize when query omits granularity', async () => { + const requests: MockRequestOptions[] = []; + let committed = 0; + const payload = Buffer.alloc(2 * GRANULARITY + 64, 0x3c); + + const auth = mockAuth(async opts => { + requests.push(opts); + const command = commandOf(opts); + if (command === 'start') { + return resumableUploadResponse(200, { + 'x-goog-upload-url': SESSION_URL, + 'x-goog-upload-status': 'active', + 'x-goog-upload-chunk-granularity': String(GRANULARITY), + }); + } + if (command === 'query') { + // Server does not repeat x-goog-upload-chunk-granularity on query. + return resumableUploadResponse(200, { + 'x-goog-upload-status': 'active', + 'x-goog-upload-size-received': String(committed), + }); + } + if (command === 'upload') { + const len = bodyLength(opts); + if (len % GRANULARITY !== 0) { + return resumableUploadResponse(400, { + 'x-goog-upload-status': 'active', + }); + } + committed += len; + return resumableUploadResponse(200, {'x-goog-upload-status': 'active'}); + } + if (command === 'upload, finalize') { + committed += bodyLength(opts); + return resumableUploadResponse( + 200, + {'x-goog-upload-status': 'final'}, + JSON.stringify({name: 'complete'}), + ); + } + throw new Error(`Unexpected command: ${command}`); + }); + + const session1 = new gax.ResumableUploadSession(buildContext(auth)); + let pausedHandle: {uploadUrl: string; chunkSize: number} | null = null; + await session1.start({ + uploadSource: bufferSource(payload).source, + chunkSize: GRANULARITY / 2, + onProgress: progress => { + if (progress.bytesUploaded >= GRANULARITY && !pausedHandle) { + pausedHandle = { + uploadUrl: progress.uploadUrl, + chunkSize: session1.chunkSize!, + }; + throw new Error('user pause'); + } + }, + }); + await assert.rejects(session1.finished(), /user pause/); + assert.ok(pausedHandle); + const handle = pausedHandle as {uploadUrl: string; chunkSize: number}; + assert.strictEqual(handle.chunkSize, GRANULARITY); + + const session2 = new gax.ResumableUploadSession(buildContext(auth)); + await session2.start({ + uploadSource: bufferSource(payload).source, + resumeUrl: handle.uploadUrl, + chunkSize: handle.chunkSize, + }); + const response = await session2.finished(); + assert.deepStrictEqual(response, {name: 'complete'}); + assert.strictEqual(committed, payload.length); + }); + + it('restarts chunk transmission with the aligned chunk size when a mid-transfer recovery query returns granularity', async () => { + const requests: MockRequestOptions[] = []; + let firstUploadFailed = false; + const auth = mockAuth(async opts => { + requests.push(opts); + const command = commandOf(opts); + if (command === 'start') { + return resumableUploadResponse(200, { + 'x-goog-upload-url': SESSION_URL, + 'x-goog-upload-status': 'active', + }); + } + if (command === 'upload') { + if (!firstUploadFailed && bodyLength(opts) < GRANULARITY) { + firstUploadFailed = true; + return resumableUploadResponse(400, { + 'x-goog-upload-status': 'active', + }); + } + return resumableUploadResponse(200, {'x-goog-upload-status': 'active'}); + } + if (command === 'query') { + return resumableUploadResponse(200, { + 'x-goog-upload-status': 'active', + 'x-goog-upload-size-received': '0', + 'x-goog-upload-chunk-granularity': String(GRANULARITY), + }); + } + if (command === 'upload, finalize') { + return resumableUploadResponse( + 200, + {'x-goog-upload-status': 'final'}, + JSON.stringify({name: 'complete'}), + ); + } + throw new Error(`Unexpected command: ${command}`); + }); + + const payload = Buffer.alloc(GRANULARITY + 32, 0x77); + const helper = new gax.ResumableUploadSession(buildContext(auth)); + await helper.start({ + uploadSource: bufferSource(payload).source, + chunkSize: GRANULARITY / 2, + }); + const response = await helper.finished(); + assert.deepStrictEqual(response, {name: 'complete'}); + assert.strictEqual(helper.chunkSize, GRANULARITY); + assert.deepStrictEqual( + requests.map(r => commandOf(r)), + ['start', 'upload', 'query', 'upload', 'upload, finalize'], + ); + assert.strictEqual(bodyLength(requests[1]), GRANULARITY / 2); + assert.strictEqual(bodyLength(requests[3]), GRANULARITY); + assert.strictEqual(bodyLength(requests[4]), 32); + }); + + it('retries the start command when the response is missing both status and url headers', async () => { + const requests: MockRequestOptions[] = []; + let startAttempts = 0; + const auth = mockAuth(async opts => { + requests.push(opts); + const command = commandOf(opts); + if (command === 'start') { + startAttempts += 1; + if (startAttempts === 1) { + return resumableUploadResponse(200, {}); + } + return resumableUploadResponse(200, { + 'x-goog-upload-url': SESSION_URL, + 'x-goog-upload-status': 'active', + }); + } + if (command === 'upload, finalize') { + return resumableUploadResponse( + 200, + {'x-goog-upload-status': 'final'}, + JSON.stringify({name: 'complete'}), + ); + } + throw new Error(`Unexpected command: ${command}`); + }); + + const helper = new gax.ResumableUploadSession(buildContext(auth)); + await helper.start({ + uploadSource: bufferSource(Buffer.alloc(64)).source, + chunkSize: GRANULARITY, + retry: { + backoffSettings: { + maxRetries: 2, + initialRetryDelayMillis: 1, + retryDelayMultiplier: 1.0, + maxRetryDelayMillis: 5, + initialRpcTimeoutMillis: 1000, + rpcTimeoutMultiplier: 1.0, + maxRpcTimeoutMillis: 1000, + totalTimeoutMillis: 10000, + }, + }, + }); + const response = await helper.finished(); + assert.deepStrictEqual(response, {name: 'complete'}); + assert.strictEqual(startAttempts, 2); + assert.deepStrictEqual( + requests.map(r => commandOf(r)), + ['start', 'start', 'upload, finalize'], + ); + }); + + it('includes the underlying failure message when retries are exhausted', async () => { + const auth = mockAuth(async () => { + throw new Error('Could not load the default credentials'); + }); + + const helper = new gax.ResumableUploadSession(buildContext(auth)); + await assert.rejects( + helper.start({ + uploadSource: bufferSource(Buffer.alloc(64)).source, + chunkSize: GRANULARITY, + retry: { + backoffSettings: { + maxRetries: 1, + initialRetryDelayMillis: 1, + retryDelayMultiplier: 1.0, + maxRetryDelayMillis: 1, + initialRpcTimeoutMillis: 1000, + rpcTimeoutMultiplier: 1.0, + maxRpcTimeoutMillis: 1000, + totalTimeoutMillis: 10000, + }, + }, + }), + /Exceeded the maximum number of retries \(1\) while sending the resumable upload command "start".*Could not load the default credentials/, + ); + }); }); From daea1c8f1fd6ba8acd92ecec0874a6140ddf0c25 Mon Sep 17 00:00:00 2001 From: feywind <57276408+feywind@users.noreply.github.com> Date: Wed, 30 Sep 2026 18:14:32 +0000 Subject: [PATCH 2/3] chore(gax): use defensive instanceof Error check when formatting retry errors --- core/packages/gax/src/resumableUpload.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/core/packages/gax/src/resumableUpload.ts b/core/packages/gax/src/resumableUpload.ts index ca957515cd0f..98aec1f814b0 100644 --- a/core/packages/gax/src/resumableUpload.ts +++ b/core/packages/gax/src/resumableUpload.ts @@ -1125,7 +1125,7 @@ export class ResumableUploadSession { throw createGoogleError( `Exceeded the maximum number of recovery retries (${retry.maxRetries}) ` + `without forward progress at byte offset ${currentOffset}: ` + - `${(err as Error).message}`, + `${err instanceof Error ? err.message : String(err)}`, Status.DEADLINE_EXCEEDED, ); } @@ -1263,7 +1263,7 @@ export class ResumableUploadSession { throw createGoogleError( `Exceeded the maximum number of recovery retries (${retry.maxRetries}) ` + `while finalizing at byte offset ${currentOffset}: ` + - `${(err as Error).message}`, + `${err instanceof Error ? err.message : String(err)}`, Status.DEADLINE_EXCEEDED, ); } @@ -1345,7 +1345,7 @@ export class ResumableUploadSession { throw createGoogleError( `Exceeded the maximum number of retries (${retry.maxRetries}) ` + `while sending the resumable upload command "${command}": ` + - `${err.message}`, + `${err instanceof Error ? err.message : String(err)}`, Status.DEADLINE_EXCEEDED, ); } From 5e7bd153b3ce5790c672f9665da04634ce534f79 Mon Sep 17 00:00:00 2001 From: feywind <57276408+feywind@users.noreply.github.com> Date: Wed, 30 Sep 2026 22:26:15 +0000 Subject: [PATCH 3/3] fix(gax): address review comments on resumable upload granularity and retry handling --- core/packages/gax/src/resumableUpload.ts | 64 ++-- .../gax/test/system-test/resumableUpload.ts | 16 +- .../packages/gax/test/unit/resumableUpload.ts | 279 ++++++++++++++++-- 3 files changed, 308 insertions(+), 51 deletions(-) diff --git a/core/packages/gax/src/resumableUpload.ts b/core/packages/gax/src/resumableUpload.ts index 98aec1f814b0..268a836c723c 100644 --- a/core/packages/gax/src/resumableUpload.ts +++ b/core/packages/gax/src/resumableUpload.ts @@ -327,6 +327,11 @@ function parseHeaderInt(value: string | null): number | null { return Number.isNaN(parsed) || parsed < 0 ? null : parsed; } +function parsePositiveHeaderInt(value: string | null): number | null { + const parsed = parseHeaderInt(value); + return parsed === null || parsed <= 0 ? null : parsed; +} + /** * Descriptor that identifies a method as a resumable upload method and * provides the caller that constructs the {@link ResumableUploadSession} @@ -439,6 +444,7 @@ export class ResumableUploadSession { private startTimeMs = 0; private globalDeadlineMs = DEFAULT_GLOBAL_DEADLINE_MS; private effectiveChunkSize_ = DEFAULT_CHUNK_SIZE; + private negotiatedGranularity_: number | null = null; private activeAbortController: AbortController | null = null; private activeStream_: NodeJS.ReadableStream | ReadableStream | null = null; private activeIterator_: AsyncIterator | null = null; @@ -548,7 +554,6 @@ export class ResumableUploadSession { this.armDeadlineTimer(); let sessionUrl: string; - let granularity: number | null = null; let finalResponse: {} | undefined; try { @@ -557,10 +562,9 @@ export class ResumableUploadSession { const queried = await this.queryOffset(params.resumeUrl); this.uploadUrl_ = params.resumeUrl; this.committedBytes_ = queried.offset; - granularity = queried.granularity; this.effectiveChunkSize_ = this.computeEffectiveChunkSize( params.chunkSize, - granularity, + queried.granularity, ); this.reportProgress(); sessionUrl = params.resumeUrl; @@ -569,11 +573,10 @@ export class ResumableUploadSession { this.state_ = ResumableUploadState.STARTING; const started = await this.sendStart(); sessionUrl = started.uploadUrl; - granularity = started.granularity; this.uploadUrl_ = sessionUrl; this.effectiveChunkSize_ = this.computeEffectiveChunkSize( params.chunkSize, - granularity, + started.granularity, ); if (started.committedOffset !== undefined) { this.committedBytes_ = started.committedOffset; @@ -641,6 +644,7 @@ export class ResumableUploadSession { if (!granularity || granularity <= 0) { return requested; } + this.negotiatedGranularity_ = granularity; const effective = Math.floor(requested / granularity) * granularity; // Non-final chunks must be a multiple of the server granularity. If the // user's requested chunk size rounds down to zero, fall back to the @@ -725,7 +729,7 @@ export class ResumableUploadSession { `${UPLOAD_URL_HEADER} header.`, ); } - const granularity = parseHeaderInt( + const granularity = parsePositiveHeaderInt( response.headers.get(UPLOAD_CHUNK_GRANULARITY_HEADER), ); // A successful start normally includes X-Goog-Upload-Status. If it is @@ -779,7 +783,7 @@ export class ResumableUploadSession { const size = parseHeaderInt( response.headers.get(UPLOAD_SIZE_RECEIVED_HEADER), ); - const granularity = parseHeaderInt( + const granularity = parsePositiveHeaderInt( response.headers.get(UPLOAD_CHUNK_GRANULARITY_HEADER), ); if (uploadStatus === 'final') { @@ -867,7 +871,6 @@ export class ResumableUploadSession { // buffer. Re-open the source at the committed offset and continue. try { this.committedBytes_ = err.offset; - this.reportProgress(); this.state_ = ResumableUploadState.TRANSMISSION; // Restart the transmission loop from the new offset. await this.runTransmission(sessionUrl); @@ -1084,13 +1087,16 @@ export class ResumableUploadSession { // Outer recovery: query the server for the exact committed offset, // align the local state, and re-enter the transmission phase. this.state_ = ResumableUploadState.RECOVERY; + const prevEffectiveChunkSize = this.effectiveChunkSize_; const queried = await this.queryOffset(sessionUrl); - if (queried.granularity !== null) { + if (queried.granularity !== null && queried.granularity > 0) { this.effectiveChunkSize_ = this.computeEffectiveChunkSize( this.params?.chunkSize, queried.granularity, ); } + const chunkSizeChanged = + this.effectiveChunkSize_ !== prevEffectiveChunkSize; const serverOffset = queried.offset; this.committedBytes_ = serverOffset; this.reportProgress(); @@ -1113,9 +1119,15 @@ export class ResumableUploadSession { } if (serverOffset === currentOffset) { // If the recovery query negotiated a different effective chunk size - // than the chunk currently in flight, restart from the server offset - // so the source is re-read with the aligned chunk size. - if (!isFinal && currentChunk.length !== this.effectiveChunkSize_) { + // or the chunk currently in flight is not a multiple of the + // negotiated granularity, restart from the server offset so the + // source is re-read with an aligned chunk size. + if ( + !isFinal && + (chunkSizeChanged || + (this.negotiatedGranularity_ !== null && + currentChunk.length % this.negotiatedGranularity_ !== 0)) + ) { throw new RestartUploadError(serverOffset); } // Nothing was committed; apply retry limit and backoff before @@ -1125,7 +1137,7 @@ export class ResumableUploadSession { throw createGoogleError( `Exceeded the maximum number of recovery retries (${retry.maxRetries}) ` + `without forward progress at byte offset ${currentOffset}: ` + - `${err instanceof Error ? err.message : String(err)}`, + `${err.message}`, Status.DEADLINE_EXCEEDED, ); } @@ -1184,8 +1196,19 @@ export class ResumableUploadSession { skipAhead: serverOffset - (currentOffset + currentChunk.length), }; } - // The server committed part of this chunk; retransmit the tail. - currentChunk = currentChunk.subarray(serverOffset - currentOffset); + // The server committed part of this chunk; retransmit the tail if + // it remains aligned to the negotiated granularity, or restart from + // the server offset otherwise. + const tail = currentChunk.subarray(serverOffset - currentOffset); + if ( + !isFinal && + (chunkSizeChanged || + (this.negotiatedGranularity_ !== null && + tail.length % this.negotiatedGranularity_ !== 0)) + ) { + throw new RestartUploadError(serverOffset); + } + currentChunk = tail; currentOffset = serverOffset; sameOffsetAttempts = 0; sameOffsetDelay = retry.initialDelayMs; @@ -1240,7 +1263,7 @@ export class ResumableUploadSession { } this.state_ = ResumableUploadState.RECOVERY; const queried = await this.queryOffset(sessionUrl); - if (queried.granularity !== null) { + if (queried.granularity !== null && queried.granularity > 0) { this.effectiveChunkSize_ = this.computeEffectiveChunkSize( this.params?.chunkSize, queried.granularity, @@ -1263,7 +1286,7 @@ export class ResumableUploadSession { throw createGoogleError( `Exceeded the maximum number of recovery retries (${retry.maxRetries}) ` + `while finalizing at byte offset ${currentOffset}: ` + - `${err instanceof Error ? err.message : String(err)}`, + `${err.message}`, Status.DEADLINE_EXCEEDED, ); } @@ -1345,7 +1368,7 @@ export class ResumableUploadSession { throw createGoogleError( `Exceeded the maximum number of retries (${retry.maxRetries}) ` + `while sending the resumable upload command "${command}": ` + - `${err instanceof Error ? err.message : String(err)}`, + `${err.message}`, Status.DEADLINE_EXCEEDED, ); } @@ -1452,8 +1475,7 @@ export class ResumableUploadSession { throw err; } throw new TransientError( - 'Transient failure while sending the resumable upload command ' + - `"${command}": ${(err as Error).message}`, + err instanceof Error ? err.message : String(err), ); } finally { this.clearStallTimer(); @@ -1484,7 +1506,7 @@ export class ResumableUploadSession { ); } if (command === COMMAND_START) { - if (response.headers.get(UPLOAD_URL_HEADER) !== null) { + if (response.headers.get(UPLOAD_URL_HEADER)) { // The start handler reconciles via the returned session URL. return; } diff --git a/core/packages/gax/test/system-test/resumableUpload.ts b/core/packages/gax/test/system-test/resumableUpload.ts index 595db0e2439f..765415aab531 100644 --- a/core/packages/gax/test/system-test/resumableUpload.ts +++ b/core/packages/gax/test/system-test/resumableUpload.ts @@ -299,16 +299,13 @@ describe('resumable upload (system)', () => { server.commands = []; const session1 = new gax.ResumableUploadSession(context); - let pausedHandle: {uploadUrl: string; chunkSize: number} | null = null; + let pausedUrl = ''; await session1.start({ uploadSource: gax.resumableSourceFromFile(file), chunkSize: GRANULARITY / 2, onProgress: status => { - if (status.bytesUploaded >= GRANULARITY && !pausedHandle) { - pausedHandle = { - uploadUrl: status.uploadUrl, - chunkSize: session1.chunkSize!, - }; + if (status.bytesUploaded >= GRANULARITY && !pausedUrl) { + pausedUrl = status.uploadUrl; throw new Error('user pause after first aligned chunk'); } }, @@ -317,14 +314,13 @@ describe('resumable upload (system)', () => { session1.finished(), /user pause after first aligned chunk/, ); - assert.ok(pausedHandle); - const handle = pausedHandle as {uploadUrl: string; chunkSize: number}; - assert.strictEqual(handle.chunkSize, GRANULARITY); + assert.ok(pausedUrl); + assert.strictEqual(session1.chunkSize, GRANULARITY); const session2 = new gax.ResumableUploadSession(context); await session2.start({ uploadSource: gax.resumableSourceFromFile(file), - resumeUrl: handle.uploadUrl, + resumeUrl: pausedUrl, chunkSize: GRANULARITY / 2, }); assert.strictEqual(session2.chunkSize, GRANULARITY); diff --git a/core/packages/gax/test/unit/resumableUpload.ts b/core/packages/gax/test/unit/resumableUpload.ts index d9014d4297d7..def383be42e6 100644 --- a/core/packages/gax/test/unit/resumableUpload.ts +++ b/core/packages/gax/test/unit/resumableUpload.ts @@ -1853,7 +1853,7 @@ describe('resumable upload', () => { assert.strictEqual(bodyLength(requests[4]), 32); }); - it('retries the start command when the response is missing both status and url headers', async () => { + it('retries the start command when the response is missing both status and url headers or has an empty url header', async () => { const requests: MockRequestOptions[] = []; let startAttempts = 0; const auth = mockAuth(async opts => { @@ -1864,6 +1864,9 @@ describe('resumable upload', () => { if (startAttempts === 1) { return resumableUploadResponse(200, {}); } + if (startAttempts === 2) { + return resumableUploadResponse(200, {'x-goog-upload-url': ''}); + } return resumableUploadResponse(200, { 'x-goog-upload-url': SESSION_URL, 'x-goog-upload-status': 'active', @@ -1898,37 +1901,273 @@ describe('resumable upload', () => { }); const response = await helper.finished(); assert.deepStrictEqual(response, {name: 'complete'}); - assert.strictEqual(startAttempts, 2); + assert.strictEqual(startAttempts, 3); assert.deepStrictEqual( requests.map(r => commandOf(r)), - ['start', 'start', 'upload, finalize'], + ['start', 'start', 'start', 'upload, finalize'], ); }); it('includes the underlying failure message when retries are exhausted', async () => { - const auth = mockAuth(async () => { + const fastRetry = { + backoffSettings: { + maxRetries: 1, + initialRetryDelayMillis: 1, + retryDelayMultiplier: 1.0, + maxRetryDelayMillis: 1, + initialRpcTimeoutMillis: 1000, + rpcTimeoutMultiplier: 1.0, + maxRpcTimeoutMillis: 1000, + totalTimeoutMillis: 10000, + }, + }; + + // 1. sendCommandWithRetry exhaustion on start + const startAuth = mockAuth(async () => { throw new Error('Could not load the default credentials'); }); - - const helper = new gax.ResumableUploadSession(buildContext(auth)); + const startHelper = new gax.ResumableUploadSession(buildContext(startAuth)); await assert.rejects( - helper.start({ + startHelper.start({ uploadSource: bufferSource(Buffer.alloc(64)).source, chunkSize: GRANULARITY, - retry: { - backoffSettings: { - maxRetries: 1, - initialRetryDelayMillis: 1, - retryDelayMultiplier: 1.0, - maxRetryDelayMillis: 1, - initialRpcTimeoutMillis: 1000, - rpcTimeoutMultiplier: 1.0, - maxRpcTimeoutMillis: 1000, - totalTimeoutMillis: 10000, - }, - }, + retry: fastRetry, }), - /Exceeded the maximum number of retries \(1\) while sending the resumable upload command "start".*Could not load the default credentials/, + /^Error: Exceeded the maximum number of retries \(1\) while sending the resumable upload command "start": Could not load the default credentials$/, + ); + + // 2. transmitChunk recovery retry exhaustion + const chunkAuth = mockAuth(async opts => { + const command = commandOf(opts); + if (command === 'start') { + return resumableUploadResponse(200, { + 'x-goog-upload-url': SESSION_URL, + 'x-goog-upload-status': 'active', + }); + } + if (command === 'upload, finalize') { + return resumableUploadResponse(412, {'x-goog-upload-status': 'active'}); + } + if (command === 'query') { + return resumableUploadResponse(200, { + 'x-goog-upload-status': 'active', + 'x-goog-upload-size-received': '0', + }); + } + throw new Error(`Unexpected command: ${command}`); + }); + const chunkHelper = new gax.ResumableUploadSession(buildContext(chunkAuth)); + await chunkHelper.start({ + uploadSource: bufferSource(Buffer.alloc(64)).source, + chunkSize: GRANULARITY, + retry: fastRetry, + }); + await assert.rejects( + chunkHelper.finished(), + /^Error: Exceeded the maximum number of recovery retries \(1\) without forward progress at byte offset 0: Resumable upload state mismatch: HTTP 412\.$/, ); + + // 3. transmitFinalize recovery retry exhaustion + const finalizeAuth = mockAuth(async opts => { + const command = commandOf(opts); + if (command === 'start') { + return resumableUploadResponse(200, { + 'x-goog-upload-url': SESSION_URL, + 'x-goog-upload-status': 'active', + }); + } + if (command === 'finalize') { + return resumableUploadResponse(412, {'x-goog-upload-status': 'active'}); + } + if (command === 'query') { + return resumableUploadResponse(200, { + 'x-goog-upload-status': 'active', + 'x-goog-upload-size-received': '0', + }); + } + throw new Error(`Unexpected command: ${command}`); + }); + const finalizeHelper = new gax.ResumableUploadSession( + buildContext(finalizeAuth), + ); + await finalizeHelper.start({ + uploadSource: bufferSource(Buffer.alloc(0)).source, + chunkSize: GRANULARITY, + retry: fastRetry, + }); + await assert.rejects( + finalizeHelper.finished(), + /^Error: Exceeded the maximum number of recovery retries \(1\) while finalizing at byte offset 0: Resumable upload state mismatch: HTTP 412\.$/, + ); + }); + + it('ignores zero x-goog-upload-chunk-granularity header values', async () => { + const auth = mockAuth(async opts => { + const command = commandOf(opts); + if (command === 'start') { + return resumableUploadResponse(200, { + 'x-goog-upload-url': SESSION_URL, + 'x-goog-upload-status': 'active', + 'x-goog-upload-chunk-granularity': String(GRANULARITY), + }); + } + if (command === 'upload') { + if (offsetOf(opts) === 0) { + return resumableUploadResponse(412, { + 'x-goog-upload-status': 'active', + }); + } + return resumableUploadResponse(200, {'x-goog-upload-status': 'active'}); + } + if (command === 'query') { + return resumableUploadResponse(200, { + 'x-goog-upload-status': 'active', + 'x-goog-upload-size-received': String(GRANULARITY), + 'x-goog-upload-chunk-granularity': '0', + }); + } + if (command === 'upload, finalize') { + return resumableUploadResponse( + 200, + {'x-goog-upload-status': 'final'}, + JSON.stringify({name: 'complete'}), + ); + } + throw new Error(`Unexpected command: ${command}`); + }); + + const payload = Buffer.alloc(GRANULARITY + 32, 0x33); + const helper = new gax.ResumableUploadSession(buildContext(auth)); + await helper.start({ + uploadSource: bufferSource(payload).source, + chunkSize: GRANULARITY / 2, + }); + const response = await helper.finished(); + assert.deepStrictEqual(response, {name: 'complete'}); + assert.strictEqual(helper.chunkSize, GRANULARITY); + }); + + it('restarts on unaligned partial commit and preserves aligned partial-commit tail across retries without duplicate progress', async () => { + // Part 1: Partial commit with newly negotiated granularity that leaves an + // unaligned tail must restart from serverOffset instead of sending the tail. + const requests1: MockRequestOptions[] = []; + const progress1: number[] = []; + let firstUpload1 = true; + const auth1 = mockAuth(async opts => { + requests1.push(opts); + const command = commandOf(opts); + if (command === 'start') { + return resumableUploadResponse(200, { + 'x-goog-upload-url': SESSION_URL, + 'x-goog-upload-status': 'active', + }); + } + if (command === 'upload') { + if (firstUpload1) { + firstUpload1 = false; + return resumableUploadResponse(412, { + 'x-goog-upload-status': 'active', + }); + } + return resumableUploadResponse(200, {'x-goog-upload-status': 'active'}); + } + if (command === 'query') { + return resumableUploadResponse(200, { + 'x-goog-upload-status': 'active', + 'x-goog-upload-size-received': '100', + 'x-goog-upload-chunk-granularity': String(GRANULARITY), + }); + } + if (command === 'upload, finalize') { + return resumableUploadResponse( + 200, + {'x-goog-upload-status': 'final'}, + JSON.stringify({name: 'complete'}), + ); + } + throw new Error(`Unexpected command: ${command}`); + }); + + const payload1 = Buffer.alloc(100 + GRANULARITY * 2 + 32, 0x44); + const src1 = bufferSource(payload1); + const helper1 = new gax.ResumableUploadSession(buildContext(auth1)); + await helper1.start({ + uploadSource: src1.source, + chunkSize: GRANULARITY * 2, + onProgress: status => { + progress1.push(status.bytesUploaded); + }, + }); + await helper1.finished(); + assert.strictEqual(src1.streams.length, 2); + assert.strictEqual(bodyLength(requests1[3]), GRANULARITY * 2); + // Progress at offset 100 is reported once (not duplicated by RestartUploadError). + assert.deepStrictEqual(progress1, [ + 100, + 100 + GRANULARITY * 2, + payload1.length, + ]); + + // Part 2: Valid aligned partial-commit tail (768 bytes with granularity 256 + // and chunkSize 1024) that subsequently fails at the same offset does NOT + // trigger a false RestartUploadError. + let uploadCalls2 = 0; + const auth2 = mockAuth(async opts => { + const command = commandOf(opts); + if (command === 'start') { + return resumableUploadResponse(200, { + 'x-goog-upload-url': SESSION_URL, + 'x-goog-upload-status': 'active', + 'x-goog-upload-chunk-granularity': String(GRANULARITY), + }); + } + if (command === 'upload') { + uploadCalls2 += 1; + if (uploadCalls2 <= 2) { + return resumableUploadResponse(412, { + 'x-goog-upload-status': 'active', + }); + } + return resumableUploadResponse(200, {'x-goog-upload-status': 'active'}); + } + if (command === 'query') { + return resumableUploadResponse(200, { + 'x-goog-upload-status': 'active', + 'x-goog-upload-size-received': String(GRANULARITY), + 'x-goog-upload-chunk-granularity': String(GRANULARITY), + }); + } + if (command === 'upload, finalize') { + return resumableUploadResponse( + 200, + {'x-goog-upload-status': 'final'}, + JSON.stringify({name: 'complete'}), + ); + } + throw new Error(`Unexpected command: ${command}`); + }); + + const payload2 = Buffer.alloc(GRANULARITY * 4 + 32, 0x55); + const src2 = bufferSource(payload2); + const helper2 = new gax.ResumableUploadSession(buildContext(auth2)); + await helper2.start({ + uploadSource: src2.source, + chunkSize: GRANULARITY * 4, + retry: { + backoffSettings: { + maxRetries: 2, + initialRetryDelayMillis: 1, + retryDelayMultiplier: 1.0, + maxRetryDelayMillis: 1, + initialRpcTimeoutMillis: 1000, + rpcTimeoutMultiplier: 1.0, + maxRpcTimeoutMillis: 1000, + totalTimeoutMillis: 10000, + }, + }, + }); + await helper2.finished(); + // Stream was only opened once (at offset 0), never reopened at offset 256. + assert.strictEqual(src2.streams.length, 1); }); });