From bab2b6048659ecaf07d624014bff0ea3bc14c8cf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Sun, 13 Sep 2026 12:23:13 +0200 Subject: [PATCH] perf(spanner): carry MetricsTracer on call context instead of global registry Pass the MetricsTracer instance directly through gRPC call options rather than registering and retrieving it via a global map keyed by Spanner request ID. - Attach MetricsTracer to call options in request() and requestStream(), resolving it directly in MetricInterceptor. - Eliminate the global _currentOperationTracers map and the background setInterval cleanup timer from MetricsTracerFactory. - Remove per-attempt metadata regex parsing for project ID and request ID in MetricInterceptor. - Record attempt completion before forwarding status to downstream listeners to prevent race conditions during operation completion. - Guard synchronous requestFn invocations to guarantee operation completion metrics on early failures. - Retain getCurrentTracer, clearCurrentTracer, and cleanup interval constants as @deprecated stubs for backwards compatibility. --- .../observability-test/context-isolation.ts | 24 +- handwritten/spanner/src/index.ts | 164 ++++++++-- handwritten/spanner/src/metrics/constants.ts | 7 +- .../spanner/src/metrics/interceptor.ts | 54 ++-- .../src/metrics/metrics-tracer-factory.ts | 111 ++----- .../spanner/src/metrics/metrics-tracer.ts | 54 ++-- handwritten/spanner/test/index.ts | 232 +++++++++++++- .../spanner/test/metrics/interceptor.ts | 230 +++++++++++++- .../test/metrics/metrics-tracer-factory.ts | 154 ++++----- .../spanner/test/metrics/metrics-tracer.ts | 294 +++++++++++++++++- handwritten/spanner/test/metrics/metrics.ts | 247 ++++++++++++++- 11 files changed, 1265 insertions(+), 306 deletions(-) diff --git a/handwritten/spanner/observability-test/context-isolation.ts b/handwritten/spanner/observability-test/context-isolation.ts index 40b5a6761dee..3d4188b9790a 100644 --- a/handwritten/spanner/observability-test/context-isolation.ts +++ b/handwritten/spanner/observability-test/context-isolation.ts @@ -165,37 +165,19 @@ describe('OpenTelemetry Context Isolation Tests', () => { await MetricsTracerFactory.resetInstance(); }); - it('should schedule MetricsTracerFactory cleanup setInterval in ROOT_CONTEXT', () => { + it('should not schedule any background cleanup setInterval', () => { const tracer = trace.getTracer('test'); + const setIntervalStub = sandbox.stub(global, 'setInterval'); - const setIntervalStub = sandbox - .stub(global, 'setInterval') - .callsFake(() => { - const activeSpan = trace.getSpan(context.active()); - - // Assert that the active context is ROOT_CONTEXT (i.e., no active span) - assert.strictEqual( - activeSpan, - undefined, - 'setInterval scheduling must be isolated within ROOT_CONTEXT and not carry any active request span', - ); - return { - unref: () => {}, - } as unknown as NodeJS.Timeout; - }); - - // Start an active request context tracer.startActiveSpan('request-span', span => { try { - // Instantiate the singleton under a request context MetricsTracerFactory.getInstance('mock-project-id'); } finally { span.end(); } }); - // Verify that the cleanup interval was scheduled - assert.strictEqual(setIntervalStub.callCount, 1); + assert.strictEqual(setIntervalStub.callCount, 0); }); }); }); diff --git a/handwritten/spanner/src/index.ts b/handwritten/spanner/src/index.ts index 17287fba98bf..03c85f3080ea 100644 --- a/handwritten/spanner/src/index.ts +++ b/handwritten/spanner/src/index.ts @@ -534,7 +534,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 +558,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,6 +615,7 @@ class Spanner extends GrpcService { if (callback) { // process.nextTick prevents Unhandled Promise Rejections if callback throws + // eslint-disable-next-line promise/catch-or-return res.then( () => process.nextTick(() => callback(null)), err => process.nextTick(() => callback(err)), @@ -1727,6 +1728,7 @@ class Spanner extends GrpcService { const clientName = config.client; try { if (!this.clients_.has(clientName)) { + // eslint-disable-next-line import/namespace this.clients_.set(clientName, new v1[clientName](this.options)); } } catch (err) { @@ -1759,6 +1761,7 @@ class Spanner extends GrpcService { }); this.projectIdReplaced_ = true; } + config.headers = extend(true, {}, config.headers); config.headers[CLOUD_RESOURCE_HEADER] = replaceProjectIdToken( config.headers[CLOUD_RESOURCE_HEADER], projectId!, @@ -1774,7 +1777,10 @@ class Spanner extends GrpcService { attributeXGoogSpannerRequestIdToActiveSpan(config); } const interceptors: any[] = []; - if (this._metricsEnabled) { + if ( + this._metricsEnabled && + (config.client === 'SpannerClient' || config.metricsTracer) + ) { interceptors.push(MetricInterceptor); } const requestFn = gaxClient[config.method].bind( @@ -1786,6 +1792,7 @@ class Spanner extends GrpcService { headers: config.headers, options: { interceptors: interceptors, + metricsTracer: config.metricsTracer, }, }, }), @@ -1843,6 +1850,26 @@ class Spanner extends GrpcService { }); } + private _getResourceName(reqOpts?: { + database?: string | object; + session?: string | object; + name?: string; + }): string { + if (!reqOpts) { + return ''; + } + if (typeof reqOpts.database === 'string') { + return reqOpts.database; + } + if (typeof reqOpts.session === 'string') { + return reqOpts.session; + } + if (typeof reqOpts.name === 'string') { + return reqOpts.name; + } + return ''; + } + /** * Funnel all API requests through this method to be sure we have a project * ID. @@ -1865,22 +1892,48 @@ class Spanner extends GrpcService { metricsTracer = MetricsTracerFactory?.getInstance(this.projectId_)?.createMetricsTracer( config.method, - config.reqOpts.database ?? config.reqOpts.session, - config.headers['x-goog-spanner-request-id'], + this._getResourceName(config.reqOpts), + config.headers?.['x-goog-spanner-request-id'], ) ?? null; } metricsTracer?.recordOperationStart(); + config.metricsTracer = metricsTracer ?? undefined; if (typeof callback === 'function') { this.prepareGapicRequest_(config, (err, requestFn) => { if (err) { callback(err); metricsTracer?.recordOperationCompletion(); } else { - const wrappedCallback = (...args) => { + let callbackInvoked = false; + let callbackThrew = false; + let callbackError: unknown; + const wrappedCallback = (...args: unknown[]) => { + if (callbackInvoked) { + return; + } + callbackInvoked = true; metricsTracer?.recordOperationCompletion(); - callback(...args); + try { + callback(...args); + } catch (error) { + callbackThrew = true; + callbackError = error; + throw error; + } }; - requestFn(wrappedCallback); + try { + requestFn(wrappedCallback); + } catch (error) { + if (callbackThrew) { + throw callbackError; + } + if (callbackInvoked) { + return; + } + callbackInvoked = true; + metricsTracer?.recordOperationCompletion(); + callback(error); + } } }); } else { @@ -1890,20 +1943,27 @@ class Spanner extends GrpcService { metricsTracer?.recordOperationCompletion(); reject(err); } else { - const result = requestFn(); - if (result && typeof result.then === 'function') { - result - .then(val => { - metricsTracer?.recordOperationCompletion(); - resolve(val); - }) - .catch(error => { - metricsTracer?.recordOperationCompletion(); - reject(error); - }); - } else { + try { + const result = requestFn(); + if (result && typeof result.then === 'function') { + result + .then(val => { + metricsTracer?.recordOperationCompletion(); + resolve(val); + return val; + }) + .catch(error => { + metricsTracer?.recordOperationCompletion(); + reject(error); + return null; + }); + } else { + metricsTracer?.recordOperationCompletion(); + resolve(result); + } + } catch (error) { metricsTracer?.recordOperationCompletion(); - resolve(result); + reject(error); } } }); @@ -1933,30 +1993,76 @@ class Spanner extends GrpcService { metricsTracer = MetricsTracerFactory?.getInstance(this.projectId_)?.createMetricsTracer( config.method, - config.reqOpts.session ?? config.reqOpts.database, - config.headers['x-goog-spanner-request-id'], + this._getResourceName(config.reqOpts), + config.headers?.['x-goog-spanner-request-id'], ) ?? null; } metricsTracer?.recordOperationStart(); + config.metricsTracer = metricsTracer ?? undefined; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + let callStream: any = null; + let cleanedUp = false; + const cleanup = () => { + if (cleanedUp) { + return; + } + cleanedUp = true; + if ( + callStream && + typeof callStream.destroy === 'function' && + !callStream.destroyed + ) { + callStream.destroy(); + } + metricsTracer?.recordOperationCompletion(); + }; + const stream = streamEvents(through.obj()); + const origDestroy = stream._destroy; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + stream._destroy = function (err: any, cb: any) { + cleanup(); + if (typeof origDestroy === 'function') { + origDestroy.call(stream, err, cb); + } else if (typeof cb === 'function') { + cb(err); + } + }; stream.once('reading', () => { this.prepareGapicRequest_(config, (err, requestFn) => { + if (stream.destroyed) { + cleanup(); + return; + } if (err) { stream.destroy(err); return; } - requestFn() - .on('error', err => { - stream.destroy(err); - }) - .pipe(stream); + try { + callStream = requestFn(); + if (stream.destroyed) { + cleanup(); + return; + } + if (callStream) { + callStream + .on('error', err => { + stream.destroy(err); + }) + .pipe(stream); + } else { + stream.destroy(new Error('Failed to initialize request stream.')); + } + } catch (error) { + stream.destroy(error as Error); + } }); }); stream.on('finish', () => { stream.destroy(); }); stream.on('close', () => { - metricsTracer?.recordOperationCompletion(); + cleanup(); }); return stream; } diff --git a/handwritten/spanner/src/metrics/constants.ts b/handwritten/spanner/src/metrics/constants.ts index 959eeb3d817b..4d51e3782353 100644 --- a/handwritten/spanner/src/metrics/constants.ts +++ b/handwritten/spanner/src/metrics/constants.ts @@ -19,8 +19,13 @@ import { export const SPANNER_METER_NAME = 'spanner-nodejs'; export const CLIENT_METRICS_PREFIX = 'spanner.googleapis.com/internal/client'; export const SPANNER_RESOURCE_TYPE = 'spanner_instance_client'; -// Maximum time to keep MetricsTracers before considering them stale, and stop tracking them. +/** + * @deprecated No longer used after eliminating the background tracer cleanup timer. + */ export const TRACER_CLEANUP_THRESHOLD_MS = 60 * 60 * 1000; // 60 minutes +/** + * @deprecated No longer used after eliminating the background tracer cleanup timer. + */ export const TRACER_CLEANUP_INTERVAL_MS = 30 * 60 * 1000; // 30 Minutes // OTel semantic conventions // See https://github.com/open-telemetry/opentelemetry-js/blob/main/semantic-conventions/README.md#unstable-semconv diff --git a/handwritten/spanner/src/metrics/interceptor.ts b/handwritten/spanner/src/metrics/interceptor.ts index c3ec57acfd44..40577fa047f4 100644 --- a/handwritten/spanner/src/metrics/interceptor.ts +++ b/handwritten/spanner/src/metrics/interceptor.ts @@ -1,4 +1,4 @@ -// Copyright 2025 Google LLC +// Copyright 2025 Google LLC // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -13,7 +13,6 @@ // limitations under the License. import {grpc} from 'google-gax'; -import {MetricsTracerFactory} from './metrics-tracer-factory'; /** * Interceptor for recording metrics on gRPC calls. @@ -29,18 +28,9 @@ import {MetricsTracerFactory} from './metrics-tracer-factory'; export const MetricInterceptor = (options, nextCall) => { return new grpc.InterceptingCall(nextCall(options), { start: function (metadata, listener, next) { - // Record attempt metric on request start - const resourcePrefix = metadata.get( - 'google-cloud-resource-prefix', - )[0] as string; - const match = resourcePrefix?.match(/^projects\/([^/]+)\//); - const projectId = match ? match[1] : undefined; - let factory; - if (projectId) { - factory = MetricsTracerFactory.getInstance(projectId); - } - const requestId = metadata.get('x-goog-spanner-request-id')[0] as string; - const metricsTracer = factory?.getCurrentTracer(requestId); + // Record attempt metric on request start. + // The tracer is carried directly on the call options. + const metricsTracer = options?.metricsTracer ?? null; metricsTracer?.recordAttemptStart(); const newListener = { onReceiveMetadata: function (metadata, next) { @@ -63,20 +53,30 @@ export const MetricInterceptor = (options, nextCall) => { next(message); }, onReceiveStatus: function (status, next) { - next(status); - - // Record attempt metric completion - metricsTracer?.recordAttemptCompletion(status.code); - if (metricsTracer?.gfeLatency) { - metricsTracer?.recordGfeLatency(status.code); - } else { - metricsTracer?.recordGfeConnectivityErrorCount(status.code); - } - if (metricsTracer?.afeLatency) { - metricsTracer?.recordAfeLatency(status.code); - } else { - metricsTracer?.recordAfeConnectivityErrorCount(status.code); + if (metricsTracer) { + // Record attempt metric completion before notifying downstream listener + metricsTracer.recordAttemptCompletion(status?.code); + if ( + typeof metricsTracer.gfeLatency === 'number' && + Number.isFinite(metricsTracer.gfeLatency) && + metricsTracer.gfeLatency >= 0 + ) { + metricsTracer.recordGfeLatency(status?.code); + } else { + metricsTracer.recordGfeConnectivityErrorCount(status?.code); + } + if ( + typeof metricsTracer.afeLatency === 'number' && + Number.isFinite(metricsTracer.afeLatency) && + metricsTracer.afeLatency >= 0 + ) { + metricsTracer.recordAfeLatency(status?.code); + } else { + metricsTracer.recordAfeConnectivityErrorCount(status?.code); + } } + + next(status); }, }; next(metadata, newListener); diff --git a/handwritten/spanner/src/metrics/metrics-tracer-factory.ts b/handwritten/spanner/src/metrics/metrics-tracer-factory.ts index d4de5f4c234b..0fcc05d7761b 100644 --- a/handwritten/spanner/src/metrics/metrics-tracer-factory.ts +++ b/handwritten/spanner/src/metrics/metrics-tracer-factory.ts @@ -16,7 +16,7 @@ import * as crypto from 'crypto'; import * as os from 'os'; import * as process from 'process'; import {MeterProvider, MetricReader} from '@opentelemetry/sdk-metrics'; -import {Counter, Histogram, context, ROOT_CONTEXT} from '@opentelemetry/api'; +import {Counter, Histogram} from '@opentelemetry/api'; import {detectResources, Resource} from '@opentelemetry/resources'; import {GcpDetectorSync} from '@google-cloud/opentelemetry-resource-util'; import * as Constants from './constants'; @@ -52,9 +52,6 @@ export class MetricsTracerFactory { private _clientUid: string; private _location = 'global'; private _projectId: string; - private _currentOperationTracers = new Map(); - private _currentOperationLastUpdatedMs = new Map(); - private _intervalTracerCleanup: NodeJS.Timeout; public static enabled = true; /** @@ -68,30 +65,20 @@ export class MetricsTracerFactory { this._clientUid = MetricsTracerFactory._generateClientUId(); this._clientName = `${Constants.SPANNER_METER_NAME}/${version}`; - // Only perform async call to retrieve location is metrics are enabled. + // Only perform async call to retrieve location if metrics are enabled. if (MetricsTracerFactory.enabled) { (async () => { const location = await MetricsTracerFactory._detectClientLocation(); this._location = location.length > 0 ? location : 'global'; })().catch(error => { - throw error; + console.warn('Unable to detect client location.', error); + this._location = 'global'; }); } this._clientHash = MetricsTracerFactory._generateClientHash( this._clientUid, ); - - // Start the Tracer cleanup task at an interval - this._intervalTracerCleanup = context.with(ROOT_CONTEXT, () => - setInterval( - this._cleanMetricsTracers.bind(this), - Constants.TRACER_CLEANUP_INTERVAL_MS, - ), - ); - // unref the interval to prevent it from blocking app termination - // in the event loop - this._intervalTracerCleanup.unref(); } /** @@ -101,17 +88,21 @@ export class MetricsTracerFactory { * @param projectId Optional GCP project ID for the factory instantiation. Does nothing for subsequent calls. * @returns The singleton MetricsTracerFactory instance or null if disabled. */ - public static getInstance(projectId: string): MetricsTracerFactory | null { + public static getInstance(projectId?: string): MetricsTracerFactory | null { if (!MetricsTracerFactory.enabled) { return null; } // Create a singleton instance, enabling/disabling metrics can only be done on the initial call if (MetricsTracerFactory._instance === null) { - MetricsTracerFactory._instance = new MetricsTracerFactory(projectId); + MetricsTracerFactory._instance = new MetricsTracerFactory( + projectId || '', + ); + } else if (projectId && !MetricsTracerFactory._instance._projectId) { + MetricsTracerFactory._instance._projectId = projectId; } - return MetricsTracerFactory!._instance; + return MetricsTracerFactory._instance; } /** @@ -143,7 +134,6 @@ export class MetricsTracerFactory { * Resets the singleton instance of the MetricsTracerFactory. */ public static async resetInstance() { - clearInterval(MetricsTracerFactory._instance?._intervalTracerCleanup); await MetricsTracerFactory._instance?.resetMeterProvider(); MetricsTracerFactory._instance = null; } @@ -153,11 +143,9 @@ export class MetricsTracerFactory { */ public async resetMeterProvider() { if (this._meterProvider !== null) { - await this._meterProvider!.shutdown(); + await this._meterProvider.shutdown(); } this._meterProvider = null; - this._currentOperationTracers = new Map(); - this._currentOperationLastUpdatedMs = new Map(); } /** @@ -226,19 +214,14 @@ export class MetricsTracerFactory { public createMetricsTracer( method: string, formattedName: string, - requestId: string, + requestId?: string, ): MetricsTracer | null { if (!MetricsTracerFactory.enabled) { return null; } - const operationRequest = this._extractOperationRequest(requestId); - - if (this._currentOperationTracers.has(operationRequest)) { - return this._currentOperationTracers.get(operationRequest); - } const {instance, database} = this.getInstanceAttributes(formattedName); - const tracer = new MetricsTracer( + return new MetricsTracer( this._instrumentAttemptCounter, this._instrumentAttemptLatency, this._instrumentOperationCounter, @@ -252,11 +235,8 @@ export class MetricsTracerFactory { instance, this._projectId, method, - operationRequest, + requestId, ); - this._currentOperationTracers.set(operationRequest, tracer); - this._currentOperationLastUpdatedMs.set(operationRequest, Date.now()); - return tracer; } /** @@ -283,53 +263,22 @@ export class MetricsTracerFactory { /** * Retrieves the current MetricsTracer for a given request id. - * Returns null if no tracer exists for the request. - * Does not implicitly create MetricsTracers as that should be done - * explicitly using the createMetricsTracer function. - * request id is expected to be as set in the gRPC metadata. + * @deprecated MetricsTracer is carried directly on the call context. * @param requestId The request id of the gRPC call set under 'x-goog-spanner-request-id'. - * @returns The MetricsTracer instance or null if not found. + * @returns null. */ + // eslint-disable-next-line @typescript-eslint/no-unused-vars public getCurrentTracer(requestId: string): MetricsTracer | null { - const operationRequest: string = this._extractOperationRequest(requestId); - if (!this._currentOperationTracers.has(operationRequest)) { - // Attempting to retrieve tracer that doesn't exist. - return null; - } - this._currentOperationLastUpdatedMs.set(operationRequest, Date.now()); - - return this._currentOperationTracers.get(operationRequest) ?? null; + return null; } /** * Removes the MetricsTracer associated with the given request id. + * @deprecated MetricsTracer is carried directly on the call context. * @param requestId The request id of the gRPC call set under 'x-goog-spanner-request-id'. */ - public clearCurrentTracer(requestId: string) { - const operationRequest = - this._extractOperationRequest(requestId) || requestId; - if (!this._currentOperationTracers.has(operationRequest)) { - return; - } - this._currentOperationTracers.delete(operationRequest); - this._currentOperationLastUpdatedMs.delete(operationRequest); - } - - private _extractOperationRequest(requestId: string): string { - if (!requestId) { - return ''; - } - - const regex = /^(\d+\.[a-z0-9]+\.\d+\.\d+\.\d+)\.\d+$/i; - const match = requestId.match(regex); - - if (!match) { - return ''; - } - - const request = match[1]; - return request; - } + // eslint-disable-next-line @typescript-eslint/no-unused-vars + public clearCurrentTracer(requestId: string): void {} /** * Creates and initializes all metric instruments (counters and histograms) for the MeterProvider. @@ -475,20 +424,4 @@ export class MetricsTracerFactory { } return defaultRegion; } - - private _cleanMetricsTracers() { - if (this._currentOperationLastUpdatedMs.size === 0) { - return; - } - - for (const [ - operationTracer, - lastUpdated, - ] of this._currentOperationLastUpdatedMs.entries()) { - if (Date.now() - lastUpdated >= Constants.TRACER_CLEANUP_THRESHOLD_MS) { - this._currentOperationTracers.delete(operationTracer); - this._currentOperationLastUpdatedMs.delete(operationTracer); - } - } - } } diff --git a/handwritten/spanner/src/metrics/metrics-tracer.ts b/handwritten/spanner/src/metrics/metrics-tracer.ts index 709df2987939..2da412280977 100644 --- a/handwritten/spanner/src/metrics/metrics-tracer.ts +++ b/handwritten/spanner/src/metrics/metrics-tracer.ts @@ -14,7 +14,7 @@ import {status as Status} from '@grpc/grpc-js'; import {Counter, Histogram} from '@opentelemetry/api'; -import {MetricsTracerFactory} from './metrics-tracer-factory'; + import { METRIC_LABEL_KEY_DATABASE, METRIC_LABEL_KEY_METHOD, @@ -165,7 +165,7 @@ export class MetricsTracer { private _instance: string, private _projectId: string, private _methodName: string, - private _request: string, + private _request?: string, ) { this._clientAttributes[METRIC_LABEL_KEY_DATABASE] = _database; this._clientAttributes[METRIC_LABEL_KEY_METHOD] = _methodName; @@ -222,8 +222,8 @@ export class MetricsTracer { * Increments the attempt count and creates a new MetricAttemptTracer. */ public recordAttemptStart() { - if (!this.enabled) return; - this.currentOperation!.createNewAttempt(); + if (!this.enabled || !this.currentOperation) return; + this.currentOperation.createNewAttempt(); } /** @@ -232,12 +232,13 @@ export class MetricsTracer { * @param status The status code of the attempt (default: Status.OK). */ public recordAttemptCompletion(statusCode: Status = Status.OK) { - if (!this.enabled) return; - this.currentOperation!.currentAttempt!.status = Status[statusCode]; + if (!this.enabled || !this.currentOperation?.currentAttempt) return; + this.currentOperation.currentAttempt.status = + Status[statusCode] ?? Status[Status.UNKNOWN]; const attemptAttributes = this._createAttemptOtelAttributes(); const endTime = performance.now(); const attemptLatencyMilliseconds = this._getMillisecondTimeDifference( - this.currentOperation!.currentAttempt!.startTime, + this.currentOperation.currentAttempt.startTime, endTime, ); this.instrumentAttemptLatency?.record( @@ -268,7 +269,7 @@ export class MetricsTracer { const endTime = performance.now(); const operationAttributes = this._createOperationOtelAttributes(); const operationLatencyMilliseconds = this._getMillisecondTimeDifference( - this.currentOperation!.startTime, + this.currentOperation.startTime, endTime, ); @@ -277,9 +278,7 @@ export class MetricsTracer { operationLatencyMilliseconds, operationAttributes, ); - MetricsTracerFactory.getInstance(this._projectId)!.clearCurrentTracer( - this._request, - ); + this.currentOperation = null; } /** @@ -319,7 +318,11 @@ export class MetricsTracer { */ public recordGfeLatency(statusCode: Status) { if (!this.enabled) return; - if (!this.gfeLatency) { + if ( + typeof this.gfeLatency !== 'number' || + !Number.isFinite(this.gfeLatency) || + this.gfeLatency < 0 + ) { console.error( 'ERROR: Attempted to record GFE metric with no latency value.', ); @@ -327,7 +330,8 @@ export class MetricsTracer { } const attributes = {...this._clientAttributes}; - attributes[METRIC_LABEL_KEY_STATUS] = Status[statusCode]; + attributes[METRIC_LABEL_KEY_STATUS] = + Status[statusCode] ?? Status[Status.UNKNOWN]; this._instrumentGfeLatency?.record(this.gfeLatency, attributes); this.gfeLatency = null; // Reset latency value @@ -339,7 +343,8 @@ export class MetricsTracer { public recordGfeConnectivityErrorCount(statusCode: Status) { if (!this.enabled) return; const attributes = {...this._clientAttributes}; - attributes[METRIC_LABEL_KEY_STATUS] = Status[statusCode]; + attributes[METRIC_LABEL_KEY_STATUS] = + Status[statusCode] ?? Status[Status.UNKNOWN]; this._instrumentGfeConnectivityErrorCount?.add(1, attributes); } @@ -349,7 +354,8 @@ export class MetricsTracer { public recordAfeConnectivityErrorCount(statusCode: Status) { if (!this.enabled || !Spanner.isAFEServerTimingEnabled()) return; const attributes = {...this._clientAttributes}; - attributes[METRIC_LABEL_KEY_STATUS] = Status[statusCode]; + attributes[METRIC_LABEL_KEY_STATUS] = + Status[statusCode] ?? Status[Status.UNKNOWN]; this._instrumentAfeConnectivityErrorCount?.add(1, attributes); } @@ -359,7 +365,11 @@ export class MetricsTracer { */ public recordAfeLatency(statusCode: Status) { if (!this.enabled || !Spanner.isAFEServerTimingEnabled()) return; - if (!this.afeLatency) { + if ( + typeof this.afeLatency !== 'number' || + !Number.isFinite(this.afeLatency) || + this.afeLatency < 0 + ) { console.error( 'ERROR: Attempted to record AFE metric with no latency value.', ); @@ -367,7 +377,8 @@ export class MetricsTracer { } const attributes = {...this._clientAttributes}; - attributes[METRIC_LABEL_KEY_STATUS] = Status[statusCode]; + attributes[METRIC_LABEL_KEY_STATUS] = + Status[statusCode] ?? Status[Status.UNKNOWN]; this._instrumentAfeLatency?.record(this.afeLatency, attributes); this.afeLatency = null; // Reset latency value @@ -381,7 +392,7 @@ export class MetricsTracer { if (!this.enabled) return {}; const attributes = {...this._clientAttributes}; attributes[METRIC_LABEL_KEY_STATUS] = - this.currentOperation!.currentAttempt?.status ?? Status[Status.UNKNOWN]; + this.currentOperation?.currentAttempt?.status ?? Status[Status.UNKNOWN]; return attributes; } @@ -394,9 +405,12 @@ export class MetricsTracer { private _createAttemptOtelAttributes() { if (!this.enabled) return {}; const attributes = {...this._clientAttributes}; - if (this.currentOperation?.currentAttempt === null) return attributes; + if (!this.currentOperation?.currentAttempt) { + attributes[METRIC_LABEL_KEY_STATUS] = Status[Status.UNKNOWN]; + return attributes; + } attributes[METRIC_LABEL_KEY_STATUS] = - this.currentOperation!.currentAttempt.status; + this.currentOperation.currentAttempt.status; return attributes; } diff --git a/handwritten/spanner/test/index.ts b/handwritten/spanner/test/index.ts index 6652d25fec63..63f5c63fdc0e 100644 --- a/handwritten/spanner/test/index.ts +++ b/handwritten/spanner/test/index.ts @@ -38,6 +38,7 @@ import { GetInstancesOptions, } from '../src'; import {Duplex} from 'stream'; +import {EventEmitter} from 'events'; import {CLOUD_RESOURCE_HEADER, AFE_SERVER_TIMING_HEADER} from '../src/common'; import {MetricsTracerFactory} from '../src/metrics/metrics-tracer-factory'; import IsolationLevel = protos.google.spanner.v1.TransactionOptions.IsolationLevel; @@ -52,12 +53,16 @@ assert.strictEqual(CLOUD_RESOURCE_HEADER, 'google-cloud-resource-prefix'); const apiConfig = require('../src/spanner_grpc_config.json'); async function disableMetrics(sandbox: sinon.SinonSandbox) { + process.env['SPANNER_DISABLE_BUILTIN_METRICS'] = + process.env['SPANNER_DISABLE_BUILTIN_METRICS'] ?? 'false'; sandbox.stub(process.env, 'SPANNER_DISABLE_BUILTIN_METRICS').value('true'); await MetricsTracerFactory.resetInstance(); MetricsTracerFactory.enabled = false; } async function enableMetrics(sandbox: sinon.SinonSandbox) { + process.env['SPANNER_DISABLE_BUILTIN_METRICS'] = + process.env['SPANNER_DISABLE_BUILTIN_METRICS'] ?? 'false'; sandbox.stub(process.env, 'SPANNER_DISABLE_BUILTIN_METRICS').value('false'); await MetricsTracerFactory.resetInstance(); } @@ -2201,12 +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); assert.deepStrictEqual(reqOpts, CONFIG.reqOpts); @@ -2223,6 +2222,30 @@ describe('Spanner', () => { requestFn(done); // (FAKE_GAPIC_CLIENT[CONFIG.method]) }); }); + + it('should not mutate caller-provided headers object', done => { + const originalHeaders = { + [CLOUD_RESOURCE_HEADER]: 'original-header', + 'custom-header': 'custom-value', + }; + const config = { + client: CONFIG.client, + method: CONFIG.method, + reqOpts: CONFIG.reqOpts, + gaxOpts: CONFIG.gaxOpts, + headers: originalHeaders, + }; + + spanner.prepareGapicRequest_(config, err => { + assert.ifError(err); + assert.notStrictEqual(config.headers, originalHeaders); + assert.deepStrictEqual(originalHeaders, { + [CLOUD_RESOURCE_HEADER]: 'original-header', + 'custom-header': 'custom-value', + }); + done(); + }); + }); }); describe('request', () => { @@ -2271,6 +2294,75 @@ describe('Spanner', () => { spanner.request(CONFIG, done); }); + + it('should not execute callback twice if requestFn calls callback and then throws', done => { + let callbackCallCount = 0; + const error = new Error('Synchronous error after callback.'); + + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, (wrappedCallback: (...args: unknown[]) => void) => { + wrappedCallback(null, 'result'); + throw error; + }); + }; + + spanner.request(CONFIG, (err: Error | null, result?: unknown) => { + callbackCallCount++; + assert.strictEqual(err, null); + assert.strictEqual(result, 'result'); + setImmediate(() => { + assert.strictEqual(callbackCallCount, 1); + done(); + }); + }); + }); + + it('should call callback with error if requestFn throws synchronously before callback is invoked', done => { + const error = new Error('Synchronous error before callback.'); + + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => { + throw error; + }); + }; + + spanner.request(CONFIG, (err: Error | null) => { + assert.strictEqual(err, error); + done(); + }); + }); + + it('should rethrow error if user callback throws synchronously', done => { + const error = new Error('User callback threw error.'); + + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, (wrappedCallback: (...args: unknown[]) => void) => { + wrappedCallback(null, 'result'); + }); + }; + + assert.throws(() => { + spanner.request(CONFIG, () => { + throw error; + }); + }, error); + done(); + }); + + it('should overwrite stale metricsTracer with undefined when metrics are not enabled for the call', done => { + const config = { + client: 'DatabaseAdminClient', + metricsTracer: {stale: true}, + }; + + spanner.prepareGapicRequest_ = (cfg, callback) => { + assert.strictEqual(cfg.metricsTracer, undefined); + callback(null, util.noop); + done(); + }; + + spanner.request(config, util.noop); + }); }); describe('promise mode', () => { @@ -2288,20 +2380,39 @@ describe('Spanner', () => { spanner.request(CONFIG); }); - it('should reject the promise', done => { + it('should reject the promise', async () => { const error = new Error('Error.'); spanner.prepareGapicRequest_ = (config, callback) => { callback(error); }; - spanner.request(CONFIG).catch(err => { - assert.strictEqual(err, error); - done(); - }); + await assert.rejects(spanner.request(CONFIG), error); }); - it('should resolve the promise with the request fn', () => { + it('should reject the promise if requestFn throws synchronously', async () => { + const error = new Error('Synchronous requestFn error.'); + + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => { + throw error; + }); + }; + + await assert.rejects(spanner.request(CONFIG), error); + }); + + it('should reject the promise if requestFn returns a rejecting promise', async () => { + const error = new Error('Async requestFn rejection.'); + + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => Promise.reject(error)); + }; + + await assert.rejects(spanner.request(CONFIG), error); + }); + + it('should resolve the promise with the request fn', async () => { const gapicRequestFnResult = {}; function gapicRequestFn() { @@ -2312,9 +2423,8 @@ 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); }); }); }); @@ -2396,6 +2506,98 @@ describe('Spanner', () => { }) .emit('reading'); }); + + it('should destroy the stream if requestFn returns nullish', done => { + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => null); + }; + + spanner + .requestStream(CONFIG) + .on('error', err => { + assert.strictEqual( + err.message, + 'Failed to initialize request stream.', + ); + done(); + }) + .emit('reading'); + }); + + it('should destroy the stream if requestFn throws an error', done => { + const error = new Error('Synchronous initialization failure.'); + + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => { + throw error; + }); + }; + + spanner + .requestStream(CONFIG) + .on('error', err => { + assert.strictEqual(err, error); + done(); + }) + .emit('reading'); + }); + + it('should destroy callStream when requestStream is destroyed', done => { + const fakeCallStream = Object.assign(new EventEmitter(), { + destroy: sinon.spy(), + pipe: sinon.spy(), + }); + + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => fakeCallStream); + }; + + const stream = spanner.requestStream(CONFIG); + stream.emit('reading'); + stream.destroy(); + + setImmediate(() => { + assert.strictEqual(fakeCallStream.destroy.calledOnce, true); + done(); + }); + }); + + it('should not call requestFn if stream is destroyed before prepareGapicRequest_ finishes', done => { + const requestFn = sinon.spy(); + + spanner.prepareGapicRequest_ = (config, callback) => { + setImmediate(() => { + callback(null, requestFn); + assert.strictEqual(requestFn.called, false); + done(); + }); + }; + + const stream = spanner.requestStream(CONFIG); + stream.emit('reading'); + stream.destroy(); + }); + + it('should destroy callStream when stream is destroyed even if removeAllListeners was called', done => { + const fakeCallStream = Object.assign(new EventEmitter(), { + destroy: sinon.spy(), + pipe: sinon.spy(), + }); + + spanner.prepareGapicRequest_ = (config, callback) => { + callback(null, () => fakeCallStream); + }; + + const stream = spanner.requestStream(CONFIG); + stream.emit('reading'); + stream.removeAllListeners(); + stream.destroy(); + + setImmediate(() => { + assert.strictEqual(fakeCallStream.destroy.calledOnce, true); + done(); + }); + }); }); describe('close', () => { diff --git a/handwritten/spanner/test/metrics/interceptor.ts b/handwritten/spanner/test/metrics/interceptor.ts index b28dd95bbfc3..5bf77a7ad837 100644 --- a/handwritten/spanner/test/metrics/interceptor.ts +++ b/handwritten/spanner/test/metrics/interceptor.ts @@ -1,4 +1,4 @@ -// Copyright 2025 Google LLC +// Copyright 2025 Google LLC // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -16,14 +16,12 @@ import * as assert from 'assert'; import * as sinon from 'sinon'; import {grpc} from 'google-gax'; import {status as Status} from '@grpc/grpc-js'; -import {MetricsTracerFactory} from '../../src/metrics/metrics-tracer-factory'; import {MetricsTracer} from '../../src/metrics/metrics-tracer'; import {MetricInterceptor} from '../../src/metrics/interceptor'; describe('MetricInterceptor', () => { let sandbox: sinon.SinonSandbox; let mockMetricsTracer: sinon.SinonStubbedInstance; - let mockFactory: sinon.SinonStubbedInstance; let mockNextCall: sinon.SinonStub; let mockInterceptingCall: any; let mockListener: any; @@ -69,21 +67,16 @@ describe('MetricInterceptor', () => { void >(); - // Mock MetricsTracerFactory - mockFactory = sandbox.createStubInstance(MetricsTracerFactory); - mockFactory.getCurrentTracer = sandbox - .stub() - .returns(mockMetricsTracer) as sinon.SinonStub< - [string], - MetricsTracer | null - >; - sandbox.stub(MetricsTracerFactory, 'getInstance').returns(mockFactory); - // Mock GRPC call components mockInterceptingCall = { start: sinon.spy((metadata: grpc.Metadata, listener: grpc.Listener) => { capturedListener = listener; }), + sendMessageWithContext: sandbox.stub(), + sendMessage: sandbox.stub(), + halfClose: sandbox.stub(), + cancel: sandbox.stub(), + cancelWithStatus: sandbox.stub(), }; mockNextCall = sinon.stub().returns(mockInterceptingCall); @@ -115,6 +108,7 @@ describe('MetricInterceptor', () => { method_definition: { path: '/google.spanner.v1.Spanner/ExecuteSql', }, + metricsTracer: mockMetricsTracer, }; testMetadata = new grpc.Metadata(); testMetadata.set( @@ -179,6 +173,38 @@ describe('MetricInterceptor', () => { ); }); + it('GFE and AFE Metrics - Latency when duration is 0ms', () => { + const zeroLatencyMetadata = new grpc.Metadata(); + zeroLatencyMetadata.set('server-timing', 'gfet4t7; dur=0, afe; dur=0'); + (mockMetricsTracer.extractGfeLatency as sinon.SinonStub).returns(0); + (mockMetricsTracer.extractAfeLatency as sinon.SinonStub).returns(0); + + const interceptingCall = MetricInterceptor(mockOptions, mockNextCall); + interceptingCall.start(testMetadata, mockListener); + + capturedListener.onReceiveMetadata(zeroLatencyMetadata); + capturedListener.onReceiveStatus(mockStatus); + + assert.strictEqual(mockMetricsTracer.recordGfeLatency.callCount, 1); + assert.strictEqual( + mockMetricsTracer.recordGfeLatency.getCall(0).args[0], + Status.OK, + ); + assert.strictEqual( + mockMetricsTracer.recordGfeConnectivityErrorCount.callCount, + 0, + ); + assert.strictEqual(mockMetricsTracer.recordAfeLatency.callCount, 1); + assert.strictEqual( + mockMetricsTracer.recordAfeLatency.getCall(0).args[0], + Status.OK, + ); + assert.strictEqual( + mockMetricsTracer.recordAfeConnectivityErrorCount.callCount, + 0, + ); + }); + it('GFE Metrics - Connectivity Error Count', () => { const interceptingCall = MetricInterceptor(mockOptions, mockNextCall); interceptingCall.start(testMetadata, mockListener); @@ -215,4 +241,182 @@ describe('MetricInterceptor', () => { ); }); }); + + describe('Tracer resolution', () => { + it('should use options.metricsTracer when provided on call options', () => { + const customTracer = sandbox.createStubInstance(MetricsTracer); + customTracer.recordAttemptStart = sandbox.stub<[], void>(); + + const optionsWithTracer = { + ...mockOptions, + metricsTracer: customTracer, + }; + const interceptingCall = MetricInterceptor( + optionsWithTracer, + mockNextCall, + ); + interceptingCall.start(testMetadata, mockListener); + + assert.strictEqual(customTracer.recordAttemptStart.callCount, 1); + assert.strictEqual(mockMetricsTracer.recordAttemptStart.callCount, 0); + }); + + it('should safely proceed with null tracer if options.metricsTracer is not provided', () => { + const optionsWithoutTracer = { + method_definition: { + path: '/google.spanner.v1.Spanner/ExecuteSql', + }, + }; + const interceptingCall = MetricInterceptor( + optionsWithoutTracer, + mockNextCall, + ); + interceptingCall.start(testMetadata, mockListener); + + assert.strictEqual(mockMetricsTracer.recordAttemptStart.callCount, 0); + }); + }); + + describe('Unhappy paths and error handling', () => { + it('should handle undefined call options and proceed cleanly without tracer', () => { + const interceptingCall = MetricInterceptor( + undefined as any, + mockNextCall, + ); + const metadata = new grpc.Metadata(); + interceptingCall.start(metadata, mockListener); + + assert.strictEqual(mockMetricsTracer.recordAttemptStart.callCount, 0); + + assert.doesNotThrow(() => { + capturedListener.onReceiveMetadata(serverTimingMetadata); + capturedListener.onReceiveMessage({data: 'payload'}); + capturedListener.onReceiveStatus(mockStatus); + }); + + assert.strictEqual(mockListener.onReceiveMetadata.callCount, 1); + assert.strictEqual(mockListener.onReceiveMessage.callCount, 1); + assert.strictEqual(mockListener.onReceiveStatus.callCount, 1); + assert.strictEqual( + mockMetricsTracer.recordGfeConnectivityErrorCount.callCount, + 0, + ); + assert.strictEqual( + mockMetricsTracer.recordAfeConnectivityErrorCount.callCount, + 0, + ); + }); + + it('should record attempt completion before invoking downstream listener in onReceiveStatus', () => { + let attemptCompletionRecorded = false; + mockMetricsTracer.recordAttemptCompletion.callsFake(() => { + attemptCompletionRecorded = true; + }); + mockListener.onReceiveStatus.callsFake(() => { + assert.strictEqual( + attemptCompletionRecorded, + true, + 'Attempt completion must be recorded before downstream next(status) is invoked', + ); + }); + + const interceptingCall = MetricInterceptor(mockOptions, mockNextCall); + interceptingCall.start(testMetadata, mockListener); + capturedListener.onReceiveStatus(mockStatus); + + assert.strictEqual(mockListener.onReceiveStatus.callCount, 1); + assert.strictEqual( + mockMetricsTracer.recordAttemptCompletion.callCount, + 1, + ); + }); + + it('should record non-OK status codes when status is PERMISSION_DENIED', () => { + const errorStatus = { + code: Status.PERMISSION_DENIED, + details: 'Permission denied on table', + metadata: new grpc.Metadata(), + }; + const interceptingCall = MetricInterceptor(mockOptions, mockNextCall); + interceptingCall.start(testMetadata, mockListener); + + capturedListener.onReceiveMetadata(emptyMetadata); + capturedListener.onReceiveStatus(errorStatus); + + assert.strictEqual( + mockMetricsTracer.recordAttemptCompletion.callCount, + 1, + ); + assert.strictEqual( + mockMetricsTracer.recordAttemptCompletion.getCall(0).args[0], + Status.PERMISSION_DENIED, + ); + assert.strictEqual( + mockMetricsTracer.recordGfeConnectivityErrorCount.callCount, + 1, + ); + assert.strictEqual( + mockMetricsTracer.recordGfeConnectivityErrorCount.getCall(0).args[0], + Status.PERMISSION_DENIED, + ); + assert.strictEqual( + mockMetricsTracer.recordAfeConnectivityErrorCount.callCount, + 1, + ); + assert.strictEqual( + mockMetricsTracer.recordAfeConnectivityErrorCount.getCall(0).args[0], + Status.PERMISSION_DENIED, + ); + }); + + it('should record non-OK status codes with latency metrics when server-timing is present', () => { + const errorStatus = { + code: Status.UNAVAILABLE, + details: 'Service unavailable', + metadata: new grpc.Metadata(), + }; + const interceptingCall = MetricInterceptor(mockOptions, mockNextCall); + interceptingCall.start(testMetadata, mockListener); + + capturedListener.onReceiveMetadata(serverTimingMetadata); + capturedListener.onReceiveStatus(errorStatus); + + assert.strictEqual( + mockMetricsTracer.recordAttemptCompletion.callCount, + 1, + ); + assert.strictEqual( + mockMetricsTracer.recordAttemptCompletion.getCall(0).args[0], + Status.UNAVAILABLE, + ); + assert.strictEqual(mockMetricsTracer.recordGfeLatency.callCount, 1); + assert.strictEqual( + mockMetricsTracer.recordGfeLatency.getCall(0).args[0], + Status.UNAVAILABLE, + ); + assert.strictEqual(mockMetricsTracer.recordAfeLatency.callCount, 1); + assert.strictEqual( + mockMetricsTracer.recordAfeLatency.getCall(0).args[0], + Status.UNAVAILABLE, + ); + }); + + it('should pass through sendMessage, halfClose, and cancel calls cleanly', () => { + const interceptingCall = MetricInterceptor(mockOptions, mockNextCall); + interceptingCall.sendMessage('test-payload'); + interceptingCall.halfClose(); + interceptingCall.cancelWithStatus(Status.CANCELLED, 'Cancelled'); + + assert.strictEqual( + mockInterceptingCall.sendMessageWithContext.callCount, + 1, + ); + assert.strictEqual( + mockInterceptingCall.sendMessageWithContext.getCall(0).args[1], + 'test-payload', + ); + assert.strictEqual(mockInterceptingCall.halfClose.callCount, 1); + assert.strictEqual(mockInterceptingCall.cancelWithStatus.callCount, 1); + }); + }); }); diff --git a/handwritten/spanner/test/metrics/metrics-tracer-factory.ts b/handwritten/spanner/test/metrics/metrics-tracer-factory.ts index 78dc1d27a132..ed0cbfd49226 100644 --- a/handwritten/spanner/test/metrics/metrics-tracer-factory.ts +++ b/handwritten/spanner/test/metrics/metrics-tracer-factory.ts @@ -158,20 +158,32 @@ describe('MetricsTracerFactory', () => { assert.ok(tracer); }); - it('should clear a MetricsTracer using an extracted operation request id', () => { + it('should create independent MetricsTracer instances directly', () => { const factory = MetricsTracerFactory.getInstance('project-id'); - factory!.createMetricsTracer( + const tracer1 = factory!.createMetricsTracer( + 'some-method', + 'method-name', + '1.1a2bc3d4.1.1.1.1', + ); + const tracer2 = factory!.createMetricsTracer( 'some-method', 'method-name', '1.1a2bc3d4.1.1.1.1', ); - assert.strictEqual((factory as any)._currentOperationTracers.size, 1); - - factory!.clearCurrentTracer('1.1a2bc3d4.1.1.1'); + assert.ok(tracer1); + assert.ok(tracer2); + assert.notStrictEqual(tracer1, tracer2); + assert.strictEqual(factory!.getCurrentTracer('1.1a2bc3d4.1.1.1.1'), null); + assert.doesNotThrow(() => { + factory!.clearCurrentTracer('1.1a2bc3d4.1.1.1.1'); + }); + }); - assert.strictEqual((factory as any)._currentOperationTracers.size, 0); - assert.strictEqual((factory as any)._currentOperationLastUpdatedMs.size, 0); + it('should not schedule background cleanup interval', () => { + const setIntervalSpy = sandbox.spy(global, 'setInterval'); + MetricsTracerFactory.getInstance('project-id'); + assert.strictEqual(setIntervalSpy.called, false); }); it('should correctly set default attributes', () => { @@ -194,6 +206,65 @@ describe('MetricsTracerFactory', () => { 'instance', ); }); + + it('should return null when createMetricsTracer is called and factory is disabled', () => { + const factory = MetricsTracerFactory.getInstance('project-id'); + MetricsTracerFactory.enabled = false; + assert.strictEqual(MetricsTracerFactory.getInstance('project-id'), null); + const tracer = factory!.createMetricsTracer( + 'some-method', + 'projects/project/instances/instance/databases/database', + '1.1a2bc3d4.1.1.1.1', + ); + assert.strictEqual(tracer, null); + MetricsTracerFactory.enabled = true; + }); + + it('should create tracer successfully when requestId is omitted', () => { + const factory = MetricsTracerFactory.getInstance('project-id'); + const tracer = factory!.createMetricsTracer( + 'some-method', + 'projects/project/instances/instance/databases/database', + ); + assert.ok(tracer); + }); + + it('should create tracer with fallback attributes when formattedName is malformed', () => { + const factory = MetricsTracerFactory.getInstance('project-id'); + const tracer = factory!.createMetricsTracer('some-method', ''); + assert.ok(tracer); + assert.strictEqual( + tracer!.clientAttributes[Constants.METRIC_LABEL_KEY_DATABASE], + 'unknown', + ); + assert.strictEqual( + tracer!.clientAttributes[Constants.MONITORED_RES_LABEL_KEY_INSTANCE], + 'unknown', + ); + }); + + it('should update singleton _projectId if initially created without one', async () => { + await MetricsTracerFactory.resetInstance(); + const initialFactory = MetricsTracerFactory.getInstance(); + assert.ok(initialFactory); + assert.strictEqual((initialFactory as any)._projectId, ''); + + const updatedFactory = MetricsTracerFactory.getInstance( + 'resolved-project-id', + ); + assert.strictEqual(updatedFactory, initialFactory); + assert.strictEqual( + (updatedFactory as any)._projectId, + 'resolved-project-id', + ); + + // Subsequent call with different project ID should not overwrite already set project ID + const thirdFactory = MetricsTracerFactory.getInstance( + 'different-project-id', + ); + assert.strictEqual(thirdFactory, initialFactory); + assert.strictEqual((thirdFactory as any)._projectId, 'resolved-project-id'); + }); }); describe('getInstanceAttributes', () => { @@ -204,7 +275,6 @@ describe('getInstanceAttributes', () => { afterEach(async () => { await factory.resetMeterProvider(); - clearInterval(factory['_intervalTracerCleanup']); }); it('should extract project, instance, and database from full resource path', () => { @@ -245,71 +315,3 @@ describe('getInstanceAttributes', () => { }); }); }); - -describe('MetricsTracerFactory with set clock', () => { - let clock: sinon.SinonFakeTimers; - - beforeEach(async () => { - MetricsTracerFactory.enabled = true; - await MetricsTracerFactory.resetInstance(); - // Use fake timers to control the clock - clock = sinon.useFakeTimers(); - }); - - afterEach(() => { - // Restore the real timers - clock.restore(); - }); - - describe('_cleanMetricTracers', () => { - it('should prune stale tracers', () => { - const factory = MetricsTracerFactory.getInstance('test-project'); - assert(factory); - - factory.createMetricsTracer( - 'method1', - 'projects/p/instances/i/databases/d', - '1.1a2b3c.1.1.1.1', - ); - - // Advance the clock to make the tracer stale - clock.tick(Constants.TRACER_CLEANUP_THRESHOLD_MS); - - // Add another tracer to trigger pruning - factory.createMetricsTracer( - 'method2', - 'projects/p/instances/i/databases/d', - '2.1a2b3c.1.1.1.1', - ); - // Only most recent tracer should remain - assert.strictEqual(factory['_currentOperationTracers'].size, 1); - assert.ok(factory['_currentOperationTracers'].has('2.1a2b3c.1.1.1')); - }); - - it('should not prune recent tracers', () => { - const factory = MetricsTracerFactory.getInstance('test-project'); - assert(factory); - - factory.createMetricsTracer( - 'method1', - 'projects/p/instances/i/databases/d', - '1.1a2b3c.1.1.1.1', - ); - - // Advance the clock, but not enough to hit the threshold - clock.tick(Constants.TRACER_CLEANUP_INTERVAL_MS); - - // Add another tracer to trigger pruning - factory.createMetricsTracer( - 'method2', - 'projects/p/instances/i/databases/d', - '2.1a2b3c.1.1.1.1', - ); - - // Both tracers should be available - assert.strictEqual(factory['_currentOperationTracers'].size, 2); - assert.ok(factory['_currentOperationTracers'].has('1.1a2b3c.1.1.1')); - assert.ok(factory['_currentOperationTracers'].has('2.1a2b3c.1.1.1')); - }); - }); -}); diff --git a/handwritten/spanner/test/metrics/metrics-tracer.ts b/handwritten/spanner/test/metrics/metrics-tracer.ts index 3cd5530f5bca..a099b3f9374d 100644 --- a/handwritten/spanner/test/metrics/metrics-tracer.ts +++ b/handwritten/spanner/test/metrics/metrics-tracer.ts @@ -17,8 +17,6 @@ import * as assert from 'assert'; import * as sinon from 'sinon'; import * as Constants from '../../src/metrics/constants'; import {MetricsTracer} from '../../src/metrics/metrics-tracer'; - -import {MetricsTracerFactory} from '../../src/metrics/metrics-tracer-factory'; import {Spanner} from '../../src'; const DATABASE = 'test-db'; @@ -134,15 +132,57 @@ describe('MetricsTracer', () => { tracer.recordAttemptCompletion(Status.OK); assert.strictEqual(fakeAttemptLatency.record.called, false); }); + + it('should record attempt error status when status is not OK', () => { + tracer.recordOperationStart(); + tracer.recordAttemptStart(); + tracer.recordAttemptCompletion(Status.PERMISSION_DENIED); + + assert.strictEqual(fakeAttemptLatency.record.calledOnce, true); + assert.strictEqual(fakeAttemptCounter.add.calledOnce, true); + const [[latency, latencyAttributes]] = fakeAttemptLatency.record.args; + const [[count, countAttributes]] = fakeAttemptCounter.add.args; + assert.strictEqual(typeof latency, 'number'); + assert.strictEqual(count, 1); + assert.strictEqual( + latencyAttributes[Constants.METRIC_LABEL_KEY_STATUS], + 'PERMISSION_DENIED', + ); + assert.strictEqual( + countAttributes[Constants.METRIC_LABEL_KEY_STATUS], + 'PERMISSION_DENIED', + ); + }); + + it('should safely do nothing if recordAttemptStart is called without active operation', () => { + // currentOperation is null + assert.doesNotThrow(() => { + tracer.recordAttemptStart(); + }); + assert.strictEqual(tracer.currentOperation, null); + }); + + it('should safely do nothing if recordAttemptCompletion is called without active attempt', () => { + // No attempt started + assert.doesNotThrow(() => { + tracer.recordAttemptCompletion(Status.OK); + }); + assert.strictEqual(fakeAttemptLatency.record.called, false); + assert.strictEqual(fakeAttemptCounter.add.called, false); + }); + + it('should safely do nothing if recordAttemptCompletion is called without active operation', () => { + tracer.currentOperation = null; + assert.doesNotThrow(() => { + tracer.recordAttemptCompletion(Status.OK); + }); + assert.strictEqual(fakeAttemptLatency.record.called, false); + assert.strictEqual(fakeAttemptCounter.add.called, false); + }); }); describe('recordOperationCompletion', () => { it('should record operation and attempt metrics when enabled', () => { - const factory = sandbox - .stub(MetricsTracerFactory, 'getInstance') - .returns({ - clearCurrentTracer: sinon.spy(), - } as any); tracer.recordOperationStart(); assert.ok(tracer.currentOperation!.startTime); tracer.recordAttemptStart(); @@ -153,8 +193,11 @@ describe('MetricsTracer', () => { assert.strictEqual(fakeAttemptCounter.add.calledOnce, true); assert.strictEqual(fakeOperationLatency.record.calledOnce, true); - const [[_, opAttrs]] = fakeOperationLatency.record.args; - assert.strictEqual(opAttrs[Constants.METRIC_LABEL_KEY_STATUS], 'OK'); + const [[, operationAttributes]] = fakeOperationLatency.record.args; + assert.strictEqual( + operationAttributes[Constants.METRIC_LABEL_KEY_STATUS], + 'OK', + ); }); it('should record fractional operation latency with sub-millisecond precision', () => { @@ -164,10 +207,6 @@ describe('MetricsTracer', () => { nowStub.onCall(2).returns(108.0); // attempt end nowStub.onCall(3).returns(110.25); // op end - sandbox.stub(MetricsTracerFactory, 'getInstance').returns({ - clearCurrentTracer: sinon.spy(), - } as any); - tracer.recordOperationStart(); tracer.recordAttemptStart(); tracer.recordAttemptCompletion(Status.OK); @@ -178,6 +217,113 @@ describe('MetricsTracer', () => { assert.strictEqual(latency, 9.75); // 110.25 - 100.5 }); + it('should record operation error status matching the failed attempt status', () => { + tracer.recordOperationStart(); + tracer.recordAttemptStart(); + tracer.recordAttemptCompletion(Status.UNAVAILABLE); + tracer.recordOperationCompletion(); + + assert.strictEqual(fakeOperationCounter.add.calledOnce, true); + assert.strictEqual(fakeOperationLatency.record.calledOnce, true); + + const [[, operationCounterAttributes]] = fakeOperationCounter.add.args; + const [[, operationLatencyAttributes]] = fakeOperationLatency.record.args; + assert.strictEqual( + operationCounterAttributes[Constants.METRIC_LABEL_KEY_STATUS], + 'UNAVAILABLE', + ); + assert.strictEqual( + operationLatencyAttributes[Constants.METRIC_LABEL_KEY_STATUS], + 'UNAVAILABLE', + ); + }); + + it('should record UNKNOWN status when operation completes without any attempts', () => { + tracer.recordOperationStart(); + // Operation completed before any attempt was started (e.g. client-side error before RPC) + tracer.recordOperationCompletion(); + + assert.strictEqual(fakeOperationCounter.add.calledOnce, true); + assert.strictEqual(fakeOperationLatency.record.calledOnce, true); + + const [[, operationCounterAttributes]] = fakeOperationCounter.add.args; + assert.strictEqual( + operationCounterAttributes[Constants.METRIC_LABEL_KEY_STATUS], + 'UNKNOWN', + ); + }); + + it('should safely do nothing if recordOperationCompletion is called without active operation', () => { + tracer.currentOperation = null; + assert.doesNotThrow(() => { + tracer.recordOperationCompletion(); + }); + assert.strictEqual(fakeOperationLatency.record.called, false); + assert.strictEqual(fakeOperationCounter.add.called, false); + }); + + it('should handle missing currentOperation in _createOperationOtelAttributes', () => { + tracer.currentOperation = null; + const attributes = (tracer as any)._createOperationOtelAttributes(); + assert.strictEqual( + attributes[Constants.METRIC_LABEL_KEY_STATUS], + 'UNKNOWN', + ); + }); + + it('should handle missing currentOperation in _createAttemptOtelAttributes', () => { + tracer.currentOperation = null; + const attributes = (tracer as any)._createAttemptOtelAttributes(); + assert.strictEqual( + attributes[Constants.METRIC_LABEL_KEY_STATUS], + 'UNKNOWN', + ); + }); + + it('should fallback to UNKNOWN status when attempt completes with unrecognized status code', () => { + tracer.recordOperationStart(); + tracer.recordAttemptStart(); + tracer.recordAttemptCompletion(999 as any); + tracer.recordOperationCompletion(); + + const [[, attemptAttributes]] = fakeAttemptLatency.record.args; + const [[, operationAttributes]] = fakeOperationLatency.record.args; + + assert.strictEqual( + attemptAttributes[Constants.METRIC_LABEL_KEY_STATUS], + 'UNKNOWN', + ); + assert.strictEqual( + operationAttributes[Constants.METRIC_LABEL_KEY_STATUS], + 'UNKNOWN', + ); + }); + + it('should not overwrite existing operation when recordOperationStart is called repeatedly', () => { + tracer.recordOperationStart(); + const initialOperation = tracer.currentOperation; + assert.ok(initialOperation); + + tracer.recordOperationStart(); + assert.strictEqual(tracer.currentOperation, initialOperation); + }); + + it('should be idempotent and not double-record if recordOperationCompletion is called multiple times', () => { + tracer.recordOperationStart(); + tracer.recordAttemptStart(); + tracer.recordAttemptCompletion(Status.OK); + tracer.recordOperationCompletion(); + + assert.strictEqual(fakeOperationCounter.add.callCount, 1); + assert.strictEqual(fakeOperationLatency.record.callCount, 1); + assert.strictEqual(tracer.currentOperation, null); + + // Subsequent call should be a no-op + tracer.recordOperationCompletion(); + assert.strictEqual(fakeOperationCounter.add.callCount, 1); + assert.strictEqual(fakeOperationLatency.record.callCount, 1); + }); + it('should do nothing if disabled', () => { tracer.enabled = false; tracer.recordOperationCompletion(); @@ -194,6 +340,50 @@ describe('MetricsTracer', () => { assert.strictEqual(fakeGfeLatency.record.calledOnce, true); }); + it('should record GFE latency when latency is 0ms', () => { + tracer.enabled = true; + tracer.gfeLatency = 0; + tracer.recordGfeLatency(Status.OK); + assert.strictEqual(fakeGfeLatency.record.calledOnce, true); + assert.strictEqual(fakeGfeLatency.record.getCall(0).args[0], 0); + assert.strictEqual(tracer.gfeLatency, null); + }); + + it('should not record and log error when gfeLatency is null', () => { + tracer.enabled = true; + tracer.gfeLatency = null; + const errorStub = sandbox.stub(console, 'error'); + tracer.recordGfeLatency(Status.OK); + assert.strictEqual(fakeGfeLatency.record.called, false); + assert.strictEqual(errorStub.calledOnce, true); + }); + + it('should not record and log error when gfeLatency is NaN or negative', () => { + tracer.enabled = true; + tracer.gfeLatency = NaN; + const errorStub = sandbox.stub(console, 'error'); + tracer.recordGfeLatency(Status.OK); + assert.strictEqual(fakeGfeLatency.record.called, false); + + tracer.gfeLatency = -1; + tracer.recordGfeLatency(Status.OK); + assert.strictEqual(fakeGfeLatency.record.called, false); + assert.strictEqual(errorStub.calledTwice, true); + }); + + it('should fallback to UNKNOWN status when called with unrecognized status code', () => { + tracer.enabled = true; + tracer.gfeLatency = 123; + tracer.recordGfeLatency(999 as any); + assert.strictEqual(fakeGfeLatency.record.calledOnce, true); + assert.strictEqual( + fakeGfeLatency.record.getCall(0).args[1][ + Constants.METRIC_LABEL_KEY_STATUS + ], + 'UNKNOWN', + ); + }); + it('should not record if disabled', () => { tracer.enabled = false; tracer.gfeLatency = 123; @@ -208,6 +398,18 @@ describe('MetricsTracer', () => { assert.strictEqual(fakeGfeCounter.add.calledOnce, true); }); + it('should fallback to UNKNOWN status when called with unrecognized status code', () => { + tracer.enabled = true; + tracer.recordGfeConnectivityErrorCount(999 as any); + assert.strictEqual(fakeGfeCounter.add.calledOnce, true); + assert.strictEqual( + fakeGfeCounter.add.getCall(0).args[1][ + Constants.METRIC_LABEL_KEY_STATUS + ], + 'UNKNOWN', + ); + }); + it('should not increment if disabled', () => { tracer.enabled = false; tracer.recordGfeConnectivityErrorCount(Status.OK); @@ -228,6 +430,50 @@ describe('MetricsTracer', () => { assert.strictEqual(fakeAfeLatency.record.calledOnce, true); }); + it('should fallback to UNKNOWN status when called with unrecognized status code', () => { + tracer.enabled = true; + tracer.afeLatency = 123; + tracer.recordAfeLatency(999 as any); + assert.strictEqual(fakeAfeLatency.record.calledOnce, true); + assert.strictEqual( + fakeAfeLatency.record.getCall(0).args[1][ + Constants.METRIC_LABEL_KEY_STATUS + ], + 'UNKNOWN', + ); + }); + + it('should record AFE latency when latency is 0ms', () => { + tracer.enabled = true; + tracer.afeLatency = 0; + tracer.recordAfeLatency(Status.OK); + assert.strictEqual(fakeAfeLatency.record.calledOnce, true); + assert.strictEqual(fakeAfeLatency.record.getCall(0).args[0], 0); + assert.strictEqual(tracer.afeLatency, null); + }); + + it('should not record and log error when afeLatency is null', () => { + tracer.enabled = true; + tracer.afeLatency = null; + const errorStub = sandbox.stub(console, 'error'); + tracer.recordAfeLatency(Status.OK); + assert.strictEqual(fakeAfeLatency.record.called, false); + assert.strictEqual(errorStub.calledOnce, true); + }); + + it('should not record and log error when afeLatency is NaN or negative', () => { + tracer.enabled = true; + tracer.afeLatency = NaN; + const errorStub = sandbox.stub(console, 'error'); + tracer.recordAfeLatency(Status.OK); + assert.strictEqual(fakeAfeLatency.record.called, false); + + tracer.afeLatency = -1; + tracer.recordAfeLatency(Status.OK); + assert.strictEqual(fakeAfeLatency.record.called, false); + assert.strictEqual(errorStub.calledTwice, true); + }); + it('should not record if AFE server timing is disabled', () => { tracer.enabled = true; Spanner._resetAFEServerTimingForTest(); @@ -245,7 +491,7 @@ describe('MetricsTracer', () => { }); }); - describe('recordGfeConnectivityErrorCount', () => { + describe('recordAfeConnectivityErrorCount', () => { afterEach(() => { Spanner._resetAFEServerTimingForTest(); process.env['SPANNER_DISABLE_AFE_SERVER_TIMING'] = 'false'; @@ -257,6 +503,18 @@ describe('MetricsTracer', () => { assert.strictEqual(fakeAfeCounter.add.calledOnce, true); }); + it('should fallback to UNKNOWN status when called with unrecognized status code', () => { + tracer.enabled = true; + tracer.recordAfeConnectivityErrorCount(999 as any); + assert.strictEqual(fakeAfeCounter.add.calledOnce, true); + assert.strictEqual( + fakeAfeCounter.add.getCall(0).args[1][ + Constants.METRIC_LABEL_KEY_STATUS + ], + 'UNKNOWN', + ); + }); + it('should not increment if metrics are disabled', () => { tracer.enabled = false; tracer.recordAfeConnectivityErrorCount(Status.OK); @@ -301,6 +559,14 @@ describe('MetricsTracer', () => { assert.strictEqual(afeLatency, 30); }); + it('should extract 0ms latency when dur=0 in server-timing header', () => { + const header = 'gfet4t7; dur=0, afe; dur=0'; + const gfeLatency = tracer.extractGfeLatency(header); + assert.strictEqual(gfeLatency, 0); + const afeLatency = tracer.extractAfeLatency(header); + assert.strictEqual(afeLatency, 0); + }); + it('should return null if header is undefined', () => { const gfeLatency = tracer.extractGfeLatency(undefined as any); assert.strictEqual(gfeLatency, null); diff --git a/handwritten/spanner/test/metrics/metrics.ts b/handwritten/spanner/test/metrics/metrics.ts index 70dbfa4e8cd4..1135a127ba78 100644 --- a/handwritten/spanner/test/metrics/metrics.ts +++ b/handwritten/spanner/test/metrics/metrics.ts @@ -18,6 +18,7 @@ import {grpc} from 'google-gax'; import * as mock from '../mockserver/mockspanner'; import {MockError, SimulatedExecutionTime} from '../mockserver/mockspanner'; import {Database, Instance, Spanner} from '../../src'; +import {CLOUD_RESOURCE_HEADER} from '../../src/common'; import {MetricsTracerFactory} from '../../src/metrics/metrics-tracer-factory'; import {MetricsTracer} from '../../src/metrics/metrics-tracer'; import {MetricReader} from '@opentelemetry/sdk-metrics'; @@ -294,7 +295,7 @@ describe('Test metrics with mock server', () => { attributes, ); // Since we only have one attempt, the attempt latency should be fairly close to the operation latency - assertApprox(operationLatency, attemptLatency, 30); + assertApprox(operationLatency, attemptLatency, 100); const gfeLatency = getAggregatedValue(gfeLatenciesData, attributes); assert.strictEqual(gfeLatency, 123); @@ -751,5 +752,249 @@ describe('Test metrics with mock server', () => { Spanner._resetAFEServerTimingForTest(); process.env['SPANNER_DISABLE_AFE_SERVER_TIMING'] = 'false'; }); + + it('should record failure metrics when streaming query fails with non-retryable error', async () => { + const database = newTestDatabase(); + const permissionDeniedError = { + message: 'Permission denied on table NUMBERS', + code: grpc.status.PERMISSION_DENIED, + } as MockError; + spannerMock.setExecutionTime( + spannerMock.executeStreamingSql, + SimulatedExecutionTime.ofError(permissionDeniedError), + ); + + await assert.rejects( + database.run(selectSql), + (error: any) => error.code === grpc.status.PERMISSION_DENIED, + ); + + const {resourceMetrics} = await reader.collect(); + const operationCountData = getMetricData( + resourceMetrics, + METRIC_NAME_OPERATION_COUNT, + ); + const attemptCountData = getMetricData( + resourceMetrics, + METRIC_NAME_ATTEMPT_COUNT, + ); + const operationLatenciesData = getMetricData( + resourceMetrics, + METRIC_NAME_OPERATION_LATENCIES, + ); + const attemptLatenciesData = getMetricData( + resourceMetrics, + METRIC_NAME_ATTEMPT_LATENCIES, + ); + + const failedAttributes = { + instance_id: 'instance', + database: `database-${dbCounter}`, + method: 'executeStreamingSql', + status: 'PERMISSION_DENIED', + }; + + assert.strictEqual( + getAggregatedValue(operationCountData, failedAttributes), + 1, + ); + assert.strictEqual( + getAggregatedValue(attemptCountData, failedAttributes), + 1, + ); + assert.ok(getAggregatedValue(operationLatenciesData, failedAttributes)); + assert.ok(getAggregatedValue(attemptLatenciesData, failedAttributes)); + + await database.close(); + }); + + it('should record failure metrics when commit fails with non-retryable error', async () => { + const database = newTestDatabase(); + const permissionDeniedError = { + message: 'Permission denied on commit', + code: grpc.status.PERMISSION_DENIED, + } as MockError; + spannerMock.setExecutionTime( + spannerMock.commit, + SimulatedExecutionTime.ofError(permissionDeniedError), + ); + + await assert.rejects( + database.runTransactionAsync(async transaction => { + await transaction.run(selectSql); + await transaction.commit(); + }), + (error: any) => error.code === grpc.status.PERMISSION_DENIED, + ); + + const {resourceMetrics} = await reader.collect(); + const operationCountData = getMetricData( + resourceMetrics, + METRIC_NAME_OPERATION_COUNT, + ); + const attemptCountData = getMetricData( + resourceMetrics, + METRIC_NAME_ATTEMPT_COUNT, + ); + const operationLatenciesData = getMetricData( + resourceMetrics, + METRIC_NAME_OPERATION_LATENCIES, + ); + + const failedCommitAttributes = { + instance_id: 'instance', + database: `database-${dbCounter}`, + method: 'commit', + status: 'PERMISSION_DENIED', + }; + + assert.strictEqual( + getAggregatedValue(operationCountData, failedCommitAttributes), + 1, + ); + assert.strictEqual( + getAggregatedValue(attemptCountData, failedCommitAttributes), + 1, + ); + assert.ok( + getAggregatedValue(operationLatenciesData, failedCommitAttributes), + ); + + await database.close(); + }); + + it('should maintain strict metrics context isolation between concurrent succeeding and failing queries', async () => { + const database = newTestDatabase(); + const failingSql = 'SELECT * FROM NON_EXISTENT_TABLE'; + const notFoundError = { + message: 'Table not found', + code: grpc.status.NOT_FOUND, + } as MockError; + spannerMock.putStatementResult( + failingSql, + mock.StatementResult.error(notFoundError), + ); + + const [successSettledResult, failureSettledResult] = + await Promise.allSettled([ + database.run(selectSql), + database.run(failingSql), + ]); + + assert.strictEqual(successSettledResult.status, 'fulfilled'); + assert.strictEqual(failureSettledResult.status, 'rejected'); + + const {resourceMetrics} = await reader.collect(); + const operationCountData = getMetricData( + resourceMetrics, + METRIC_NAME_OPERATION_COUNT, + ); + const attemptCountData = getMetricData( + resourceMetrics, + METRIC_NAME_ATTEMPT_COUNT, + ); + + const successAttributes = { + instance_id: 'instance', + database: `database-${dbCounter}`, + method: 'executeStreamingSql', + status: 'OK', + }; + const failureAttributes = { + instance_id: 'instance', + database: `database-${dbCounter}`, + method: 'executeStreamingSql', + status: 'NOT_FOUND', + }; + + assert.strictEqual( + getAggregatedValue(operationCountData, successAttributes), + 1, + ); + assert.strictEqual( + getAggregatedValue(attemptCountData, successAttributes), + 1, + ); + + assert.strictEqual( + getAggregatedValue(operationCountData, failureAttributes), + 1, + ); + assert.strictEqual( + getAggregatedValue(attemptCountData, failureAttributes), + 1, + ); + + await database.close(); + }); + + it('should record metrics when Spanner.request is called in callback mode', async () => { + const databaseName = `projects/${PROJECT_ID}/instances/instance/databases/database-${++dbCounter}`; + + await new Promise((resolve, reject) => { + (spanner as any).request( + { + client: 'SpannerClient', + method: 'createSession', + headers: { + [CLOUD_RESOURCE_HEADER]: databaseName, + 'x-goog-spanner-request-id': 'test-callback-req-id', + }, + reqOpts: { + database: databaseName, + }, + }, + (error: any) => { + if (error) { + reject(error); + } else { + resolve(); + } + }, + ); + }); + + const {resourceMetrics} = await reader.collect(); + const operationCountData = getMetricData( + resourceMetrics, + METRIC_NAME_OPERATION_COUNT, + ); + const attemptCountData = getMetricData( + resourceMetrics, + METRIC_NAME_ATTEMPT_COUNT, + ); + + const callbackAttributes = { + instance_id: 'instance', + database: `database-${dbCounter}`, + method: 'createSession', + status: 'OK', + }; + + assert.strictEqual( + getAggregatedValue(operationCountData, callbackAttributes), + 1, + ); + assert.strictEqual( + getAggregatedValue(attemptCountData, callbackAttributes), + 1, + ); + }); + + it('should safely handle request when reqOpts is omitted', async () => { + await new Promise(resolve => { + (spanner as any).request( + { + client: 'SpannerClient', + method: 'createSession', + headers: {}, + // reqOpts intentionally omitted + }, + () => { + resolve(); + }, + ); + }); + }); }); });