Skip to content
Merged
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
4 changes: 4 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -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/
12 changes: 12 additions & 0 deletions packages/grpc-js/src/call-credentials.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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]));
}

Expand Down Expand Up @@ -197,6 +200,9 @@ class SingleCallCredentials extends CallCredentials {
}

compose(other: CallCredentials): CallCredentials {
if (other instanceof EmptyCallCredentials) {
return this;
}
return new ComposedCallCredentials([this, other]);
}

Expand Down Expand Up @@ -225,3 +231,9 @@ class EmptyCallCredentials extends CallCredentials {
return other instanceof EmptyCallCredentials;
}
}

export function isEmptyCallCredentials(
callCredentials: CallCredentials
): boolean {
return callCredentials instanceof EmptyCallCredentials;
}
3 changes: 2 additions & 1 deletion packages/grpc-js/src/channel-credentials.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand All @@ -195,7 +196,7 @@ class InsecureChannelCredentialsImpl extends ChannelCredentials {
return Promise.resolve();
},
getCallCredentials: () => {
return callCredentials ?? CallCredentials.createEmpty();
return credentials;
},
destroy() {}
}
Expand Down
22 changes: 12 additions & 10 deletions packages/grpc-js/src/channelz.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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,
};
}
Expand Down
14 changes: 10 additions & 4 deletions packages/grpc-js/src/deadline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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';
}
58 changes: 50 additions & 8 deletions packages/grpc-js/src/internal-channel.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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;
Expand All @@ -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<string, Map<string, string>>();

constructor(
target: string,
private readonly credentials: ChannelCredentials,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 ' +
Expand Down Expand Up @@ -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<string, string>();
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,
Expand Down Expand Up @@ -820,6 +861,7 @@ export class InternalChannel {
this.subchannelPool.unrefUnusedSubchannels();
this.configSelector?.unref();
this.configSelector = null;
this.serviceUrlCache.clear();
}

getTarget() {
Expand All @@ -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;
Expand Down
51 changes: 27 additions & 24 deletions packages/grpc-js/src/load-balancing-call.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
*
*/

import { CallCredentials } from './call-credentials';
import { CallCredentials, isEmptyCallCredentials } from './call-credentials';
import {
Call,
DeadlineInfoProvider,
Expand All @@ -31,14 +31,15 @@ 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';
import { AuthContext } from './auth-context';
import { SubchannelInterface } from './subchannel-interface';

const TRACER_NAME = 'load_balancing_call';
const RESOLVED_EMPTY_METADATA: Promise<Metadata | undefined> =
Promise.resolve(undefined);

export type RpcProgress = 'NOT_STARTED' | 'DROP' | 'REFUSED' | 'PROCESSED';

Expand All @@ -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,
Expand All @@ -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');
Expand Down Expand Up @@ -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 };
Expand All @@ -152,6 +142,17 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider {
}
}

private generateCallCredentialsMetadata(
callCredentials: CallCredentials
): Promise<Metadata | undefined> {
return isEmptyCallCredentials(callCredentials)
? RESOLVED_EMPTY_METADATA
: callCredentials.generateMetadata({
method_name: this.methodName,
service_url: this.serviceUrl,
});
}

doPick() {
if (this.ended) {
return;
Expand Down Expand Up @@ -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
Expand All @@ -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(
{
Expand Down Expand Up @@ -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(
Expand Down
Loading