From e245d07021d310743f06b1fc792446aa49b4878b Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Thu, 10 Sep 2026 15:10:43 +0800 Subject: [PATCH 1/2] Prevent abort listener leaks in async client requests --- lib/client.js | 63 +++++++++++++------------- test/test-client.js | 98 ++++++++++++++++++++++++++++++++++++++++ test/test-destruction.js | 32 +++++++++++++ 3 files changed, 162 insertions(+), 31 deletions(-) diff --git a/lib/client.js b/lib/client.js index 1b3980d0d..624646365 100644 --- a/lib/client.js +++ b/lib/client.js @@ -136,60 +136,61 @@ class Client extends Entity { let sequenceNumber = null; let isResolved = false; let isTimeout = false; + let effectiveSignal = options.signal; + let timeoutSignal; + + const onTimeout = () => { + isTimeout = true; + }; const cleanup = () => { if (sequenceNumber !== null) { this._sequenceNumberToCallbackMap.delete(sequenceNumber); } + if (effectiveSignal) { + effectiveSignal.removeEventListener('abort', onAbort); + } + if (timeoutSignal) { + timeoutSignal.removeEventListener('abort', onTimeout); + } isResolved = true; }; - let effectiveSignal = options.signal; + const onAbort = () => { + if (!isResolved) { + cleanup(); + const error = isTimeout + ? new TimeoutError('Service request', options.timeout, { + entityType: 'service', + entityName: this._serviceName, + }) + : new AbortError('Service request', undefined, { + entityType: 'service', + entityName: this._serviceName, + }); + reject(error); + } + }; if (options.timeout !== undefined && options.timeout >= 0) { - const timeoutSignal = AbortSignal.timeout(options.timeout); - - timeoutSignal.addEventListener('abort', () => { - isTimeout = true; - }); + timeoutSignal = AbortSignal.timeout(options.timeout); if (options.signal) { effectiveSignal = AbortSignal.any([options.signal, timeoutSignal]); } else { effectiveSignal = timeoutSignal; } + + timeoutSignal.addEventListener('abort', onTimeout, { once: true }); } if (effectiveSignal) { if (effectiveSignal.aborted) { - const error = isTimeout - ? new TimeoutError('Service request', options.timeout, { - entityType: 'service', - entityName: this._serviceName, - }) - : new AbortError('Service request', undefined, { - entityType: 'service', - entityName: this._serviceName, - }); - reject(error); + onAbort(); return; } - effectiveSignal.addEventListener('abort', () => { - if (!isResolved) { - cleanup(); - const error = isTimeout - ? new TimeoutError('Service request', options.timeout, { - entityType: 'service', - entityName: this._serviceName, - }) - : new AbortError('Service request', undefined, { - entityType: 'service', - entityName: this._serviceName, - }); - reject(error); - } - }); + effectiveSignal.addEventListener('abort', onAbort, { once: true }); } try { diff --git a/test/test-client.js b/test/test-client.js index 6261f0186..1a527ad93 100644 --- a/test/test-client.js +++ b/test/test-client.js @@ -1,4 +1,5 @@ import assert from 'assert'; +import { getEventListeners } from 'events'; import sinon from 'sinon'; import rclnodejsBinding from '../lib/native_loader.js'; import Client from '../lib/client.js'; @@ -191,6 +192,103 @@ describe('Client coverage testing', function () { assert.deepStrictEqual(result, new MockTypeClass.Response()); }); + it('sendRequestAsync releases abort listeners after successful requests', async function () { + const client = new Client( + mockHandle, + mockNodeHandle, + 'test_service', + MockTypeClass, + {} + ); + const controller = new AbortController(); + const callerListener = sinon.spy(); + controller.signal.addEventListener('abort', callerListener); + + for (let requestIndex = 0; requestIndex < 3; requestIndex++) { + const promise = client.sendRequestAsync( + { a: requestIndex }, + { signal: controller.signal } + ); + client.processResponse(12345, new MockTypeClass.Response()); + assert.deepStrictEqual(await promise, { sum: 3 }); + assert.strictEqual(client._sequenceNumberToCallbackMap.size, 0); + assert.deepStrictEqual(getEventListeners(controller.signal, 'abort'), [ + callerListener, + ]); + } + + controller.abort(); + assert.ok(callerListener.calledOnce); + }); + + for (const outcome of [ + 'success', + 'manual abort', + 'timeout', + 'serialization error', + 'send error', + 'pre-aborted signal', + ]) { + it(`sendRequestAsync cleans up timeout and abort listeners on ${outcome}`, async function () { + const client = new Client( + mockHandle, + mockNodeHandle, + 'test_service', + MockTypeClass, + {} + ); + const controller = new AbortController(); + const timeoutController = new AbortController(); + const timeoutStub = sandbox + .stub(AbortSignal, 'timeout') + .returns(timeoutController.signal); + const combinedSignalSpy = sandbox.spy(AbortSignal, 'any'); + const failure = new Error('Request failed'); + + if (outcome === 'pre-aborted signal') { + controller.abort(); + } else if (outcome === 'serialization error') { + sandbox + .stub(MockTypeClass.Request.prototype, 'serialize') + .throws(failure); + } else if (outcome === 'send error') { + rclnodejsBinding.sendRequest.throws(failure); + } + + const promise = client.sendRequestAsync( + { a: 1 }, + { signal: controller.signal, timeout: 1000 } + ); + const effectiveSignal = combinedSignalSpy.firstCall.returnValue; + + if (outcome === 'success') { + client.processResponse(12345, new MockTypeClass.Response()); + assert.deepStrictEqual(await promise, { sum: 3 }); + } else { + const expectedError = outcome.endsWith('error') + ? failure + : { name: outcome === 'timeout' ? 'TimeoutError' : 'AbortError' }; + const rejected = assert.rejects(promise, expectedError); + if (outcome === 'manual abort') controller.abort(); + if (outcome === 'timeout') timeoutController.abort(); + await rejected; + } + + assert.ok(timeoutStub.calledOnceWithExactly(1000)); + assert.strictEqual(client._sequenceNumberToCallbackMap.size, 0); + for (const signal of [ + controller.signal, + timeoutController.signal, + effectiveSignal, + ]) { + assert.deepStrictEqual(getEventListeners(signal, 'abort'), []); + } + if (outcome === 'pre-aborted signal') { + assert.ok(rclnodejsBinding.sendRequest.notCalled); + } + }); + } + it('sendRequestAsync handles timeout', async function () { const client = new Client( mockHandle, diff --git a/test/test-destruction.js b/test/test-destruction.js index 70ab2f58d..13916ba1a 100644 --- a/test/test-destruction.js +++ b/test/test-destruction.js @@ -17,6 +17,7 @@ import childProcess from 'child_process'; import rclnodejs from '../index.js'; import { fileURLToPath } from 'url'; import { dirname } from 'path'; +import { createDelay } from './utils.js'; const __filename = fileURLToPath(import.meta.url); const __dirname = dirname(__filename); @@ -36,6 +37,37 @@ describe('Node & Entity destroy testing', function () { node.destroy(); }); + it('node.destroy() removes the node from the graph while referenced', async function () { + const observer = new rclnodejs.Node(`graph_observer_${process.pid}`); + const nodeName = `graph_destroyed_${process.pid}`; + const node = new rclnodejs.Node(nodeName); + + const waitForPresence = async (expectedPresence) => { + const deadline = Date.now() + 10000; + let present = observer.getNodeNames().includes(nodeName); + while (present !== expectedPresence && Date.now() < deadline) { + await createDelay(50); + present = observer.getNodeNames().includes(nodeName); + } + assert.strictEqual( + present, + expectedPresence, + expectedPresence + ? 'Node never appeared in the graph' + : 'Destroyed node is still advertised in the graph' + ); + }; + + try { + await waitForPresence(true); + node.destroy(); + await waitForPresence(false); + } finally { + node.destroy(); + observer.destroy(); + } + }); + it('node.destroy() twice', function () { var node = rclnodejs.createNode('my_node2'); node.destroy(); From 61102810d16f8f825ab333d79126be746e9b34a3 Mon Sep 17 00:00:00 2001 From: Minggang Wang Date: Thu, 10 Sep 2026 15:25:14 +0800 Subject: [PATCH 2/2] Remove unit test --- test/test-destruction.js | 32 -------------------------------- 1 file changed, 32 deletions(-) diff --git a/test/test-destruction.js b/test/test-destruction.js index 13916ba1a..70ab2f58d 100644 --- a/test/test-destruction.js +++ b/test/test-destruction.js @@ -17,7 +17,6 @@ import childProcess from 'child_process'; import rclnodejs from '../index.js'; import { fileURLToPath } from 'url'; import { dirname } from 'path'; -import { createDelay } from './utils.js'; const __filename = fileURLToPath(import.meta.url); const __dirname = dirname(__filename); @@ -37,37 +36,6 @@ describe('Node & Entity destroy testing', function () { node.destroy(); }); - it('node.destroy() removes the node from the graph while referenced', async function () { - const observer = new rclnodejs.Node(`graph_observer_${process.pid}`); - const nodeName = `graph_destroyed_${process.pid}`; - const node = new rclnodejs.Node(nodeName); - - const waitForPresence = async (expectedPresence) => { - const deadline = Date.now() + 10000; - let present = observer.getNodeNames().includes(nodeName); - while (present !== expectedPresence && Date.now() < deadline) { - await createDelay(50); - present = observer.getNodeNames().includes(nodeName); - } - assert.strictEqual( - present, - expectedPresence, - expectedPresence - ? 'Node never appeared in the graph' - : 'Destroyed node is still advertised in the graph' - ); - }; - - try { - await waitForPresence(true); - node.destroy(); - await waitForPresence(false); - } finally { - node.destroy(); - observer.destroy(); - } - }); - it('node.destroy() twice', function () { var node = rclnodejs.createNode('my_node2'); node.destroy();