diff --git a/core/packages/gax/samples/observability.js b/core/packages/gax/samples/observability.js index 6d31636144c..ab2318a803a 100644 --- a/core/packages/gax/samples/observability.js +++ b/core/packages/gax/samples/observability.js @@ -17,17 +17,12 @@ // [START gax_observability] 'use strict'; -// 1. INITIALIZE OPENTELEMETRY BEFORE IMPORTING ANY CLIENT LIBRARIES -// In Node.js, instrumentations must patch the networking modules (http, grpc) -// before any Google Cloud client libraries are loaded into the module cache. +// 1. IMPORT OPENTELEMETRY MODULES const {NodeTracerProvider} = require('@opentelemetry/sdk-trace-node'); const {BatchSpanProcessor} = require('@opentelemetry/sdk-trace-base'); const { TraceExporter, } = require('@google-cloud/opentelemetry-cloud-trace-exporter'); -const {registerInstrumentations} = require('@opentelemetry/instrumentation'); -const {HttpInstrumentation} = require('@opentelemetry/instrumentation-http'); -const {GrpcInstrumentation} = require('@opentelemetry/instrumentation-grpc'); // 2. CONFIGURE TRACING: SET UP A TRACER PROVIDER AND EXPORTER const cloudTraceExporter = new TraceExporter(); @@ -38,19 +33,11 @@ const provider = new NodeTracerProvider({ }); provider.register(); -// 3. ENABLE LOW-LEVEL NETWORK TRACING SPANS USING INSTRUMENTATION LIBRARIES -registerInstrumentations({ - instrumentations: [ - new HttpInstrumentation(), - new GrpcInstrumentation(), - ], -}); - -// 4. ENABLE CLIENT REQUEST TRACING SPANS WITH ENV VARIABLE +// 3. ENABLE TRACING SPANS WITH ENV VARIABLE // Sets the flag before client libraries or RPC callers initialize process.env.GOOGLE_SDK_NODE_ENABLE_TRACING = 'true'; -// 5. IMPORT CLIENT LIBRARIES AFTER OPENTELEMETRY SETUP +// 4. IMPORT CLIENT LIBRARIES AFTER OPENTELEMETRY SETUP // Replace with your Google Cloud client library, for example: // const { SecretManagerServiceClient } = require('@google-cloud/secret-manager'); @@ -58,7 +45,7 @@ async function main() { // const client = new SecretManagerServiceClient(); // await client.listSecrets({parent: 'projects/my-project'}); - // 6. FLUSH SPANS BEFORE PROCESS EXIT + // 5. FLUSH SPANS BEFORE PROCESS EXIT // Ensures all buffered spans in BatchSpanProcessor are exported to Cloud Trace await provider.forceFlush(); console.log('Tracing initialized successfully.'); diff --git a/core/packages/gax/src/createApiCall.ts b/core/packages/gax/src/createApiCall.ts index ac73d9aa8fb..9d016891774 100644 --- a/core/packages/gax/src/createApiCall.ts +++ b/core/packages/gax/src/createApiCall.ts @@ -18,6 +18,7 @@ * Provides function wrappers that implement page streaming and retrying. */ +import {context} from '@opentelemetry/api'; import {createAPICaller} from './apiCaller'; import { APICallback, @@ -26,6 +27,7 @@ import { GRPCCallOtherArgs, RequestType, SimpleCallbackFunction, + UnaryCall, } from './apitypes'; import {Descriptor} from './descriptor'; import {CallOptions, CallSettings, convertRetryOptions} from './gax'; @@ -36,8 +38,10 @@ import {StreamProxy} from './streamingCalls/streaming'; import {warn} from './warnings'; import { traceCall, + traceAttempt, StaticTraceContext, DynamicTraceContext, + AttemptTraceContext, ResendRecorder, } from './observability/TracerHelper'; import {resolveStaticTraceContext} from './observability/metadataResolver'; @@ -77,6 +81,26 @@ export function createApiCall( const apiCaller = createAPICaller(settings, descriptor); const tracingEnabled = checkTelemetryEnabled(settings); + const staticArgs: StaticTraceContext | undefined = tracingEnabled + ? resolveStaticTraceContext(settings) + : undefined; + const serviceName = settings.apiName?.split('.').pop() ?? ''; + const isFallback = Boolean(_fallback); + const dynamicArgs: DynamicTraceContext | undefined = tracingEnabled + ? { + clientName: serviceName ? `${serviceName}Client` : '', + methodName: settings.otherArgs?.internalMethodName ?? '', + rpcType: isFallback ? 'http' : 'grpc', + } + : undefined; + // Context for per-attempt T4 CLIENT spans. + const attemptDynamicArgs: AttemptTraceContext | undefined = + tracingEnabled && dynamicArgs + ? { + ...dynamicArgs, + apiName: settings.apiName ?? '', + } + : undefined; const invokeCall = ( request: RequestType, @@ -84,6 +108,8 @@ export function createApiCall( callback?: APICallback, recordResend?: ResendRecorder, ) => { + // Capture the active T3 call span context to parent async T4 attempt spans. + const parentContext = tracingEnabled ? context.active() : undefined; let currentApiCaller = apiCaller; let thisSettings: CallSettings; @@ -119,11 +145,39 @@ export function createApiCall( funcPromise .then((func: GRPCCall) => { // Initially, the function is just what gRPC server stub contains. - func = currentApiCaller.wrap(func); + let wrappedFunc = currentApiCaller.wrap(func); const streaming = (currentApiCaller as StreamingApiCaller).descriptor ?.streaming; + // Wrap the transport call so each attempt (initial send and retries) emits a T4 CLIENT span. + if (tracingEnabled && attemptDynamicArgs && staticArgs) { + const callerWrappedFunc = wrappedFunc; + let attemptCount = 0; + wrappedFunc = (( + argument: {}, + metadata: {}, + options: {}, + attemptCallback: APICallback, + ) => { + const resendCount = attemptCount++; + return traceAttempt( + {...attemptDynamicArgs, resendCount}, + staticArgs, + tracedAttemptCallback => + (callerWrappedFunc as UnaryCall)( + argument, + metadata, + options, + tracedAttemptCallback ?? attemptCallback, + ), + Boolean(streaming), + attemptCallback, + parentContext, + ); + }) as GRPCCall; + } + const retry = thisSettings.retry; if (streaming && retry) { @@ -159,7 +213,7 @@ export function createApiCall( retry.backoffSettings.initialRpcTimeoutMillis ??= thisSettings.timeout; return retryable( - func, + wrappedFunc, thisSettings.retry!, thisSettings.otherArgs as GRPCCallOtherArgs, thisSettings.apiName, @@ -168,7 +222,7 @@ export function createApiCall( } } return addTimeoutArg( - func, + wrappedFunc, thisSettings.timeout, thisSettings.otherArgs as GRPCCallOtherArgs, ); @@ -192,16 +246,7 @@ export function createApiCall( return currentApiCaller.result(ongoingCall); }; - if (tracingEnabled) { - const staticArgs: StaticTraceContext = resolveStaticTraceContext(settings); - - const serviceName = settings.apiName?.split('.').pop() ?? ''; - const isFallback = Boolean(_fallback); - const dynamicArgs: DynamicTraceContext = { - clientName: serviceName ? `${serviceName}Client` : '', - methodName: settings.otherArgs?.internalMethodName ?? '', - rpcType: isFallback ? 'http' : 'grpc', - }; + if (tracingEnabled && dynamicArgs && staticArgs) { const isStreamingCall = apiCaller instanceof StreamingApiCaller; return ( request: RequestType, diff --git a/core/packages/gax/src/fallback.ts b/core/packages/gax/src/fallback.ts index 9f66c8bef75..1d2c3afabcc 100644 --- a/core/packages/gax/src/fallback.ts +++ b/core/packages/gax/src/fallback.ts @@ -108,6 +108,8 @@ export class GrpcClient { httpRules?: Array; numericEnums: boolean; minifyJson: boolean; + private _servicePath?: string; + private _port?: number; /** * In rare cases users might need to deallocate all memory consumed by loaded protos. @@ -156,6 +158,10 @@ export class GrpcClient { this.httpRules = (options as GrpcClientOptions).httpRules; this.numericEnums = (options as GrpcClientOptions).numericEnums ?? false; this.minifyJson = (options as GrpcClientOptions).minifyJson ?? false; + this._servicePath = + (options as GrpcClientOptions).servicePath || + (options as GrpcClientOptions).apiEndpoint; + this._port = (options as GrpcClientOptions).port; } /** @@ -279,12 +285,21 @@ export class GrpcClient { } return metadata; } + const otherArgs: Record = { + metadataBuilder: buildMetadata, + }; + if (this._servicePath) { + otherArgs.servicePath = this._servicePath; + } + if (this._port !== undefined) { + otherArgs.port = this._port; + } return gax.constructSettings( serviceName, clientConfig, configOverrides, Status, - {metadataBuilder: buildMetadata}, + otherArgs, enableTelemetryTracing, internalTelemetryInfo, ); diff --git a/core/packages/gax/src/fallbackServiceStub.ts b/core/packages/gax/src/fallbackServiceStub.ts index ad196dc1932..2816076d1b2 100644 --- a/core/packages/gax/src/fallbackServiceStub.ts +++ b/core/packages/gax/src/fallbackServiceStub.ts @@ -25,6 +25,7 @@ import {isNodeJS} from './featureDetection'; import {StreamArrayParser} from './streamArrayParser'; import {defaultToObjectOptions} from './fallback'; import {GoogleError} from './googleError'; +import {setAttemptHttpMethod} from './observability/TracerHelper'; import {rpcCodeFromHttpStatusCode, Status} from './status'; import {pipeline, PipelineSource} from 'stream'; import type {Agent as HttpAgent} from 'http'; @@ -302,6 +303,8 @@ export function generateServiceStub( }; } + setAttemptHttpMethod(fetchParameters.method); + const cancelController = new AbortController(); const cancelSignal = cancelController.signal as AbortSignal; let cancelRequested = false; diff --git a/core/packages/gax/src/grpc.ts b/core/packages/gax/src/grpc.ts index 14d35398d51..5c68f256456 100644 --- a/core/packages/gax/src/grpc.ts +++ b/core/packages/gax/src/grpc.ts @@ -53,6 +53,9 @@ export interface GrpcClientOptions extends GoogleAuthOptions { httpRules?: Array; numericEnums?: boolean; universeDomain?: string; + servicePath?: string; + apiEndpoint?: string; + port?: number; } export interface MetadataValue { @@ -121,6 +124,8 @@ export class GrpcClient { fallback: boolean | 'rest' | 'proto'; private static protoCache = new Map(); httpRules?: Array; + private _servicePath?: string; + private _port?: number; /** * Base directory for resolving client certificates. * @@ -174,6 +179,8 @@ export class GrpcClient { constructor(options: GrpcClientOptions = {}) { this.auth = options.auth || new GoogleAuth(options); this.fallback = false; + this._servicePath = options.servicePath || options.apiEndpoint; + this._port = options.port; const minimumVersion = 10; const major = Number(process.version.match(/^v(\d+)/)?.[1]); @@ -388,12 +395,21 @@ export class GrpcClient { enableTelemetryTracing?: boolean, internalTelemetryInfo?: StaticTraceContext, ) { + const otherArgs: Record = { + metadataBuilder: this.metadataBuilder(headers), + }; + if (this._servicePath) { + otherArgs.servicePath = this._servicePath; + } + if (this._port !== undefined) { + otherArgs.port = this._port; + } return gax.constructSettings( serviceName, clientConfig, configOverrides, this.grpc.status, - {metadataBuilder: this.metadataBuilder(headers)}, + otherArgs, enableTelemetryTracing, internalTelemetryInfo, ); diff --git a/core/packages/gax/src/observability/TracerHelper.ts b/core/packages/gax/src/observability/TracerHelper.ts index 81b28d044cd..438b1510992 100644 --- a/core/packages/gax/src/observability/TracerHelper.ts +++ b/core/packages/gax/src/observability/TracerHelper.ts @@ -18,7 +18,10 @@ import {EventEmitter} from 'events'; import { Attributes, context, + Context, + createContextKey, Span, + SpanKind, SpanStatusCode, trace, Tracer, @@ -66,6 +69,10 @@ export interface StaticTraceContext { * Server port number for the RPC call. */ serverPort?: number; + /** + * Target service domain (e.g. 'cloudkms.googleapis.com'). + */ + urlDomain?: string; } /** @@ -93,6 +100,42 @@ export interface DynamicTraceContext { * Server port number for the RPC call. */ serverPort?: number; + /** + * Target service domain (e.g. 'cloudkms.googleapis.com'). + */ + urlDomain?: string; +} + +/** + * Dynamic metadata specific to an individual RPC transport attempt (T4 span). + */ +export interface AttemptTraceContext extends DynamicTraceContext { + /** + * The fully-qualified protobuf service name (e.g. 'google.cloud.kms.v1.KeyManagementService'). + */ + apiName?: string; + /** + * The ordinal resend count for this attempt (0 for the initial attempt, 1 for the first retry, etc.). + * Omitted from span attributes when 0 or undefined. + */ + resendCount?: number; + /** + * The HTTP request method for REST fallback attempts (e.g. 'GET', 'POST', 'PUT', 'PATCH', 'DELETE'). + */ + httpMethod?: string; +} + +const ATTEMPT_SPAN_KEY = createContextKey('google-gax-attempt-span'); + +/** + * Updates the `http.request.method` attribute on the currently active T4 attempt span, if any. + */ +export function setAttemptHttpMethod(httpMethod: string): void { + const attemptSpan = context.active().getValue(ATTEMPT_SPAN_KEY) as + Span | undefined; + if (attemptSpan) { + attemptSpan.setAttribute('http.request.method', httpMethod); + } } /** @@ -124,15 +167,16 @@ export function resolveErrorInfoReason(e: unknown): string | undefined { } // Decode binary gRPC status details if present and not yet parsed. + const errWithMeta = e as GoogleError; if ( - e instanceof GoogleError && - e.metadata && - typeof e.metadata.get === 'function' && - (e.metadata.get('grpc-status-details-bin') as unknown[])?.length > 0 && - !e.reason + errWithMeta.metadata && + typeof errWithMeta.metadata.get === 'function' && + (errWithMeta.metadata.get('grpc-status-details-bin') as unknown[])?.length > + 0 && + !errWithMeta.reason ) { try { - GoogleError.parseGRPCStatusDetails(e); + GoogleError.parseGRPCStatusDetails(errWithMeta); } catch { // Ignore decoding errors. } @@ -696,16 +740,17 @@ export function resolveServerExceptionDetails(e: Error): { }) : undefined; - // If e is a GoogleError with gRPC metadata that hasn't decoded statusDetails yet, parse it: + // If e has gRPC metadata that hasn't decoded statusDetails yet, parse it: + const errWithMeta = e as GoogleError; if ( !errObj.statusDetails && - e instanceof GoogleError && - e.metadata && - typeof e.metadata.get === 'function' && - (e.metadata.get('grpc-status-details-bin') as unknown[])?.length > 0 + errWithMeta.metadata && + typeof errWithMeta.metadata.get === 'function' && + (errWithMeta.metadata.get('grpc-status-details-bin') as unknown[])?.length > + 0 ) { try { - GoogleError.parseGRPCStatusDetails(e); + GoogleError.parseGRPCStatusDetails(errWithMeta); } catch { // Ignore decoding errors } @@ -952,6 +997,30 @@ export function handleStream( } } +/** + * Resolves the target service domain (`url.domain`) from dynamic and static trace contexts. + */ +function resolveUrlDomain( + dynamicArgs: DynamicTraceContext, + staticArgs: StaticTraceContext, +): string | undefined { + const explicit = dynamicArgs.urlDomain ?? staticArgs.urlDomain; + if (explicit) { + return explicit; + } + const rawAddress = dynamicArgs.serverAddress ?? staticArgs.serverAddress; + if (rawAddress) { + const match = rawAddress.match(/^(\[[^\]]+\]|[^:]+):(\d+)$/); + return match ? match[1] : rawAddress; + } + if (staticArgs.gcpClientService) { + return staticArgs.gcpClientService.includes('.') + ? staticArgs.gcpClientService + : `${staticArgs.gcpClientService}.googleapis.com`; + } + return undefined; +} + /** * Executes a function within an active OpenTelemetry span, populating standard * GCP telemetry attributes and recording errors/exceptions if thrown. @@ -1005,14 +1074,19 @@ export function traceCall( ): GaxCallResult { const spanName = `${dynamicArgs.clientName}.${dynamicArgs.methodName}`; return getGaxTracer().startActiveSpan(spanName, {}, (span: Span) => { - span.setAttributes({ + const urlDomain = resolveUrlDomain(dynamicArgs, staticArgs); + const initialAttributes: Attributes = { 'gcp.client.service': staticArgs.gcpClientService, 'gcp.client.version': staticArgs.gcpVersion, 'gcp.repo': staticArgs.gcpRepo, 'gcp.artifact': staticArgs.gcpArtifact, 'gcp.method.name': dynamicArgs.methodName, 'gcp.method.type': dynamicArgs.rpcType, - }); + }; + if (urlDomain !== undefined) { + initialAttributes['url.domain'] = urlDomain; + } + span.setAttributes(initialAttributes); let rawAddress = dynamicArgs.serverAddress ?? staticArgs.serverAddress; let rawPort = dynamicArgs.serverPort ?? staticArgs.serverPort; @@ -1068,11 +1142,8 @@ export function traceCall( const setStatusAttributes = () => { const attributes: Attributes = {}; - if (rpcStatusName !== undefined) { + if (dynamicArgs.rpcType === 'grpc' && rpcStatusName !== undefined) { attributes['rpc.response.status_code'] = rpcStatusName; - if (dynamicArgs.rpcType === 'grpc') { - attributes['grpc.response.status_code'] = rpcStatusName; - } } if (dynamicArgs.rpcType === 'http' && httpStatusCode !== undefined) { attributes['http.response.status_code'] = httpStatusCode; @@ -1136,7 +1207,8 @@ export function traceCall( : undefined; try { - const result = context.with(trace.setSpan(context.active(), span), () => + const activeContext = trace.setSpan(context.active(), span); + const result = context.with(activeContext, () => fn(tracedCallback, recordResend), ); const promiseTarget = !isStreamCall ? getPromiseTarget(result) : null; @@ -1158,3 +1230,172 @@ export function traceCall( } }); } + +/** + * Executes an individual RPC transport attempt within an active OpenTelemetry + * CLIENT span (T4 span), parenting it to the active T3 client request span + * and recording per-attempt network, status, and error attributes without + * injecting span context into outgoing headers. + * + * @param {AttemptTraceContext} dynamicArgs - Dynamic trace context for the RPC attempt. + * @param {StaticTraceContext} staticArgs - Static trace context for the client library. + * @param {function} fn - The transport attempt operation to trace. + * @param {boolean} [isStreamCall=false] - Whether the operation is a stream call. + * @param {APICallback} [callback] - The attempt callback. + * @param {Context} [parentContext] - Optional parent OpenTelemetry context (e.g. T3 span context). + * @returns {GaxCallResult} The result of the traced attempt. + */ +export function traceAttempt( + dynamicArgs: AttemptTraceContext, + staticArgs: StaticTraceContext, + fn: (tracedCallback?: APICallback) => T, + isStreamCall = false, + callback?: APICallback, + parentContext?: Context, +): T { + const spanName = dynamicArgs.apiName + ? `${dynamicArgs.apiName}/${dynamicArgs.methodName}` + : dynamicArgs.methodName; + const baseContext = parentContext ?? context.active(); + return getGaxTracer().startActiveSpan( + spanName, + {kind: SpanKind.CLIENT}, + baseContext, + (span: Span) => { + // Populate initial transport, method, domain, and retry attributes. + const urlDomain = resolveUrlDomain(dynamicArgs, staticArgs); + const initialAttributes: Attributes = { + 'gcp.client.service': staticArgs.gcpClientService, + 'rpc.system': dynamicArgs.rpcType, + }; + if (dynamicArgs.rpcType === 'grpc') { + initialAttributes['rpc.method'] = spanName; + } else { + initialAttributes['http.request.method'] = + dynamicArgs.httpMethod ?? 'POST'; + } + if (urlDomain !== undefined) { + initialAttributes['url.domain'] = urlDomain; + } + if ( + dynamicArgs.resendCount !== undefined && + dynamicArgs.resendCount > 0 + ) { + const resendCountAttribute = + dynamicArgs.rpcType === 'grpc' + ? 'gcp.grpc.resend_count' + : 'http.request.resend_count'; + initialAttributes[resendCountAttribute] = dynamicArgs.resendCount; + } + span.setAttributes(initialAttributes); + + // Parse server address and port, defaulting to port 443. + let rawAddress = + dynamicArgs.serverAddress ?? staticArgs.serverAddress ?? urlDomain; + let rawPort = dynamicArgs.serverPort ?? staticArgs.serverPort; + if (rawAddress) { + const match = rawAddress.match(/^(\[[^\]]+\]|[^:]+):(\d+)$/); + if (match) { + rawAddress = match[1]; + rawPort = rawPort ?? Number(match[2]); + } + rawPort = rawPort ?? 443; + } + + let spanEnded = false; + let errorRecorded = false; + let recordedError: unknown; + let rpcStatusName: string | undefined; + let httpStatusCode: number | undefined; + + const setErrorStatus = (message: string) => { + errorRecorded = true; + span.setStatus({code: SpanStatusCode.ERROR, message}); + }; + + // Record response status and server endpoint (omitted on pre-connection failures). + const setStatusAttributes = () => { + const attributes: Attributes = {}; + if (dynamicArgs.rpcType === 'grpc' && rpcStatusName !== undefined) { + attributes['rpc.response.status_code'] = rpcStatusName; + } + if (dynamicArgs.rpcType === 'http' && httpStatusCode !== undefined) { + attributes['http.response.status_code'] = httpStatusCode; + } + if ( + rawAddress !== undefined && + (!errorRecorded || !isPreConnectionFailure(recordedError)) + ) { + attributes['server.address'] = rawAddress; + if (rawPort !== undefined) { + attributes['server.port'] = rawPort; + } + } + span.setAttributes(attributes); + }; + + // Finalize status attributes and end the span once. + const endSpan = () => { + if (!spanEnded) { + spanEnded = true; + if (!errorRecorded) { + rpcStatusName = Status[Status.OK]; + httpStatusCode = 200; + } + setStatusAttributes(); + span.end(); + } + }; + + // Capture error type, exception event, and span error status on failure. + const recordError = (e: unknown) => { + recordedError = e; + rpcStatusName = resolveRpcStatusName(e); + httpStatusCode = resolveHttpStatusCode(e); + span.setAttributes({ + 'error.type': resolveErrorType(e, dynamicArgs.rpcType), + }); + if (e instanceof Error) { + recordExceptionEvent(span, e, dynamicArgs.rpcType); + setErrorStatus(e.message); + } else { + setErrorStatus(resolveErrorMessage(e)); + } + }; + + const tracedCallback: APICallback | undefined = callback + ? function (this: unknown, ...args: Parameters) { + const err = args[0]; + if (err) { + recordError(err); + } + endSpan(); + callback.apply(this, args); + } + : undefined; + + try { + // Expose the attempt span in context so the HTTP transport can update http.request.method. + const attemptContext = context + .active() + .setValue(ATTEMPT_SPAN_KEY, span); + const result = context.with(attemptContext, () => fn(tracedCallback)); + const promiseTarget = !isStreamCall ? getPromiseTarget(result) : null; + if (isStreamCall && result instanceof EventEmitter) { + handleStream(result, recordError, endSpan, !!callback); + } else if (promiseTarget) { + handlePromise(promiseTarget, recordError, endSpan); + } else if (tracedCallback) { + // Span stays open; tracedCallback ends it when the attempt completes. + } else { + endSpan(); + } + return result; + } catch (e) { + recordError(e); + endSpan(); + throw e; + } + }, + ); +} diff --git a/core/packages/gax/test/unit/apiCallable.ts b/core/packages/gax/test/unit/apiCallable.ts index 20465bf1b3e..6182c82890b 100644 --- a/core/packages/gax/test/unit/apiCallable.ts +++ b/core/packages/gax/test/unit/apiCallable.ts @@ -17,12 +17,22 @@ import assert from 'assert'; import {PassThrough} from 'stream'; import {status} from '@grpc/grpc-js'; +import {SpanKind, SpanStatusCode} from '@opentelemetry/api'; import {afterEach, beforeEach, describe, it} from 'mocha'; import * as sinon from 'sinon'; -import {CancellableStream, GRPCCall, RequestType} from '../../src/apitypes'; +import { + CancellableStream, + GRPCCall, + GRPCCallResult, + RequestType, +} from '../../src/apitypes'; import {createApiCall as gaxCreateApiCall} from '../../src/createApiCall'; -import {createApiCall as fallbackCreateApiCall} from '../../src/fallback'; +import { + createApiCall as fallbackCreateApiCall, + GrpcClient as FallbackGrpcClient, +} from '../../src/fallback'; +import {GrpcClient} from '../../src/grpc'; import {StreamDescriptor} from '../../src/descriptor'; import {StreamType} from '../../src/streamingCalls/streaming'; import * as gax from '../../src/gax'; @@ -583,9 +593,17 @@ describe('createApiCall', () => { assert.deepStrictEqual(response, {data: 'hello'}); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + const attemptSpan = spans[0]; + const span = spans[1]; + assert.strictEqual(attemptSpan.name, 'google.example.v1.Echo/Echo'); + assert.strictEqual(attemptSpan.kind, SpanKind.CLIENT); + assert.strictEqual( + attemptSpan.parentSpanContext?.spanId, + span.spanContext().spanId, + ); assert.strictEqual(span.name, 'EchoClient.Echo'); + assert.strictEqual(span.kind, SpanKind.INTERNAL); assert.strictEqual(span.ended, true); assert.strictEqual( span.attributes['gcp.client.service'], @@ -599,6 +617,7 @@ describe('createApiCall', () => { assert.strictEqual(span.attributes['gcp.artifact'], '@google-cloud/echo'); assert.strictEqual(span.attributes['gcp.method.name'], 'Echo'); assert.strictEqual(span.attributes['gcp.method.type'], 'grpc'); + assert.strictEqual(span.attributes['url.domain'], 'echo.googleapis.com'); }); it('enables tracing purely through GOOGLE_SDK_NODE_ENABLE_TRACING and resolves static metadata dynamically at runtime', async () => { @@ -631,8 +650,14 @@ describe('createApiCall', () => { assert.deepStrictEqual(response, {data: 'hello'}); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + const attemptSpan = spans[0]; + const span = spans[1]; + assert.strictEqual( + attemptSpan.name, + 'google.cloud.redis.v1.CloudRedis/GetInstance', + ); + assert.strictEqual(attemptSpan.kind, SpanKind.CLIENT); assert.strictEqual(span.name, 'CloudRedisClient.GetInstance'); assert.strictEqual(span.ended, true); assert.strictEqual(span.attributes['gcp.client.service'], 'redis'); @@ -672,8 +697,8 @@ describe('createApiCall', () => { await apiCall({}, undefined); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + const span = spans[1]; assert.strictEqual(span.attributes['gcp.client.service'], 'env-service'); assert.strictEqual(span.attributes['gcp.client.version'], '9.9.9'); assert.strictEqual( @@ -709,8 +734,13 @@ describe('createApiCall', () => { await apiCall({}, undefined); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + const attemptSpan = spans[0]; + const span = spans[1]; + assert.strictEqual(attemptSpan.name, 'google.example.v1.Echo/Echo'); + assert.strictEqual(attemptSpan.kind, SpanKind.CLIENT); + assert.strictEqual(attemptSpan.attributes['rpc.system'], 'http'); + assert.strictEqual(attemptSpan.attributes['gcp.method.type'], undefined); assert.strictEqual(span.name, 'EchoClient.Echo'); assert.strictEqual(span.ended, true); assert.strictEqual(span.attributes['gcp.method.type'], 'http'); @@ -748,8 +778,11 @@ describe('createApiCall', () => { assert.strictEqual(dynamicArgs.rpcType, 'http'); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + const attemptSpan = spans[0]; + const span = spans[1]; + assert.strictEqual(attemptSpan.name, 'google.example.v1.Echo/Echo'); + assert.strictEqual(attemptSpan.attributes['rpc.system'], 'http'); assert.strictEqual(span.name, 'EchoClient.Echo'); assert.strictEqual(span.ended, true); assert.strictEqual(span.attributes['gcp.method.type'], 'http'); @@ -790,8 +823,8 @@ describe('createApiCall', () => { assert.strictEqual(dynamicArgs.rpcType, 'http'); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + const span = spans[1]; assert.strictEqual(span.attributes['gcp.method.type'], 'http'); }); @@ -845,20 +878,26 @@ describe('createApiCall', () => { ); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; - assert.strictEqual(span.ended, true); - assert.strictEqual(span.attributes['gcp.method.type'], 'http'); - // On the fallback transport error.type reports the HTTP status the - // server sent. A deadline expires before any response arrives, so there - // is none, and the attribute resolves to CLIENT_TIMEOUT per Tier 3. - assert.strictEqual(span.attributes['error.type'], 'CLIENT_TIMEOUT'); - assert.strictEqual( - span.attributes['rpc.response.status_code'], - 'DEADLINE_EXCEEDED', - ); - assert.strictEqual(span.events.length, 1); - assert.strictEqual(span.events[0].name, 'exception'); + assert.strictEqual(spans.length, 2); + assert.strictEqual(spans[0].attributes['rpc.system'], 'http'); + assert.strictEqual(spans[1].attributes['gcp.method.type'], 'http'); + for (const span of spans) { + assert.strictEqual(span.ended, true); + // On the fallback transport error.type reports the HTTP status the + // server sent. A deadline expires before any response arrives, so there + // is none, and the attribute resolves to CLIENT_TIMEOUT per Tier 3. + assert.strictEqual(span.attributes['error.type'], 'CLIENT_TIMEOUT'); + assert.strictEqual( + span.attributes['rpc.response.status_code'], + undefined, + ); + assert.strictEqual( + span.attributes['http.response.status_code'], + undefined, + ); + assert.strictEqual(span.events.length, 1); + assert.strictEqual(span.events[0].name, 'exception'); + } }); it('ends the span and preserves system error codes like ECONNREFUSED on a fallback call', async () => { @@ -906,28 +945,30 @@ describe('createApiCall', () => { ); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; - assert.strictEqual(span.ended, true); - assert.strictEqual(span.attributes['gcp.method.type'], 'http'); - assert.strictEqual( - span.attributes['error.type'], - 'CLIENT_CONNECTION_ERROR', - ); - assert.strictEqual( - span.attributes['rpc.response.status_code'], - 'UNAVAILABLE', - ); - assert.strictEqual( - span.attributes['http.response.status_code'], - undefined, - ); - assert.strictEqual(span.events.length, 1); - assert.strictEqual(span.events[0].name, 'exception'); - assert.strictEqual( - span.events[0].attributes?.['exception.type'], - 'GoogleError', - ); + assert.strictEqual(spans.length, 2); + assert.strictEqual(spans[0].attributes['rpc.system'], 'http'); + assert.strictEqual(spans[1].attributes['gcp.method.type'], 'http'); + for (const span of spans) { + assert.strictEqual(span.ended, true); + assert.strictEqual( + span.attributes['error.type'], + 'CLIENT_CONNECTION_ERROR', + ); + assert.strictEqual( + span.attributes['rpc.response.status_code'], + undefined, + ); + assert.strictEqual( + span.attributes['http.response.status_code'], + undefined, + ); + assert.strictEqual(span.events.length, 1); + assert.strictEqual(span.events[0].name, 'exception'); + assert.strictEqual( + span.events[0].attributes?.['exception.type'], + 'GoogleError', + ); + } }); it('passes fallback flag and isStreamingCall as true for server-streaming fallback calls', () => { @@ -987,9 +1028,9 @@ describe('createApiCall', () => { await apiCall({}, undefined); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; - assert.strictEqual(span.attributes['gcp.method.type'], 'grpc'); + assert.strictEqual(spans.length, 2); + assert.strictEqual(spans[0].attributes['rpc.system'], 'grpc'); + assert.strictEqual(spans[1].attributes['gcp.method.type'], 'grpc'); }); it('sets rpcType to http when _fallback is "rest"', async () => { @@ -1018,9 +1059,9 @@ describe('createApiCall', () => { await apiCall({}, undefined); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; - assert.strictEqual(span.attributes['gcp.method.type'], 'http'); + assert.strictEqual(spans.length, 2); + assert.strictEqual(spans[0].attributes['rpc.system'], 'http'); + assert.strictEqual(spans[1].attributes['gcp.method.type'], 'http'); }); it('sets rpcType to http when _fallback is "proto"', async () => { @@ -1049,9 +1090,9 @@ describe('createApiCall', () => { await apiCall({}, undefined); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; - assert.strictEqual(span.attributes['gcp.method.type'], 'http'); + assert.strictEqual(spans.length, 2); + assert.strictEqual(spans[0].attributes['rpc.system'], 'http'); + assert.strictEqual(spans[1].attributes['gcp.method.type'], 'http'); }); it('pipes telemetry information configured via constructSettings', async () => { @@ -1090,8 +1131,11 @@ describe('createApiCall', () => { await apiCall({}, undefined); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + const attemptSpan = spans[0]; + const span = spans[1]; + assert.strictEqual(attemptSpan.name, 'google.example.v1.Echo/Echo'); + assert.strictEqual(attemptSpan.kind, SpanKind.CLIENT); assert.strictEqual(span.name, 'EchoClient.Echo'); assert.strictEqual(span.ended, true); assert.strictEqual( @@ -1148,12 +1192,13 @@ describe('createApiCall', () => { ); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; - assert.strictEqual(span.ended, true); - assert.strictEqual(span.status.message, 'RPC test failure'); - assert.strictEqual(span.events.length, 1); - assert.strictEqual(span.events[0].name, 'exception'); + assert.strictEqual(spans.length, 2); + for (const span of spans) { + assert.strictEqual(span.ended, true); + assert.strictEqual(span.status.message, 'RPC test failure'); + assert.strictEqual(span.events.length, 1); + assert.strictEqual(span.events[0].name, 'exception'); + } }); it('does not end span prematurely for successful asynchronous API calls', async () => { @@ -1189,10 +1234,12 @@ describe('createApiCall', () => { const [response] = (await promise) as [{data: string}, unknown, unknown]; assert.deepStrictEqual(response, {data: 'hello'}); - // Span must only be ended after completion + // Spans must only be ended after completion const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + assert.strictEqual(spans[0].ended, true); + assert.strictEqual(spans[0].name, 'google.example.v1.Echo/Echo'); + const span = spans[1]; assert.strictEqual(span.ended, true); assert.strictEqual(span.name, 'EchoClient.Echo'); }); @@ -1308,8 +1355,12 @@ describe('createApiCall', () => { try { assert.strictEqual(received.length, 2); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; + assert.strictEqual(spans.length, 2); + const attemptSpan = spans[0]; + const span = spans[1]; + assert.strictEqual(attemptSpan.ended, true); + assert.strictEqual(attemptSpan.name, 'google.example.v1.Echo/Echo'); + assert.strictEqual(attemptSpan.kind, SpanKind.CLIENT); assert.strictEqual(span.ended, true); assert.strictEqual(span.name, 'EchoClient.Echo'); assert.strictEqual(span.attributes['gcp.method.type'], 'grpc'); @@ -1352,12 +1403,13 @@ describe('createApiCall', () => { try { assert.strictEqual(err.message, 'streaming test failure'); const spans = harness.getSpans('google-gax'); - assert.strictEqual(spans.length, 1); - const span = spans[0]; - assert.strictEqual(span.ended, true); - assert.strictEqual(span.status.message, 'streaming test failure'); - assert.strictEqual(span.events.length, 1); - assert.strictEqual(span.events[0].name, 'exception'); + assert.strictEqual(spans.length, 2); + for (const span of spans) { + assert.strictEqual(span.ended, true); + assert.strictEqual(span.status.message, 'streaming test failure'); + assert.strictEqual(span.events.length, 1); + assert.strictEqual(span.events[0].name, 'exception'); + } done(); } catch (e) { done(e); @@ -1415,7 +1467,10 @@ describe('createApiCall', () => { await apiCall({}, undefined); assert.strictEqual(attempts, 1); - const span = harness.requireSingleSpan('google-gax'); + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 1 + attempts); + const span = spans.find(s => s.name === 'EchoClient.Echo')!; + assert.ok(span); // Asserting the transport too, so that a span accidentally // produced by the other one cannot satisfy this test. assert.strictEqual( @@ -1423,6 +1478,11 @@ describe('createApiCall', () => { transport.rpcType, ); harness.assertResendCount(0, {span}); + const attemptSpans = spans.filter( + s => s.name === 'google.example.v1.Echo/Echo', + ); + assert.strictEqual(attemptSpans.length, attempts); + harness.assertResendCount(0, {span: attemptSpans[0]}); }); it('reports one resend per retry', async () => { @@ -1448,7 +1508,10 @@ describe('createApiCall', () => { await apiCall({}, undefined); assert.strictEqual(attempts, 3); - const span = harness.requireSingleSpan('google-gax'); + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 1 + attempts); + const span = spans.find(s => s.name === 'EchoClient.Echo')!; + assert.ok(span); assert.strictEqual( span.attributes['gcp.method.type'], transport.rpcType, @@ -1458,10 +1521,19 @@ describe('createApiCall', () => { // off-by-one between the two is exactly what the attribute // defines. harness.assertResendCount(attempts - 1, {span}); - harness.assertResponseStatus({ - rpcStatus: 'OK', - ...(transport.rpcType === 'http' ? {httpStatus: 200} : {}), + const attemptSpans = spans.filter( + s => s.name === 'google.example.v1.Echo/Echo', + ); + assert.strictEqual(attemptSpans.length, attempts); + attemptSpans.forEach((attemptSpan, idx) => { + harness.assertResendCount(idx, {span: attemptSpan}); }); + harness.assertResponseStatus( + transport.rpcType === 'http' + ? {httpStatus: 200} + : {rpcStatus: 'OK'}, + {span}, + ); }); it('reports resends correctly when retries are exhausted by maxRetries', async () => { @@ -1505,13 +1577,23 @@ describe('createApiCall', () => { }); assert.strictEqual(attempts, 2); - const span = harness.requireSingleSpan('google-gax'); + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 1 + attempts); + const span = spans.find(s => s.name === 'EchoClient.Echo')!; + assert.ok(span); assert.strictEqual( span.attributes['gcp.method.type'], transport.rpcType, ); // 2 attempts made: initial send + 1 resend. The 2nd retry was not sent because maxRetries was reached. harness.assertResendCount(1, {span}); + const attemptSpans = spans.filter( + s => s.name === 'google.example.v1.Echo/Echo', + ); + assert.strictEqual(attemptSpans.length, attempts); + attemptSpans.forEach((attemptSpan, idx) => { + harness.assertResendCount(idx, {span: attemptSpan}); + }); }); it('reports resends correctly when retries are exhausted by totalTimeoutMillis', async () => { @@ -1553,7 +1635,10 @@ describe('createApiCall', () => { await apiCall({}, undefined); }); - const span = harness.requireSingleSpan('google-gax'); + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 1 + attempts); + const span = spans.find(s => s.name === 'EchoClient.Echo')!; + assert.ok(span); assert.strictEqual( span.attributes['gcp.method.type'], transport.rpcType, @@ -1562,6 +1647,13 @@ describe('createApiCall', () => { // resend count should match the number of retries actually made // (attempts - 1), without counting the attempt aborted by the deadline. harness.assertResendCount(attempts - 1, {span}); + const attemptSpans = spans.filter( + s => s.name === 'google.example.v1.Echo/Echo', + ); + assert.strictEqual(attemptSpans.length, attempts); + attemptSpans.forEach((attemptSpan, idx) => { + harness.assertResendCount(idx, {span: attemptSpan}); + }); }); }); } @@ -1622,9 +1714,19 @@ describe('createApiCall', () => { try { // Three attempts: the initial send plus the two allowed resends. assert.strictEqual(attempts, 3); - const span = harness.requireSingleSpan('google-gax'); + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 1 + attempts); + const span = spans.find(s => s.name === 'EchoClient.Echo')!; + assert.ok(span); assert.strictEqual(span.attributes['gcp.method.type'], 'grpc'); harness.assertResendCount(2, {span}); + const attemptSpans = spans.filter( + s => s.name === 'google.example.v1.Echo/Echo', + ); + assert.strictEqual(attemptSpans.length, attempts); + attemptSpans.forEach((attemptSpan, idx) => { + harness.assertResendCount(idx, {span: attemptSpan}); + }); done(); } catch (e) { done(e); @@ -1675,9 +1777,19 @@ describe('createApiCall', () => { stream.on('error', () => { try { assert.strictEqual(attempts, 3); - const span = harness.requireSingleSpan('google-gax'); + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 1 + attempts); + const span = spans.find(s => s.name === 'EchoClient.Echo')!; + assert.ok(span); assert.strictEqual(span.attributes['gcp.method.type'], 'grpc'); harness.assertResendCount(2, {span}); + const attemptSpans = spans.filter( + s => s.name === 'google.example.v1.Echo/Echo', + ); + assert.strictEqual(attemptSpans.length, attempts); + attemptSpans.forEach((attemptSpan, idx) => { + harness.assertResendCount(idx, {span: attemptSpan}); + }); done(); } catch (e) { done(e); @@ -1685,6 +1797,358 @@ describe('createApiCall', () => { }); }); }); + + describe('T3 client request span to T4 per-attempt span correlation', () => { + it('emits a T4 gRPC attempt span parented to its T3 client request span', async () => { + const grpcClient = new GrpcClient({ + servicePath: 'echo.googleapis.com', + port: 443, + }); + const defaults = grpcClient.constructSettings( + 'google.example.v1.Echo', + { + interfaces: { + 'google.example.v1.Echo': { + methods: { + Echo: {timeout_millis: 5000}, + }, + }, + }, + }, + {}, + {'x-goog-api-client': 'test'}, + true, + telemetryInfo, + ); + + const stubFunc = ( + argument: {}, + metadata: {}, + options: {}, + callback: Function, + ): GRPCCallResult => { + callback(null, {echo: 'ok'}); + return {cancel: () => {}}; + }; + + const apiCall = gaxCreateApiCall( + stubFunc as unknown as GRPCCall, + defaults.echo, + ); + await apiCall({message: 'hello'}, undefined); + + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 2); + + const t4Span = spans.find( + s => s.name === 'google.example.v1.Echo/Echo', + )!; + const t3Span = spans.find(s => s.name === 'EchoClient.Echo')!; + assert.ok(t4Span); + assert.ok(t3Span); + + assert.strictEqual(t3Span.kind, SpanKind.INTERNAL); + assert.strictEqual(t4Span.kind, SpanKind.CLIENT); + assert.strictEqual( + t4Span.spanContext().traceId, + t3Span.spanContext().traceId, + ); + assert.strictEqual( + t4Span.parentSpanContext?.spanId, + t3Span.spanContext().spanId, + ); + assert.strictEqual( + t4Span.attributes['url.domain'], + 'echo.googleapis.com', + ); + assert.strictEqual( + t4Span.attributes['server.address'], + 'echo.googleapis.com', + ); + assert.strictEqual(t4Span.attributes['server.port'], 443); + assert.strictEqual( + t4Span.attributes['rpc.method'], + 'google.example.v1.Echo/Echo', + ); + assert.strictEqual(t4Span.attributes['http.request.method'], undefined); + assert.strictEqual(t4Span.attributes['rpc.response.status_code'], 'OK'); + assert.strictEqual( + t4Span.attributes['grpc.response.status_code'], + undefined, + ); + assert.strictEqual( + t4Span.attributes['http.response.status_code'], + undefined, + ); + }); + + it('emits a T4 HTTP/REST attempt span parented to its T3 client request span', async () => { + const fallbackClient = new FallbackGrpcClient({ + servicePath: 'echo.googleapis.com', + port: 443, + }); + const defaults = fallbackClient.constructSettings( + 'google.example.v1.Echo', + { + interfaces: { + 'google.example.v1.Echo': { + methods: { + Echo: {timeout_millis: 5000}, + }, + }, + }, + }, + {}, + {'x-goog-api-client': 'test'}, + true, + telemetryInfo, + ); + + const stubFunc = ( + argument: {}, + metadata: {}, + options: {}, + callback: Function, + ): GRPCCallResult => { + callback(null, {echo: 'ok'}); + return {cancel: () => {}}; + }; + + const apiCall = gaxCreateApiCall( + stubFunc as unknown as GRPCCall, + defaults.echo, + undefined, + 'rest', + ); + await apiCall({message: 'hello'}, undefined); + + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 2); + + const t4Span = spans.find( + s => s.name === 'google.example.v1.Echo/Echo', + )!; + const t3Span = spans.find(s => s.name === 'EchoClient.Echo')!; + assert.ok(t4Span); + assert.ok(t3Span); + + assert.strictEqual(t3Span.kind, SpanKind.INTERNAL); + assert.strictEqual(t4Span.kind, SpanKind.CLIENT); + assert.strictEqual( + t4Span.spanContext().traceId, + t3Span.spanContext().traceId, + ); + assert.strictEqual( + t4Span.parentSpanContext?.spanId, + t3Span.spanContext().spanId, + ); + assert.strictEqual( + t4Span.attributes['url.domain'], + 'echo.googleapis.com', + ); + assert.strictEqual( + t4Span.attributes['server.address'], + 'echo.googleapis.com', + ); + assert.strictEqual(t4Span.attributes['server.port'], 443); + assert.strictEqual(t4Span.attributes['http.request.method'], 'POST'); + assert.strictEqual(t4Span.attributes['rpc.method'], undefined); + assert.strictEqual( + t4Span.attributes['rpc.response.status_code'], + undefined, + ); + assert.strictEqual(t4Span.attributes['http.response.status_code'], 200); + assert.strictEqual( + t3Span.attributes['rpc.response.status_code'], + undefined, + ); + assert.strictEqual(t3Span.attributes['http.response.status_code'], 200); + }); + + it('ties concurrent T4 attempt spans to their respective T3 client request spans without cross-talk', async () => { + const echoSettings = new gax.CallSettings({ + apiName: 'google.example.v1.Echo', + enableTelemetryTracing: true, + otherArgs: { + internalTelemetryInfo: telemetryInfo, + internalMethodName: 'Echo', + }, + }); + + const expandSettings = new gax.CallSettings({ + apiName: 'google.example.v1.Echo', + enableTelemetryTracing: true, + otherArgs: { + internalTelemetryInfo: telemetryInfo, + internalMethodName: 'Expand', + }, + }); + + const makeStub = (delayMs: number): GRPCCall => { + return (( + argument: {}, + metadata: {}, + options: {}, + callback: Function, + ): GRPCCallResult => { + setTimeout(() => { + callback(null, {ok: true}); + }, delayMs); + return {cancel: () => {}}; + }) as unknown as GRPCCall; + }; + + const echoCall = gaxCreateApiCall(makeStub(15), echoSettings); + const expandCall = gaxCreateApiCall(makeStub(5), expandSettings); + + await Promise.all([ + echoCall({id: 1}, undefined), + expandCall({id: 2}, undefined), + ]); + + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 4); + + const t3Echo = spans.find(s => s.name === 'EchoClient.Echo')!; + const t3Expand = spans.find(s => s.name === 'EchoClient.Expand')!; + const t4Echo = spans.find( + s => s.name === 'google.example.v1.Echo/Echo', + )!; + const t4Expand = spans.find( + s => s.name === 'google.example.v1.Echo/Expand', + )!; + + assert.ok(t3Echo && t3Expand && t4Echo && t4Expand); + assert.notStrictEqual( + t3Echo.spanContext().spanId, + t3Expand.spanContext().spanId, + ); + + assert.strictEqual( + t4Echo.spanContext().traceId, + t3Echo.spanContext().traceId, + ); + assert.strictEqual( + t4Echo.parentSpanContext?.spanId, + t3Echo.spanContext().spanId, + ); + + assert.strictEqual( + t4Expand.spanContext().traceId, + t3Expand.spanContext().traceId, + ); + assert.strictEqual( + t4Expand.parentSpanContext?.spanId, + t3Expand.spanContext().spanId, + ); + }); + + it('emits one T4 attempt span per retry attempt, all parented to the single T3 client request span', async () => { + const retryOptions = gax.createRetryOptions( + [status.UNAVAILABLE], + gax.createBackoffSettings(1, 1.1, 5, 100, 1.0, 100, 1000), + ); + + const settings = new gax.CallSettings({ + apiName: 'google.example.v1.Echo', + retry: retryOptions, + enableTelemetryTracing: true, + otherArgs: { + internalTelemetryInfo: telemetryInfo, + internalMethodName: 'Echo', + }, + }); + + let attempt = 0; + const stubFunc = ( + argument: {}, + metadata: {}, + options: {}, + callback: Function, + ): GRPCCallResult => { + attempt++; + if (attempt === 1) { + const err = new GoogleError('transient failure'); + err.code = status.UNAVAILABLE; + callback(err); + } else { + callback(null, {echo: 'recovered'}); + } + return {cancel: () => {}}; + }; + + const apiCall = gaxCreateApiCall( + stubFunc as unknown as GRPCCall, + settings, + ); + await apiCall({message: 'retry-me'}, undefined); + + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 3); + + const t3Span = spans.find(s => s.name === 'EchoClient.Echo')!; + const t4Spans = spans.filter( + s => s.name === 'google.example.v1.Echo/Echo', + ); + assert.ok(t3Span); + assert.strictEqual(t4Spans.length, 2); + + // First attempt failed with UNAVAILABLE (resend_count omitted on initial attempt) + assert.strictEqual(t4Spans[0].kind, SpanKind.CLIENT); + assert.strictEqual(t4Spans[0].status.code, SpanStatusCode.ERROR); + assert.strictEqual(t4Spans[0].status.message, 'transient failure'); + assert.strictEqual( + t4Spans[0].attributes['gcp.grpc.resend_count'], + undefined, + ); + assert.strictEqual(t4Spans[0].attributes['error.type'], 'UNAVAILABLE'); + assert.strictEqual( + t4Spans[0].attributes['rpc.response.status_code'], + 'UNAVAILABLE', + ); + assert.strictEqual( + t4Spans[0].attributes['grpc.response.status_code'], + undefined, + ); + assert.strictEqual(t4Spans[0].events.length, 1); + assert.strictEqual(t4Spans[0].events[0].name, 'exception'); + assert.strictEqual( + t4Spans[0].events[0].attributes?.['exception.type'], + 'GoogleError', + ); + assert.strictEqual( + t4Spans[0].parentSpanContext?.spanId, + t3Span.spanContext().spanId, + ); + + // Second attempt succeeded with OK (resend_count = 1) + assert.strictEqual(t4Spans[1].kind, SpanKind.CLIENT); + assert.strictEqual(t4Spans[1].status.code, SpanStatusCode.UNSET); + assert.strictEqual(t4Spans[1].attributes['gcp.grpc.resend_count'], 1); + assert.strictEqual( + t4Spans[1].attributes['rpc.response.status_code'], + 'OK', + ); + assert.strictEqual( + t4Spans[1].attributes['grpc.response.status_code'], + undefined, + ); + assert.strictEqual( + t4Spans[1].parentSpanContext?.spanId, + t3Span.spanContext().spanId, + ); + + // Overall T3 call span succeeded with resend_count = 1 + assert.strictEqual(t3Span.kind, SpanKind.INTERNAL); + assert.strictEqual(t3Span.status.code, SpanStatusCode.UNSET); + assert.strictEqual(t3Span.attributes['gcp.grpc.resend_count'], 1); + assert.strictEqual(t3Span.attributes['rpc.response.status_code'], 'OK'); + assert.strictEqual( + t3Span.attributes['grpc.response.status_code'], + undefined, + ); + }); + }); }); }); diff --git a/core/packages/gax/test/unit/otelHarness.ts b/core/packages/gax/test/unit/otelHarness.ts index 21a6fae262f..dd53cbcc090 100644 --- a/core/packages/gax/test/unit/otelHarness.ts +++ b/core/packages/gax/test/unit/otelHarness.ts @@ -296,7 +296,8 @@ export class OtelHarness { ): void { const target = options.span ?? this.requireSingleSpan(options.tracerName); const actual = this.responseStatus(target); - const transport = target.attributes['gcp.method.type']; + const transport = + target.attributes['gcp.method.type'] ?? target.attributes['rpc.system']; const where = `span '${target.name}'`; if ('serverAddress' in expected || 'serverPort' in expected) { @@ -314,26 +315,21 @@ export class OtelHarness { ); assert.strictEqual( - actual.rpc, - expected.rpcStatus, - expected.rpcStatus === undefined - ? `expected ${where} to report no rpc.response.status_code, got ` + - `${JSON.stringify(actual.rpc)}. Response status is omitted when there is no server response.` - : `expected ${where} to report rpc.response.status_code ` + - `${JSON.stringify(expected.rpcStatus)}, got ${JSON.stringify(actual.rpc)}. ` + - 'This attribute is reported on every call with a server response, on both transports.', + actual.grpc, + undefined, + `${where} reported grpc.response.status_code ${JSON.stringify(actual.grpc)}; ` + + 'only rpc.response.status_code or http.response.status_code should be set.', ); if (transport === 'grpc') { assert.strictEqual( - actual.grpc, + actual.rpc, expected.rpcStatus, expected.rpcStatus === undefined - ? `expected ${where} to report no grpc.response.status_code, got ` + - `${JSON.stringify(actual.grpc)}. Response status is omitted when there is no server response.` - : `expected ${where} to report grpc.response.status_code ` + - `${JSON.stringify(expected.rpcStatus)}, got ${JSON.stringify(actual.grpc)}. ` + - 'On a gRPC span it mirrors rpc.response.status_code.', + ? `expected ${where} to report no rpc.response.status_code, got ` + + `${JSON.stringify(actual.rpc)}. Response status is omitted when there is no server response.` + : `expected ${where} to report rpc.response.status_code ` + + `${JSON.stringify(expected.rpcStatus)}, got ${JSON.stringify(actual.rpc)}.`, ); assert.strictEqual( actual.http, @@ -353,11 +349,16 @@ export class OtelHarness { } assert.strictEqual( - actual.grpc, + actual.rpc, + undefined, + `${where} is a fallback span but reported rpc.response.status_code ` + + `${JSON.stringify(actual.rpc)}. HTTP spans only report http.response.status_code.`, + ); + assert.strictEqual( + expected.rpcStatus, undefined, - `${where} is a fallback span but reported grpc.response.status_code ` + - `${JSON.stringify(actual.grpc)}. The gRPC status is reported as ` + - 'rpc.response.status_code there, not under the grpc.* name.', + 'assertResponseStatus was given an expected rpcStatus for an HTTP ' + + 'span, which can never hold one. Use httpStatus instead.', ); assert.strictEqual( actual.http, @@ -431,7 +432,8 @@ export class OtelHarness { options: {tracerName?: string; span?: ReadableSpan} = {}, ): void { const target = options.span ?? this.requireSingleSpan(options.tracerName); - const transport = target.attributes['gcp.method.type']; + const transport = + target.attributes['gcp.method.type'] ?? target.attributes['rpc.system']; const where = `span '${target.name}'`; assert.ok( diff --git a/core/packages/gax/test/unit/tracerHelper.ts b/core/packages/gax/test/unit/tracerHelper.ts index 647fae36b01..4a1af5769dc 100644 --- a/core/packages/gax/test/unit/tracerHelper.ts +++ b/core/packages/gax/test/unit/tracerHelper.ts @@ -18,14 +18,17 @@ import * as assert from 'assert'; import * as vm from 'vm'; import {EventEmitter} from 'events'; import {Duplex, Writable} from 'stream'; -import {SpanStatusCode, trace} from '@opentelemetry/api'; +import {SpanKind, SpanStatusCode, trace} from '@opentelemetry/api'; import {describe, it, beforeEach, afterEach} from 'mocha'; import * as grpc from '@grpc/grpc-js'; import { getGaxTracer, traceCall, + traceAttempt, + setAttemptHttpMethod, handlePromise, handleStream, + AttemptTraceContext, DynamicTraceContext, StaticTraceContext, resolveErrorInfoReason, @@ -119,6 +122,10 @@ describe('TracerHelper', () => { ); assert.strictEqual(span.attributes['gcp.method.name'], 'GetObject'); assert.strictEqual(span.attributes['gcp.method.type'], 'grpc'); + assert.strictEqual( + span.attributes['url.domain'], + 'storage.googleapis.com', + ); // A successful call reports no error.type, and leaves the status unset // rather than claiming OK on the application's behalf. assert.strictEqual(span.attributes['error.type'], undefined); @@ -415,7 +422,7 @@ describe('TracerHelper', () => { await errorTypeOf(error, httpDynamicArgs), 'CLIENT_TIMEOUT', ); - harness.assertResponseStatus({rpcStatus: 'DEADLINE_EXCEEDED'}); + harness.assertResponseStatus({}); }); it('prefers a system error code to the class on either transport', async () => { @@ -448,7 +455,7 @@ describe('TracerHelper', () => { await errorTypeOf(error, httpDynamicArgs), 'CLIENT_CONNECTION_ERROR', ); - harness.assertResponseStatus({rpcStatus: 'UNAVAILABLE'}); + harness.assertResponseStatus({}); }); it('checks e.cause when the outer error is a GoogleError', async () => { @@ -1992,10 +1999,10 @@ describe('TracerHelper', () => { harness.assertResponseStatus({rpcStatus: 'OK'}); }); - it('reports OK and 200 on a successful http call', async () => { + it('reports 200 on a successful http call', async () => { await traceCall(httpDynamicArgs, staticArgs, async () => 'ok'); - harness.assertResponseStatus({rpcStatus: 'OK', httpStatus: 200}); + harness.assertResponseStatus({httpStatus: 200}); }); it('reports the gRPC status name on a failed grpc call', async () => { @@ -2013,7 +2020,7 @@ describe('TracerHelper', () => { assert.strictEqual(spans[0].status.message, '5 NOT_FOUND'); }); - it('reports the received http status alongside the mapped gRPC status', async () => { + it('reports the received http status on a failed http call', async () => { // 418 is unmapped, so rpcCodeFromHttpStatusCode collapses it to // FAILED_PRECONDITION. The received status is therefore not // recoverable from `code`, which is why it is carried separately. @@ -2028,7 +2035,6 @@ describe('TracerHelper', () => { }); harness.assertResponseStatus({ - rpcStatus: 'FAILED_PRECONDITION', httpStatus: 418, }); const spans = harness.getSpans('google-gax'); @@ -2048,7 +2054,7 @@ describe('TracerHelper', () => { }); }); - harness.assertResponseStatus({rpcStatus: 'DEADLINE_EXCEEDED'}); + harness.assertResponseStatus({}); const spans = harness.getSpans('google-gax'); assert.strictEqual(spans[0].status.message, 'Deadline exceeded'); }); @@ -2143,7 +2149,7 @@ describe('TracerHelper', () => { }); }); - harness.assertResponseStatus({rpcStatus: undefined}); + harness.assertResponseStatus({}); const span = harness.requireSingleSpan('google-gax'); assert.strictEqual(span.status.message, 'client-side http error'); assert.strictEqual( @@ -2253,7 +2259,7 @@ describe('TracerHelper', () => { const span = harness.requireSingleSpan('google-gax'); assert.strictEqual(span.attributes['error.type'], 'NOT_FOUND'); assert.strictEqual( - span.attributes['grpc.response.status_code'], + span.attributes['rpc.response.status_code'], 'NOT_FOUND', ); }); @@ -3825,4 +3831,173 @@ describe('TracerHelper', () => { assertListenersRestored(baseline, 'after stream end'); }); }); + + describe('traceAttempt', () => { + const telemetryInfo: StaticTraceContext = { + gcpClientService: 'echo.googleapis.com', + gcpVersion: '1.2.3', + gcpRepo: 'googleapis/google-cloud-node', + gcpArtifact: '@google-cloud/echo', + }; + + it('creates a CLIENT T4 span in traceAttempt with url.domain, server.address, server.port, and status_code', async () => { + const attemptArgs: AttemptTraceContext = { + apiName: 'google.example.v1.Echo', + clientName: 'EchoClient', + methodName: 'Echo', + rpcType: 'grpc', + }; + + await traceAttempt(attemptArgs, telemetryInfo, async () => { + return [{echo: 'ok'}, undefined, undefined]; + }); + + const span = harness.requireSingleSpan('google-gax'); + assert.strictEqual(span.name, 'google.example.v1.Echo/Echo'); + assert.strictEqual(span.kind, SpanKind.CLIENT); + assert.strictEqual(span.attributes['url.domain'], 'echo.googleapis.com'); + assert.strictEqual( + span.attributes['server.address'], + 'echo.googleapis.com', + ); + assert.strictEqual(span.attributes['server.port'], 443); + assert.strictEqual(span.attributes['rpc.system'], 'grpc'); + assert.strictEqual( + span.attributes['rpc.method'], + 'google.example.v1.Echo/Echo', + ); + assert.strictEqual(span.attributes['http.request.method'], undefined); + assert.strictEqual(span.attributes['rpc.response.status_code'], 'OK'); + assert.strictEqual( + span.attributes['grpc.response.status_code'], + undefined, + ); + assert.strictEqual(span.attributes['gcp.repo'], undefined); + assert.strictEqual(span.attributes['gcp.method.type'], undefined); + assert.strictEqual(span.attributes['gcp.method.name'], undefined); + assert.strictEqual(span.attributes['gcp.client.version'], undefined); + assert.strictEqual(span.attributes['gcp.artifact'], undefined); + }); + + it('sets http.request.method on HTTP T4 spans and allows setAttemptHttpMethod to update it', async () => { + await traceAttempt( + { + apiName: 'google.example.v1.Echo', + clientName: 'EchoClient', + methodName: 'Echo', + rpcType: 'http', + }, + telemetryInfo, + async () => [{echo: 'ok'}, undefined, undefined], + ); + + await traceAttempt( + { + apiName: 'google.example.v1.Echo', + clientName: 'EchoClient', + methodName: 'Echo', + rpcType: 'http', + }, + telemetryInfo, + async () => { + setAttemptHttpMethod('GET'); + return [{echo: 'ok'}, undefined, undefined]; + }, + ); + + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 2); + assert.strictEqual(spans[0].attributes['http.request.method'], 'POST'); + assert.strictEqual(spans[0].attributes['rpc.method'], undefined); + assert.strictEqual(spans[0].attributes['http.response.status_code'], 200); + assert.strictEqual( + spans[0].attributes['rpc.response.status_code'], + undefined, + ); + assert.strictEqual(spans[1].attributes['http.request.method'], 'GET'); + assert.strictEqual(spans[1].attributes['rpc.method'], undefined); + assert.strictEqual(spans[1].attributes['http.response.status_code'], 200); + assert.strictEqual( + spans[1].attributes['rpc.response.status_code'], + undefined, + ); + }); + + it('omits server.address and server.port on T4 span for pre-connection failures while preserving url.domain', async () => { + const attemptArgs: AttemptTraceContext = { + apiName: 'google.example.v1.Echo', + clientName: 'EchoClient', + methodName: 'Echo', + rpcType: 'grpc', + }; + + const dnsError = Object.assign(new Error('getaddrinfo ENOTFOUND'), { + code: 'ENOTFOUND', + }); + + await assert.rejects(async () => { + await traceAttempt(attemptArgs, telemetryInfo, async () => { + throw dnsError; + }); + }); + + const span = harness.requireSingleSpan('google-gax'); + assert.strictEqual(span.name, 'google.example.v1.Echo/Echo'); + assert.strictEqual(span.kind, SpanKind.CLIENT); + assert.strictEqual(span.attributes['url.domain'], 'echo.googleapis.com'); + assert.strictEqual(span.attributes['server.address'], undefined); + assert.strictEqual(span.attributes['server.port'], undefined); + assert.strictEqual( + span.attributes['error.type'], + 'CLIENT_CONNECTION_ERROR', + ); + assert.strictEqual(span.status.code, SpanStatusCode.ERROR); + assert.strictEqual(span.status.message, 'getaddrinfo ENOTFOUND'); + assert.strictEqual(span.events.length, 1); + assert.strictEqual(span.events[0].name, 'exception'); + assert.strictEqual( + span.events[0].attributes?.['exception.type'], + 'Error', + ); + }); + + it('sets gcp.grpc.resend_count and http.request.resend_count on T4 spans when resendCount > 0', async () => { + await traceAttempt( + { + apiName: 'google.example.v1.Echo', + clientName: 'EchoClient', + methodName: 'Echo', + rpcType: 'grpc', + resendCount: 2, + }, + telemetryInfo, + async () => [{echo: 'ok'}, undefined, undefined], + ); + + await traceAttempt( + { + apiName: 'google.example.v1.Echo', + clientName: 'EchoClient', + methodName: 'Echo', + rpcType: 'http', + resendCount: 3, + }, + telemetryInfo, + async () => [{echo: 'ok'}, undefined, undefined], + ); + + const spans = harness.getSpans('google-gax'); + assert.strictEqual(spans.length, 2); + assert.strictEqual(spans[0].attributes['gcp.grpc.resend_count'], 2); + assert.strictEqual( + spans[0].attributes['http.request.resend_count'], + undefined, + ); + assert.strictEqual(spans[1].attributes['http.request.resend_count'], 3); + assert.strictEqual( + spans[1].attributes['gcp.grpc.resend_count'], + undefined, + ); + }); + }); });