diff --git a/handwritten/spanner/src/helper.ts b/handwritten/spanner/src/helper.ts index 3f4e87e05d2..ac67ae3739d 100644 --- a/handwritten/spanner/src/helper.ts +++ b/handwritten/spanner/src/helper.ts @@ -365,6 +365,51 @@ export function replaceProjectIdToken(value: any, projectId: string): any { return value; } +/** + * Checks whether an input value contains the `{{projectId}}` placeholder. + * + * @param {*} value - The value to inspect. + * @return {boolean} - `true` if any placeholder is found, otherwise `false`. + */ +// eslint-disable-next-line @typescript-eslint/no-explicit-any +export function hasProjectIdToken(value: any): boolean { + if (typeof value === 'string') { + return value.includes(PROJECT_ID_TOKEN); + } + + if ( + value === null || + typeof value !== 'object' || + value instanceof Buffer || + value instanceof Stream || + isDate(value) + ) { + return false; + } + + if (Array.isArray(value)) { + for (let i = 0; i < value.length; i++) { + if (hasProjectIdToken(value[i])) { + return true; + } + } + return false; + } + + for (const key in value) { + if (Object.prototype.hasOwnProperty.call(value, key)) { + if (!KEYS_TO_SCAN.has(key)) { + continue; + } + if (hasProjectIdToken(value[key])) { + return true; + } + } + } + + return false; +} + /** * Custom error type for missing project ID errors. */ diff --git a/handwritten/spanner/src/index.ts b/handwritten/spanner/src/index.ts index 17287fba98b..84c348e48ed 100644 --- a/handwritten/spanner/src/index.ts +++ b/handwritten/spanner/src/index.ts @@ -16,7 +16,7 @@ import {GrpcService, GrpcServiceConfig} from './common-grpc/service'; import {PreciseDate} from '@google-cloud/precise-date'; -import {replaceProjectIdToken} from './helper'; +import {hasProjectIdToken, replaceProjectIdToken} from './helper'; import {promisifyAll} from '@google-cloud/promisify'; import * as extend from 'extend'; import {GoogleAuth, GoogleAuthOptions} from 'google-auth-library'; @@ -331,6 +331,9 @@ class Spanner extends GrpcService { private _metricsEnabled = false; private static _isAFEServerTimingEnabled: boolean | undefined; readonly _nthClientId: number; + private _pendingProjectIdCallbacks?: Array< + (error: Error | null, projectId?: string) => void + >; /** * Placeholder used to auto populate a column with the commit timestamp. @@ -534,7 +537,7 @@ class Spanner extends GrpcService { if (!this.clients_.has(clientName)) { this.clients_.set( clientName, - new v1[clientName](this.options as ClientOptions), + new v1.InstanceAdminClient(this.options as ClientOptions), ); } return this.clients_.get(clientName)! as v1.InstanceAdminClient; @@ -558,7 +561,7 @@ class Spanner extends GrpcService { if (!this.clients_.has(clientName)) { this.clients_.set( clientName, - new v1[clientName](this.options as ClientOptions), + new v1.DatabaseAdminClient(this.options as ClientOptions), ); } return this.clients_.get(clientName)! as v1.DatabaseAdminClient; @@ -615,10 +618,12 @@ class Spanner extends GrpcService { if (callback) { // process.nextTick prevents Unhandled Promise Rejections if callback throws - res.then( - () => process.nextTick(() => callback(null)), - err => process.nextTick(() => callback(err)), - ); + res + .then( + () => process.nextTick(() => callback(null)), + err => process.nextTick(() => callback(err)), + ) + .catch(() => {}); } else { return res; } @@ -1718,51 +1723,136 @@ class Spanner extends GrpcService { * @param {object} config Request config * @param {function} callback Callback function */ - prepareGapicRequest_(config, callback) { - this.auth.getProjectId((err, projectId) => { - if (err) { - callback(err); + prepareGapicRequest_( + config: RequestConfig, + // eslint-disable-next-line @typescript-eslint/no-explicit-any + callback: (err: Error | null, requestFn?: any) => void, + ): void { + if (this.projectId && this.projectIdReplaced_) { + this._prepareGapicRequestWithProjectId(config, this.projectId, callback); + return; + } + if (this._pendingProjectIdCallbacks) { + this._pendingProjectIdCallbacks.push((error, projectId) => { + if (error) { + callback(error); + return; + } + this._prepareGapicRequestWithProjectId(config, projectId!, callback); + }); + return; + } + this._pendingProjectIdCallbacks = []; + this.auth.getProjectId((error, projectId) => { + const pendingCallbacks = this._pendingProjectIdCallbacks || []; + this._pendingProjectIdCallbacks = undefined; + if (error) { + try { + callback(error); + } finally { + for (const pendingCallback of pendingCallbacks) { + try { + pendingCallback(error); + } catch { + // Prevent one failing user callback from stranding subsequent pending callers. + } + } + } return; } - const clientName = config.client; try { - if (!this.clients_.has(clientName)) { - this.clients_.set(clientName, new v1[clientName](this.options)); + this._prepareGapicRequestWithProjectId(config, projectId!, callback); + } finally { + for (const pendingCallback of pendingCallbacks) { + try { + pendingCallback(null, projectId!); + } catch { + // Prevent one failing user callback from stranding subsequent pending callers. + } } - } catch (err) { - callback(err, null); + } + }); + } + + private _prepareGapicRequestWithProjectId( + config: RequestConfig, + projectId: string, + // eslint-disable-next-line @typescript-eslint/no-explicit-any + callback: (err: Error | null, requestFn?: any) => void, + ): void { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + let wrappedRequestFn: any; + try { + const clientName = config.client; + if (!this.clients_.has(clientName)) { + // eslint-disable-next-line import/namespace + this.clients_.set(clientName, new v1[clientName](this.options)); } const gaxClient = this.clients_.get(clientName)!; - let reqOpts = extend(true, {}, config.reqOpts); - reqOpts = replaceProjectIdToken(reqOpts, projectId!); - // It would have been preferable to replace the projectId already in the - // constructor of Spanner, but that is not possible as auth.getProjectId - // is an async method. This is therefore the first place where we have - // access to the value that should be used instead of the placeholder. + let reqOpts = config.reqOpts; + if (!this.projectIdReplaced_ || hasProjectIdToken(reqOpts)) { + reqOpts = extend(true, {}, config.reqOpts); + reqOpts = replaceProjectIdToken(reqOpts, projectId); + } if (!this.projectIdReplaced_) { - this.projectId = replaceProjectIdToken(this.projectId, projectId!); + this.projectId = replaceProjectIdToken(this.projectId, projectId); this.projectFormattedName_ = replaceProjectIdToken( this.projectFormattedName_, - projectId!, + projectId, ); + if ( + this.commonHeaders_[CLOUD_RESOURCE_HEADER]?.includes('{{projectId}}') + ) { + this.commonHeaders_[CLOUD_RESOURCE_HEADER] = replaceProjectIdToken( + this.commonHeaders_[CLOUD_RESOURCE_HEADER], + projectId, + ); + } this.instances_.forEach(instance => { instance.formattedName_ = replaceProjectIdToken( instance.formattedName_, - projectId!, + projectId, ); + if ( + instance.commonHeaders_?.[CLOUD_RESOURCE_HEADER]?.includes( + '{{projectId}}', + ) + ) { + instance.commonHeaders_[CLOUD_RESOURCE_HEADER] = + replaceProjectIdToken( + instance.commonHeaders_[CLOUD_RESOURCE_HEADER], + projectId, + ); + } instance.databases_.forEach(database => { database.formattedName_ = replaceProjectIdToken( database.formattedName_, - projectId!, + projectId, ); + if ( + database.commonHeaders_?.[CLOUD_RESOURCE_HEADER]?.includes( + '{{projectId}}', + ) + ) { + database.commonHeaders_[CLOUD_RESOURCE_HEADER] = + replaceProjectIdToken( + database.commonHeaders_[CLOUD_RESOURCE_HEADER], + projectId, + ); + } }); }); this.projectIdReplaced_ = true; } - config.headers[CLOUD_RESOURCE_HEADER] = replaceProjectIdToken( - config.headers[CLOUD_RESOURCE_HEADER], - projectId!, - ); + if (!config.headers) { + config.headers = {}; + } + if (config.headers[CLOUD_RESOURCE_HEADER]?.includes('{{projectId}}')) { + config.headers[CLOUD_RESOURCE_HEADER] = replaceProjectIdToken( + config.headers[CLOUD_RESOURCE_HEADER], + projectId, + ); + } if (isTracingEnabled(this._observabilityOptions)) { // Do context propagation propagation.inject(context.active(), config.headers, { @@ -1773,26 +1863,34 @@ class Spanner extends GrpcService { // Attach the x-goog-spanner-request-id to the currently active span. attributeXGoogSpannerRequestIdToActiveSpan(config); } - const interceptors: any[] = []; - if (this._metricsEnabled) { - interceptors.push(MetricInterceptor); - } + // eslint-disable-next-line @typescript-eslint/no-explicit-any + const customInterceptors: any[] = + config.gaxOpts?.otherArgs?.options?.interceptors ?? []; + const interceptors = this._metricsEnabled + ? [...customInterceptors, MetricInterceptor] + : customInterceptors; + const headers = Object.assign( + {}, + config.gaxOpts?.otherArgs?.headers, + config.headers, + ); + const options = Object.assign({}, config.gaxOpts?.otherArgs?.options, { + interceptors, + }); + const gaxOpts = extend(true, {}, config.gaxOpts, { + otherArgs: { + headers, + options, + }, + }); const requestFn = gaxClient[config.method].bind( gaxClient, reqOpts, - // Add headers to `gaxOpts` - extend(true, {}, config.gaxOpts, { - otherArgs: { - headers: config.headers, - options: { - interceptors: interceptors, - }, - }, - }), + gaxOpts, ); // Wrap requestFn to inject the spanner request id into every returned error. - const wrappedRequestFn = (...args) => { + wrappedRequestFn = (...args) => { const hasCallback = args && args.length > 0 && @@ -1814,33 +1912,38 @@ class Spanner extends GrpcService { } case false: { - const res = requestFn(...args); - const stream = res as EventEmitter; - if (stream) { - stream.on('error', err => { + let res; + try { + res = requestFn(...args); + } catch (err) { + injectRequestIDIntoError(config, err as Error); + throw err; + } + + if (res instanceof Promise) { + return res.catch(err => { injectRequestIDIntoError(config, err as Error); + throw err; }); } - const originallyPromise = res instanceof Promise; - if (!originallyPromise) { - return res; + const stream = res as EventEmitter; + if (stream && typeof stream.on === 'function') { + stream.on('error', err => { + injectRequestIDIntoError(config, err as Error); + }); } - return new Promise((resolve, reject) => { - requestFn(...args) - .then(resolve) - .catch(err => { - injectRequestIDIntoError(config, err as Error); - reject(err); - }); - }); + return res; } } }; + } catch (error) { + callback(error as Error, null); + return; + } - callback(null, wrappedRequestFn); - }); + callback(null, wrappedRequestFn); } /** @@ -1866,7 +1969,7 @@ class Spanner extends GrpcService { MetricsTracerFactory?.getInstance(this.projectId_)?.createMetricsTracer( config.method, config.reqOpts.database ?? config.reqOpts.session, - config.headers['x-goog-spanner-request-id'], + config.headers?.['x-goog-spanner-request-id'], ) ?? null; } metricsTracer?.recordOperationStart(); @@ -1896,6 +1999,7 @@ class Spanner extends GrpcService { .then(val => { metricsTracer?.recordOperationCompletion(); resolve(val); + return val; }) .catch(error => { metricsTracer?.recordOperationCompletion(); @@ -1934,7 +2038,7 @@ class Spanner extends GrpcService { MetricsTracerFactory?.getInstance(this.projectId_)?.createMetricsTracer( config.method, config.reqOpts.session ?? config.reqOpts.database, - config.headers['x-goog-spanner-request-id'], + config.headers?.['x-goog-spanner-request-id'], ) ?? null; } metricsTracer?.recordOperationStart(); diff --git a/handwritten/spanner/test/helper.ts b/handwritten/spanner/test/helper.ts index 0d9a4cc0a08..668df900db9 100644 --- a/handwritten/spanner/test/helper.ts +++ b/handwritten/spanner/test/helper.ts @@ -16,7 +16,7 @@ import * as assert from 'assert'; import {describe, it} from 'mocha'; -import {replaceProjectIdToken} from '../src/helper'; +import {hasProjectIdToken, replaceProjectIdToken} from '../src/helper'; import {Stream} from 'stream'; describe('helper', () => { @@ -209,4 +209,138 @@ describe('helper', () => { }, /Sorry, we cannot connect to Cloud Services/); }); }); + + describe('hasProjectIdToken', () => { + it('should return true for strings containing placeholder', () => { + assert.strictEqual( + hasProjectIdToken('projects/{{projectId}}/instances'), + true, + ); + assert.strictEqual(hasProjectIdToken('{{projectId}}'), true); + assert.strictEqual( + hasProjectIdToken('prefix-{{projectId}}-suffix'), + true, + ); + }); + + it('should return false for strings without placeholder', () => { + assert.strictEqual( + hasProjectIdToken('projects/my-project/instances'), + false, + ); + assert.strictEqual(hasProjectIdToken(''), false); + }); + + it('should return false for primitive non-string values', () => { + assert.strictEqual(hasProjectIdToken(null), false); + assert.strictEqual(hasProjectIdToken(undefined), false); + assert.strictEqual(hasProjectIdToken(12345), false); + assert.strictEqual(hasProjectIdToken(true), false); + assert.strictEqual(hasProjectIdToken(false), false); + }); + + it('should return false for Buffers, Streams, and Dates', () => { + const buffer = Buffer.from('projects/{{projectId}}'); + const stream = new Stream(); + Object.assign(stream, {session: 'projects/{{projectId}}'}); + const date = new Date(); + Object.assign(date, {session: 'projects/{{projectId}}'}); + + assert.strictEqual(hasProjectIdToken(buffer), false); + assert.strictEqual(hasProjectIdToken(stream), false); + assert.strictEqual(hasProjectIdToken(date), false); + }); + + it('should return true for arrays containing placeholder', () => { + assert.strictEqual( + hasProjectIdToken(['normal-string', 'projects/{{projectId}}']), + true, + ); + assert.strictEqual( + hasProjectIdToken([['nested', 'projects/{{projectId}}']]), + true, + ); + }); + + it('should return false for arrays without placeholder', () => { + assert.strictEqual( + hasProjectIdToken(['normal-string', 'another-string']), + false, + ); + assert.strictEqual(hasProjectIdToken([]), false); + }); + + it('should return true for objects with placeholder in KEYS_TO_SCAN', () => { + assert.strictEqual( + hasProjectIdToken({session: 'projects/{{projectId}}/sessions/123'}), + true, + ); + assert.strictEqual( + hasProjectIdToken({database: 'projects/{{projectId}}/databases/db'}), + true, + ); + assert.strictEqual( + hasProjectIdToken({parent: 'projects/{{projectId}}'}), + true, + ); + assert.strictEqual( + hasProjectIdToken({name: 'projects/{{projectId}}/instances/inst'}), + true, + ); + assert.strictEqual( + hasProjectIdToken({instance: 'projects/{{projectId}}/instances/inst'}), + true, + ); + assert.strictEqual( + hasProjectIdToken({backup: 'projects/{{projectId}}/backups/b1'}), + true, + ); + assert.strictEqual( + hasProjectIdToken({ + config: {parent: 'projects/{{projectId}}'}, + }), + true, + ); + assert.strictEqual( + hasProjectIdToken({ + encryptionConfig: {kmsKeyName: 'projects/{{projectId}}/keys/k1'}, + }), + true, + ); + }); + + it('should return false for objects with placeholder only in non-scanned keys', () => { + const input = { + sql: 'SELECT * FROM users WHERE col = "{{projectId}}"', + query: '{{projectId}}', + params: { + param1: '{{projectId}}', + deep: { + token: '{{projectId}}', + }, + }, + mutations: [ + { + insert: { + table: 'users', + values: ['{{projectId}}'], + }, + }, + ], + }; + + assert.strictEqual(hasProjectIdToken(input), false); + }); + + it('should return false for objects without placeholder', () => { + const input = { + session: 'projects/my-project/sessions/123', + database: 'projects/my-project/databases/db', + parent: 'projects/my-project', + }; + + assert.strictEqual(hasProjectIdToken(input), false); + assert.strictEqual(hasProjectIdToken({}), false); + }); + }); }); diff --git a/handwritten/spanner/test/index.ts b/handwritten/spanner/test/index.ts index 6652d25fec6..d71129b83e4 100644 --- a/handwritten/spanner/test/index.ts +++ b/handwritten/spanner/test/index.ts @@ -24,7 +24,7 @@ import * as proxyquire from 'proxyquire'; import * as through from 'through2'; import {util} from '@google-cloud/common'; import {PreciseDate} from '@google-cloud/precise-date'; -import {replaceProjectIdToken} from '../src/helper'; +import {hasProjectIdToken, replaceProjectIdToken} from '../src/helper'; import * as pfy from '@google-cloud/promisify'; import {grpc} from 'google-gax'; import * as sinon from 'sinon'; @@ -51,6 +51,10 @@ assert.strictEqual(CLOUD_RESOURCE_HEADER, 'google-cloud-resource-prefix'); // eslint-disable-next-line @typescript-eslint/no-var-requires const apiConfig = require('../src/spanner_grpc_config.json'); +if (!('SPANNER_DISABLE_BUILTIN_METRICS' in process.env)) { + process.env.SPANNER_DISABLE_BUILTIN_METRICS = 'false'; +} + async function disableMetrics(sandbox: sinon.SinonSandbox) { sandbox.stub(process.env, 'SPANNER_DISABLE_BUILTIN_METRICS').value('true'); await MetricsTracerFactory.resetInstance(); @@ -193,6 +197,7 @@ describe('Spanner', () => { '@google-cloud/promisify': fakePfy, './helper.js': { replaceProjectIdToken: fakeReplaceProjectIdToken, + hasProjectIdToken, }, 'google-auth-library': { GoogleAuth: fakeGoogleAuth, @@ -2201,11 +2206,6 @@ describe('Spanner', () => { replaceProjectIdTokenOverride = reqOpts => { return reqOpts; }; - const expectedGaxOpts = extend(true, {}, CONFIG.gaxOpts, { - otherArgs: { - headers: CONFIG.headers, - }, - }); FAKE_GAPIC_CLIENT[CONFIG.method] = function (reqOpts, gaxOpts, arg) { assert.strictEqual(this, FAKE_GAPIC_CLIENT); @@ -2223,6 +2223,560 @@ describe('Spanner', () => { requestFn(done); // (FAKE_GAPIC_CLIENT[CONFIG.method]) }); }); + + it('should synchronously return requestFn when project ID is already cached and replaced', done => { + spanner.projectId = PROJECT_ID; + spanner.projectIdReplaced_ = true; + asAny(spanner).auth.getProjectId = sinon.stub().callsFake(() => { + done( + new Error( + 'auth.getProjectId should not be called when project ID is replaced', + ), + ); + }); + + let called = false; + spanner.prepareGapicRequest_(CONFIG, (err, requestFn) => { + assert.ifError(err); + assert.strictEqual(typeof requestFn, 'function'); + called = true; + }); + + assert.strictEqual(called, true); + assert.strictEqual(asAny(spanner).auth.getProjectId.called, false); + done(); + }); + + it('should not clone reqOpts when projectIdReplaced_ is already true', done => { + spanner.projectId = PROJECT_ID; + spanner.projectIdReplaced_ = true; + + FAKE_GAPIC_CLIENT[CONFIG.method] = function (reqOpts) { + assert.strictEqual(reqOpts, CONFIG.reqOpts); + done(); + }; + + spanner.prepareGapicRequest_(CONFIG, (err, requestFn) => { + assert.ifError(err); + requestFn(); + }); + }); + + it('should not call requestFn twice in promise mode', async () => { + let callCount = 0; + FAKE_GAPIC_CLIENT[CONFIG.method] = sinon.spy(async () => { + callCount++; + return 'response-data'; + }); + + return new Promise((resolve, reject) => { + spanner.prepareGapicRequest_(CONFIG, async (err, requestFn) => { + if (err) { + reject(err); + return; + } + try { + const result = await requestFn(); + assert.strictEqual(result, 'response-data'); + assert.strictEqual(callCount, 1); + assert.strictEqual(FAKE_GAPIC_CLIENT[CONFIG.method].callCount, 1); + resolve(); + } catch (error) { + reject(error); + } + }); + }); + }); + + it('should inject request ID and re-throw on promise rejection', async () => { + const apiError = new Error('API failure'); + FAKE_GAPIC_CLIENT[CONFIG.method] = sinon.spy(async () => { + throw apiError; + }); + + return new Promise((resolve, reject) => { + spanner.prepareGapicRequest_(CONFIG, async (err, requestFn) => { + if (err) { + reject(err); + return; + } + try { + await requestFn(); + reject(new Error('Expected requestFn to reject')); + } catch (caughtError) { + assert.strictEqual(caughtError, apiError); + assert.strictEqual(FAKE_GAPIC_CLIENT[CONFIG.method].callCount, 1); + resolve(); + } + }); + }); + }); + + it('should inject request ID and re-throw on synchronous exception in requestFn', done => { + const syncError = new Error('Synchronous failure'); + FAKE_GAPIC_CLIENT[CONFIG.method] = sinon.spy(() => { + throw syncError; + }); + + spanner.prepareGapicRequest_(CONFIG, (err, requestFn) => { + assert.ifError(err); + assert.throws( + () => { + requestFn(); + }, + caughtError => caughtError === syncError, + ); + done(); + }); + }); + + it('should preserve existing headers and options in gaxOpts without deep cloning', done => { + const customConfig = { + ...CONFIG, + gaxOpts: { + timeout: 5000, + otherArgs: { + headers: { + 'x-custom-header': 'custom-val', + }, + options: { + customOption: true, + }, + }, + }, + }; + + FAKE_GAPIC_CLIENT[CONFIG.method] = function (reqOpts, gaxOpts) { + assert.strictEqual(gaxOpts.timeout, 5000); + assert.strictEqual( + gaxOpts.otherArgs.headers['x-custom-header'], + 'custom-val', + ); + assert.strictEqual( + gaxOpts.otherArgs.headers[CLOUD_RESOURCE_HEADER], + 'header', + ); + assert.strictEqual(gaxOpts.otherArgs.options.customOption, true); + assert.ok(Array.isArray(gaxOpts.otherArgs.options.interceptors)); + done(); + }; + + spanner.prepareGapicRequest_(customConfig, (err, requestFn) => { + assert.ifError(err); + requestFn(); + }); + }); + + it('should coalesce concurrent auth.getProjectId calls and replace tokens for both requests', done => { + let getProjectIdCalls = 0; + let authCallback: Function; + asAny(spanner).auth.getProjectId = (callback: Function) => { + getProjectIdCalls++; + authCallback = callback; + }; + + const requestConfig1 = { + client: 'SpannerClient', + method: 'methodName', + reqOpts: { + database: 'projects/{{projectId}}/instances/inst/databases/db1', + }, + }; + + const requestConfig2 = { + client: 'SpannerClient', + method: 'methodName', + reqOpts: { + database: 'projects/{{projectId}}/instances/inst/databases/db2', + }, + }; + + let completed = 0; + const onComplete = () => { + completed++; + if (completed === 2) { + assert.strictEqual(getProjectIdCalls, 1); + done(); + } + }; + + FAKE_GAPIC_CLIENT[CONFIG.method] = function (reqOpts: { + database: string; + }) { + assert.ok(!reqOpts.database.includes('{{projectId}}')); + assert.ok(reqOpts.database.includes(PROJECT_ID)); + }; + + spanner.prepareGapicRequest_( + requestConfig1 as unknown as spnr.RequestConfig, + (err, requestFn) => { + assert.ifError(err); + requestFn(); + onComplete(); + }, + ); + + spanner.prepareGapicRequest_( + requestConfig2 as unknown as spnr.RequestConfig, + (err, requestFn) => { + assert.ifError(err); + requestFn(); + onComplete(); + }, + ); + + assert.strictEqual(getProjectIdCalls, 1); + authCallback!(null, PROJECT_ID); + }); + + it('should propagate auth error to both initial and coalesced requests', done => { + const authError = new Error('Auth failed'); + let authCallback: Function; + asAny(spanner).auth.getProjectId = (callback: Function) => { + authCallback = callback; + }; + + let errorCount = 0; + const onError = (error: Error | null) => { + assert.strictEqual(error, authError); + errorCount++; + if (errorCount === 2) { + done(); + } + }; + + spanner.prepareGapicRequest_(CONFIG, onError); + spanner.prepareGapicRequest_(CONFIG, onError); + + authCallback!(authError); + }); + + it('should not call callback a second time if the callback throws synchronously', done => { + spanner.projectId = PROJECT_ID; + spanner.projectIdReplaced_ = true; + + let callCount = 0; + const throwingCallback = () => { + callCount++; + throw new Error('Callback throws synchronously'); + }; + + assert.throws(() => { + spanner.prepareGapicRequest_(CONFIG, throwingCallback); + }, /Callback throws synchronously/); + + assert.strictEqual(callCount, 1); + done(); + }); + + it('should isolate exceptions in pending callbacks so all pending callbacks are called on success', done => { + let authCallback: Function; + asAny(spanner).auth.getProjectId = (callback: Function) => { + authCallback = callback; + }; + + let secondCallbackCalled = false; + let thirdCallbackCalled = false; + + // First callback (initiator) + spanner.prepareGapicRequest_(CONFIG, () => { + // Initiator succeeds + }); + + // Second callback (throws synchronously) + spanner.prepareGapicRequest_(CONFIG, () => { + secondCallbackCalled = true; + throw new Error('Explosion in user callback'); + }); + + // Third callback (must still be called!) + spanner.prepareGapicRequest_(CONFIG, (err, requestFn) => { + assert.ifError(err); + assert.strictEqual(typeof requestFn, 'function'); + thirdCallbackCalled = true; + assert.ok(secondCallbackCalled); + assert.ok(thirdCallbackCalled); + done(); + }); + + authCallback!(null, PROJECT_ID); + }); + + it('should isolate exceptions in pending callbacks so all pending callbacks are called on auth error', done => { + const authError = new Error('Auth failure'); + let authCallback: Function; + asAny(spanner).auth.getProjectId = (callback: Function) => { + authCallback = callback; + }; + + let firstCallbackCalled = false; + let secondCallbackCalled = false; + let thirdCallbackCalled = false; + + // First callback (throws synchronously) + spanner.prepareGapicRequest_(CONFIG, () => { + firstCallbackCalled = true; + throw new Error('First callback throws'); + }); + + // Second callback (throws synchronously) + spanner.prepareGapicRequest_(CONFIG, () => { + secondCallbackCalled = true; + throw new Error('Second callback throws'); + }); + + // Third callback (must still be called with authError!) + spanner.prepareGapicRequest_(CONFIG, error => { + assert.strictEqual(error, authError); + thirdCallbackCalled = true; + assert.ok(firstCallbackCalled); + assert.ok(secondCallbackCalled); + assert.ok(thirdCallbackCalled); + done(); + }); + + authCallback!(authError); + }); + + it('should replace tokens if reqOpts contains placeholder even when projectIdReplaced_ is true', done => { + spanner.projectId = PROJECT_ID; + spanner.projectIdReplaced_ = true; + + const configWithLateToken = { + client: 'SpannerClient', + method: 'methodName', + reqOpts: { + database: 'projects/{{projectId}}/instances/inst/databases/db-late', + }, + }; + + FAKE_GAPIC_CLIENT[CONFIG.method] = function (reqOpts: { + database: string; + }) { + assert.strictEqual( + reqOpts.database, + `projects/${PROJECT_ID}/instances/inst/databases/db-late`, + ); + done(); + }; + + spanner.prepareGapicRequest_( + configWithLateToken as unknown as spnr.RequestConfig, + (err, requestFn) => { + assert.ifError(err); + requestFn(); + }, + ); + }); + + it('should preserve custom interceptors in gaxOpts options', done => { + const customInterceptor = () => {}; + const configWithInterceptors = { + ...CONFIG, + gaxOpts: { + otherArgs: { + options: { + interceptors: [customInterceptor], + }, + }, + }, + }; + + FAKE_GAPIC_CLIENT[CONFIG.method] = function ( + reqOpts: unknown, + gaxOpts: {otherArgs: {options: {interceptors: Function[]}}}, + ) { + assert.ok( + gaxOpts.otherArgs.options.interceptors.includes(customInterceptor), + ); + done(); + }; + + spanner.prepareGapicRequest_( + configWithInterceptors as unknown as spnr.RequestConfig, + (err, requestFn) => { + assert.ifError(err); + requestFn(); + }, + ); + }); + + it('should handle config without headers without error', done => { + const configWithoutHeaders = { + client: 'SpannerClient', + method: 'methodName', + reqOpts: {a: 'b'}, + }; + + spanner.prepareGapicRequest_( + configWithoutHeaders as unknown as spnr.RequestConfig, + (err, requestFn) => { + assert.ifError(err); + assert.strictEqual(typeof requestFn, 'function'); + done(); + }, + ); + }); + + it('should replace {{projectId}} in CLOUD_RESOURCE_HEADER', done => { + const configWithHeaderToken = { + client: 'SpannerClient', + method: 'methodName', + reqOpts: {}, + headers: { + [CLOUD_RESOURCE_HEADER]: + 'projects/{{projectId}}/instances/inst/databases/db', + }, + }; + + FAKE_GAPIC_CLIENT[CONFIG.method] = function ( + reqOpts: unknown, + gaxOpts: {otherArgs: {headers: {[k: string]: string}}}, + ) { + assert.strictEqual( + gaxOpts.otherArgs.headers[CLOUD_RESOURCE_HEADER], + `projects/${PROJECT_ID}/instances/inst/databases/db`, + ); + done(); + }; + + spanner.prepareGapicRequest_(configWithHeaderToken, (err, requestFn) => { + assert.ifError(err); + requestFn(); + }); + }); + + it('should shallow copy headers so modifying gaxOpts headers does not mutate commonHeaders_ or config.headers', done => { + const initialCommonHeaders = { + 'x-goog-spanner-route-to-leader': 'true', + }; + const config = { + ...CONFIG, + headers: initialCommonHeaders, + }; + + FAKE_GAPIC_CLIENT[CONFIG.method] = function ( + reqOpts: unknown, + gaxOpts: {otherArgs: {headers: {[k: string]: string}}}, + ) { + // Mutate the headers passed to GAPIC + gaxOpts.otherArgs.headers['x-goog-api-client'] = 'gax/1.0.0'; + gaxOpts.otherArgs.headers['x-goog-spanner-route-to-leader'] = 'false'; + + // Verify initialCommonHeaders was not mutated + assert.strictEqual( + initialCommonHeaders['x-goog-spanner-route-to-leader'], + 'true', + ); + assert.strictEqual( + (initialCommonHeaders as Record)['x-goog-api-client'], + undefined, + ); + done(); + }; + + spanner.prepareGapicRequest_(config, (err, requestFn) => { + assert.ifError(err); + requestFn(); + }); + }); + + it('should attach request ID to stream errors', done => { + const {EventEmitter} = require('events'); + const fakeStream = new EventEmitter(); + FAKE_GAPIC_CLIENT[CONFIG.method] = () => fakeStream; + + const config = { + ...CONFIG, + headers: { + 'x-goog-spanner-request-id': 'req-12345', + }, + }; + + spanner.prepareGapicRequest_(config, (err, requestFn) => { + assert.ifError(err); + const stream = requestFn(); + stream.on('error', (error: Error & {requestID?: string}) => { + assert.strictEqual(error.message, 'Stream failed'); + assert.strictEqual(error.requestID, 'req-12345'); + done(); + }); + fakeStream.emit('error', new Error('Stream failed')); + }); + }); + + it('should attach request ID to callback errors', done => { + const apiError = new Error('Callback failed'); + FAKE_GAPIC_CLIENT[CONFIG.method] = ( + reqOpts: unknown, + gaxOpts: unknown, + callback: Function, + ) => { + callback(apiError); + }; + + const config = { + ...CONFIG, + headers: { + 'x-goog-spanner-request-id': 'req-67890', + }, + }; + + spanner.prepareGapicRequest_(config, (err, requestFn) => { + assert.ifError(err); + requestFn((error: Error & {requestID?: string}) => { + assert.strictEqual(error, apiError); + assert.strictEqual(error.requestID, 'req-67890'); + done(); + }); + }); + }); + + it('should update formattedName_ and commonHeaders_ on cached instances and databases', done => { + const fakeDatabase = { + formattedName_: 'projects/{{projectId}}/instances/inst/databases/db', + commonHeaders_: { + [CLOUD_RESOURCE_HEADER]: + 'projects/{{projectId}}/instances/inst/databases/db', + }, + }; + const fakeInstance = { + formattedName_: 'projects/{{projectId}}/instances/inst', + commonHeaders_: { + [CLOUD_RESOURCE_HEADER]: 'projects/{{projectId}}/instances/inst', + }, + databases_: new Map([['db', fakeDatabase]]), + }; + spanner.instances_.set('inst', fakeInstance as unknown as spnr.Instance); + spanner.commonHeaders_ = { + [CLOUD_RESOURCE_HEADER]: 'projects/{{projectId}}', + }; + + spanner.prepareGapicRequest_(CONFIG, err => { + assert.ifError(err); + assert.strictEqual( + spanner.commonHeaders_[CLOUD_RESOURCE_HEADER], + `projects/${PROJECT_ID}`, + ); + assert.strictEqual( + fakeInstance.formattedName_, + `projects/${PROJECT_ID}/instances/inst`, + ); + assert.strictEqual( + fakeInstance.commonHeaders_[CLOUD_RESOURCE_HEADER], + `projects/${PROJECT_ID}/instances/inst`, + ); + assert.strictEqual( + fakeDatabase.formattedName_, + `projects/${PROJECT_ID}/instances/inst/databases/db`, + ); + assert.strictEqual( + fakeDatabase.commonHeaders_[CLOUD_RESOURCE_HEADER], + `projects/${PROJECT_ID}/instances/inst/databases/db`, + ); + done(); + }); + }); }); describe('request', () => { @@ -2301,7 +2855,7 @@ describe('Spanner', () => { }); }); - it('should resolve the promise with the request fn', () => { + it('should resolve the promise with the request fn', async () => { const gapicRequestFnResult = {}; function gapicRequestFn() { @@ -2312,9 +2866,24 @@ describe('Spanner', () => { callback(null, gapicRequestFn); }; - return spanner.request(CONFIG).then(result => { - assert.strictEqual(result, gapicRequestFnResult); - }); + const result = await spanner.request(CONFIG); + assert.strictEqual(result, gapicRequestFnResult); + }); + + it('should handle config without headers when metrics are enabled', async () => { + asAny(spanner)._metricsEnabled = true; + asAny(spanner).projectId_ = 'project-id'; + const configWithoutHeaders = { + client: 'SpannerClient', + method: 'executeSql', + reqOpts: {database: 'db'}, + }; + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => 'ok'); + }; + + const result = await spanner.request(configWithoutHeaders); + assert.strictEqual(result, 'ok'); }); }); }); @@ -2396,6 +2965,26 @@ describe('Spanner', () => { }) .emit('reading'); }); + + it('should handle config without headers when metrics are enabled', done => { + asAny(spanner)._metricsEnabled = true; + asAny(spanner).projectId_ = 'project-id'; + const configWithoutHeaders = { + client: 'SpannerClient', + method: 'executeStreamingSql', + reqOpts: {session: 'session-name'}, + }; + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => through.obj()); + }; + + spanner + .requestStream(configWithoutHeaders) + .on('pipe', () => { + done(); + }) + .emit('reading'); + }); }); describe('close', () => {