From 1ec453bcebf37c1c65290d544c77f675a43b16cf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Thu, 24 Sep 2026 10:02:28 -0700 Subject: [PATCH 01/12] perf(grpc-js): avoid throwaway timer allocation in ResolvingCall (#3092) Make ResolvingCall.deadlineTimer nullable and initialize it to null instead of allocating a throwaway setTimeout(() => {}, 0) on every RPC. - Initialize deadlineTimer to null and add null checks before clearTimeout. - Set deadlineTimer to null once cleared or when the deadline handler fires. - Add a test verifying that calls without deadlines do not schedule timers. --- packages/grpc-js/src/resolving-call.ts | 13 +++-- packages/grpc-js/test/test-deadline.ts | 67 ++++++++++++++++++++++++++ 2 files changed, 77 insertions(+), 3 deletions(-) diff --git a/packages/grpc-js/src/resolving-call.ts b/packages/grpc-js/src/resolving-call.ts index d7ec6ddcd..8b3b7ea20 100644 --- a/packages/grpc-js/src/resolving-call.ts +++ b/packages/grpc-js/src/resolving-call.ts @@ -56,7 +56,7 @@ 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; @@ -117,7 +117,10 @@ export class ResolvingCall implements Call { } private runDeadlineTimer() { - clearTimeout(this.deadlineTimer); + if (this.deadlineTimer) { + clearTimeout(this.deadlineTimer); + this.deadlineTimer = null; + } this.deadlineStartTime = new Date(); if (this.traceEnabled) { this.trace('Deadline: ' + deadlineToString(this.deadline)); @@ -128,6 +131,7 @@ export class ResolvingCall implements Call { this.trace('Deadline will be reached in ' + timeout + 'ms'); } const handleDeadline = () => { + this.deadlineTimer = null; if (!this.deadlineStartTime) { this.cancelWithStatus(Status.DEADLINE_EXCEEDED, 'Deadline exceeded'); return; @@ -168,7 +172,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( diff --git a/packages/grpc-js/test/test-deadline.ts b/packages/grpc-js/test/test-deadline.ts index 24aebd4d7..2a7f3c700 100644 --- a/packages/grpc-js/test/test-deadline.ts +++ b/packages/grpc-js/test/test-deadline.ts @@ -91,3 +91,70 @@ describe('Client with configured timeout', () => { }); }); }); + +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(); + }); + }); +}); From 555126ce6a604e39de964c640fde693cfe9e744a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Thu, 24 Sep 2026 10:05:11 -0700 Subject: [PATCH 02/12] perf(grpc-js): avoid building metadata map when parsing trailers (#3093) Read grpc-status and grpc-message directly from trailers using metadata.get() instead of calling metadata.getMap(). Calling metadata.getMap() on every completed RPC allocates a map object, iterates over all trailers, and clones binary buffers (such as grpc-status-details-bin) that are immediately thrown away. Reading the two needed headers directly avoids these allocations on the response hot path. --- packages/grpc-js/src/subchannel-call.ts | 20 ++-- packages/grpc-js/test/test-client.ts | 116 +++++++++++++++++++++++- 2 files changed, 125 insertions(+), 11 deletions(-) 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/test/test-client.ts b/packages/grpc-js/test/test-client.ts index e19c87bb6..587dc0a92 100644 --- a/packages/grpc-js/test/test-client.ts +++ b/packages/grpc-js/test/test-client.ts @@ -20,8 +20,7 @@ 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 { ConnectivityState } from '../src/connectivity-state'; import { Http2SubchannelCall } from '../src/subchannel-call'; @@ -240,6 +239,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(() => { From c3d46d5c0e72add63cd9f77a4ce6218126b4bf91 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Thu, 24 Sep 2026 13:47:54 -0700 Subject: [PATCH 03/12] perf(grpc-js): avoid array allocation for single-value metadata headers (#3099) Most gRPC metadata headers contain only a single value (such as authorization tokens, tracing headers, and request IDs). Previously, `toHttp2Headers` always called `.map()`, allocating a temporary array for every header key on every RPC. For single-value headers, we now assign the formatted value directly to the header object instead of wrapping it in an array. This eliminates unnecessary allocations and reduces garbage collection pressure on the RPC hot path. --- packages/grpc-js/src/metadata.ts | 9 ++-- packages/grpc-js/test/test-metadata.ts | 57 +++++++++++++++++++++++++- 2 files changed, 60 insertions(+), 6 deletions(-) diff --git a/packages/grpc-js/src/metadata.ts b/packages/grpc-js/src/metadata.ts index 7ae68ba32..0c00facde 100644 --- a/packages/grpc-js/src/metadata.ts +++ b/packages/grpc-js/src/metadata.ts @@ -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; diff --git a/packages/grpc-js/test/test-metadata.ts b/packages/grpc-js/test/test-metadata.ts index 44182ef39..7faddf5c5 100644 --- a/packages/grpc-js/test/test-metadata.ts +++ b/packages/grpc-js/test/test-metadata.ts @@ -268,14 +268,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 +286,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 +315,7 @@ describe('Metadata', () => { key2: ['value2'], key3: ['value3a', 'value3b'], key4: ['part1, part2'], + 'single-bin': Buffer.from('hello').toString('base64'), 'key-bin': [ 'AAECAwQFBgcICQoLDA0ODw==', 'EBESExQVFhcYGRobHB0eHw==', @@ -309,6 +329,7 @@ describe('Metadata', () => { ['key2', ['value2']], ['key3', ['value3a', 'value3b']], ['key4', ['part1, part2']], + ['single-bin', [Buffer.from('hello')]], [ 'key-bin', [ @@ -327,4 +348,36 @@ describe('Metadata', () => { assert.deepStrictEqual(internalRepr, new Map()); }); }); + + 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]), + ]); + }); + }); }); From 9ddfe5b66bf335d5bc69e8e0c7dbd5cfb02ead64 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Fri, 25 Sep 2026 09:09:28 -0700 Subject: [PATCH 04/12] chore(grpc-js): ignore copied proto dependencies (#3096) The copy-protos script (run during `npm run prepare`) copies external proto definitions from grpc-js-xds into `packages/grpc-js/proto/protoc-gen-validate` and `packages/grpc-js/proto/xds`. Add both directories to .gitignore so these generated dependencies do not show up as untracked files or get accidentally committed. --- .gitignore | 4 ++++ 1 file changed, 4 insertions(+) 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/ From df27d5acf5a0eef04a7b2527e00c3e1da3d6ac54 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Fri, 25 Sep 2026 10:48:22 -0700 Subject: [PATCH 05/12] perf(grpc-js): avoid Date allocations on request and channelz tracking paths (#3090) Store timestamps as epoch milliseconds (Date.now()) instead of allocating Date objects on the request path and channelz counters: - Call layers (ResolvingCall, RetryingCall, LoadBalancingCall): store call start and resolution times as numbers; construct Date instances only when formatting deadline error messages or when trace logging is enabled. - Channelz call tracker and HTTP/2 transport: record last call started, message sent, and message received timestamps as numbers. - Channel activity tracking: record channel idle timestamp as a number. - Server timeout handling: calculate incoming call deadlines with Date.now() + timeout instead of allocating a Date object. - dateToProtoTimestamp: support numeric millisecond timestamps in addition to Date objects, and use Math.floor() to avoid 32-bit integer overflow beyond year 2038. - formatDateDifference: accept numeric timestamps as well as Date instances. - Update channelz and deadline tests for numeric timestamps. --- packages/grpc-js/src/channelz.ts | 22 +-- packages/grpc-js/src/deadline.ts | 14 +- packages/grpc-js/src/internal-channel.ts | 13 +- packages/grpc-js/src/load-balancing-call.ts | 12 +- packages/grpc-js/src/resolving-call.ts | 50 ++++--- packages/grpc-js/src/retrying-call.ts | 10 +- packages/grpc-js/src/server-interceptors.ts | 3 +- packages/grpc-js/src/server.ts | 8 +- packages/grpc-js/src/transport.ts | 8 +- packages/grpc-js/test/test-channelz.ts | 103 +++++++++++++ packages/grpc-js/test/test-deadline.ts | 157 ++++++++++++++++++++ 11 files changed, 339 insertions(+), 61 deletions(-) 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..7b395d575 100644 --- a/packages/grpc-js/src/internal-channel.ts +++ b/packages/grpc-js/src/internal-channel.ts @@ -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; @@ -461,7 +461,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 +645,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,7 +690,7 @@ export class InternalChannel { } } this.callCount -= 1; - this.lastActivityTimestamp = new Date(); + this.lastActivityTimestamp = Date.now(); this.maybeStartIdleTimer(); } @@ -830,7 +829,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..3b29c32ae 100644 --- a/packages/grpc-js/src/load-balancing-call.ts +++ b/packages/grpc-js/src/load-balancing-call.ts @@ -62,8 +62,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, @@ -85,11 +85,11 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { /* 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.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 +142,7 @@ export class LoadBalancingCall implements Call, DeadlineInfoProvider { ' details="' + status.details + '" start time=' + - this.startTime.toISOString() + new Date(this.startTime).toISOString() ); } const finalStatus = { ...status, progress }; @@ -261,7 +261,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/resolving-call.ts b/packages/grpc-js/src/resolving-call.ts index 8b3b7ea20..13aae2191 100644 --- a/packages/grpc-js/src/resolving-call.ts +++ b/packages/grpc-js/src/resolving-call.ts @@ -59,9 +59,9 @@ export class ResolvingCall implements Call { 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 @@ -121,7 +121,7 @@ export class ResolvingCall implements Call { clearTimeout(this.deadlineTimer); this.deadlineTimer = null; } - this.deadlineStartTime = new Date(); + this.deadlineStartTime = Date.now(); if (this.traceEnabled) { this.trace('Deadline: ' + deadlineToString(this.deadline)); } @@ -132,20 +132,35 @@ export class ResolvingCall implements Call { } const handleDeadline = () => { this.deadlineTimer = null; - if (!this.deadlineStartTime) { + 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'); @@ -235,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( @@ -251,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(); } @@ -278,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/retrying-call.ts b/packages/grpc-js/src/retrying-call.ts index 9023b5241..eaa609b40 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,7 +715,7 @@ 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(); 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/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/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-deadline.ts b/packages/grpc-js/test/test-deadline.ts index 2a7f3c700..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,6 +99,154 @@ 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', () => { From 1ef18575ec3cd1662e9e57220f397930d1317817 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Fri, 25 Sep 2026 10:54:33 -0700 Subject: [PATCH 06/12] perf(grpc-js): lazily allocate Metadata opaqueData map (#3095) Every Metadata instance previously allocated a Map for opaqueData in its constructor, even though opaqueData is only used in specific scenarios like caching parsed ORCA load reports on trailers. On high-throughput RPC paths where multiple Metadata instances are created per call, this eager allocation resulted in unnecessary heap allocation and GC pressure. Changes: - Make Metadata.prototype.opaqueData optional and lazily instantiate the Map on the first setOpaque() call. - Use optional chaining in getOpaque() when reading from opaqueData. - Preserve existing clone() and constructor behavior to maintain exact compatibility. - Add unit tests covering lazy allocation, key updates, value types, and cloned instance behavior. --- packages/grpc-js/src/metadata.ts | 7 ++- packages/grpc-js/test/test-metadata.ts | 84 ++++++++++++++++++++++++++ 2 files changed, 89 insertions(+), 2 deletions(-) diff --git a/packages/grpc-js/src/metadata.ts b/packages/grpc-js/src/metadata.ts index 0c00facde..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; @@ -256,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); } @@ -265,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/test/test-metadata.ts b/packages/grpc-js/test/test-metadata.ts index 7faddf5c5..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 = @@ -349,6 +354,85 @@ describe('Metadata', () => { }); }); + 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(); From 17c61a5e53fd375d510c035c6a83501a355666dc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Fri, 25 Sep 2026 11:06:57 -0700 Subject: [PATCH 07/12] test(grpc-js): fix Windows CI test failures (#3094) Fix two test failures when running on Windows: - Increase the timeout for the nonexistent domain test on Windows to allow sequential DNS lookup failures to complete. - Skip the Unix Domain Socket (UDS) idle timer test on Windows where UDS file paths are not supported. --- packages/grpc-js/test/test-client.ts | 2 +- packages/grpc-js/test/test-idle-timer.ts | 7 +++++-- packages/grpc-reflection/package.json | 2 +- 3 files changed, 7 insertions(+), 4 deletions(-) diff --git a/packages/grpc-js/test/test-client.ts b/packages/grpc-js/test/test-client.ts index 587dc0a92..cb80443c4 100644 --- a/packages/grpc-js/test/test-client.ts +++ b/packages/grpc-js/test/test-client.ts @@ -413,7 +413,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', 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-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" From e4ccd66bdc6bdce0cf16847f2e15964bebbc7217 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Tue, 29 Sep 2026 10:45:53 -0700 Subject: [PATCH 08/12] perf(grpc-js): cache precomputed serviceUrl in InternalChannel (#3097) Cache the precomputed serviceUrl used by call credentials in a bounded two-level cache on InternalChannel instead of re-parsing the method path, splitting the host, and constructing a new template string on every RPC. - Add zero-allocation `extractServiceName` and `computeServiceUrl` helpers in `uri-parser.ts`. - Introduce a bounded two-level cache (`serviceUrlCache: Map>`) on `InternalChannel`, capped at 5 hosts and 100 methods per host to strictly bound memory while achieving zero-allocation lookups on cache hits. - Clear `serviceUrlCache` when the channel is closed. - Update `LoadBalancingCall` and `SingleSubchannelChannel` to use the shared URL computation and caching. --- packages/grpc-js/src/internal-channel.ts | 45 ++++++- packages/grpc-js/src/load-balancing-call.ts | 14 +-- .../grpc-js/src/single-subchannel-channel.ts | 15 +-- packages/grpc-js/src/uri-parser.ts | 26 ++++ packages/grpc-js/test/test-client.ts | 116 ++++++++++++++++++ packages/grpc-js/test/test-uri-parser.ts | 64 ++++++++++ 6 files changed, 253 insertions(+), 27 deletions(-) diff --git a/packages/grpc-js/src/internal-channel.ts b/packages/grpc-js/src/internal-channel.ts index 7b395d575..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'; @@ -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, @@ -694,6 +709,33 @@ export class InternalChannel { 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, @@ -819,6 +861,7 @@ export class InternalChannel { this.subchannelPool.unrefUnusedSubchannels(); this.configSelector?.unref(); this.configSelector = null; + this.serviceUrlCache.clear(); } getTarget() { diff --git a/packages/grpc-js/src/load-balancing-call.ts b/packages/grpc-js/src/load-balancing-call.ts index 3b29c32ae..382d644a9 100644 --- a/packages/grpc-js/src/load-balancing-call.ts +++ b/packages/grpc-js/src/load-balancing-call.ts @@ -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'; @@ -73,18 +72,7 @@ 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.serviceUrl = this.channel.getServiceUrl(this.host, this.methodName); this.startTime = Date.now(); } getDeadlineInfo(): string[] { 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/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-client.ts b/packages/grpc-js/test/test-client.ts index cb80443c4..48fc23610 100644 --- a/packages/grpc-js/test/test-client.ts +++ b/packages/grpc-js/test/test-client.ts @@ -21,6 +21,7 @@ import * as http2 from 'http2'; import * as grpc 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'; @@ -453,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-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 + ); + }); + } + }); }); From 2aa8b8836f263809a9cb568e4844f445154017fb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Tue, 29 Sep 2026 10:50:30 -0700 Subject: [PATCH 09/12] perf(grpc-js): memoize method config lookup in default config selector (#3098) On every RPC, the default config selector parsed the method name by splitting the string, filtering empty segments, and looping through method config match levels to find the appropriate configuration. In high-throughput clients, this performed repetitive string splitting, array allocations, and lookups on the hot path. This change: 1. Caches the resolved MethodConfig by method name in a Map within getDefaultConfigSelector, avoiding repeated string splitting and search loops on cache hits. 2. Bounds the cache to 100 method entries to keep memory usage fixed. 3. Clears the cache in unref() when channel resolution updates or the channel closes. 4. Returns fresh CallConfig containers for each RPC, preserving call isolation for pickInformation and dynamic filters. --- .../grpc-js/src/resolving-load-balancer.ts | 71 +++--- .../test/test-resolving-load-balancer.ts | 241 ++++++++++++++++++ 2 files changed, 281 insertions(+), 31 deletions(-) create mode 100644 packages/grpc-js/test/test-resolving-load-balancer.ts 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/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 + ); + }); +}); From 8db8a06638bb8e4293d8ed82490d2bcb4d954ed3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Tue, 29 Sep 2026 11:00:54 -0700 Subject: [PATCH 10/12] perf(grpc-js): avoid redundant metadata clone in RetryingCall on initial attempt (#3101) RetryingCall.startNewAttempt() previously cloned this.initialMetadata on every attempt, even when previousAttempts was 0 and no grpc-previous-rpc-attempts header needed to be added. Because LoadBalancingCall never mutates this.metadata and always clones it in doPick(), the extra clone in RetryingCall on the initial attempt was redundant. Only clone this.initialMetadata in RetryingCall.startNewAttempt() when previousAttempts > 0. --- packages/grpc-js/src/retrying-call.ts | 3 +- packages/grpc-js/test/test-retry.ts | 154 ++++++++++++++++++++++++++ 2 files changed, 156 insertions(+), 1 deletion(-) diff --git a/packages/grpc-js/src/retrying-call.ts b/packages/grpc-js/src/retrying-call.ts index eaa609b40..d84779ab2 100644 --- a/packages/grpc-js/src/retrying-call.ts +++ b/packages/grpc-js/src/retrying-call.ts @@ -718,8 +718,9 @@ export class RetryingCall implements Call, DeadlineInfoProvider { 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/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(); + } + ); + }); }); }); From aed174522e4b104c30c9541927108765b2637872 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Thu, 1 Oct 2026 11:16:11 -0700 Subject: [PATCH 11/12] test(grpc-js): stabilize weighted round robin tests on slow runners (#3104) In test-weighted-round-robin, test traffic was issued immediately against freshly created clients with aggressive 100ms per-call deadlines. On slower or heavily loaded CI environments (such as Windows Kokoro runners), cold-start connection establishment, name resolution, and picker readiness frequently exceeded 100ms, triggering intermittent DEADLINE_EXCEEDED ("Waiting for LB pick") failures. Additionally, blind initial traffic could reach only the first backend before the second finished connecting, leaving the second backend's metrics unrecorded. Furthermore, executing 50 sequential RPC round-trips combined with arbitrary sleep timers pushed total execution time beyond Mocha's default 2-second timeout on Windows VMs (often failing at ~2.005s). This change stabilizes the tests by: 1. Ensuring the channel reaches readiness and that both backends have responded to initial traffic before measurement begins, guaranteeing that connections and initial metric tracking are established across all endpoints. 2. Synchronizing deterministically against picker updates rather than relying on fixed sleep delays, and executing measurement RPCs concurrently over HTTP/2 so the suite runs in ~1s while preserving the expected weight distribution. 3. Increasing the test timeout to provide sufficient headroom against VM scheduling jitter on Windows CI. --- .../grpc-js/test/test-weighted-round-robin.ts | 108 +++++++++++++++--- 1 file changed, 89 insertions(+), 19 deletions(-) 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']}`); }); From 40b61a00726b62ef032cd432d46d5cf47e204d22 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Knut=20Olav=20L=C3=B8ite?= Date: Thu, 1 Oct 2026 11:23:50 -0700 Subject: [PATCH 12/12] perf(grpc-js): avoid allocations for empty CallCredentials on the RPC hot path (#3100) Previously, every RPC executed through LoadBalancingCall composed per-call and subchannel CallCredentials and called generateMetadata(), even when both sides were EmptyCallCredentials. In addition: - Composing any CallCredentials with EmptyCallCredentials on the right-hand side wrapped them in a new ComposedCallCredentials instance. - InsecureChannelCredentialsImpl allocated a new EmptyCallCredentials instance on every getCallCredentials() call when no call credentials were provided. - EmptyCallCredentials.generateMetadata() allocated a new Metadata instance and Promise on every RPC only to merge an empty map. This change: - Treats EmptyCallCredentials as a two-sided identity in compose(), returning `this` directly when composing with an empty instance. - Evaluates the fallback CallCredentials.createEmpty() once per connector in InsecureChannelCredentialsImpl._createSecureConnector(). - Short-circuits metadata generation and merging in LoadBalancingCall when the composed CallCredentials is empty, reusing a pre-resolved Promise. --- packages/grpc-js/src/call-credentials.ts | 12 ++ packages/grpc-js/src/channel-credentials.ts | 3 +- packages/grpc-js/src/load-balancing-call.ts | 25 ++- .../grpc-js/test/test-call-credentials.ts | 182 ++++++++++++++++++ .../grpc-js/test/test-channel-credentials.ts | 29 +++ 5 files changed, 245 insertions(+), 6 deletions(-) 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/load-balancing-call.ts b/packages/grpc-js/src/load-balancing-call.ts index 382d644a9..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, @@ -38,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'; @@ -140,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; @@ -168,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 @@ -182,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( { 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(); + }); }); });