Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
101 changes: 80 additions & 21 deletions core/packages/gax/src/resumableUpload.ts
Original file line number Diff line number Diff line change
Expand Up @@ -291,6 +291,7 @@ interface TransmitResult {

interface QueryOffsetResult {
offset: number;
granularity: number | null;
finalResponse?: {};
}

Expand Down Expand Up @@ -326,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}
Expand Down Expand Up @@ -438,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<unknown> | null = null;
Expand Down Expand Up @@ -547,7 +554,6 @@ export class ResumableUploadSession {
this.armDeadlineTimer();

let sessionUrl: string;
let granularity: number | null = null;
let finalResponse: {} | undefined;

try {
Expand All @@ -556,15 +562,22 @@ export class ResumableUploadSession {
const queried = await this.queryOffset(params.resumeUrl);
this.uploadUrl_ = params.resumeUrl;
this.committedBytes_ = queried.offset;
this.effectiveChunkSize_ = this.computeEffectiveChunkSize(
params.chunkSize,
queried.granularity,
);
this.reportProgress();
sessionUrl = params.resumeUrl;
finalResponse = queried.finalResponse;
} else {
this.state_ = ResumableUploadState.STARTING;
const started = await this.sendStart();
sessionUrl = started.uploadUrl;
granularity = started.granularity;
this.uploadUrl_ = sessionUrl;
this.effectiveChunkSize_ = this.computeEffectiveChunkSize(
params.chunkSize,
started.granularity,
);
if (started.committedOffset !== undefined) {
this.committedBytes_ = started.committedOffset;
this.reportProgress();
Expand All @@ -580,10 +593,6 @@ export class ResumableUploadSession {
throw err;
}

this.effectiveChunkSize_ = this.computeEffectiveChunkSize(
params.chunkSize,
granularity,
);
if (finalResponse !== undefined) {
this.state_ = ResumableUploadState.FINALIZING;
this.done_ = true;
Expand Down Expand Up @@ -635,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
Expand Down Expand Up @@ -719,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
Expand All @@ -730,7 +740,7 @@ export class ResumableUploadSession {
const queried = await this.queryOffset(uploadUrl);
return {
uploadUrl,
granularity,
granularity: granularity ?? queried.granularity,
committedOffset: queried.offset,
finalResponse: queried.finalResponse,
};
Expand Down Expand Up @@ -773,21 +783,24 @@ export class ResumableUploadSession {
const size = parseHeaderInt(
response.headers.get(UPLOAD_SIZE_RECEIVED_HEADER),
);
const granularity = parsePositiveHeaderInt(
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(
'The resumable upload recovery query did not include a ' +
`${UPLOAD_SIZE_RECEIVED_HEADER} header.`,
);
}
return {offset: size};
return {offset: size, granularity};
}

private async runTransmission(sessionUrl: string): Promise<void> {
Expand Down Expand Up @@ -858,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);
Expand Down Expand Up @@ -1075,7 +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 && 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();
Expand All @@ -1097,13 +1118,26 @@ export class ResumableUploadSession {
throw new RestartUploadError(serverOffset);
}
if (serverOffset === currentOffset) {
// If the recovery query negotiated a different effective chunk size
// 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);
Comment thread
feywind marked this conversation as resolved.
}
// 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.message}`,
Status.DEADLINE_EXCEEDED,
);
}
Expand Down Expand Up @@ -1162,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;
Expand Down Expand Up @@ -1218,6 +1263,12 @@ export class ResumableUploadSession {
}
this.state_ = ResumableUploadState.RECOVERY;
const queried = await this.queryOffset(sessionUrl);
if (queried.granularity !== null && queried.granularity > 0) {
this.effectiveChunkSize_ = this.computeEffectiveChunkSize(
this.params?.chunkSize,
queried.granularity,
);
}
const serverOffset = queried.offset;
this.committedBytes_ = serverOffset;
this.reportProgress();
Expand All @@ -1234,7 +1285,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.message}`,
Status.DEADLINE_EXCEEDED,
);
}
Expand Down Expand Up @@ -1315,7 +1367,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,
);
}
Expand Down Expand Up @@ -1416,13 +1469,13 @@ export class ResumableUploadSession {
if (
err instanceof Category2Error ||
err instanceof Category3Error ||
err instanceof GoogleError
err instanceof GoogleError ||
err instanceof TransientError
) {
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();
Expand Down Expand Up @@ -1453,8 +1506,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)) {
// 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.
Expand Down
49 changes: 49 additions & 0 deletions core/packages/gax/test/showcase-resumable-upload/sample.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
40 changes: 40 additions & 0 deletions core/packages/gax/test/system-test/resumableUpload.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
});

Expand All @@ -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(
Expand Down Expand Up @@ -289,4 +293,40 @@ 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 pausedUrl = '';
await session1.start({
uploadSource: gax.resumableSourceFromFile(file),
chunkSize: GRANULARITY / 2,
onProgress: status => {
if (status.bytesUploaded >= GRANULARITY && !pausedUrl) {
pausedUrl = status.uploadUrl;
throw new Error('user pause after first aligned chunk');
}
},
});
await assert.rejects(
session1.finished(),
/user pause after first aligned chunk/,
);
assert.ok(pausedUrl);
assert.strictEqual(session1.chunkSize, GRANULARITY);

const session2 = new gax.ResumableUploadSession(context);
await session2.start({
uploadSource: gax.resumableSourceFromFile(file),
resumeUrl: pausedUrl,
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));
});
});
Loading
Loading