Skip to content
Closed
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
12 changes: 9 additions & 3 deletions handwritten/storage/src/bucket.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4505,13 +4505,19 @@ class Bucket extends ServiceObject<Bucket, BucketMetadata> {
if (options.onUploadProgress) {
writable.on('progress', options.onUploadProgress);
}
fs.createReadStream(pathString)
.on('error', bail)
const readStream = fs.createReadStream(pathString);
readStream
.on('error', err => {
readStream.destroy();
writable.destroy();
bail(err);
})
.pipe(writable)
.on('error', err => {
readStream.destroy();
if (
this.storage.retryOptions.autoRetry &&
this.storage.retryOptions.retryableErrorFn!(err)
this.storage.retryOptions.retryableErrorFn!(err as ApiError)
) {
return reject(err);
} else {
Expand Down
76 changes: 75 additions & 1 deletion handwritten/storage/src/file.ts
Original file line number Diff line number Diff line change
Expand Up @@ -148,21 +148,94 @@ export interface SignedPostPolicyV4Output {
url: string;
fields: PolicyFields;
}

export interface GetSignedUrlConfig
extends Pick<SignerGetSignedUrlConfig, 'host' | 'signingEndpoint'> {
/**
* The action to permit with the signed URL.
* - `'read'`: Allows downloading/viewing the file (HTTP GET).
* - `'write'`: Allows uploading/overwriting the file (HTTP PUT).
* - `'delete'`: Allows removing the file (HTTP DELETE).
* - `'resumable'`: Allows resumable uploads (HTTP POST).
* Note: When using `'resumable'`, the header `X-Goog-Resumable: start` must be sent in the client request.
*/
action: 'read' | 'write' | 'delete' | 'resumable';

/**
* The signing version to use.
* @default 'v2'
*/
version?: 'v2' | 'v4';

/**
* Determines the URL structure for accessing bucket resources.
* - `true`: Uses virtual hosted-style URLs (e.g., `https://mybucket.storage.googleapis.com/...`)
* - `false`: Uses path-style URLs (e.g., `https://storage.googleapis.com/mybucket/...`).
* Virtual hosted-style URLs are generally preferred.
* @default false
*/
virtualHostedStyle?: boolean;

/**
* The custom domain name (CNAME) mapped to this bucket (e.g., `"https://cdn.example.com"`).
*/
cname?: string;

/**
* The MD5 digest value in base64. If provided, the client request **must**
* include an identical `Content-MD5` HTTP header.
* If omitted, the client request must not include this header.
*/
contentMd5?: string;

/**
* The expected Content-Type of the file. If provided, the client request **must**
* include an identical `Content-Type` HTTP header.
* If omitted, the client request must not include this header.
*/
contentType?: string;

/**
* The expiration timestamp for the link. Any provided value is passed directly to `new Date()`.
* @throws {Error} If an expiration timestamp from the past is given.
* Note: `'v4'` signing supports a maximum duration of 7 days (604,800 seconds) from the creation time.
*/
expires: string | number | Date;

/**
* The timestamp when this link becomes usable. Any provided value is passed directly to `new Date()`.
* @default Date.now()
* Note: Only supported/applicable when `version` is set to `'v4'`.
*/
accessibleAt?: string | number | Date;

/**
* Canonical extension headers that the server will validate against the client's request.
* Requirements:
* - Header names must be prefixed with `x-goog-` and must be entirely lowercase.
* - Multi-valued headers passed as an array are converted into a comma-separated string (no spaces).
* The client must format them identically to prevent signature mismatches.
*/
extensionHeaders?: http.OutgoingHttpHeaders;

/**
* The filename to prompt the browser/user to save the file as upon access.
* Note: This option is ignored if `responseDisposition` is explicitly set.
*/
promptSaveAs?: string;

/**
* Maps to the `response-content-disposition` query parameter in the signed URL.
*/
responseDisposition?: string;

/**
* Maps to the `response-content-type` query parameter in the signed URL.
*/
responseType?: string;

/**
* Additional query parameters to include natively in the generated signed URL.
*/
queryParams?: Query;
}

Expand Down Expand Up @@ -3236,6 +3309,7 @@ class File extends ServiceObject<File, FileMetadata> {
contentMd5: cfg.contentMd5,
contentType: cfg.contentType,
host: cfg.host,
signingEndpoint: cfg.signingEndpoint,
};

if (cfg.cname) {
Expand Down
71 changes: 65 additions & 6 deletions handwritten/storage/src/resumable-upload.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1386,9 +1386,7 @@ export class Upload extends Writable {

if (retryDelay <= 0) {
this.destroy(
new Error(
`Retry total time limit exceeded - ${JSON.stringify(resp.data)}`,
),
buildRetryError('Retry total time limit exceeded', resp),
);
return;
}
Expand All @@ -1409,9 +1407,7 @@ export class Upload extends Writable {
}
this.numRetries++;
} else {
this.destroy(
new Error(`Retry limit exceeded - ${JSON.stringify(resp.data)}`),
);
this.destroy(buildRetryError('Retry limit exceeded', resp));
}
}

Expand Down Expand Up @@ -1456,6 +1452,69 @@ export class Upload extends Writable {
}
}

function buildRetryError(
prefix: string,
resp: Pick<GaxiosResponse, 'data' | 'status'>,
): Error {
const parts: string[] = [];

if (typeof resp.status === 'number' && !isNaN(resp.status)) {
parts.push(`status: ${resp.status}`);
}

const err = resp.data;
if (err !== undefined && err !== null) {
if (typeof err === 'object') {
const gaxiosErrLike = err as any;
const errParts: string[] = [];
if (gaxiosErrLike.message) {
errParts.push(String(gaxiosErrLike.message));
}
const status = gaxiosErrLike.status ?? gaxiosErrLike.response?.status;
if (typeof status === 'number' && !isNaN(status) && status !== resp.status) {
errParts.push(`status: ${status}`);
}
const statusText = gaxiosErrLike.response?.statusText;
if (statusText) {
errParts.push(`statusText: ${statusText}`);
}
const responseData = gaxiosErrLike.response?.data;
if (responseData !== undefined && responseData !== null && responseData !== '') {
errParts.push(
`response: ${
typeof responseData === 'object'
? JSON.stringify(responseData)
: responseData
}`,
);
}
if (gaxiosErrLike.code) {
errParts.push(`code: ${String(gaxiosErrLike.code)}`);
}

if (errParts.length > 0) {
parts.push(...errParts);
} else if (err instanceof Error) {
parts.push(err.toString() || err.name || 'Unknown Error');
} else {
const stringified = JSON.stringify(err);
if (stringified && stringified !== '{}') {
parts.push(stringified);
}
}
} else if (typeof err === 'string') {
if (err !== '') {
parts.push(err);
}
} else {
parts.push(String(err));
}
}

const suffix = parts.join(' - ');
return new Error(`${prefix} - ${suffix || 'Unknown Error'}`);
}

export function upload(cfg: UploadConfig) {
return new Upload(cfg);
}
Expand Down
2 changes: 1 addition & 1 deletion handwritten/storage/src/signer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ type GoogleAuthLike = Pick<GoogleAuth, 'getCredentials' | 'sign'>;
* @deprecated Use {@link GoogleAuth} instead
*/
export interface AuthClient {
sign(blobToSign: string): Promise<string>;
sign(blobToSign: string, signingEndpoint?: string): Promise<string>;
getCredentials(): Promise<{
client_email?: string;
}>;
Expand Down
46 changes: 46 additions & 0 deletions handwritten/storage/test/bucket.ts
Original file line number Diff line number Diff line change
Expand Up @@ -100,11 +100,18 @@ class FakeNotification {
}

let fsStatOverride: Function | null;
let fsCreateReadStreamOverride: Function | null;
const fakeFs = {
...fs,
stat: (filePath: string, callback: Function) => {
return (fsStatOverride || fs.stat)(filePath, callback);
},
createReadStream: (filePath: string, options?: Parameters<typeof fs.createReadStream>[1]) => {
return (fsCreateReadStreamOverride || fs.createReadStream)(
filePath,
options
);
},
};

let pLimitOverride: Function | null;
Expand Down Expand Up @@ -234,6 +241,7 @@ describe('Bucket', () => {

beforeEach(() => {
fsStatOverride = null;
fsCreateReadStreamOverride = null;
pLimitOverride = null;
bucket = new Bucket(STORAGE, BUCKET_NAME);
});
Expand Down Expand Up @@ -3231,6 +3239,44 @@ describe('Bucket', () => {
});
});

it('should destroy the local read stream if write stream fails', done => {
const fakeFile = new FakeFile(bucket, 'file-name');
const options = {destination: fakeFile, resumable: false};
const originalCreateReadStream = fs.createReadStream;
let readStream: fs.ReadStream;
fsCreateReadStreamOverride = (path: string, opts: any) => {
readStream = originalCreateReadStream(path, opts);
return readStream;
};

fakeFile.createWriteStream = (options_: CreateWriteStreamOptions) => {
const ws = new stream.Writable({
write(chunk, encoding, callback) {
callback(new Error('write error'));
},
});
return ws;
};

const textfilepath = path.join(
getDirName(),
'../../../test/testdata/textfile.txt'
);

bucket.upload(textfilepath, options, (err: Error) => {
try {
assert.strictEqual(err.message, 'write error');
assert.ok(readStream);
assert.ok(readStream.destroyed);
done();
} catch (e) {
done(e);
} finally {
fsCreateReadStreamOverride = null;
}
});
});

it('should allow overriding content type', done => {
const fakeFile = new FakeFile(bucket, 'file-name');
const metadata = {contentType: 'made-up-content-type'};
Expand Down
19 changes: 19 additions & 0 deletions handwritten/storage/test/file.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3794,11 +3794,30 @@ describe('File', () => {
contentType: config.contentType,
cname: CNAME,
virtualHostedStyle: true,
signingEndpoint: undefined,
});
done();
});
});

it('should pass signingEndpoint to URLSigner', done => {
const signingEndpoint = 'https://my-endpoint.com';
const config = {
...SIGNED_URL_CONFIG,
signingEndpoint,
};

file.getSignedUrl(config, (err: Error | null) => {
assert.ifError(err);
const getSignedUrlArgs = signerGetSignedUrlStub.getCall(0).args;
assert.strictEqual(
getSignedUrlArgs[0]['signingEndpoint'],
signingEndpoint
);
done();
});
});

it('should add "x-goog-resumable: start" header if action is resumable', done => {
SIGNED_URL_CONFIG.action = 'resumable';
SIGNED_URL_CONFIG.extensionHeaders = {
Expand Down
Loading
Loading