diff --git a/.gitignore b/.gitignore index 8e63d8bf0..f35fc5cf0 100644 --- a/.gitignore +++ b/.gitignore @@ -23,3 +23,7 @@ coverage # Node's bash completion file .node_bash_completion + +# Copied proto dependencies +packages/grpc-js/proto/protoc-gen-validate/ +packages/grpc-js/proto/xds/ diff --git a/packages/grpc-js/src/call-credentials.ts b/packages/grpc-js/src/call-credentials.ts index a9afe4ae3..fa1c8b92d 100644 --- a/packages/grpc-js/src/call-credentials.ts +++ b/packages/grpc-js/src/call-credentials.ts @@ -162,6 +162,9 @@ class ComposedCallCredentials extends CallCredentials { } compose(other: CallCredentials): CallCredentials { + if (other instanceof EmptyCallCredentials) { + return this; + } return new ComposedCallCredentials(this.creds.concat([other])); } @@ -197,6 +200,9 @@ class SingleCallCredentials extends CallCredentials { } compose(other: CallCredentials): CallCredentials { + if (other instanceof EmptyCallCredentials) { + return this; + } return new ComposedCallCredentials([this, other]); } @@ -225,3 +231,9 @@ class EmptyCallCredentials extends CallCredentials { return other instanceof EmptyCallCredentials; } } + +export function isEmptyCallCredentials( + callCredentials: CallCredentials +): boolean { + return callCredentials instanceof EmptyCallCredentials; +} diff --git a/packages/grpc-js/src/channel-credentials.ts b/packages/grpc-js/src/channel-credentials.ts index a6ded81ea..04b162899 100644 --- a/packages/grpc-js/src/channel-credentials.ts +++ b/packages/grpc-js/src/channel-credentials.ts @@ -184,6 +184,7 @@ class InsecureChannelCredentialsImpl extends ChannelCredentials { return other instanceof InsecureChannelCredentialsImpl; } _createSecureConnector(channelTarget: GrpcUri, options: ChannelOptions, callCredentials?: CallCredentials): SecureConnector { + const credentials = callCredentials ?? CallCredentials.createEmpty(); return { connect(socket) { return Promise.resolve({ @@ -195,7 +196,7 @@ class InsecureChannelCredentialsImpl extends ChannelCredentials { return Promise.resolve(); }, getCallCredentials: () => { - return callCredentials ?? CallCredentials.createEmpty(); + return credentials; }, destroy() {} } diff --git a/packages/grpc-js/src/channelz.ts b/packages/grpc-js/src/channelz.ts index 7242d7160..cb00b4e41 100644 --- a/packages/grpc-js/src/channelz.ts +++ b/packages/grpc-js/src/channelz.ts @@ -263,11 +263,11 @@ export class ChannelzCallTracker { callsStarted = 0; callsSucceeded = 0; callsFailed = 0; - lastCallStartedTimestamp: Date | null = null; + lastCallStartedTimestamp: Date | number | null = null; addCallStarted() { this.callsStarted += 1; - this.lastCallStartedTimestamp = new Date(); + this.lastCallStartedTimestamp = Date.now(); } addCallSucceeded() { this.callsSucceeded += 1; @@ -324,10 +324,10 @@ export interface SocketInfo { messagesSent: number; messagesReceived: number; keepAlivesSent: number; - lastLocalStreamCreatedTimestamp: Date | null; - lastRemoteStreamCreatedTimestamp: Date | null; - lastMessageSentTimestamp: Date | null; - lastMessageReceivedTimestamp: Date | null; + lastLocalStreamCreatedTimestamp: Date | number | null; + lastRemoteStreamCreatedTimestamp: Date | number | null; + lastMessageSentTimestamp: Date | number | null; + lastMessageReceivedTimestamp: Date | number | null; localFlowControlWindow: number | null; remoteFlowControlWindow: number | null; } @@ -542,13 +542,15 @@ function connectivityStateToMessage( } } -function dateToProtoTimestamp(date?: Date | null): Timestamp | null { - if (!date) { +export function dateToProtoTimestamp( + date?: Date | number | null +): Timestamp | null { + if (date === null || date === undefined) { return null; } - const millisSinceEpoch = date.getTime(); + const millisSinceEpoch = typeof date === 'number' ? date : date.getTime(); return { - seconds: (millisSinceEpoch / 1000) | 0, + seconds: Math.floor(millisSinceEpoch / 1000), nanos: (millisSinceEpoch % 1000) * 1_000_000, }; } diff --git a/packages/grpc-js/src/deadline.ts b/packages/grpc-js/src/deadline.ts index de05e381e..c20dee06d 100644 --- a/packages/grpc-js/src/deadline.ts +++ b/packages/grpc-js/src/deadline.ts @@ -37,7 +37,7 @@ const units: Array<[string, number]> = [ ]; export function getDeadlineTimeoutString(deadline: Deadline) { - const now = new Date().getTime(); + const now = Date.now(); if (deadline instanceof Date) { deadline = deadline.getTime(); } @@ -70,7 +70,7 @@ const MAX_TIMEOUT_TIME = 2147483647; */ export function getRelativeTimeout(deadline: Deadline) { const deadlineMs = deadline instanceof Date ? deadline.getTime() : deadline; - const now = new Date().getTime(); + const now = Date.now(); const timeout = deadlineMs - now; if (timeout < 0) { return 0; @@ -101,6 +101,12 @@ export function deadlineToString(deadline: Deadline): string { * @param endDate * @returns */ -export function formatDateDifference(startDate: Date, endDate: Date): string { - return ((endDate.getTime() - startDate.getTime()) / 1000).toFixed(3) + 's'; +export function formatDateDifference( + startDate: Date | number, + endDate: Date | number +): string { + const startMs = + typeof startDate === 'number' ? startDate : startDate.getTime(); + const endMs = typeof endDate === 'number' ? endDate : endDate.getTime(); + return ((endMs - startMs) / 1000).toFixed(3) + 's'; } diff --git a/packages/grpc-js/src/internal-channel.ts b/packages/grpc-js/src/internal-channel.ts index 9c05b91eb..d269ae1bd 100644 --- a/packages/grpc-js/src/internal-channel.ts +++ b/packages/grpc-js/src/internal-channel.ts @@ -34,7 +34,7 @@ import { import { trace, isTracerEnabled } from './logging'; import { SubchannelAddress } from './subchannel-address'; import { mapProxyName } from './http_proxy'; -import { GrpcUri, parseUri, uriToString } from './uri-parser'; +import { GrpcUri, computeServiceUrl, parseUri, uriToString } from './uri-parser'; import { ServerSurfaceCall } from './server-call'; import { ConnectivityState } from './connectivity-state'; @@ -224,7 +224,7 @@ export class InternalChannel { private callCount = 0; private idleTimer: NodeJS.Timeout | null = null; private readonly idleTimeoutMs: number; - private lastActivityTimestamp: Date; + private lastActivityTimestamp: number; // Channelz info private readonly channelzEnabled: boolean = true; @@ -240,6 +240,21 @@ export class InternalChannel { Math.random() * Number.MAX_SAFE_INTEGER ); + /** + * Maximum number of distinct hosts to cache service URLs for on this channel. + */ + private static readonly MAX_CACHED_HOSTS = 5; + /** + * Maximum number of distinct method paths to cache per host. + */ + private static readonly MAX_METHODS_PER_HOST = 100; + /** + * Two-level cache mapping host to method name to precomputed service URL + * (`https://${hostname}/${serviceName}`) used by call credentials. Nested + * maps avoid intermediate composite key string allocations on the hot path. + */ + private readonly serviceUrlCache = new Map>(); + constructor( target: string, private readonly credentials: ChannelCredentials, @@ -461,7 +476,7 @@ export class InternalChannel { error.stack?.substring(error.stack.indexOf('\n') + 1) ); } - this.lastActivityTimestamp = new Date(); + this.lastActivityTimestamp = Date.now(); } private get traceEnabled(): boolean { @@ -645,9 +660,8 @@ export class InternalChannel { this.startIdleTimeout(this.idleTimeoutMs); return; } - const now = new Date(); - const timeSinceLastActivity = - now.valueOf() - this.lastActivityTimestamp.valueOf(); + const now = Date.now(); + const timeSinceLastActivity = now - this.lastActivityTimestamp; if (timeSinceLastActivity >= this.idleTimeoutMs) { this.trace( 'Idle timer triggered after ' + @@ -691,10 +705,37 @@ export class InternalChannel { } } this.callCount -= 1; - this.lastActivityTimestamp = new Date(); + this.lastActivityTimestamp = Date.now(); this.maybeStartIdleTimer(); } + /** + * Returns the precomputed service URL (`https://${hostname}/${serviceName}`) + * for the given host and method, using a bounded cache to avoid parsing on + * every RPC. + */ + getServiceUrl(host: string, methodName: string): string { + let methodMap = this.serviceUrlCache.get(host); + if (methodMap === undefined) { + if (this.serviceUrlCache.size < InternalChannel.MAX_CACHED_HOSTS) { + methodMap = new Map(); + this.serviceUrlCache.set(host, methodMap); + } + } + if (methodMap !== undefined) { + const serviceUrl = methodMap.get(methodName); + if (serviceUrl !== undefined) { + return serviceUrl; + } + const computedServiceUrl = computeServiceUrl(host, methodName); + if (methodMap.size < InternalChannel.MAX_METHODS_PER_HOST) { + methodMap.set(methodName, computedServiceUrl); + } + return computedServiceUrl; + } + return computeServiceUrl(host, methodName); + } + createLoadBalancingCall( callConfig: CallConfig, method: string, @@ -820,6 +861,7 @@ export class InternalChannel { this.subchannelPool.unrefUnusedSubchannels(); this.configSelector?.unref(); this.configSelector = null; + this.serviceUrlCache.clear(); } getTarget() { @@ -830,7 +872,7 @@ export class InternalChannel { const connectivityState = this.connectivityState; if (tryToConnect) { this.resolvingLoadBalancer.exitIdle(); - this.lastActivityTimestamp = new Date(); + this.lastActivityTimestamp = Date.now(); this.maybeStartIdleTimer(); } return connectivityState; diff --git a/packages/grpc-js/src/load-balancing-call.ts b/packages/grpc-js/src/load-balancing-call.ts index a3a7495da..68aa3dfba 100644 --- a/packages/grpc-js/src/load-balancing-call.ts +++ b/packages/grpc-js/src/load-balancing-call.ts @@ -15,7 +15,7 @@ * */ -import { CallCredentials } from './call-credentials'; +import { CallCredentials, isEmptyCallCredentials } from './call-credentials'; import { Call, DeadlineInfoProvider, @@ -31,7 +31,6 @@ import { InternalChannel } from './internal-channel'; import { Metadata } from './metadata'; import { OnCallEnded, PickResultType } from './picker'; import { CallConfig } from './resolver'; -import { splitHostPort } from './uri-parser'; import * as logging from './logging'; import { restrictControlPlaneStatusCode } from './control-plane-status'; import * as http2 from 'http2'; @@ -39,6 +38,8 @@ import { AuthContext } from './auth-context'; import { SubchannelInterface } from './subchannel-interface'; const TRACER_NAME = 'load_balancing_call'; +const RESOLVED_EMPTY_METADATA: Promise = + Promise.resolve(undefined); export type RpcProgress = 'NOT_STARTED' | 'DROP' | 'REFUSED' | 'PROCESSED'; @@ -62,8 +63,8 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { private metadata: Metadata | null = null; private listener: InterceptingListener | null = null; private onCallEnded: OnCallEnded | null = null; - private startTime: Date; - private childStartTime: Date | null = null; + private startTime: number; + private childStartTime: number | null = null; constructor( private readonly channel: InternalChannel, private readonly callConfig: CallConfig, @@ -73,23 +74,12 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { private readonly deadline: Deadline, private readonly callNumber: number ) { - const splitPath: string[] = this.methodName.split('/'); - let serviceName = ''; - /* The standard path format is "/{serviceName}/{methodName}", so if we split - * by '/', the first item should be empty and the second should be the - * service name */ - if (splitPath.length >= 2) { - serviceName = splitPath[1]; - } - const hostname = splitHostPort(this.host)?.host ?? 'localhost'; - /* Currently, call credentials are only allowed on HTTPS connections, so we - * can assume that the scheme is "https" */ - this.serviceUrl = `https://${hostname}/${serviceName}`; - this.startTime = new Date(); + this.serviceUrl = this.channel.getServiceUrl(this.host, this.methodName); + this.startTime = Date.now(); } getDeadlineInfo(): string[] { const deadlineInfo: string[] = []; - if (this.childStartTime) { + if (this.childStartTime !== null) { if (this.childStartTime > this.startTime) { if (this.metadata?.getOptions().waitForReady) { deadlineInfo.push('wait_for_ready'); @@ -142,7 +132,7 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { ' details="' + status.details + '" start time=' + - this.startTime.toISOString() + new Date(this.startTime).toISOString() ); } const finalStatus = { ...status, progress }; @@ -152,6 +142,17 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { } } + private generateCallCredentialsMetadata( + callCredentials: CallCredentials + ): Promise { + return isEmptyCallCredentials(callCredentials) + ? RESOLVED_EMPTY_METADATA + : callCredentials.generateMetadata({ + method_name: this.methodName, + service_url: this.serviceUrl, + }); + } + doPick() { if (this.ended) { return; @@ -180,9 +181,9 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { } switch (pickResult.pickResultType) { case PickResultType.COMPLETE: - const combinedCallCredentials = this.credentials.compose(pickResult.subchannel!.getCallCredentials()); - combinedCallCredentials - .generateMetadata({ method_name: this.methodName, service_url: this.serviceUrl }) + this.generateCallCredentialsMetadata( + this.credentials.compose(pickResult.subchannel!.getCallCredentials()) + ) .then( credsMetadata => { /* If this call was cancelled (e.g. by the deadline) before @@ -194,7 +195,9 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { ); return; } - finalMetadata.merge(credsMetadata); + if (credsMetadata) { + finalMetadata.merge(credsMetadata); + } if (finalMetadata.get('authorization').length > 1) { this.outputStatus( { @@ -261,7 +264,7 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { }, this.callNumber ); - this.childStartTime = new Date(); + this.childStartTime = Date.now(); } catch (error) { if (this.traceEnabled) { this.trace( diff --git a/packages/grpc-js/src/metadata.ts b/packages/grpc-js/src/metadata.ts index 7ae68ba32..9d00ffbc7 100644 --- a/packages/grpc-js/src/metadata.ts +++ b/packages/grpc-js/src/metadata.ts @@ -89,7 +89,7 @@ export interface MetadataOptions { export class Metadata { protected internalRepr: MetadataObject = new Map(); private options: MetadataOptions; - private opaqueData: Map = new Map(); + private opaqueData?: Map; constructor(options: MetadataOptions = {}) { this.options = options; @@ -219,16 +219,17 @@ export class Metadata { * Creates an OutgoingHttpHeaders object that can be used with the http2 API. */ toHttp2Headers(): http2.OutgoingHttpHeaders { - // NOTE: Node <8.9 formats http2 headers incorrectly. const result: http2.OutgoingHttpHeaders = {}; for (const [key, values] of this.internalRepr) { if (key.startsWith(':')) { continue; } - // We assume that the user's interaction with this object is limited to - // through its public API (i.e. keys and values are already validated). - result[key] = values.map(bufToString); + if (values.length === 1) { + result[key] = bufToString(values[0]); + } else { + result[key] = values.map(bufToString); + } } return result; @@ -255,6 +256,9 @@ export class Metadata { * @param value */ setOpaque(key: string, value: unknown) { + if (!this.opaqueData) { + this.opaqueData = new Map(); + } this.opaqueData.set(key, value); } @@ -264,7 +268,7 @@ export class Metadata { * @returns */ getOpaque(key: string) { - return this.opaqueData.get(key); + return this.opaqueData?.get(key); } /** diff --git a/packages/grpc-js/src/resolving-call.ts b/packages/grpc-js/src/resolving-call.ts index d7ec6ddcd..13aae2191 100644 --- a/packages/grpc-js/src/resolving-call.ts +++ b/packages/grpc-js/src/resolving-call.ts @@ -56,12 +56,12 @@ export class ResolvingCall implements Call { private deadline: Deadline; private host: string; private statusWatchers: ((status: StatusObject) => void)[] = []; - private deadlineTimer: NodeJS.Timeout = setTimeout(() => {}, 0); + private deadlineTimer: NodeJS.Timeout | null = null; private filterStack: FilterStack | null = null; - private deadlineStartTime: Date | null = null; - private configReceivedTime: Date | null = null; - private childStartTime: Date | null = null; + private deadlineStartTime: number | null = null; + private configReceivedTime: number | null = null; + private childStartTime: number | null = null; /** * Credentials configured for this specific call. Does not include @@ -117,8 +117,11 @@ export class ResolvingCall implements Call { } private runDeadlineTimer() { - clearTimeout(this.deadlineTimer); - this.deadlineStartTime = new Date(); + if (this.deadlineTimer) { + clearTimeout(this.deadlineTimer); + this.deadlineTimer = null; + } + this.deadlineStartTime = Date.now(); if (this.traceEnabled) { this.trace('Deadline: ' + deadlineToString(this.deadline)); } @@ -128,20 +131,36 @@ export class ResolvingCall implements Call { this.trace('Deadline will be reached in ' + timeout + 'ms'); } const handleDeadline = () => { - if (!this.deadlineStartTime) { + this.deadlineTimer = null; + if (this.deadlineStartTime === null) { this.cancelWithStatus(Status.DEADLINE_EXCEEDED, 'Deadline exceeded'); return; } const deadlineInfo: string[] = []; - const deadlineEndTime = new Date(); - deadlineInfo.push(`Deadline exceeded after ${formatDateDifference(this.deadlineStartTime, deadlineEndTime)}`); - if (this.configReceivedTime) { + const deadlineEndTime = Date.now(); + deadlineInfo.push( + `Deadline exceeded after ${formatDateDifference( + this.deadlineStartTime, + deadlineEndTime + )}` + ); + if (this.configReceivedTime !== null) { if (this.configReceivedTime > this.deadlineStartTime) { - deadlineInfo.push(`name resolution: ${formatDateDifference(this.deadlineStartTime, this.configReceivedTime)}`); + deadlineInfo.push( + `name resolution: ${formatDateDifference( + this.deadlineStartTime, + this.configReceivedTime + )}` + ); } - if (this.childStartTime) { + if (this.childStartTime !== null) { if (this.childStartTime > this.configReceivedTime) { - deadlineInfo.push(`metadata filters: ${formatDateDifference(this.configReceivedTime, this.childStartTime)}`); + deadlineInfo.push( + `metadata filters: ${formatDateDifference( + this.configReceivedTime, + this.childStartTime + )}` + ); } } else { deadlineInfo.push('waiting for metadata filters'); @@ -168,7 +187,10 @@ export class ResolvingCall implements Call { if (!this.filterStack) { this.filterStack = this.filterStackFactory.createFilter(); } - clearTimeout(this.deadlineTimer); + if (this.deadlineTimer) { + clearTimeout(this.deadlineTimer); + this.deadlineTimer = null; + } const filteredStatus = this.filterStack.receiveTrailers(status); if (this.traceEnabled) { this.trace( @@ -228,7 +250,7 @@ export class ResolvingCall implements Call { return; } // configResult.type === 'SUCCESS' - this.configReceivedTime = new Date(); + this.configReceivedTime = Date.now(); const config = configResult.config; if (config.status !== Status.OK) { const { code, details } = restrictControlPlaneStatusCode( @@ -244,14 +266,11 @@ export class ResolvingCall implements Call { } if (config.methodConfig.timeout) { - const configDeadline = new Date(); - configDeadline.setSeconds( - configDeadline.getSeconds() + config.methodConfig.timeout.seconds - ); - configDeadline.setMilliseconds( - configDeadline.getMilliseconds() + + const timeoutMs = Math.floor( + config.methodConfig.timeout.seconds * 1000 + config.methodConfig.timeout.nanos / 1_000_000 ); + const configDeadline = Date.now() + timeoutMs; this.deadline = minDeadline(this.deadline, configDeadline); this.runDeadlineTimer(); } @@ -271,7 +290,7 @@ export class ResolvingCall implements Call { if (this.traceEnabled) { this.trace('Created child [' + this.child.getCallNumber() + ']'); } - this.childStartTime = new Date(); + this.childStartTime = Date.now(); this.child.start(filteredMetadata, { onReceiveMetadata: metadata => { this.trace('Received metadata'); diff --git a/packages/grpc-js/src/resolving-load-balancer.ts b/packages/grpc-js/src/resolving-load-balancer.ts index c117e9455..adb45a2fe 100644 --- a/packages/grpc-js/src/resolving-load-balancer.ts +++ b/packages/grpc-js/src/resolving-load-balancer.ts @@ -27,7 +27,7 @@ import { validateServiceConfig, } from './service-config'; import { ConnectivityState } from './connectivity-state'; -import { CHANNEL_ARGS_CONFIG_SELECTOR_KEY, ConfigSelector, createResolver, Resolver } from './resolver'; +import { CallConfig, CHANNEL_ARGS_CONFIG_SELECTOR_KEY, ConfigSelector, createResolver, Resolver } from './resolver'; import { Picker, UnavailablePicker, QueuePicker } from './picker'; import { BackoffOptions, BackoffTimeout } from './backoff-timeout'; import { Status } from './constants'; @@ -99,49 +99,58 @@ function findMatchingConfig( return null; } -function getDefaultConfigSelector( +const MAX_CACHED_METHOD_CONFIGS = 100; + +export function getDefaultConfigSelector( serviceConfig: ServiceConfig | null ): ConfigSelector { + const methodConfigCache = new Map(); return { - invoke( + invoke( methodName: string, - metadata: Metadata - ) { - const splitName = methodName.split('/').filter(x => x.length > 0); - const service = splitName[0] ?? ''; - const method = splitName[1] ?? ''; - if (serviceConfig && serviceConfig.methodConfig) { - /* Check for the following in order, and return the first method - * config that matches: - * 1. A name that exactly matches the service and method - * 2. A name with no method set that matches the service - * 3. An empty name - */ - for (const matchLevel of NAME_MATCH_LEVEL_ORDER) { - const matchingConfig = findMatchingConfig( - service, - method, - serviceConfig.methodConfig, - matchLevel - ); - if (matchingConfig) { - return { - methodConfig: matchingConfig, - pickInformation: {}, - status: Status.OK, - dynamicFilterFactories: [], - }; + metadata: Metadata, + channelId: number + ): CallConfig { + let matchingConfig = methodConfigCache.get(methodName); + if (matchingConfig === undefined) { + const splitName = methodName.split('/').filter(x => x.length > 0); + const service = splitName[0] ?? ''; + const method = splitName[1] ?? ''; + matchingConfig = { name: [] }; + if (serviceConfig && serviceConfig.methodConfig) { + /* Check for the following in order, and return the first method + * config that matches: + * 1. A name that exactly matches the service and method + * 2. A name with no method set that matches the service + * 3. An empty name + */ + for (const matchLevel of NAME_MATCH_LEVEL_ORDER) { + const foundConfig = findMatchingConfig( + service, + method, + serviceConfig.methodConfig, + matchLevel + ); + if (foundConfig) { + matchingConfig = foundConfig; + break; + } } } + if (methodConfigCache.size < MAX_CACHED_METHOD_CONFIGS) { + methodConfigCache.set(methodName, matchingConfig); + } } return { - methodConfig: { name: [] }, + methodConfig: matchingConfig, pickInformation: {}, status: Status.OK, dynamicFilterFactories: [], }; }, - unref() {} + unref() { + methodConfigCache.clear(); + } }; } diff --git a/packages/grpc-js/src/retrying-call.ts b/packages/grpc-js/src/retrying-call.ts index 9023b5241..d84779ab2 100644 --- a/packages/grpc-js/src/retrying-call.ts +++ b/packages/grpc-js/src/retrying-call.ts @@ -124,7 +124,7 @@ interface UnderlyingCall { state: UnderlyingCallState; call: LoadBalancingCall; nextMessageToSend: number; - startTime: Date; + startTime: number; } /** @@ -210,7 +210,7 @@ export class RetryingCall implements Call, DeadlineInfoProvider { private committedCallIndex: number | null = null; private initialRetryBackoffSec = 0; private nextRetryBackoffSec = 0; - private startTime: Date; + private startTime: number; private maxAttempts: number; constructor( private readonly channel: InternalChannel, @@ -249,7 +249,7 @@ export class RetryingCall implements Call, DeadlineInfoProvider { this.state = 'TRANSPARENT_ONLY'; this.maxAttempts = 1; } - this.startTime = new Date(); + this.startTime = Date.now(); } getDeadlineInfo(): string[] { if (this.underlyingCalls.length === 0) { @@ -299,7 +299,7 @@ export class RetryingCall implements Call, DeadlineInfoProvider { ' details="' + statusObject.details + '" start time=' + - this.startTime.toISOString() + new Date(this.startTime).toISOString() ); } this.bufferTracker.freeAll(this.callNumber); @@ -715,11 +715,12 @@ export class RetryingCall implements Call, DeadlineInfoProvider { state: 'ACTIVE', call: child, nextMessageToSend: 0, - startTime: new Date(), + startTime: Date.now(), }); const previousAttempts = this.attempts - 1; - const initialMetadata = this.initialMetadata!.clone(); + let initialMetadata = this.initialMetadata!; if (previousAttempts > 0) { + initialMetadata = initialMetadata.clone(); initialMetadata.set( PREVIONS_RPC_ATTEMPTS_METADATA_KEY, `${previousAttempts}` diff --git a/packages/grpc-js/src/server-interceptors.ts b/packages/grpc-js/src/server-interceptors.ts index afadc9d90..56e597827 100644 --- a/packages/grpc-js/src/server-interceptors.ts +++ b/packages/grpc-js/src/server-interceptors.ts @@ -654,8 +654,7 @@ export class BaseServerInterceptingCall const timeout = (+match[1] * deadlineUnitsToMs[match[2]]) | 0; - const now = new Date(); - this.deadline = now.setMilliseconds(now.getMilliseconds() + timeout); + this.deadline = Date.now() + timeout; this.deadlineTimer = setTimeout(() => { const status: PartialStatusObject = { code: Status.DEADLINE_EXCEEDED, diff --git a/packages/grpc-js/src/server.ts b/packages/grpc-js/src/server.ts index 82aa9fef8..c44ac5161 100644 --- a/packages/grpc-js/src/server.ts +++ b/packages/grpc-js/src/server.ts @@ -193,8 +193,8 @@ interface ChannelzSessionInfo { messagesSent: number; messagesReceived: number; keepAlivesSent: number; - lastMessageSentTimestamp: Date | null; - lastMessageReceivedTimestamp: Date | null; + lastMessageSentTimestamp: number | null; + lastMessageReceivedTimestamp: number | null; } /** @@ -1376,13 +1376,13 @@ export class Server { addMessageSent: () => { if (channelzSessionInfo) { channelzSessionInfo.messagesSent += 1; - channelzSessionInfo.lastMessageSentTimestamp = new Date(); + channelzSessionInfo.lastMessageSentTimestamp = Date.now(); } }, addMessageReceived: () => { if (channelzSessionInfo) { channelzSessionInfo.messagesReceived += 1; - channelzSessionInfo.lastMessageReceivedTimestamp = new Date(); + channelzSessionInfo.lastMessageReceivedTimestamp = Date.now(); } }, onCallEnd: status => { diff --git a/packages/grpc-js/src/single-subchannel-channel.ts b/packages/grpc-js/src/single-subchannel-channel.ts index 4e286795c..df6016b0c 100644 --- a/packages/grpc-js/src/single-subchannel-channel.ts +++ b/packages/grpc-js/src/single-subchannel-channel.ts @@ -32,7 +32,7 @@ import { Metadata } from "./metadata"; import { getDefaultAuthority } from "./resolver"; import { Subchannel } from "./subchannel"; import { SubchannelCall } from "./subchannel-call"; -import { GrpcUri, splitHostPort, uriToString } from "./uri-parser"; +import { GrpcUri, computeServiceUrl, uriToString } from "./uri-parser"; class SubchannelCallWrapper implements Call { private childCall: SubchannelCall | null = null; @@ -46,18 +46,7 @@ class SubchannelCallWrapper implements Call { private readFilterPending = false; private writeFilterPending = false; constructor(private subchannel: Subchannel, private method: string, filterStackFactory: FilterStackFactory, private options: CallStreamOptions, private callNumber: number) { - const splitPath: string[] = this.method.split('/'); - let serviceName = ''; - /* The standard path format is "/{serviceName}/{methodName}", so if we split - * by '/', the first item should be empty and the second should be the - * service name */ - if (splitPath.length >= 2) { - serviceName = splitPath[1]; - } - const hostname = splitHostPort(this.options.host)?.host ?? 'localhost'; - /* Currently, call credentials are only allowed on HTTPS connections, so we - * can assume that the scheme is "https" */ - this.serviceUrl = `https://${hostname}/${serviceName}`; + this.serviceUrl = computeServiceUrl(this.options.host, this.method); const timeout = getRelativeTimeout(options.deadline); if (timeout !== Infinity) { if (timeout <= 0) { diff --git a/packages/grpc-js/src/subchannel-call.ts b/packages/grpc-js/src/subchannel-call.ts index a9cbe08d8..0b738ab4c 100644 --- a/packages/grpc-js/src/subchannel-call.ts +++ b/packages/grpc-js/src/subchannel-call.ts @@ -500,20 +500,22 @@ export class Http2SubchannelCall implements SubchannelCall { } catch (e) { metadata = new Metadata(); } - const metadataMap = metadata.getMap(); + const statusValues = metadata.get('grpc-status'); let status: StatusObject; - if (typeof metadataMap['grpc-status'] === 'string') { - const receivedStatus: Status = Number(metadataMap['grpc-status']); + if (statusValues.length > 0 && typeof statusValues[0] === 'string') { + const receivedStatus: Status = Number(statusValues[0]); if (this.traceEnabled) { this.trace('received status code ' + receivedStatus + ' from server'); } metadata.remove('grpc-status'); let details = ''; - if (typeof metadataMap['grpc-message'] === 'string') { + const messageValues = metadata.get('grpc-message'); + if (messageValues.length > 0 && typeof messageValues[0] === 'string') { + const messageString = messageValues[0]; try { - details = decodeURI(metadataMap['grpc-message']); - } catch (e) { - details = metadataMap['grpc-message']; + details = decodeURI(messageString); + } catch (error) { + details = messageString; } metadata.remove('grpc-message'); if (this.traceEnabled) { @@ -525,7 +527,7 @@ export class Http2SubchannelCall implements SubchannelCall { status = { code: receivedStatus, details: details, - metadata: metadata + metadata: metadata, }; } else if (this.httpStatusCode) { status = mapHttpStatusCode(this.httpStatusCode); @@ -534,7 +536,7 @@ export class Http2SubchannelCall implements SubchannelCall { status = { code: Status.UNKNOWN, details: 'No status information received', - metadata: metadata + metadata: metadata, }; } // This is a no-op if the call was already ended when handling headers. diff --git a/packages/grpc-js/src/transport.ts b/packages/grpc-js/src/transport.ts index bc19ca2a6..154e3ffd5 100644 --- a/packages/grpc-js/src/transport.ts +++ b/packages/grpc-js/src/transport.ts @@ -141,8 +141,8 @@ class Http2Transport implements Transport { private keepalivesSent = 0; private messagesSent = 0; private messagesReceived = 0; - private lastMessageSentTimestamp: Date | null = null; - private lastMessageReceivedTimestamp: Date | null = null; + private lastMessageSentTimestamp: number | null = null; + private lastMessageReceivedTimestamp: number | null = null; constructor( private session: http2.ClientHttp2Session, @@ -616,12 +616,12 @@ class Http2Transport implements Transport { eventTracker = { addMessageSent: () => { this.messagesSent += 1; - this.lastMessageSentTimestamp = new Date(); + this.lastMessageSentTimestamp = Date.now(); subchannelCallStatsTracker.addMessageSent?.(); }, addMessageReceived: () => { this.messagesReceived += 1; - this.lastMessageReceivedTimestamp = new Date(); + this.lastMessageReceivedTimestamp = Date.now(); subchannelCallStatsTracker.addMessageReceived?.(); }, onCallEnd: status => { diff --git a/packages/grpc-js/src/uri-parser.ts b/packages/grpc-js/src/uri-parser.ts index 2b2efeca0..0ceac640b 100644 --- a/packages/grpc-js/src/uri-parser.ts +++ b/packages/grpc-js/src/uri-parser.ts @@ -125,3 +125,29 @@ export function uriToString(uri: GrpcUri): string { result += uri.path; return result; } + +/** + * Extracts the service name from a gRPC method path (e.g. + * "/package.Service/Method" -> "package.Service"). + */ +export function extractServiceName(methodName: string): string { + const firstSlashIndex = methodName.indexOf('/'); + if (firstSlashIndex === -1) { + return ''; + } + const secondSlashIndex = methodName.indexOf('/', firstSlashIndex + 1); + if (secondSlashIndex === -1) { + return methodName.substring(firstSlashIndex + 1); + } + return methodName.substring(firstSlashIndex + 1, secondSlashIndex); +} + +/** + * Computes the service URL used by call credentials (e.g. + * "https://hostname/serviceName"). + */ +export function computeServiceUrl(host: string, methodName: string): string { + const serviceName = extractServiceName(methodName); + const hostname = splitHostPort(host)?.host ?? 'localhost'; + return `https://${hostname}/${serviceName}`; +} diff --git a/packages/grpc-js/test/test-call-credentials.ts b/packages/grpc-js/test/test-call-credentials.ts index 0a6827274..145d52a1f 100644 --- a/packages/grpc-js/test/test-call-credentials.ts +++ b/packages/grpc-js/test/test-call-credentials.ts @@ -20,8 +20,14 @@ import * as assert from 'assert'; import { CallCredentials, CallMetadataGenerator, + isEmptyCallCredentials, } from '../src/call-credentials'; +import { ConnectivityState } from '../src/connectivity-state'; +import { Status } from '../src/constants'; +import { LoadBalancingCall } from '../src/load-balancing-call'; import { Metadata } from '../src/metadata'; +import { PickResultType } from '../src/picker'; +import { computeServiceUrl } from '../src/uri-parser'; // Metadata generators @@ -52,6 +58,37 @@ describe('CallCredentials', () => { }); }); + describe('isEmptyCallCredentials', () => { + it('should return true for empty call credentials', () => { + const emptyCredentials = CallCredentials.createEmpty(); + assert.strictEqual(isEmptyCallCredentials(emptyCredentials), true); + assert.strictEqual( + isEmptyCallCredentials(emptyCredentials.compose(emptyCredentials)), + true + ); + }); + + it('should return false for non-empty call credentials', () => { + const callCredentials1 = CallCredentials.createFromMetadataGenerator( + generateFromServiceURL + ); + const callCredentials2 = CallCredentials.createFromMetadataGenerator( + generateFromServiceURL + ); + assert.strictEqual(isEmptyCallCredentials(callCredentials1), false); + assert.strictEqual( + isEmptyCallCredentials( + callCredentials1.compose(CallCredentials.createEmpty()) + ), + false + ); + assert.strictEqual( + isEmptyCallCredentials(callCredentials1.compose(callCredentials2)), + false + ); + }); + }); + describe('compose', () => { it('should accept a CallCredentials object and return a new object', () => { const callCredentials1 = CallCredentials.createFromMetadataGenerator( @@ -65,6 +102,44 @@ describe('CallCredentials', () => { assert.notStrictEqual(combinedCredentials, callCredentials2); }); + it('should return the same object when composed with empty credentials', async () => { + const callCredentials1 = CallCredentials.createFromMetadataGenerator( + generateFromServiceURL + ); + const callCredentials2 = CallCredentials.createFromMetadataGenerator( + generateFromServiceURL + ); + const emptyCredentials = CallCredentials.createEmpty(); + assert.strictEqual( + emptyCredentials.compose(emptyCredentials), + emptyCredentials + ); + assert.strictEqual( + callCredentials1.compose(emptyCredentials), + callCredentials1 + ); + assert.strictEqual( + emptyCredentials.compose(callCredentials1), + callCredentials1 + ); + const combinedCredentials = callCredentials1.compose(callCredentials2); + assert.strictEqual( + combinedCredentials.compose(emptyCredentials), + combinedCredentials + ); + assert.strictEqual( + emptyCredentials.compose(combinedCredentials), + combinedCredentials + ); + const metadata = await callCredentials1 + .compose(emptyCredentials) + .generateMetadata({ + method_name: 'bar', + service_url: 'foo', + }); + assert.deepStrictEqual(metadata.get('service_url'), ['foo']); + }); + it('should be chainable', () => { const callCredentials1 = CallCredentials.createFromMetadataGenerator( generateFromServiceURL @@ -149,4 +224,111 @@ describe('CallCredentials', () => { ); }); }); + + describe('LoadBalancingCall integration', () => { + function runLoadBalancingCallPick( + callCredentials: CallCredentials, + subchannelCredentials: CallCredentials, + initialMetadata: Metadata + ): Promise { + return new Promise((resolve, reject) => { + const mockSubchannel: any = { + getCallCredentials: () => subchannelCredentials, + getConnectivityState: () => ConnectivityState.READY, + getRealSubchannel: () => mockSubchannel, + getChannelzRef: () => ({ id: 1 }), + getAddress: () => 'localhost:12345', + createCall: (metadata: Metadata) => { + resolve(metadata); + return { + getCallNumber: () => 1, + startRead: () => {}, + sendMessageWithContext: () => {}, + halfClose: () => {}, + cancelWithStatus: () => {}, + }; + }, + }; + const mockChannel: any = { + getServiceUrl: (host: string, methodName: string) => + computeServiceUrl(host, methodName), + doPick: () => ({ + pickResultType: PickResultType.COMPLETE, + subchannel: mockSubchannel, + status: null, + onCallStarted: null, + onCallEnded: null, + }), + }; + const callConfig: any = { + methodConfig: { name: [] }, + pickInformation: {}, + status: Status.OK, + dynamicFilterFactories: [], + }; + const call = new LoadBalancingCall( + mockChannel, + callConfig, + '/service/method', + 'localhost:12345', + callCredentials, + Infinity, + 1 + ); + call.start(initialMetadata, { + onReceiveMetadata: () => {}, + onReceiveMessage: () => {}, + onReceiveStatus: status => { + reject(new Error(`Unexpected status: ${status.details}`)); + }, + }); + }); + } + + it('should bypass generateMetadata when credentials are empty', async () => { + const emptyPrototype = Object.getPrototypeOf( + CallCredentials.createEmpty() + ); + const originalGenerateMetadata = emptyPrototype.generateMetadata; + let generateMetadataCalls = 0; + emptyPrototype.generateMetadata = function ( + ...args: Parameters + ) { + generateMetadataCalls += 1; + return originalGenerateMetadata.apply(this, args); + }; + try { + const initialMetadata = new Metadata(); + initialMetadata.set('custom-key', 'custom-value'); + const finalMetadata = await runLoadBalancingCallPick( + CallCredentials.createEmpty(), + CallCredentials.createEmpty(), + initialMetadata + ); + assert.strictEqual(generateMetadataCalls, 0); + assert.deepStrictEqual(finalMetadata.get('custom-key'), [ + 'custom-value', + ]); + } finally { + emptyPrototype.generateMetadata = originalGenerateMetadata; + } + }); + + it('should call generateMetadata and merge metadata when credentials are non-empty', async () => { + const initialMetadata = new Metadata(); + initialMetadata.set('custom-key', 'custom-value'); + const callCredentials = CallCredentials.createFromMetadataGenerator( + generateFromServiceURL + ); + const finalMetadata = await runLoadBalancingCallPick( + callCredentials, + CallCredentials.createEmpty(), + initialMetadata + ); + assert.deepStrictEqual(finalMetadata.get('custom-key'), ['custom-value']); + assert.deepStrictEqual(finalMetadata.get('service_url'), [ + 'https://localhost/service', + ]); + }); + }); }); diff --git a/packages/grpc-js/test/test-channel-credentials.ts b/packages/grpc-js/test/test-channel-credentials.ts index 00ab08778..c439821d7 100644 --- a/packages/grpc-js/test/test-channel-credentials.ts +++ b/packages/grpc-js/test/test-channel-credentials.ts @@ -91,6 +91,35 @@ describe('ChannelCredentials Implementation', () => { const composedChannelCreds = channelCreds.compose(callCreds); assert.ok(composedChannelCreds instanceof ChannelCredentials); }); + + it('should preserve unwrapped CallCredentials when creating a connector', () => { + const channelCreds = ChannelCredentials.createSsl(); + const callCreds = CallCredentials.createFromMetadataGenerator( + (options, cb) => cb(null, new grpc.Metadata()) + ); + const composedChannelCreds = channelCreds.compose(callCreds); + const connector = composedChannelCreds._createSecureConnector( + { scheme: 'dns', path: 'localhost' }, + {} + ); + assert.strictEqual(connector.getCallCredentials(), callCreds); + connector.destroy(); + }); + }); + + describe('createInsecure', () => { + it('should return the same default CallCredentials instance from a connector', () => { + const insecureCreds = ChannelCredentials.createInsecure(); + const connector = insecureCreds._createSecureConnector( + { scheme: 'dns', path: 'localhost' }, + {} + ); + assert.strictEqual( + connector.getCallCredentials(), + connector.getCallCredentials() + ); + connector.destroy(); + }); }); }); diff --git a/packages/grpc-js/test/test-channelz.ts b/packages/grpc-js/test/test-channelz.ts index 5c3dc3d7d..f7940297c 100644 --- a/packages/grpc-js/test/test-channelz.ts +++ b/packages/grpc-js/test/test-channelz.ts @@ -18,6 +18,7 @@ import * as assert from 'assert'; import * as protoLoader from '@grpc/proto-loader'; import * as grpc from '../src'; +import { dateToProtoTimestamp } from '../src/channelz'; import { ProtoGrpcType } from '../src/generated/channelz'; import { ChannelzClient } from '../src/generated/grpc/channelz/v1/Channelz'; @@ -58,6 +59,25 @@ const testServiceImpl: grpc.UntypedServiceImplementation = { }, }; +function assertValidTimestamp( + timestamp: + | { seconds?: string | number | null; nanos?: number | null } + | null + | undefined, + expectedTimeMs: number +) { + assert(timestamp, 'Timestamp should be defined'); + const expectedSeconds = Math.floor(expectedTimeMs / 1000); + assert( + Math.abs(+timestamp.seconds! - expectedSeconds) <= 10, + `Timestamp seconds (${timestamp.seconds}) should be close to expected (${expectedSeconds})` + ); + assert( + +timestamp.nanos! >= 0 && +timestamp.nanos! < 1e9, + `Timestamp nanos (${timestamp.nanos}) should be in range [0, 1e9)` + ); +} + describe('Channelz', () => { let channelzServer: grpc.Server; let channelzClient: ChannelzClient; @@ -132,6 +152,10 @@ describe('Channelz', () => { +result.channel.ref.channel_id, testClient.getChannel().getChannelzRef().id ); + assert.strictEqual( + result.channel.data?.last_call_started_timestamp, + null + ); // Test that the channel is in the list of top channels channelzClient.getTopChannels( { @@ -167,6 +191,10 @@ describe('Channelz', () => { +result.server.ref.server_id, testServer.getChannelzRef().id ); + assert.strictEqual( + result.server.data?.last_call_started_timestamp, + null + ); // Test that the server is in the list of servers channelzClient.getServers( { start_server_id: testServer.getChannelzRef().id, max_results: 1 }, @@ -187,6 +215,7 @@ describe('Channelz', () => { }); it('should count successful calls', done => { + const callStartTimeMs = Date.now(); testClient.unary({}, (error: grpc.ServiceError, value: unknown) => { assert.ifError(error); // Channel data tests @@ -201,6 +230,10 @@ describe('Channelz', () => { assert.strictEqual(+channelResult.channel.data.calls_started, 1); assert.strictEqual(+channelResult.channel.data.calls_succeeded, 1); assert.strictEqual(+channelResult.channel.data.calls_failed, 0); + assertValidTimestamp( + channelResult.channel.data.last_call_started_timestamp, + callStartTimeMs + ); assert.strictEqual(channelResult.channel.subchannel_ref.length, 1); channelzClient.getSubchannel( { @@ -229,6 +262,10 @@ describe('Channelz', () => { +subchannelResult.subchannel.data.calls_failed, 0 ); + assertValidTimestamp( + subchannelResult.subchannel.data.last_call_started_timestamp, + callStartTimeMs + ); assert.strictEqual( subchannelResult.subchannel.socket_ref.length, 1 @@ -268,6 +305,19 @@ describe('Channelz', () => { +socketResult.socket.data.messages_sent, 1 ); + assertValidTimestamp( + socketResult.socket.data + .last_local_stream_created_timestamp, + callStartTimeMs + ); + assertValidTimestamp( + socketResult.socket.data.last_message_sent_timestamp, + callStartTimeMs + ); + assertValidTimestamp( + socketResult.socket.data.last_message_received_timestamp, + callStartTimeMs + ); // Server data tests channelzClient.getServer( { server_id: testServer.getChannelzRef().id }, @@ -293,6 +343,10 @@ describe('Channelz', () => { +serverResult.server.data.calls_failed, 0 ); + assertValidTimestamp( + serverResult.server.data.last_call_started_timestamp, + callStartTimeMs + ); channelzClient.getServerSockets( { server_id: testServer.getChannelzRef().id }, (error, socketsResult) => { @@ -338,6 +392,21 @@ describe('Channelz', () => { +serverSocketResult.socket.data.messages_sent, 1 ); + assertValidTimestamp( + serverSocketResult.socket.data + .last_remote_stream_created_timestamp, + callStartTimeMs + ); + assertValidTimestamp( + serverSocketResult.socket.data + .last_message_sent_timestamp, + callStartTimeMs + ); + assertValidTimestamp( + serverSocketResult.socket.data + .last_message_received_timestamp, + callStartTimeMs + ); done(); } ); @@ -572,3 +641,37 @@ describe('Disabling channelz', () => { ); }); }); + +describe('dateToProtoTimestamp', () => { + it('returns null for null or undefined', () => { + assert.strictEqual(dateToProtoTimestamp(null), null); + assert.strictEqual(dateToProtoTimestamp(undefined), null); + }); + + it('correctly serializes epoch 0 numbers and Dates', () => { + assert.deepStrictEqual(dateToProtoTimestamp(0), { seconds: 0, nanos: 0 }); + assert.deepStrictEqual(dateToProtoTimestamp(new Date(0)), { + seconds: 0, + nanos: 0, + }); + }); + + it('correctly serializes millisecond numbers and Date objects', () => { + const timestampMs = 1_500_250; + const expected = { seconds: 1500, nanos: 250_000_000 }; + assert.deepStrictEqual(dateToProtoTimestamp(timestampMs), expected); + assert.deepStrictEqual( + dateToProtoTimestamp(new Date(timestampMs)), + expected + ); + }); + + it('handles post-2038 timestamps without 32-bit integer overflow', () => { + // 2040-01-01T00:00:00.123Z => seconds: 2208988800 (exceeds 2^31 - 1) + const post2038Ms = 2_208_988_800_123; + const result = dateToProtoTimestamp(post2038Ms); + assert(result !== null); + assert.strictEqual(result.seconds, 2_208_988_800); + assert.strictEqual(result.nanos, 123_000_000); + }); +}); diff --git a/packages/grpc-js/test/test-client.ts b/packages/grpc-js/test/test-client.ts index e19c87bb6..48fc23610 100644 --- a/packages/grpc-js/test/test-client.ts +++ b/packages/grpc-js/test/test-client.ts @@ -20,8 +20,8 @@ import { EventEmitter } from 'events'; import * as http2 from 'http2'; import * as grpc from '../src'; -import { Server, ServerCredentials } from '../src'; -import { Client } from '../src'; +import { Client, Metadata, Server, ServerCredentials } from '../src'; +import { InternalChannel } from '../src/internal-channel'; import { ConnectivityState } from '../src/connectivity-state'; import { Http2SubchannelCall } from '../src/subchannel-call'; @@ -240,6 +240,119 @@ describe('Client HTTP/2 stream lifecycle', () => { }); }); +describe('Http2SubchannelCall trailers handling', () => { + let originalGetMap: typeof Metadata.prototype.getMap; + let getMapCallCount = 0; + + beforeEach(() => { + originalGetMap = Metadata.prototype.getMap; + getMapCallCount = 0; + Metadata.prototype.getMap = function (...args) { + getMapCallCount++; + return originalGetMap.apply(this, args); + }; + }); + + afterEach(() => { + Metadata.prototype.getMap = originalGetMap; + }); + + function createMockSubchannelCall( + onReceiveStatus: (status: grpc.StatusObject) => void + ) { + const mockHttp2Stream = Object.assign(new EventEmitter(), { + destroyed: false, + writableEnded: false, + end() {}, + rstCode: 0, + close() {}, + resume() {}, + pause() {}, + }); + const mockTransport = { + getOptions: () => ({}), + getPeerName: () => 'localhost', + getAuthContext: () => null, + }; + const mockTracker = { + addMessageReceived: () => {}, + addMessageSent: () => {}, + onStreamEnd: () => {}, + onCallEnd: () => {}, + }; + const mockListener = { + onReceiveMetadata: () => {}, + onReceiveMessage: () => {}, + onReceiveStatus, + }; + + // eslint-disable-next-line @typescript-eslint/no-explicit-any + new Http2SubchannelCall( + mockHttp2Stream as any, + mockTracker as any, + mockListener as any, + mockTransport as any, + 1 + ); + + return mockHttp2Stream; + } + + it('should parse grpc-status and grpc-message from trailers without calling getMap()', done => { + const mockHttp2Stream = createMockSubchannelCall(status => { + assert.strictEqual(status.code, grpc.status.OK); + assert.strictEqual(status.details, 'All Good'); + assert.strictEqual(getMapCallCount, 0); + assert.strictEqual(status.metadata.get('grpc-status').length, 0); + assert.strictEqual(status.metadata.get('grpc-message').length, 0); + assert.deepStrictEqual(status.metadata.get('custom-trailer'), [ + 'custom-val', + ]); + done(); + }); + + mockHttp2Stream.emit('trailers', { + 'grpc-status': '0', + 'grpc-message': 'All%20Good', + 'custom-trailer': 'custom-val', + }); + mockHttp2Stream.emit('end'); + }); + + it('should handle trailers without grpc-message', done => { + const mockHttp2Stream = createMockSubchannelCall(status => { + assert.strictEqual(status.code, grpc.status.OK); + assert.strictEqual(status.details, ''); + assert.strictEqual(getMapCallCount, 0); + assert.strictEqual(status.metadata.get('grpc-status').length, 0); + assert.strictEqual(status.metadata.get('grpc-message').length, 0); + done(); + }); + + mockHttp2Stream.emit('trailers', { + 'grpc-status': '0', + }); + mockHttp2Stream.emit('end'); + }); + + it('should fall back to raw string when grpc-message has invalid percent-encoding', done => { + const mockHttp2Stream = createMockSubchannelCall(status => { + assert.strictEqual(status.code, grpc.status.INTERNAL); + assert.strictEqual(status.details, 'Invalid%2'); + assert.strictEqual(getMapCallCount, 0); + assert.strictEqual(status.metadata.get('grpc-status').length, 0); + assert.strictEqual(status.metadata.get('grpc-message').length, 0); + done(); + }); + + mockHttp2Stream.emit('trailers', { + 'grpc-status': String(grpc.status.INTERNAL), + 'grpc-message': 'Invalid%2', + }); + mockHttp2Stream.emit('end'); + }); +}); + describe('Client without a server', () => { let client: Client; before(() => { @@ -301,7 +414,7 @@ describe('Client with a nonexistent target domain', () => { client.close(); }); it('should fail multiple calls', function (done) { - this.timeout(5000); + this.timeout(process.platform === 'win32' ? 15000 : 5000); // Regression test for https://github.com/grpc/grpc-node/issues/1411 client.makeUnaryRequest( '/service/method', @@ -341,3 +454,118 @@ describe('Client with a nonexistent target domain', () => { client.close(); }); }); + +describe('InternalChannel serviceUrl cache', () => { + it('returns and caches serviceUrl for a given host and method', () => { + const channel = new InternalChannel( + 'localhost:1234', + grpc.credentials.createInsecure(), + {} + ); + const firstServiceUrl = channel.getServiceUrl( + 'spanner.googleapis.com:443', + '/google.spanner.v1.Spanner/ExecuteSql' + ); + assert.strictEqual( + firstServiceUrl, + 'https://spanner.googleapis.com/google.spanner.v1.Spanner' + ); + const cache = ( + channel as unknown as { + serviceUrlCache: Map>; + } + ).serviceUrlCache; + const cachedServiceUrl = cache + .get('spanner.googleapis.com:443') + ?.get('/google.spanner.v1.Spanner/ExecuteSql'); + assert.strictEqual(firstServiceUrl, cachedServiceUrl); + + const secondServiceUrl = channel.getServiceUrl( + 'spanner.googleapis.com:443', + '/google.spanner.v1.Spanner/ExecuteSql' + ); + assert.strictEqual(secondServiceUrl, cachedServiceUrl); + channel.close(); + }); + + it('bounds the number of cached hosts', () => { + const channel = new InternalChannel( + 'localhost:1234', + grpc.credentials.createInsecure(), + {} + ); + for (let hostIndex = 0; hostIndex < 10; hostIndex++) { + const serviceUrl = channel.getServiceUrl( + `host${hostIndex}.example.com:443`, + '/service/method' + ); + assert.strictEqual( + serviceUrl, + `https://host${hostIndex}.example.com/service` + ); + } + const cache = ( + channel as unknown as { + serviceUrlCache: Map>; + } + ).serviceUrlCache; + assert.strictEqual(cache.size, 5); + assert.strictEqual(cache.has('host0.example.com:443'), true); + assert.strictEqual(cache.has('host4.example.com:443'), true); + assert.strictEqual(cache.has('host5.example.com:443'), false); + assert.strictEqual(cache.has('host9.example.com:443'), false); + + // Calling an uncached host still computes the correct URL + const uncachedHostServiceUrl = channel.getServiceUrl( + 'host5.example.com:443', + '/service/method' + ); + assert.strictEqual( + uncachedHostServiceUrl, + 'https://host5.example.com/service' + ); + assert.strictEqual(cache.size, 5); + channel.close(); + assert.strictEqual(cache.size, 0); + }); + + it('bounds the number of cached methods per host', () => { + const channel = new InternalChannel( + 'localhost:1234', + grpc.credentials.createInsecure(), + {} + ); + for (let methodIndex = 0; methodIndex < 150; methodIndex++) { + const serviceUrl = channel.getServiceUrl( + 'localhost:50051', + `/service${methodIndex}/method` + ); + assert.strictEqual(serviceUrl, `https://localhost/service${methodIndex}`); + } + const cache = ( + channel as unknown as { + serviceUrlCache: Map>; + } + ).serviceUrlCache; + const hostMethodMap = cache.get('localhost:50051'); + assert(hostMethodMap !== undefined); + assert.strictEqual(hostMethodMap.size, 100); + assert.strictEqual(hostMethodMap.has('/service0/method'), true); + assert.strictEqual(hostMethodMap.has('/service99/method'), true); + assert.strictEqual(hostMethodMap.has('/service100/method'), false); + + // Calling an uncached method still computes the correct URL + const uncachedMethodServiceUrl = channel.getServiceUrl( + 'localhost:50051', + '/service100/method' + ); + assert.strictEqual( + uncachedMethodServiceUrl, + 'https://localhost/service100' + ); + assert.strictEqual(hostMethodMap.size, 100); + + channel.close(); + assert.strictEqual(cache.size, 0); + }); +}); diff --git a/packages/grpc-js/test/test-deadline.ts b/packages/grpc-js/test/test-deadline.ts index 24aebd4d7..2b60eec21 100644 --- a/packages/grpc-js/test/test-deadline.ts +++ b/packages/grpc-js/test/test-deadline.ts @@ -18,6 +18,12 @@ import * as assert from 'assert'; import * as grpc from '../src'; +import { + formatDateDifference, + getDeadlineTimeoutString, + getRelativeTimeout, + minDeadline, +} from '../src/deadline'; import { ServiceClient, ServiceClientConstructor } from '../src/make-client'; import { loadProtoFile } from './common'; @@ -38,6 +44,7 @@ describe('Client with configured timeout', () => { let server: grpc.Server; let Client: ServiceClientConstructor; let client: ServiceClient; + let serverPort: number; before(done => { Client = loadProtoFile(__dirname + '/fixtures/test_service.proto') @@ -57,6 +64,7 @@ describe('Client with configured timeout', () => { done(error); return; } + serverPort = port; server.start(); client = new Client( `localhost:${port}`, @@ -77,6 +85,7 @@ describe('Client with configured timeout', () => { client.unary({}, (error: grpc.ServiceError, value: unknown) => { assert(error); assert.strictEqual(error.code, grpc.status.DEADLINE_EXCEEDED); + assert.match(error.details, /Deadline exceeded after \d+\.\d{3}s/); done(); }); }); @@ -90,4 +99,219 @@ describe('Client with configured timeout', () => { done(); }); }); + + it('Should handle sub-millisecond fractional nanos in methodConfig timeout', done => { + const fractionalTimeoutConfig: grpc.ServiceConfig = { + loadBalancingConfig: [], + methodConfig: [ + { + name: [{ service: 'TestService' }], + timeout: { + seconds: 0, + nanos: 50_500_000, + }, + }, + ], + }; + const fractionalClient = new Client( + `localhost:${serverPort}`, + grpc.credentials.createInsecure(), + { 'grpc.service_config': JSON.stringify(fractionalTimeoutConfig) } + ); + fractionalClient.unary({}, (error: grpc.ServiceError) => { + assert(error); + assert.strictEqual(error.code, grpc.status.DEADLINE_EXCEEDED); + assert.match(error.details, /Deadline exceeded after \d+\.\d{3}s/); + fractionalClient.close(); + done(); + }); + }); +}); + +describe('deadline utility functions', () => { + describe('formatDateDifference', () => { + it('formats difference between Date objects', () => { + const startDate = new Date(1000); + const endDate = new Date(2500); + assert.strictEqual(formatDateDifference(startDate, endDate), '1.500s'); + }); + + it('formats difference between number timestamps', () => { + const startTimestamp = 1000; + const endTimestamp = 3250; + assert.strictEqual( + formatDateDifference(startTimestamp, endTimestamp), + '2.250s' + ); + }); + + it('formats difference with mixed Date and number timestamps', () => { + const startDate = new Date(5000); + const endTimestamp = 6125; + assert.strictEqual( + formatDateDifference(startDate, endTimestamp), + '1.125s' + ); + assert.strictEqual( + formatDateDifference(endTimestamp, startDate), + '-1.125s' + ); + }); + }); + + describe('minDeadline', () => { + it('returns the minimum deadline across Date and number values', () => { + const firstDeadline = new Date(5000); + const secondDeadline = 2000; + const thirdDeadline = new Date(8000); + const fourthDeadline = 3000; + assert.strictEqual( + minDeadline( + firstDeadline, + secondDeadline, + thirdDeadline, + fourthDeadline + ), + 2000 + ); + }); + + it('returns Infinity when given an empty list', () => { + assert.strictEqual(minDeadline(), Infinity); + }); + + it('handles 0 timestamp correctly', () => { + assert.strictEqual(minDeadline(0, 1000), 0); + assert.strictEqual(minDeadline(1000, 0), 0); + assert.strictEqual(minDeadline(new Date(0), 1000), 0); + }); + }); + + describe('getRelativeTimeout', () => { + it('returns 0 for deadlines in the past', () => { + const pastDate = new Date(Date.now() - 5000); + const pastTimestamp = Date.now() - 5000; + assert.strictEqual(getRelativeTimeout(pastDate), 0); + assert.strictEqual(getRelativeTimeout(pastTimestamp), 0); + }); + + it('returns positive timeout for deadlines in the future', () => { + const futureTimestamp = Date.now() + 5000; + const futureDate = new Date(futureTimestamp); + const timeoutFromDate = getRelativeTimeout(futureDate); + const timeoutFromNumber = getRelativeTimeout(futureTimestamp); + assert(timeoutFromDate > 0 && timeoutFromDate <= 5000); + assert(timeoutFromNumber > 0 && timeoutFromNumber <= 5000); + }); + + it('returns Infinity for deadlines exceeding MAX_TIMEOUT_TIME', () => { + const distantFuture = Date.now() + 3_000_000_000; + assert.strictEqual(getRelativeTimeout(distantFuture), Infinity); + }); + }); + + describe('getDeadlineTimeoutString', () => { + it('formats timeout string correctly for milliseconds', () => { + const now = Date.now(); + const deadline = now + 500; + const timeoutString = getDeadlineTimeoutString(deadline); + assert( + timeoutString.endsWith('m'), + `Expected ${timeoutString} to end with m` + ); + }); + + it('formats timeout string for deadlines as Date instances', () => { + const now = Date.now(); + const deadlineDate = new Date(now + 2500); + const timeoutString = getDeadlineTimeoutString(deadlineDate); + assert.match( + timeoutString, + /^\d+m$/, + `Expected ${timeoutString} to match format e.g. 2500m` + ); + }); + + it('handles epoch 0 correctly', () => { + assert.strictEqual(formatDateDifference(0, 1500), '1.500s'); + assert.strictEqual(formatDateDifference(1500, 0), '-1.500s'); + assert.strictEqual(getRelativeTimeout(0), 0); + }); + + it('throws when deadline is too far in the future', () => { + const excessivelyFarDeadline = + Date.now() + 100_000_000 * 60 * 60 * 1000 + 1000; + assert.throws( + () => getDeadlineTimeoutString(excessivelyFarDeadline), + /Deadline is too far in the future/ + ); + }); + }); +}); + +describe('Calls without deadlines', () => { + let server: grpc.Server; + let Client: ServiceClientConstructor; + let client: ServiceClient; + let originalSetTimeout: typeof setTimeout; + + before(done => { + Client = loadProtoFile(__dirname + '/fixtures/test_service.proto') + .TestService as ServiceClientConstructor; + server = new grpc.Server(); + server.addService(Client.service, { + unary: ( + call: grpc.ServerUnaryCall, + callback: grpc.sendUnaryData> + ) => { + callback(null, {}); + }, + }); + server.bindAsync( + 'localhost:0', + grpc.ServerCredentials.createInsecure(), + (error, port) => { + if (error) { + done(error); + return; + } + server.start(); + client = new Client( + `localhost:${port}`, + grpc.credentials.createInsecure() + ); + client.waitForReady(Date.now() + 2000, done); + } + ); + }); + + after(done => { + client.close(); + server.tryShutdown(done); + }); + + afterEach(() => { + if (originalSetTimeout) { + global.setTimeout = originalSetTimeout; + } + }); + + it('Should not schedule any timer for calls with infinite deadline', done => { + let timerCount = 0; + originalSetTimeout = global.setTimeout; + global.setTimeout = (( + handler: (...args: any[]) => void, + timeout?: number, + ...args: any[] + ) => { + timerCount++; + return originalSetTimeout(handler as any, timeout, ...args); + }) as any; + + client.unary({}, (error: grpc.ServiceError | null) => { + assert.ifError(error); + assert.strictEqual(timerCount, 0); + done(); + }); + }); }); diff --git a/packages/grpc-js/test/test-idle-timer.ts b/packages/grpc-js/test/test-idle-timer.ts index ed6af2cf7..6429a79e4 100644 --- a/packages/grpc-js/test/test-idle-timer.ts +++ b/packages/grpc-js/test/test-idle-timer.ts @@ -132,7 +132,10 @@ describe('Channel idle timer', () => { describe('Channel idle timer with UDS', () => { let server: TestServer; let client: TestClient | null = null; - before(() => { + before(function () { + if (process.platform === 'win32') { + this.skip(); + } server = new TestServer(false); return server.startUds(); }); @@ -143,7 +146,7 @@ describe('Channel idle timer with UDS', () => { } }); after(() => { - server.shutdown(); + server?.shutdown(); }); it('Should be able to make a request after going idle', function (done) { this.timeout(5000); diff --git a/packages/grpc-js/test/test-metadata.ts b/packages/grpc-js/test/test-metadata.ts index 44182ef39..b40e43803 100644 --- a/packages/grpc-js/test/test-metadata.ts +++ b/packages/grpc-js/test/test-metadata.ts @@ -25,6 +25,11 @@ class TestMetadata extends Metadata { return this.internalRepr; } + getOpaqueData() { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + return (this as any).opaqueData; + } + static fromHttp2Headers(headers: http2.IncomingHttpHeaders): TestMetadata { const result = Metadata.fromHttp2Headers(headers) as TestMetadata; result.getInternalRepresentation = @@ -268,14 +273,16 @@ describe('Metadata', () => { metadata.add('Key2', 'value2'); metadata.add('KEY3', 'value3a'); metadata.add('key3', 'value3b'); + metadata.add('single-bin', Buffer.from('hello')); metadata.add('key-bin', Buffer.from(range(0, 16))); metadata.add('key-bin', Buffer.from(range(16, 32))); metadata.add('key-bin', Buffer.from(range(0, 32))); const headers = metadata.toHttp2Headers(); assert.deepStrictEqual(headers, { - key1: ['value1'], - key2: ['value2'], + key1: 'value1', + key2: 'value2', key3: ['value3a', 'value3b'], + 'single-bin': Buffer.from('hello').toString('base64'), 'key-bin': [ 'AAECAwQFBgcICQoLDA0ODw==', 'EBESExQVFhcYGRobHB0eHw==', @@ -284,6 +291,23 @@ describe('Metadata', () => { }); }); + it('omits pseudo-headers starting with colon', () => { + metadata.getInternalRepresentation().set(':path', ['/service/method']); + metadata.getInternalRepresentation().set(':status', ['200']); + metadata.add('custom-key', 'custom-value'); + assert.deepStrictEqual(metadata.toHttp2Headers(), { + 'custom-key': 'custom-value', + }); + }); + + it('handles empty strings and strings with commas', () => { + metadata.add('empty', ''); + metadata.add('comma-str', 'a, b, c'); + const headers = metadata.toHttp2Headers(); + assert.strictEqual(headers['empty'], ''); + assert.strictEqual(headers['comma-str'], 'a, b, c'); + }); + it('creates an empty header object from empty Metadata', () => { assert.deepStrictEqual(metadata.toHttp2Headers(), {}); }); @@ -296,6 +320,7 @@ describe('Metadata', () => { key2: ['value2'], key3: ['value3a', 'value3b'], key4: ['part1, part2'], + 'single-bin': Buffer.from('hello').toString('base64'), 'key-bin': [ 'AAECAwQFBgcICQoLDA0ODw==', 'EBESExQVFhcYGRobHB0eHw==', @@ -309,6 +334,7 @@ describe('Metadata', () => { ['key2', ['value2']], ['key3', ['value3a', 'value3b']], ['key4', ['part1, part2']], + ['single-bin', [Buffer.from('hello')]], [ 'key-bin', [ @@ -327,4 +353,115 @@ describe('Metadata', () => { assert.deepStrictEqual(internalRepr, new Map()); }); }); + + describe('opaqueData', () => { + it('does not allocate opaqueData map until setOpaque is called', () => { + const testMetadata = new TestMetadata(); + assert.strictEqual(testMetadata.getOpaqueData(), undefined); + assert.strictEqual(testMetadata.getOpaque('key'), undefined); + + testMetadata.setOpaque('key', 'value'); + assert.notStrictEqual(testMetadata.getOpaqueData(), undefined); + assert.strictEqual(testMetadata.getOpaque('key'), 'value'); + assert.strictEqual(testMetadata.getOpaque('missing'), undefined); + }); + + it('handles multiple keys, key overwrite, and various value types without reallocating map', () => { + const testMetadata = new TestMetadata(); + testMetadata.setOpaque('stringKey', 'firstValue'); + const initialMap = testMetadata.getOpaqueData(); + + // Subsequent setOpaque calls exercise the existing map branch + testMetadata.setOpaque('stringKey', 'overwrittenValue'); + testMetadata.setOpaque('numberKey', 42); + testMetadata.setOpaque('objectKey', { metric: 123 }); + testMetadata.setOpaque('undefinedKey', undefined); + + // Map reference remains identical (no reallocation) + assert.strictEqual(testMetadata.getOpaqueData(), initialMap); + assert.strictEqual( + testMetadata.getOpaque('stringKey'), + 'overwrittenValue' + ); + assert.strictEqual(testMetadata.getOpaque('numberKey'), 42); + assert.deepStrictEqual(testMetadata.getOpaque('objectKey'), { + metric: 123, + }); + assert.strictEqual(testMetadata.getOpaque('undefinedKey'), undefined); + }); + + it('does not copy opaqueData to cloned instances', () => { + const originalMetadata = new TestMetadata(); + originalMetadata.setOpaque('key', 'value'); + const clonedMetadata = originalMetadata.clone(); + + assert.strictEqual(clonedMetadata.getOpaque('key'), undefined); + // eslint-disable-next-line @typescript-eslint/no-explicit-any + assert.strictEqual((clonedMetadata as any).opaqueData, undefined); + }); + + it('does not allocate opaqueData on clone if original has no opaqueData', () => { + const originalMetadata = new Metadata(); + const clonedMetadata = originalMetadata.clone(); + + assert.strictEqual(clonedMetadata.getOpaque('key'), undefined); + // eslint-disable-next-line @typescript-eslint/no-explicit-any + assert.strictEqual((clonedMetadata as any).opaqueData, undefined); + }); + }); + + describe('options', () => { + it('defaults to empty options object when none are passed', () => { + const defaultMetadata = new Metadata(); + assert.deepStrictEqual(defaultMetadata.getOptions(), {}); + }); + + it('preserves custom options passed to constructor', () => { + const customMetadata = new Metadata({ waitForReady: true, corked: true }); + assert.deepStrictEqual(customMetadata.getOptions(), { + waitForReady: true, + corked: true, + }); + }); + + it('allows updating options via setOptions', () => { + const customMetadata = new Metadata(); + customMetadata.setOptions({ cacheableRequest: true }); + assert.deepStrictEqual(customMetadata.getOptions(), { + cacheableRequest: true, + }); + }); + }); + + describe('roundtrip toHttp2Headers and fromHttp2Headers', () => { + it('preserves single and multi-value string and binary metadata', () => { + const originalMetadata = new Metadata(); + originalMetadata.add('single-str', 'value1'); + originalMetadata.add('multi-str', 'value2a'); + originalMetadata.add('multi-str', 'value2b'); + originalMetadata.add('single-bin', Buffer.from('hello')); + originalMetadata.add('multi-bin', Buffer.from([1, 2, 3])); + originalMetadata.add('multi-bin', Buffer.from([4, 5, 6])); + + const roundtripMetadata = Metadata.fromHttp2Headers( + originalMetadata.toHttp2Headers() as http2.IncomingHttpHeaders + ); + assert.deepStrictEqual( + roundtripMetadata.getMap(), + originalMetadata.getMap() + ); + assert.deepStrictEqual(roundtripMetadata.get('single-str'), ['value1']); + assert.deepStrictEqual(roundtripMetadata.get('single-bin'), [ + Buffer.from('hello'), + ]); + assert.deepStrictEqual(roundtripMetadata.get('multi-str'), [ + 'value2a', + 'value2b', + ]); + assert.deepStrictEqual(roundtripMetadata.get('multi-bin'), [ + Buffer.from([1, 2, 3]), + Buffer.from([4, 5, 6]), + ]); + }); + }); }); diff --git a/packages/grpc-js/test/test-resolving-load-balancer.ts b/packages/grpc-js/test/test-resolving-load-balancer.ts new file mode 100644 index 000000000..2aa0660a0 --- /dev/null +++ b/packages/grpc-js/test/test-resolving-load-balancer.ts @@ -0,0 +1,241 @@ +/* + * Copyright 2026 gRPC authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +import * as assert from 'assert'; +import { getDefaultConfigSelector } from '../src/resolving-load-balancer'; +import { ServiceConfig } from '../src/service-config'; +import { Metadata } from '../src/metadata'; +import { Status } from '../src/constants'; + +describe('getDefaultConfigSelector', () => { + const dummyMetadata = new Metadata(); + const dummyChannelId = 0; + + const sampleServiceConfigWithoutWildcard: ServiceConfig = { + loadBalancingConfig: [], + methodConfig: [ + { + name: [{ service: 'TestService', method: 'SpecificMethod' }], + timeout: { seconds: 10, nanos: 0 }, + }, + { + name: [{ service: 'TestService' }], + timeout: { seconds: 20, nanos: 0 }, + }, + ], + }; + + const sampleServiceConfigWithWildcard: ServiceConfig = { + loadBalancingConfig: [], + methodConfig: [ + { + name: [{ service: 'TestService', method: 'SpecificMethod' }], + timeout: { seconds: 10, nanos: 0 }, + }, + { + name: [{ service: 'TestService' }], + timeout: { seconds: 20, nanos: 0 }, + }, + { + name: [{}], + timeout: { seconds: 30, nanos: 0 }, + }, + ], + }; + + it('matches exact service and method', () => { + const selector = getDefaultConfigSelector( + sampleServiceConfigWithoutWildcard + ); + const callConfig = selector.invoke( + '/TestService/SpecificMethod', + dummyMetadata, + dummyChannelId + ); + assert.strictEqual(callConfig.status, Status.OK); + assert.strictEqual(callConfig.methodConfig.timeout?.seconds, 10); + }); + + it('matches service-level config when method does not match', () => { + const selector = getDefaultConfigSelector( + sampleServiceConfigWithoutWildcard + ); + const callConfig = selector.invoke( + '/TestService/OtherMethod', + dummyMetadata, + dummyChannelId + ); + assert.strictEqual(callConfig.status, Status.OK); + assert.strictEqual(callConfig.methodConfig.timeout?.seconds, 20); + }); + + it('returns empty name default config when serviceConfig is present but neither method nor service matches', () => { + const selector = getDefaultConfigSelector( + sampleServiceConfigWithoutWildcard + ); + const callConfig = selector.invoke( + '/OtherService/AnyMethod', + dummyMetadata, + dummyChannelId + ); + assert.strictEqual(callConfig.status, Status.OK); + assert.deepStrictEqual(callConfig.methodConfig.name, []); + assert.strictEqual(callConfig.methodConfig.timeout, undefined); + }); + + it('matches empty name default config when service does not match and wildcard is present', () => { + const selector = getDefaultConfigSelector(sampleServiceConfigWithWildcard); + const callConfig = selector.invoke( + '/OtherService/AnyMethod', + dummyMetadata, + dummyChannelId + ); + assert.strictEqual(callConfig.status, Status.OK); + assert.strictEqual(callConfig.methodConfig.timeout?.seconds, 30); + }); + + it('returns default config when serviceConfig is null', () => { + const selector = getDefaultConfigSelector(null); + const callConfig = selector.invoke( + '/TestService/TestMethod', + dummyMetadata, + dummyChannelId + ); + assert.strictEqual(callConfig.status, Status.OK); + assert.deepStrictEqual(callConfig.methodConfig.name, []); + assert.strictEqual(callConfig.methodConfig.timeout, undefined); + }); + + it('memoizes MethodConfig and returns isolated CallConfig containers', () => { + const selector = getDefaultConfigSelector( + sampleServiceConfigWithoutWildcard + ); + const firstCallConfig = selector.invoke( + '/TestService/SpecificMethod', + dummyMetadata, + dummyChannelId + ); + const secondCallConfig = selector.invoke( + '/TestService/SpecificMethod', + dummyMetadata, + dummyChannelId + ); + + // CallConfig containers must be isolated + assert.notStrictEqual(firstCallConfig, secondCallConfig); + assert.notStrictEqual( + firstCallConfig.pickInformation, + secondCallConfig.pickInformation + ); + assert.notStrictEqual( + firstCallConfig.dynamicFilterFactories, + secondCallConfig.dynamicFilterFactories + ); + + // Underlying resolved MethodConfig must be memoized by reference + assert.strictEqual( + firstCallConfig.methodConfig, + secondCallConfig.methodConfig + ); + + // Mutations on one call's pickInformation must not leak to the other + (firstCallConfig.pickInformation as Record)['testKey'] = + 'testValue'; + assert.strictEqual( + (secondCallConfig.pickInformation as Record)['testKey'], + undefined + ); + }); + + it('bounds the cache size to 100 entries', () => { + const selector = getDefaultConfigSelector( + sampleServiceConfigWithoutWildcard + ); + for (let methodIndex = 0; methodIndex < 150; methodIndex++) { + const callConfig = selector.invoke( + `/UnmatchedService/Method${methodIndex}`, + dummyMetadata, + dummyChannelId + ); + assert.strictEqual(callConfig.status, Status.OK); + } + + // Method0 was in the first 100 calls, so its resolved MethodConfig is cached by reference + const firstCachedCallConfig = selector.invoke( + '/UnmatchedService/Method0', + dummyMetadata, + dummyChannelId + ); + const secondCachedCallConfig = selector.invoke( + '/UnmatchedService/Method0', + dummyMetadata, + dummyChannelId + ); + assert.strictEqual( + firstCachedCallConfig.methodConfig, + secondCachedCallConfig.methodConfig + ); + + // Method149 was beyond 100, so it was not cached; each invoke generates a fresh { name: [] } + const firstUncachedCallConfig = selector.invoke( + '/UnmatchedService/Method149', + dummyMetadata, + dummyChannelId + ); + const secondUncachedCallConfig = selector.invoke( + '/UnmatchedService/Method149', + dummyMetadata, + dummyChannelId + ); + assert.strictEqual(firstUncachedCallConfig.status, Status.OK); + assert.deepStrictEqual( + firstUncachedCallConfig.methodConfig, + secondUncachedCallConfig.methodConfig + ); + assert.notStrictEqual( + firstUncachedCallConfig.methodConfig, + secondUncachedCallConfig.methodConfig + ); + }); + + it('clears the cache on unref()', () => { + const selector = getDefaultConfigSelector( + sampleServiceConfigWithoutWildcard + ); + const firstCallConfig = selector.invoke( + '/UnmatchedService/MethodX', + dummyMetadata, + dummyChannelId + ); + selector.unref(); + + // After unref, a subsequent invoke will recompute a fresh MethodConfig object + const secondCallConfig = selector.invoke( + '/UnmatchedService/MethodX', + dummyMetadata, + dummyChannelId + ); + assert.deepStrictEqual( + firstCallConfig.methodConfig, + secondCallConfig.methodConfig + ); + assert.notStrictEqual( + firstCallConfig.methodConfig, + secondCallConfig.methodConfig + ); + }); +}); diff --git a/packages/grpc-js/test/test-retry.ts b/packages/grpc-js/test/test-retry.ts index d7e5a09cd..d72470d4b 100644 --- a/packages/grpc-js/test/test-retry.ts +++ b/packages/grpc-js/test/test-retry.ts @@ -53,6 +53,8 @@ const serviceImpl = { describe('Retries', () => { let server: grpc.Server; let port: number; + const originalClone = grpc.Metadata.prototype.clone; + before(done => { server = new grpc.Server(); server.addService(EchoService.service, serviceImpl); @@ -71,6 +73,10 @@ describe('Retries', () => { ); }); + afterEach(() => { + grpc.Metadata.prototype.clone = originalClone; + }); + after(() => { server.forceShutdown(); }); @@ -343,6 +349,108 @@ describe('Retries', () => { } ); }); + + it('Should not clone metadata in RetryingCall on the initial attempt', done => { + const metadata = new grpc.Metadata(); + metadata.set('custom-header', 'custom-value'); + const clonedInstances: grpc.Metadata[] = []; + grpc.Metadata.prototype.clone = function (this: grpc.Metadata) { + const cloned = originalClone.call(this); + clonedInstances.push(cloned); + return cloned; + }; + client.echo( + { value: 'test value', value2: 3 }, + metadata, + { deadline: Date.now() + 10000 }, + (error: grpc.ServiceError, response: any) => { + grpc.Metadata.prototype.clone = originalClone; + assert.ifError(error); + assert.deepStrictEqual(response, { value: 'test value', value2: 3 }); + // Cloned only in ResolvingCall.start and LoadBalancingCall.doPick + assert.strictEqual(clonedInstances.length, 2); + // Caller metadata remains unmutated + assert.deepStrictEqual(metadata.getMap(), { + 'custom-header': 'custom-value', + }); + // First clone (RetryingCall.initialMetadata / LoadBalancingCall.metadata) + // must not be mutated by LoadBalancingCall.doPick (e.g. grpc-timeout) + const retryingCallInitialMetadata = clonedInstances[0]; + assert.strictEqual( + retryingCallInitialMetadata.get('grpc-timeout').length, + 0 + ); + // Second clone (finalMetadata in LoadBalancingCall.doPick) receives grpc-timeout + const loadBalancingFinalMetadata = clonedInstances[1]; + assert.strictEqual( + loadBalancingFinalMetadata.get('grpc-timeout').length, + 1 + ); + done(); + } + ); + }); + + it('Should clone metadata on retry attempts without mutating RetryingCall initialMetadata', done => { + const metadata = new grpc.Metadata(); + metadata.set('succeed-on-retry-attempt', '2'); + metadata.set('respond-with-status', `${grpc.status.RESOURCE_EXHAUSTED}`); + const callCredentials = grpc.credentials.createFromMetadataGenerator( + (_options, callback) => { + const credentialsMetadata = new grpc.Metadata(); + credentialsMetadata.set('authorization', 'Bearer test-token'); + callback(null, credentialsMetadata); + } + ); + const clonedInstances: grpc.Metadata[] = []; + grpc.Metadata.prototype.clone = function (this: grpc.Metadata) { + const cloned = originalClone.call(this); + clonedInstances.push(cloned); + return cloned; + }; + client.echo( + { value: 'test value', value2: 3 }, + metadata, + { credentials: callCredentials, deadline: Date.now() + 10000 }, + (error: grpc.ServiceError, response: any) => { + grpc.Metadata.prototype.clone = originalClone; + assert.ifError(error); + assert.deepStrictEqual(response, { value: 'test value', value2: 3 }); + // 1 in ResolvingCall.start + 2 in RetryingCall (attempts 2 & 3) + 3 in LoadBalancingCall.doPick + assert.strictEqual(clonedInstances.length, 6); + assert.strictEqual( + metadata.get('grpc-previous-rpc-attempts').length, + 0 + ); + assert.strictEqual(metadata.get('authorization').length, 0); + assert.strictEqual(metadata.get('grpc-timeout').length, 0); + // First clone (RetryingCall.initialMetadata, passed uncloned to LoadBalancingCall on attempt 1) + // must not be polluted by RetryingCall, LoadBalancingCall, or post-start caller mutations + const retryingCallInitialMetadata = clonedInstances[0]; + assert.strictEqual( + retryingCallInitialMetadata.get('grpc-previous-rpc-attempts') + .length, + 0 + ); + assert.strictEqual( + retryingCallInitialMetadata.get('authorization').length, + 0 + ); + assert.strictEqual( + retryingCallInitialMetadata.get('grpc-timeout').length, + 0 + ); + assert.strictEqual( + retryingCallInitialMetadata.get('mutated-after-start').length, + 0 + ); + done(); + } + ); + // Mutate caller metadata while the RPC and its retry attempts are in flight + metadata.set('mutated-after-start', 'true'); + metadata.set('succeed-on-retry-attempt', '99'); + }); }); describe('Client with hedging configured', () => { @@ -466,5 +574,51 @@ describe('Retries', () => { } ); }); + + it('Should clone metadata on hedged attempts without mutating RetryingCall initialMetadata', done => { + const metadata = new grpc.Metadata(); + metadata.set('succeed-on-retry-attempt', '2'); + metadata.set('respond-with-status', `${grpc.status.RESOURCE_EXHAUSTED}`); + const callCredentials = grpc.credentials.createFromMetadataGenerator( + (_options, callback) => { + const credentialsMetadata = new grpc.Metadata(); + credentialsMetadata.set('authorization', 'Bearer test-token'); + callback(null, credentialsMetadata); + } + ); + const clonedInstances: grpc.Metadata[] = []; + grpc.Metadata.prototype.clone = function (this: grpc.Metadata) { + const cloned = originalClone.call(this); + clonedInstances.push(cloned); + return cloned; + }; + client.echo( + { value: 'test value', value2: 3 }, + metadata, + { credentials: callCredentials, deadline: Date.now() + 10000 }, + (error: grpc.ServiceError, response: any) => { + grpc.Metadata.prototype.clone = originalClone; + assert.ifError(error); + assert.deepStrictEqual(response, { value: 'test value', value2: 3 }); + // 1 in ResolvingCall.start + 2 in RetryingCall (hedged attempts 2 & 3) + 3 in LoadBalancingCall.doPick + assert.strictEqual(clonedInstances.length, 6); + const retryingCallInitialMetadata = clonedInstances[0]; + assert.strictEqual( + retryingCallInitialMetadata.get('grpc-previous-rpc-attempts') + .length, + 0 + ); + assert.strictEqual( + retryingCallInitialMetadata.get('authorization').length, + 0 + ); + assert.strictEqual( + retryingCallInitialMetadata.get('grpc-timeout').length, + 0 + ); + done(); + } + ); + }); }); }); diff --git a/packages/grpc-js/test/test-uri-parser.ts b/packages/grpc-js/test/test-uri-parser.ts index 1e20e5e26..18d9e3db1 100644 --- a/packages/grpc-js/test/test-uri-parser.ts +++ b/packages/grpc-js/test/test-uri-parser.ts @@ -144,4 +144,68 @@ describe('URI Parser', function () { }); } }); + + describe('extractServiceName', function () { + const expectationList: { methodName: string; expected: string }[] = [ + { + methodName: '/google.spanner.v1.Spanner/ExecuteSql', + expected: 'google.spanner.v1.Spanner', + }, + { methodName: '/EchoService/Echo', expected: 'EchoService' }, + { methodName: '/serviceOnly', expected: 'serviceOnly' }, + { methodName: '/a/b/c', expected: 'a' }, + { methodName: 'noLeading/second/third', expected: 'second' }, + { methodName: 'noSlash', expected: '' }, + { methodName: '', expected: '' }, + { methodName: '/', expected: '' }, + { methodName: '//', expected: '' }, + ]; + for (const { methodName, expected } of expectationList) { + it(methodName, function () { + assert.strictEqual(uriParser.extractServiceName(methodName), expected); + }); + } + }); + + describe('computeServiceUrl', function () { + const expectationList: { + host: string; + methodName: string; + expected: string; + }[] = [ + { + host: 'spanner.googleapis.com:443', + methodName: '/google.spanner.v1.Spanner/ExecuteSql', + expected: 'https://spanner.googleapis.com/google.spanner.v1.Spanner', + }, + { + host: 'localhost:50051', + methodName: '/EchoService/Echo', + expected: 'https://localhost/EchoService', + }, + { + host: '[::1]:50051', + methodName: '/EchoService/Echo', + expected: 'https://::1/EchoService', + }, + { + host: '127.0.0.1:8080', + methodName: '/EchoService/Echo', + expected: 'https://127.0.0.1/EchoService', + }, + { + host: '[', + methodName: '/EchoService/Echo', + expected: 'https://localhost/EchoService', + }, + ]; + for (const { host, methodName, expected } of expectationList) { + it(`${host} + ${methodName}`, function () { + assert.strictEqual( + uriParser.computeServiceUrl(host, methodName), + expected + ); + }); + } + }); }); diff --git a/packages/grpc-js/test/test-weighted-round-robin.ts b/packages/grpc-js/test/test-weighted-round-robin.ts index 3016cb8e4..3dc7baa9f 100644 --- a/packages/grpc-js/test/test-weighted-round-robin.ts +++ b/packages/grpc-js/test/test-weighted-round-robin.ts @@ -28,20 +28,19 @@ const protoFile = path.join(__dirname, 'fixtures', 'echo_service.proto'); const EchoService = (loadProtoFile(protoFile) as unknown as ProtoGrpcType).EchoService; function makeNCalls(client: EchoServiceClient, count: number): Promise<{[serverId: string]: number}> { - return new Promise((resolve, reject) => { - const result: {[serverId: string]: number} = {}; - function makeOneCall(callsLeft: number) { - if (callsLeft <= 0) { - resolve(result); - } else { + const result: {[serverId: string]: number} = {}; + const promises: Promise[] = []; + for (let index = 0; index < count; index++) { + promises.push( + new Promise((resolve, reject) => { const deadline = new Date(); - deadline.setMilliseconds(deadline.getMilliseconds() + 100); - const call= client.echo({}, {deadline}, (error, value) => { + deadline.setSeconds(deadline.getSeconds() + 2); + const call = client.echo({}, {deadline}, (error, value) => { if (error) { reject(error); return; } - makeOneCall(callsLeft - 1); + resolve(); }); call.on('metadata', metadata => { const serverEntry = metadata.get('server'); @@ -53,9 +52,60 @@ function makeNCalls(client: EchoServiceClient, count: number): Promise<{[serverI result[serverId] += 1; } }); + }) + ); + } + return Promise.all(promises).then(() => result); +} + +function warmupAllServers(client: EchoServiceClient, serverCount = 2): Promise { + const seenServers = new Set(); + const deadline = new Date(); + deadline.setSeconds(deadline.getSeconds() + 5); + + return new Promise((resolve, reject) => { + function sendNextCall() { + if (seenServers.size >= serverCount) { + resolve(); + return; } + const call = client.echo({}, {deadline}, error => { + if (error) { + reject(error); + return; + } + sendNextCall(); + }); + call.on('metadata', metadata => { + const serverEntry = metadata.get('server'); + if (serverEntry.length > 0) { + const serverId = serverEntry[0] as string; + if (serverId) { + seenServers.add(serverId); + } + } + }); } - makeOneCall(count); + sendNextCall(); + }); +} + +function waitForPickerUpdate(client: EchoServiceClient, timeoutMs = 2000): Promise { + const internalChannel = (client.getChannel() as any).internalChannel; + const initialPicker = internalChannel.currentPicker; + const startTime = Date.now(); + + return new Promise((resolve, reject) => { + const checkPicker = () => { + if (internalChannel.currentPicker !== initialPicker) { + resolve(); + } else if (Date.now() - startTime > timeoutMs) { + reject(new Error('Timed out waiting for WeightedRoundRobinPicker to update')); + } else { + setTimeout(checkPicker, 10); + } + }; + checkPicker(); }); } @@ -72,13 +122,28 @@ function createClient(ports: number[], serviceConfig: grpc.ServiceConfig) { return new EchoService(`ipv4:${ports.map(port => `127.0.0.1:${port}`).join(',')}`, grpc.credentials.createInsecure(), {'grpc.service_config': JSON.stringify(serviceConfig)}); } +function waitForClientReady(client: EchoServiceClient): Promise { + return new Promise((resolve, reject) => { + const deadline = new Date(); + deadline.setSeconds(deadline.getSeconds() + 5); + client.waitForReady(deadline, error => { + if (error) { + reject(error); + } else { + resolve(); + } + }); + }); +} + function asyncTimeout(delay: number): Promise { return new Promise(resolve => { setTimeout(resolve, delay); }); } -describe('Weighted round robin LB policy', () => { +describe('Weighted round robin LB policy', function () { + this.timeout(process.platform === 'win32' ? 10000 : 5000); describe('Config parsing', () => { it('Should have default values with an empty object', () => { const config = WeightedRoundRobinLoadBalancingConfig.createFromJson({}); @@ -261,7 +326,8 @@ describe('Weighted round robin LB policy', () => { it('Should evenly balance among endpoints with no weight', async () => { const serviceConfig = createServiceConfig({}); client = createClient([port1, port2], serviceConfig); - await makeNCalls(client, 10); + await waitForClientReady(client); + await warmupAllServers(client); const result = await makeNCalls(client, 30); assert(Math.abs(result['1'] - result['2']) < 3, `server1: ${result['1']}, server2: ${result[2]}`); }); @@ -271,12 +337,13 @@ describe('Weighted round robin LB policy', () => { weight_update_period: '0.1s' }); client = createClient([port1, port2], serviceConfig); + await waitForClientReady(client); server1Metrics.qps = 3; server1Metrics.utilization = 1; server2Metrics.qps = 1; server2Metrics.utilization = 1; - await makeNCalls(client, 10); - await asyncTimeout(200); + await warmupAllServers(client); + await waitForPickerUpdate(client); const result = await makeNCalls(client, 40); assert(Math.abs(result['1'] - 30) < 3, `server1: ${result['1']}, server2: ${result['2']}`); }); @@ -324,14 +391,15 @@ describe('Weighted round robin LB policy', () => { error_utilization_penalty: 1 }); client = createClient([port1, port2], serviceConfig); + await waitForClientReady(client); server1Metrics.qps = 2; server1Metrics.utilization = 1; server1Metrics.eps = 0; server2Metrics.qps = 2; server2Metrics.utilization = 1; server2Metrics.eps = 2; - await makeNCalls(client, 10); - await asyncTimeout(100); + await warmupAllServers(client); + await waitForPickerUpdate(client); const result = await makeNCalls(client, 30); assert(Math.abs(result['1'] - 20) < 3, `server1: ${result['1']}, server2: ${result['2']}`); }); @@ -411,7 +479,8 @@ describe('Weighted round robin LB policy', () => { blackout_period: '0.01s' }); client = createClient([port1, port2], serviceConfig); - await makeNCalls(client, 10); + await waitForClientReady(client); + await warmupAllServers(client); const result = await makeNCalls(client, 30); assert(Math.abs(result['1'] - result['2']) < 3, `server1: ${result['1']}, server2: ${result[2]}`); }); @@ -423,12 +492,13 @@ describe('Weighted round robin LB policy', () => { weight_update_period: '0.1s' }); client = createClient([port1, port2], serviceConfig); + await waitForClientReady(client); server1MetricRecorder.setQpsMetric(3); server1MetricRecorder.setApplicationUtilizationMetric(1); server2MetricRecorder.setQpsMetric(1); server2MetricRecorder.setApplicationUtilizationMetric(1); - await makeNCalls(client, 10); - await asyncTimeout(200); + await warmupAllServers(client); + await waitForPickerUpdate(client); const result = await makeNCalls(client, 40); assert(Math.abs(result['1'] - 30) < 3, `server1: ${result['1']}, server2: ${result['2']}`); }); diff --git a/packages/grpc-reflection/package.json b/packages/grpc-reflection/package.json index d423b0d44..cc9c183d9 100644 --- a/packages/grpc-reflection/package.json +++ b/packages/grpc-reflection/package.json @@ -25,7 +25,7 @@ "license": "Apache-2.0", "scripts": { "compile": "tsc -p .", - "postcompile": "copyfiles './proto/**/*.proto' build/", + "postcompile": "copyfiles \"./proto/**/*.proto\" build/", "prepare": "npm run generate-types && npm run compile", "test": "mocha --require ts-node/register test/**.ts", "generate-types": "proto-loader-gen-types --longs String --enums String --bytes Array --defaults --oneofs --includeComments --includeDirs proto/ -O src/generated grpc/reflection/v1/reflection.proto grpc/reflection/v1alpha/reflection.proto"