From a738dfb3c872a1b2d9b25ca636afd47629639430 Mon Sep 17 00:00:00 2001 From: Philip Youssef Date: Sun, 27 Sep 2026 08:34:31 -0700 Subject: [PATCH] fix(grpc-js): stop retry attempts after the call is cancelled cancelWithStatus() reported the status and cancelled the child calls, but left the retry state unchanged and the backoff timer running. When a call was cancelled while a retry was waiting out its backoff, the timer still started a new attempt, and the call kept retrying through its remaining attempts after the application had cancelled it. Those attempts are sent after reportStatus() has cleared the write buffer, so a unary attempt carries headers but no message. The server never responds, and the stream stays open until the connection closes, occupying one of the server's concurrent streams. Enough of them block every new call on the channel. Track the retry timer, and on cancellation disable further attempts and clear the retry and hedging timers. Co-Authored-By: Claude Opus 5.5 --- packages/grpc-js/src/retrying-call.ts | 16 ++- packages/grpc-js/test/test-retry.ts | 142 ++++++++++++++++++++++++++ 2 files changed, 157 insertions(+), 1 deletion(-) diff --git a/packages/grpc-js/src/retrying-call.ts b/packages/grpc-js/src/retrying-call.ts index eaa609b40..4065f28f4 100644 --- a/packages/grpc-js/src/retrying-call.ts +++ b/packages/grpc-js/src/retrying-call.ts @@ -207,6 +207,7 @@ export class RetryingCall implements Call, DeadlineInfoProvider { */ private attempts = 0; private hedgingTimer: NodeJS.Timeout | null = null; + private retryTimer: NodeJS.Timeout | null = null; private committedCallIndex: number | null = null; private initialRetryBackoffSec = 0; private nextRetryBackoffSec = 0; @@ -321,6 +322,18 @@ export class RetryingCall implements Call, DeadlineInfoProvider { 'cancelWithStatus code: ' + status + ' details: "' + details + '"' ); } + /* No new attempt may start after cancellation. Otherwise a pending retry + * or hedging timer would start one, opening a stream that nothing will + * ever end. */ + this.state = 'NO_RETRY'; + if (this.retryTimer) { + clearTimeout(this.retryTimer); + this.retryTimer = null; + } + if (this.hedgingTimer) { + clearTimeout(this.hedgingTimer); + this.hedgingTimer = null; + } this.reportStatus({ code: status, details, metadata: new Metadata() }); for (const { call } of this.underlyingCalls) { call.cancelWithStatus(status, details); @@ -487,7 +500,8 @@ export class RetryingCall implements Call, DeadlineInfoProvider { retryDelayMs = pushback; this.nextRetryBackoffSec = this.initialRetryBackoffSec; } - setTimeout(() => { + this.retryTimer = setTimeout(() => { + this.retryTimer = null; if (this.state !== 'RETRY') { callback(false); return; diff --git a/packages/grpc-js/test/test-retry.ts b/packages/grpc-js/test/test-retry.ts index d7e5a09cd..e6321ac84 100644 --- a/packages/grpc-js/test/test-retry.ts +++ b/packages/grpc-js/test/test-retry.ts @@ -468,3 +468,145 @@ describe('Retries', () => { }); }); }); + +describe('Retries after cancellation', () => { + /* Counts the attempts started for each call-id. Client streaming handlers + * run as soon as an attempt's headers arrive, so every attempt is counted. */ + const attemptCounts = new Map(); + const cancellationServiceImpl = { + echo: ( + call: grpc.ServerUnaryCall, + callback: grpc.sendUnaryData + ) => { + const id = call.metadata.get('call-id')[0] as string | undefined; + if (id) { + attemptCounts.set(id, (attemptCounts.get(id) ?? 0) + 1); + callback({ code: grpc.status.UNAVAILABLE, details: 'Unavailable' }); + } else { + callback(null, call.request); + } + }, + echoClientStream: ( + call: grpc.ServerReadableStream, + callback: grpc.sendUnaryData + ) => { + const id = call.metadata.get('call-id')[0] as string; + attemptCounts.set(id, (attemptCounts.get(id) ?? 0) + 1); + callback({ code: grpc.status.UNAVAILABLE, details: 'Unavailable' }); + }, + }; + const serviceConfig = { + loadBalancingConfig: [], + methodConfig: [ + { + name: [{ service: 'EchoService' }], + retryPolicy: { + maxAttempts: 3, + initialBackoff: '0.3s', + maxBackoff: '0.3s', + backoffMultiplier: 1, + retryableStatusCodes: ['UNAVAILABLE'], + }, + }, + ], + }; + let server: grpc.Server; + let client: InstanceType; + before(done => { + server = new grpc.Server({ 'grpc.max_concurrent_streams': 3 }); + server.addService(EchoService.service, cancellationServiceImpl); + server.bindAsync( + 'localhost:0', + grpc.ServerCredentials.createInsecure(), + (error, port) => { + if (error) { + done(error); + return; + } + client = new EchoService( + `localhost:${port}`, + grpc.credentials.createInsecure(), + { 'grpc.service_config': JSON.stringify(serviceConfig) } + ); + done(); + } + ); + }); + + after(() => { + client.close(); + server.forceShutdown(); + }); + + function cancelDuringBackoff( + id: string, + streaming: boolean, + done: () => void + ) { + const metadata = new grpc.Metadata(); + metadata.set('call-id', id); + let call: grpc.ClientWritableStream | grpc.ClientUnaryCall; + if (streaming) { + const stream = client.echoClientStream(metadata, () => {}); + stream.write({ value: 'test value', value2: 3 }); + call = stream; + } else { + call = client.echo( + { value: 'test value', value2: 3 }, + metadata, + () => {} + ); + } + const waitForFirstAttempt = () => { + if (!attemptCounts.get(id)) { + setTimeout(waitForFirstAttempt, 5); + return; + } + // The first attempt failed; the retry is waiting out its backoff. + setTimeout(() => { + call.cancel(); + done(); + }, 50); + }; + waitForFirstAttempt(); + } + + it('Should not start new attempts after the call is cancelled during backoff', done => { + cancelDuringBackoff('cancelled', true, () => { + setTimeout(() => { + assert.strictEqual(attemptCounts.get('cancelled'), 1); + done(); + }, 1000); + }); + }); + + it('Should not leak streams when calls are cancelled during backoff', done => { + let cancelled = 0; + for (let i = 0; i < 3; i++) { + /* A unary retry attempt started after cancellation has no message to + * send, so the server never responds and its stream stays open. */ + cancelDuringBackoff(`leak-${i}`, false, () => { + cancelled += 1; + if (cancelled < 3) { + return; + } + // The server allows 3 concurrent streams. Leaked retry attempts would + // hold all of them, so this call could never start. + setTimeout(() => { + client.echo( + { value: 'test value', value2: 3 }, + { deadline: Date.now() + 2000 }, + (error: grpc.ServiceError, response: any) => { + assert.ifError(error); + assert.deepStrictEqual(response, { + value: 'test value', + value2: 3, + }); + done(); + } + ); + }, 500); + }); + } + }); +});