Skip to content
Draft
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
2 changes: 1 addition & 1 deletion handwritten/storage/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@
"duplexify": "^4.1.3",
"fast-xml-parser": "^5.3.4",
"gaxios": "^6.0.2",
"google-auth-library": "^9.6.3",
"google-auth-library": "^11.0.0",
"html-entities": "^2.5.2",
"mime": "^3.0.0",
"p-limit": "^3.0.1",
Expand Down
2 changes: 1 addition & 1 deletion handwritten/storage/src/nodejs-common/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ export class Service {
private projectIdRequired: boolean;
providedUserAgent?: string;
makeAuthenticatedRequest: MakeAuthenticatedRequest;
authClient: GoogleAuth<AuthClient>;
authClient: GoogleAuth;
apiEndpoint: string;
timeout?: number;
universeDomain: string;
Expand Down
7 changes: 4 additions & 3 deletions handwritten/storage/src/nodejs-common/util.ts
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ export interface MakeAuthenticatedRequest {
getCredentials: (
callback: (err?: Error | null, credentials?: CredentialBody) => void
) => void;
authClient: GoogleAuth<AuthClient>;
authClient: GoogleAuth;
}

export interface Abortable {
Expand Down Expand Up @@ -643,7 +643,7 @@ export class Util {
delete googleAutoAuthConfig.projectId;
}

let authClient: GoogleAuth<AuthClient>;
let authClient: GoogleAuth;

if (googleAutoAuthConfig.authClient instanceof GoogleAuth) {
// Use an existing `GoogleAuth`
Expand All @@ -652,7 +652,8 @@ export class Util {
// Pass an `AuthClient` & `clientOptions` to `GoogleAuth`, if available
authClient = new GoogleAuth({
...googleAutoAuthConfig,
authClient: googleAutoAuthConfig.authClient,
// eslint-disable-next-line @typescript-eslint/no-explicit-any
authClient: googleAutoAuthConfig.authClient as any,
clientOptions: googleAutoAuthConfig.clientOptions,
});
}
Expand Down
55 changes: 33 additions & 22 deletions handwritten/storage/src/resumable-upload.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,6 @@
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

import AbortController from 'abort-controller';
import {createHash} from 'crypto';
import {
GaxiosOptions,
Expand Down Expand Up @@ -99,9 +97,8 @@ export interface UploadConfig extends Pick<WritableOptions, 'highWaterMark'> {
* emulator context is detected.
*/
authClient?: {
request: <T>(
opts: GaxiosOptions
) => Promise<GaxiosResponse<T>> | GaxiosPromise<T>;
// eslint-disable-next-line @typescript-eslint/no-explicit-any, @typescript-eslint/no-unused-vars
request<T = any>(opts: any): Promise<any>;
};

/**
Expand Down Expand Up @@ -301,9 +298,8 @@ export class Upload extends Writable {
* emulator context is detected.
*/
authClient: {
request: <T>(
opts: GaxiosOptions
) => Promise<GaxiosResponse<T>> | GaxiosPromise<T>;
// eslint-disable-next-line @typescript-eslint/no-explicit-any, @typescript-eslint/no-unused-vars
request<T = any>(opts: any): Promise<any>;
};
cacheKey: string;
chunkSize?: number;
Expand Down Expand Up @@ -633,8 +629,16 @@ export class Upload extends Writable {
checksums.push(`md5=${this.#clientMd5Hash}`);
}

if (checksums.length > 0) {
headers!['X-Goog-Hash'] = checksums.join(',');
if (checksums.length > 0 && headers) {
const value = checksums.join(',');

if (headers instanceof Headers) {
headers.set('X-Goog-Hash', value);
} else if (Array.isArray(headers)) {
headers.push(['X-Goog-Hash', value]);
} else {
(headers as Record<string, string>)['X-Goog-Hash'] = value;
}
}
}

Expand Down Expand Up @@ -877,7 +881,8 @@ export class Upload extends Writable {
const res = await this.makeRequest(reqOpts);
// We have successfully got a URI we can now create a new invocation id
this.currentInvocationId.uri = crypto.randomUUID();
return res.headers.location;
const respHeaders = new Headers(res.headers);
return respHeaders.get('location');
} catch (err) {
const e = err as GaxiosError;
const apiError = {
Expand Down Expand Up @@ -908,13 +913,13 @@ export class Upload extends Writable {
}
);

this.uri = uri;
this.uri = uri!;
this.offset = 0;

// emit the newly generated URI for future reuse, if necessary.
this.emit('uri', uri);

return uri;
return uri!;
}

private async continueUploading() {
Expand Down Expand Up @@ -1111,6 +1116,7 @@ export class Upload extends Writable {
return;
}

const respHeaders = new Headers(resp.headers);
// At this point we can safely create a new id for the chunk
this.currentInvocationId.chunk = crypto.randomUUID();

Expand All @@ -1119,7 +1125,7 @@ export class Upload extends Writable {
const shouldContinueWithNextMultiChunkRequest =
this.chunkSize &&
resp.status === RESUMABLE_INCOMPLETE_STATUS_CODE &&
resp.headers.range &&
respHeaders.get('range') &&
moreDataToUpload;

/**
Expand All @@ -1135,7 +1141,7 @@ export class Upload extends Writable {
// Use the upper value in this header to determine where to start the next chunk.
// We should not assume that the server received all bytes sent in the request.
// https://cloud.google.com/storage/docs/performing-resumable-uploads#chunked-upload
const range: string = resp.headers.range;
const range: string = respHeaders.get('range')!;
this.offset = Number(range.split('-')[1]) + 1;

// We should not assume that the server received all bytes sent in the request.
Expand Down Expand Up @@ -1271,8 +1277,9 @@ export class Upload extends Writable {
const resp = await this.checkUploadStatus({retry: false});

if (resp.status === RESUMABLE_INCOMPLETE_STATUS_CODE) {
if (typeof resp.headers.range === 'string') {
this.offset = Number(resp.headers.range.split('-')[1]) + 1;
const respHeaders = new Headers(resp.headers);
if (typeof respHeaders.get('range') === 'string') {
this.offset = Number(respHeaders.get('range')!.split('-')[1]) + 1;
return;
}
}
Expand All @@ -1295,10 +1302,14 @@ export class Upload extends Writable {
private async makeRequest(reqOpts: GaxiosOptions): GaxiosPromise {
if (this.encryption) {
reqOpts.headers = reqOpts.headers || {};
reqOpts.headers['x-goog-encryption-algorithm'] = 'AES256';
reqOpts.headers['x-goog-encryption-key'] = this.encryption.key.toString();
reqOpts.headers['x-goog-encryption-key-sha256'] =
this.encryption.hash.toString();
(reqOpts.headers as Record<string, string>)[
'x-goog-encryption-algorithm'
] = 'AES256';
(reqOpts.headers as Record<string, string>)['x-goog-encryption-key'] =
this.encryption.key.toString();
(reqOpts.headers as Record<string, string>)[
'x-goog-encryption-key-sha256'
] = this.encryption.hash.toString();
}

if (this.userProject) {
Expand Down Expand Up @@ -1353,7 +1364,7 @@ export class Upload extends Writable {
reqOpts.params = reqOpts.params || {};
reqOpts.params.userProject = this.userProject;
}
reqOpts.signal = controller.signal;
reqOpts.signal = controller.signal as AbortSignal;
reqOpts.validateStatus = () => true;

const combinedReqOpts: GaxiosOptions = {
Expand Down
102 changes: 60 additions & 42 deletions handwritten/storage/src/transfer-manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ import {GoogleAuth} from 'google-auth-library';
import {XMLParser, XMLBuilder} from 'fast-xml-parser';
import AsyncRetry from 'async-retry';
import {ApiError} from './nodejs-common/index.js';
import {GaxiosResponse, Headers} from 'gaxios';
import {GaxiosResponse} from 'gaxios';
import {createHash} from 'crypto';
import {GCCL_GCS_CMD_KEY} from './nodejs-common/util.js';
import {getRuntimeTrackingString, getUserAgentString} from './util.js';
Expand Down Expand Up @@ -133,6 +133,10 @@ export interface UploadFileInChunksOptions {
headers?: {[key: string]: string};
}

interface MultiPartUploadErrorResponse {
error?: object;
}

export interface MultiPartUploadHelper {
bucket: Bucket;
fileName: string;
Expand Down Expand Up @@ -220,34 +224,39 @@ class XMLMultiPartUploadHelper implements MultiPartUploadHelper {
};
}

#setGoogApiClientHeaders(headers: Headers = {}): Headers {
#setGoogApiClientHeaders(headers = new Headers()): Headers {
Comment thread
thiyaguk09 marked this conversation as resolved.
let headerFound = false;
let userAgentFound = false;

for (const [key, value] of Object.entries(headers)) {
headers.forEach((value, key) => {
if (key.toLocaleLowerCase().trim() === 'x-goog-api-client') {
headerFound = true;

// Prepend command feature to value, if not already there
if (!value.includes(GCCL_GCS_CMD_FEATURE.UPLOAD_SHARDED)) {
headers[key] =
`${value} gccl-gcs-cmd/${GCCL_GCS_CMD_FEATURE.UPLOAD_SHARDED}`;
headers.set(
key,
`${value} gccl-gcs-cmd/${GCCL_GCS_CMD_FEATURE.UPLOAD_SHARDED}`
);
}
} else if (key.toLocaleLowerCase().trim() === 'user-agent') {
userAgentFound = true;
}
}
});

// If the header isn't present, add it
if (!headerFound) {
headers['x-goog-api-client'] = `${getRuntimeTrackingString()} gccl/${
packageJson.version
} gccl-gcs-cmd/${GCCL_GCS_CMD_FEATURE.UPLOAD_SHARDED}`;
headers.set(
'x-goog-api-client',
`${getRuntimeTrackingString()} gccl/${
packageJson.version
} gccl-gcs-cmd/${GCCL_GCS_CMD_FEATURE.UPLOAD_SHARDED}`
);
}

// If the User-Agent isn't present, add it
if (!userAgentFound) {
headers['User-Agent'] = getUserAgentString();
headers.set('User-Agent', getUserAgentString());
}

return headers;
Expand All @@ -258,21 +267,26 @@ class XMLMultiPartUploadHelper implements MultiPartUploadHelper {
*
* @returns {Promise<void>}
*/
async initiateUpload(headers: Headers = {}): Promise<void> {
async initiateUpload(headers?: {[key: string]: string}): Promise<void> {
const headersObject = new Headers(headers);
const url = `${this.baseUrl}?uploads`;
return AsyncRetry(async bail => {
try {
const res = await this.authClient.request({
headers: this.#setGoogApiClientHeaders(headers),
const res = await this.authClient.request<
string | MultiPartUploadErrorResponse
>({
headers: this.#setGoogApiClientHeaders(headersObject),
method: 'POST',
url,
});

if (res.data && res.data.error) {
throw res.data.error;
if ((res?.data as MultiPartUploadErrorResponse)?.error) {
throw (res.data as MultiPartUploadErrorResponse).error;
}
if (typeof res.data === 'string') {
const parsedXML = this.xmlParser.parse(res.data);
this.uploadId = parsedXML.InitiateMultipartUploadResult.UploadId;
}
const parsedXML = this.xmlParser.parse(res.data);
this.uploadId = parsedXML.InitiateMultipartUploadResult.UploadId;
} catch (e) {
this.#handleErrorResponse(e as Error, bail);
}
Expand All @@ -294,31 +308,32 @@ class XMLMultiPartUploadHelper implements MultiPartUploadHelper {
validation?: 'md5' | 'crc32c' | false
): Promise<void> {
const url = `${this.baseUrl}?partNumber=${partNumber}&uploadId=${this.uploadId}`;
let headers: Headers = this.#setGoogApiClientHeaders();
const headers: Headers = this.#setGoogApiClientHeaders();

if (validation === 'md5') {
const hash = createHash('md5').update(chunk).digest('base64');
headers = {
'Content-MD5': hash,
};
headers.set('Content-MD5', hash);
} else if (validation === 'crc32c') {
const crc = new CRC32C();
crc.update(chunk);
headers['x-goog-hash'] = `crc32c=${crc.toString()}`;
headers.set('x-goog-hash', `crc32c=${crc.toString()}`);
}

return AsyncRetry(async bail => {
try {
const res = await this.authClient.request({
url,
method: 'PUT',
body: chunk,
headers,
});
const res = await this.authClient.request<MultiPartUploadErrorResponse>(
{
url,
method: 'PUT',
body: chunk,
headers,
}
);
if (res.data && res.data.error) {
throw res.data.error;
}
this.partsMap.set(partNumber, res.headers['etag']);
const resHeaders = new Headers(res.headers);
this.partsMap.set(partNumber, resHeaders.get('etag')!);
} catch (e) {
this.#handleErrorResponse(e as Error, bail);
}
Expand All @@ -344,16 +359,18 @@ class XMLMultiPartUploadHelper implements MultiPartUploadHelper {
)}</CompleteMultipartUpload>`;
return AsyncRetry(async bail => {
try {
const res = await this.authClient.request({
headers: this.#setGoogApiClientHeaders(),
url,
method: 'POST',
body,
});
const res = await this.authClient.request<MultiPartUploadErrorResponse>(
{
headers: this.#setGoogApiClientHeaders(),
url,
method: 'POST',
body,
}
);
if (res.data && res.data.error) {
throw res.data.error;
}
return res;
return res as unknown as GaxiosResponse;
} catch (e) {
this.#handleErrorResponse(e as Error, bail);
return;
Expand All @@ -371,16 +388,17 @@ class XMLMultiPartUploadHelper implements MultiPartUploadHelper {
const url = `${this.baseUrl}?uploadId=${this.uploadId}`;
return AsyncRetry(async bail => {
try {
const res = await this.authClient.request({
url,
method: 'DELETE',
});
const res = await this.authClient.request<MultiPartUploadErrorResponse>(
{
url,
method: 'DELETE',
}
);
if (res.data && res.data.error) {
throw res.data.error;
}
} catch (e) {
this.#handleErrorResponse(e as Error, bail);
return;
}
}, this.retryOptions);
}
Expand All @@ -398,7 +416,7 @@ class XMLMultiPartUploadHelper implements MultiPartUploadHelper {
) {
throw err;
} else {
bail(err as Error);
bail(err);
}
}
}
Expand Down
Loading
Loading