Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 4 additions & 17 deletions core/packages/gax/samples/observability.js
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand All @@ -38,27 +33,19 @@ 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');

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.');
Expand Down
71 changes: 58 additions & 13 deletions core/packages/gax/src/createApiCall.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
* Provides function wrappers that implement page streaming and retrying.
*/

import {context} from '@opentelemetry/api';
import {createAPICaller} from './apiCaller';
import {
APICallback,
Expand All @@ -26,6 +27,7 @@ import {
GRPCCallOtherArgs,
RequestType,
SimpleCallbackFunction,
UnaryCall,
} from './apitypes';
import {Descriptor} from './descriptor';
import {CallOptions, CallSettings, convertRetryOptions} from './gax';
Expand All @@ -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';
Expand Down Expand Up @@ -77,13 +81,35 @@ 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,
callOptions?: CallOptions,
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;
Expand Down Expand Up @@ -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;
}
Comment thread
shivanee-p marked this conversation as resolved.

const retry = thisSettings.retry;

if (streaming && retry) {
Expand Down Expand Up @@ -159,7 +213,7 @@ export function createApiCall(
retry.backoffSettings.initialRpcTimeoutMillis ??=
thisSettings.timeout;
return retryable(
func,
wrappedFunc,
thisSettings.retry!,
thisSettings.otherArgs as GRPCCallOtherArgs,
thisSettings.apiName,
Expand All @@ -168,7 +222,7 @@ export function createApiCall(
}
}
return addTimeoutArg(
func,
wrappedFunc,
thisSettings.timeout,
thisSettings.otherArgs as GRPCCallOtherArgs,
);
Expand All @@ -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,
Expand Down
17 changes: 16 additions & 1 deletion core/packages/gax/src/fallback.ts
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,8 @@ export class GrpcClient {
httpRules?: Array<google.api.IHttpRule>;
numericEnums: boolean;
minifyJson: boolean;
private _servicePath?: string;
private _port?: number;

/**
* In rare cases users might need to deallocate all memory consumed by loaded protos.
Expand Down Expand Up @@ -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;
}

/**
Expand Down Expand Up @@ -279,12 +285,21 @@ export class GrpcClient {
}
return metadata;
}
const otherArgs: Record<string, unknown> = {
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,
);
Expand Down
3 changes: 3 additions & 0 deletions core/packages/gax/src/fallbackServiceStub.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
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';
Expand Down Expand Up @@ -302,6 +303,8 @@
};
}

setAttemptHttpMethod(fetchParameters.method);

const cancelController = new AbortController();
const cancelSignal = cancelController.signal as AbortSignal;
let cancelRequested = false;
Expand Down Expand Up @@ -433,7 +436,7 @@
// state, as the handlers below do.
if (err && (timedOut || !cancelRequested)) {
if (callback) {
callback(err);

Check warning on line 439 in core/packages/gax/src/fallbackServiceStub.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid calling back inside of a promise
}
streamArrayParser.emit('error', err);
}
Expand All @@ -448,7 +451,7 @@
Promise.resolve(response.ok),
response.arrayBuffer(),
])
.then(([ok, buffer]: [boolean, Buffer | ArrayBuffer]) => {

Check warning on line 454 in core/packages/gax/src/fallbackServiceStub.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid nesting promises
const response = responseDecoder(
rpc,
ok,
Expand All @@ -458,7 +461,7 @@
callback!(null, response);
return;
})
.catch((err: Error) => {

Check warning on line 464 in core/packages/gax/src/fallbackServiceStub.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid nesting promises
// The deadline can expire after the response headers arrive but
// before the body is fully read, which rejects here rather than
// in the outer handler.
Expand All @@ -481,7 +484,7 @@
// state we recorded.
if (timedOut || !cancelRequested) {
if (callback) {
callback(callErr);

Check warning on line 487 in core/packages/gax/src/fallbackServiceStub.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid calling back inside of a promise
}
streamArrayParser.emit('error', callErr);
}
Expand Down Expand Up @@ -539,12 +542,12 @@
// nobody is listening to any more.
if (timedOut || !cancelRequested) {
if (callback) {
callback(err);

Check warning on line 545 in core/packages/gax/src/fallbackServiceStub.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid calling back inside of a promise
}
streamArrayParser.emit('error', err);
}
} else if (callback) {
callback(err);

Check warning on line 550 in core/packages/gax/src/fallbackServiceStub.ts

View workflow job for this annotation

GitHub Actions / lint

Avoid calling back inside of a promise
} else {
throw err;
}
Expand Down
18 changes: 17 additions & 1 deletion core/packages/gax/src/grpc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,9 @@ export interface GrpcClientOptions extends GoogleAuthOptions {
httpRules?: Array<google.api.IHttpRule>;
numericEnums?: boolean;
universeDomain?: string;
servicePath?: string;
apiEndpoint?: string;
port?: number;
}

export interface MetadataValue {
Expand Down Expand Up @@ -121,6 +124,8 @@ export class GrpcClient {
fallback: boolean | 'rest' | 'proto';
private static protoCache = new Map<string, grpc.GrpcObject>();
httpRules?: Array<google.api.IHttpRule>;
private _servicePath?: string;
private _port?: number;
/**
* Base directory for resolving client certificates.
*
Expand Down Expand Up @@ -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]);
Expand Down Expand Up @@ -388,12 +395,21 @@ export class GrpcClient {
enableTelemetryTracing?: boolean,
internalTelemetryInfo?: StaticTraceContext,
) {
const otherArgs: Record<string, unknown> = {
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,
);
Expand Down
Loading
Loading