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); + }); + } + }); +});