From 9260f9529d86aa52036b4c94e87568eec87d15fd Mon Sep 17 00:00:00 2001 From: Jason Marshall Date: Mon, 27 Jul 2026 12:27:17 -0700 Subject: [PATCH 1/6] Rework cluster and worker lifecycle logic to be more similar. In particular, cluster now uses a similar announcement system to filter workers. This fixes #181. Signed-off-by: Jason Marshall --- lib/cluster.js | 173 ++++++++++++++++++++++++++++---------------- lib/worker.js | 36 ++++----- test/clusterTest.js | 18 ++++- 3 files changed, 138 insertions(+), 89 deletions(-) diff --git a/lib/cluster.js b/lib/cluster.js index 237c9a81..96bec70c 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -22,6 +22,7 @@ * cluster master. */ +const { debuglog } = require('node:util'); const Registry = require('./registry'); // We need to lazy-load the 'cluster' module as some application servers - // namely Passenger - crash when it is imported. @@ -31,20 +32,71 @@ let cluster = () => { return data; }; +const debug = debuglog('prom:metrics:cluster'); +const ANNOUNCEMENT = '@prometheus-io/client:announcement'; const GET_METRICS_REQ = '@prometheus-io/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes'; let registries = [Registry.globalRegistry]; let requestCtr = 0; // Concurrency control -let listenersAdded = false; const requests = new Map(); // Pending requests for workers' local metrics. class AggregatorRegistry extends Registry { + /** + * Create a Registry. + * @param regContentType + */ constructor(regContentType = Registry.PROMETHEUS_CONTENT_TYPE) { super(regContentType); - addListeners(); + + if (cluster().isPrimary) { + this.workers = new Map(); + } + + addListeners(this); } + addWorker(id) { + if (this.workers.has(id)) { + debug('duplicate worker announcement', id); + return; + } + + const worker = cluster().workers[id]; + + worker.on('disconnect', () => { + debug('worker disconnected', id); + this.workers.delete(id); + }); + + worker.on('message', message => { + if (message.type === GET_METRICS_RES) { + const request = requests.get(message.requestId); + + if (request === undefined) { + debug('unexpected results from worker', id); + return; + } + + const response = request.responseHandlers.get(id); + if (response === undefined) { + return; + } + request.responseHandlers.delete(id); + + if (message.error) { + response.reject(new Error(message.error)); + } else { + response.resolve({ + threadId: id, + metrics: message.metrics, + }); + } + } + }); + + this.workers.set(id, worker); + } /** * Gets aggregated metrics for all workers. The optional callback and * returned Promise resolve with the same value; either may be used. @@ -53,9 +105,9 @@ class AggregatorRegistry extends Registry { */ clusterMetrics() { const requestId = requestCtr++; - const workers = Object.values(cluster().workers) - .filter(worker => worker.isConnected()) - .sort((left, right) => left.id - right.id); + const orderedWorkers = [...this.workers.values()].sort( + (left, right) => left.id - right.id, + ); return new Promise((resolve, reject) => { let settled = false; @@ -78,37 +130,37 @@ class AggregatorRegistry extends Registry { responseHandlers, done, errorTimeout: setTimeout(() => { - const err = new Error('Operation timed out.'); + const err = new Error( + `Operation timed out. ${request.responseHandlers.size} outstanding responses.`, + ); request.done(err); - }, 5000), + }, 5_000), }; requests.set(requestId, request); - - const message = { - type: GET_METRICS_REQ, - requestId, - }; - - if (workers.length === 0) { - // No workers were up - process.nextTick(() => done(undefined, '')); - return; - } - - const responsePromises = workers.map( + const responsePromises = orderedWorkers.map( worker => new Promise((resolveResponse, rejectResponse) => { responseHandlers.set(worker.id, { resolve: resolveResponse, reject: rejectResponse, }); - worker.send(message); + + worker.send({ + type: GET_METRICS_REQ, + requestId, + }); }), ); - Promise.all(responsePromises) - .then(metrics => Registry.aggregate(metrics.flat()).metrics()) - .then(result => done(undefined, result), done); + if (responsePromises.length === 0) { + debug('No workers found for requestId', requestId); + process.nextTick(() => done(undefined, '')); + } else { + Promise.all(responsePromises) + .then(responses => responses.flatMap(response => response.metrics)) + .then(metrics => Registry.aggregate(metrics).metrics()) + .then(result => done(undefined, result), done); + } }); } @@ -157,54 +209,51 @@ class AggregatorRegistry extends Registry { * than once). * @returns {void} */ -function addListeners() { - if (listenersAdded) return; - listenersAdded = true; - +function addListeners(registry) { if (cluster().isPrimary) { // Listen for worker responses to requests for local metrics cluster().on('message', (worker, message) => { - if (message.type === GET_METRICS_RES) { - const request = requests.get(message.requestId); - - if (request === undefined) { - return; - } - - const response = request.responseHandlers.get(worker.id); - if (response === undefined) { - return; - } - request.responseHandlers.delete(worker.id); - - if (message.error) { - response.reject(new Error(message.error)); - } else { - response.resolve(message.metrics); - } + if (message.type === ANNOUNCEMENT) { + registry.addWorker(worker.id); } }); + + announce(); } else { // Respond to master's requests for worker's local metrics. - process.on('message', message => { - if (message.type === GET_METRICS_REQ) { - Promise.all(registries.map(r => r.getMetricsAsJSON())) - .then(metrics => { - process.send({ - type: GET_METRICS_RES, - requestId: message.requestId, - metrics, - }); - }) - .catch(error => { - process.send({ - type: GET_METRICS_RES, - requestId: message.requestId, - error: error.message, - }); + process.on('message', async message => { + if (message.type === ANNOUNCEMENT) { + process.send({ type: ANNOUNCEMENT }); + } else if (message.type === GET_METRICS_REQ) { + const metrics = await Promise.all( + registries.map(r => r.getMetricsAsJSON()), + ); + + try { + process.send({ + type: GET_METRICS_RES, + requestId: message.requestId, + metrics, }); + } catch (error) { + process.send({ + type: GET_METRICS_RES, + requestId: message.requestId, + error: error.message, + }); + } } }); + + process.send({ type: ANNOUNCEMENT }); + } +} + +function announce() { + for (const worker of Object.values(cluster().workers)) { + if (worker.isConnected()) { + worker.send({ type: ANNOUNCEMENT }); + } } } diff --git a/lib/worker.js b/lib/worker.js index 03ad2c93..681afe3c 100644 --- a/lib/worker.js +++ b/lib/worker.js @@ -22,22 +22,21 @@ * main thread. */ +const { debuglog } = require('node:util'); const Registry = require('./registry'); const worker = require('node:worker_threads'); const { isMainThread, threadId, BroadcastChannel } = worker; +const debug = debuglog('prom:metrics:worker'); const ANNOUNCEMENT = '@prometheus-io/client:announcement'; const GET_METRICS_REQ = '@prometheus-io/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes'; const ANNOUNCEMENT_CHANNEL = new BroadcastChannel( '@prometheus-io/client:announce', -); - -ANNOUNCEMENT_CHANNEL.unref(); +).unref(); let registries = [Registry.globalRegistry]; let requestCtr = 0; // Concurrency control -let listenersAdded = false; const requests = new Map(); // Pending requests for workers' local metrics. class WorkerRegistry extends Registry { @@ -70,6 +69,7 @@ class WorkerRegistry extends Registry { */ addWorker(name) { if (this.channels.has(name)) { + debug('duplicate worker announcement', name); return; } @@ -85,6 +85,7 @@ class WorkerRegistry extends Registry { const request = requests.get(message.requestId); if (request === undefined) { + debug('unexpected results from worker', name); return; } @@ -115,8 +116,8 @@ class WorkerRegistry extends Registry { * metrics. */ workerMetrics() { - const requestId = requestCtr++; //TODO: We should be able to collect metrics for the collector thread. + const requestId = requestCtr++; return new Promise((resolve, reject) => { let settled = false; @@ -146,7 +147,7 @@ class WorkerRegistry extends Registry { }, 5_000), }; requests.set(requestId, request); - const responsePromises = [...this.channels.keys()].map( + const responsePromises = [...this.channels.keys()].sort().map( name => new Promise((resolveResponse, rejectResponse) => { responseHandlers.set(name, { @@ -163,19 +164,14 @@ class WorkerRegistry extends Registry { }); if (responsePromises.length === 0) { - // No workers were up + debug('No workers found for requestId', requestId); process.nextTick(() => done(undefined, '')); - return; + } else { + Promise.all(responsePromises) + .then(responses => responses.flatMap(response => response.metrics)) + .then(metrics => Registry.aggregate(metrics).metrics()) + .then(result => done(undefined, result), done); } - - Promise.all(responsePromises) - .then(responses => - responses - .sort((left, right) => left.threadId - right.threadId) - .flatMap(response => response.metrics), - ) - .then(metrics => Registry.aggregate(metrics).metrics()) - .then(result => done(undefined, result), done); }); } @@ -223,12 +219,6 @@ class WorkerRegistry extends Registry { * Watch for metrics collection events. */ function addListeners(registry) { - if (listenersAdded) { - return; - } - - listenersAdded = true; - const name = `@prometheus-io/client:worker:${threadId}`; const channel = new BroadcastChannel(name); diff --git a/test/clusterTest.js b/test/clusterTest.js index 64e1db60..8e2a690b 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -18,6 +18,7 @@ const cluster = require('cluster'); const process = require('process'); const Registry = require('../lib/cluster'); +const ANNOUNCEMENT = '@prometheus-io/client:announcement'; const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes'; function metric(value) { @@ -75,6 +76,9 @@ describe.each([ }); it('aggregates worker responses in worker id order', async () => { + jest.resetModules(); + + const registry = new Registry(regType); const originalWorkers = cluster.workers; const workers = Object.fromEntries( [1, 2, 3].map(id => [ @@ -82,23 +86,29 @@ describe.each([ { id, isConnected: () => true, + on: jest.fn(), send: jest.fn(), }, ]), ); cluster.workers = workers; + Object.keys(workers).forEach(id => { + cluster.emit('message', workers[id], { type: ANNOUNCEMENT }); + }); + try { - const registry = new Registry(regType); const result = registry.clusterMetrics(); - const requestId = workers[1].send.mock.calls[0][0].requestId; + const calls = workers[1].send.mock.calls; + const requestId = calls[0][0].requestId; for (const [id, value] of [ [3, 0.3437699], [1, 0.5848208], [2, 0.5479198], ]) { - cluster.emit('message', workers[id], { + const listener = workers[id].on.mock.calls.at(-1)[1]; + listener({ type: GET_METRICS_RES, requestId, metrics: [[metric(value)]], @@ -109,7 +119,7 @@ describe.each([ } finally { cluster.workers = originalWorkers; } - }); + }, 6_000); }); describe('message handling', () => { From 580ca69ecd31b62890e4ac5dbc79ff5fd3232929 Mon Sep 17 00:00:00 2001 From: Jason Marshall Date: Mon, 27 Jul 2026 15:15:43 -0700 Subject: [PATCH 2/6] Fixing missing unref for BroadcastChannel. Caught by jest warnings. Signed-off-by: Jason Marshall --- lib/worker.js | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/lib/worker.js b/lib/worker.js index 681afe3c..081094b6 100644 --- a/lib/worker.js +++ b/lib/worker.js @@ -73,7 +73,7 @@ class WorkerRegistry extends Registry { return; } - const channel = new BroadcastChannel(name); + const channel = new BroadcastChannel(name).unref(); channel.addEventListener('close', () => { this.channels.delete(name); }); @@ -220,9 +220,7 @@ class WorkerRegistry extends Registry { */ function addListeners(registry) { const name = `@prometheus-io/client:worker:${threadId}`; - const channel = new BroadcastChannel(name); - - channel.unref(); + const channel = new BroadcastChannel(name).unref(); ANNOUNCEMENT_CHANNEL.addEventListener('message', async event => { const message = event.data; From 6336df8ce1312da45a7d0b6d8970e18377e4e319 Mon Sep 17 00:00:00 2001 From: Jason Marshall Date: Mon, 27 Jul 2026 15:52:05 -0700 Subject: [PATCH 3/6] Export stats from the primary thread. Fixes #183 Signed-off-by: Jason Marshall --- example/cluster.js | 4 ++++ lib/cluster.js | 23 +++++++++++++++-------- test/clusterTest.js | 23 ++++++++++++++++++++++- 3 files changed, 41 insertions(+), 9 deletions(-) diff --git a/example/cluster.js b/example/cluster.js index 2686a27e..3472cdb3 100644 --- a/example/cluster.js +++ b/example/cluster.js @@ -22,6 +22,10 @@ const metricsServer = express(); const clusterRegistry = new ClusterRegistry(); if (cluster.isPrimary) { + require('../').collectDefaultMetrics({ + gcDurationBuckets: [0.001, 0.01, 0.1, 1, 2, 5], // These are the default buckets. + }); + for (let i = 1; i <= 4; i++) { cluster.fork({ ...process.env, PORT: 3000 + i }); } diff --git a/lib/cluster.js b/lib/cluster.js index 96bec70c..0c9a84a4 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -137,7 +137,7 @@ class AggregatorRegistry extends Registry { }, 5_000), }; requests.set(requestId, request); - const responsePromises = orderedWorkers.map( + const workerMetrics = orderedWorkers.map( worker => new Promise((resolveResponse, rejectResponse) => { responseHandlers.set(worker.id, { @@ -152,15 +152,22 @@ class AggregatorRegistry extends Registry { }), ); - if (responsePromises.length === 0) { + const myMetrics = Promise.all( + registries.map(r => r.getMetricsAsJSON()), + ).then(metrics => { + return { metrics }; + }); + + if (workerMetrics.length === 0) { debug('No workers found for requestId', requestId); - process.nextTick(() => done(undefined, '')); - } else { - Promise.all(responsePromises) - .then(responses => responses.flatMap(response => response.metrics)) - .then(metrics => Registry.aggregate(metrics).metrics()) - .then(result => done(undefined, result), done); } + + const allMetrics = [myMetrics, ...workerMetrics]; + + Promise.all(allMetrics) + .then(responses => responses.flatMap(response => response.metrics)) + .then(metrics => Registry.aggregate(metrics).metrics()) + .then(result => done(undefined, result), done); }); } diff --git a/test/clusterTest.js b/test/clusterTest.js index 8e2a690b..fd532ec6 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -72,7 +72,7 @@ describe.each([ const AggregatorRegistry = require('../lib/cluster'); const ar = new AggregatorRegistry(regType); const metrics = await ar.clusterMetrics(); - expect(metrics).toEqual(''); + expect(metrics.trim()).toEqual(''); }); it('aggregates worker responses in worker id order', async () => { @@ -120,6 +120,27 @@ describe.each([ cluster.workers = originalWorkers; } }, 6_000); + + it('aggregates telemetry from primary thread', async () => { + jest.resetModules(); + + require('../lib/cluster'); + const { Gauge } = require('../index'); + + const gauge = new Gauge({ name: 'primary_gauge_test', help: 'help' }); + + try { + const AggregatorRegistry = require('../lib/cluster'); + const ar = new AggregatorRegistry(regType); + + gauge.set(10); + + const result = ar.clusterMetrics(); + await expect(result).resolves.toContain('primary_gauge_test 10\n'); + } finally { + gauge.remove(); + } + }); }); describe('message handling', () => { From a652f072f40d5c1468f78751d80f86c53f73c9d3 Mon Sep 17 00:00:00 2001 From: Jason Marshall Date: Mon, 27 Jul 2026 18:55:01 -0700 Subject: [PATCH 4/6] Rework state tracking to support bot #155 and #788 Signed-off-by: Jason Marshall --- CHANGELOG.md | 2 + example/cluster.js | 1 + index.d.ts | 1 - lib/cluster.js | 198 +++++++++++++++++++++++++++----------------- lib/worker.js | 145 ++++++++++++++++++-------------- test/clusterTest.js | 30 ++++--- test/workerTest.js | 46 +++++++--- 7 files changed, 260 insertions(+), 163 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4507479f..acbe5d74 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,6 +20,7 @@ This release marks our first release under the Prometheus umbrella. runtime object. Value-style uses such as `MetricType.Counter` (which threw at runtime) no longer compile; compare against the string literals instead. Under `verbatimModuleSyntax`, import it with `import type`. +- The cluster primary now reports metrics ### Changed @@ -46,6 +47,7 @@ This release marks our first release under the Prometheus umbrella. - chore: Add copyright license headers and test - Make cluster and worker-thread metric aggregation order deterministic - Export `MetricObject`, `MetricObjectWithValues`, `MetricValue` and `MetricValueWithName` from the TypeScript definitions +- Improve cluster support to allow workers to opt out ### Added diff --git a/example/cluster.js b/example/cluster.js index 3472cdb3..40b77421 100644 --- a/example/cluster.js +++ b/example/cluster.js @@ -36,6 +36,7 @@ if (cluster.isPrimary) { res.set('Content-Type', clusterRegistry.contentType); res.send(metrics); } catch (ex) { + console.error(ex); res.statusCode = 500; res.send(ex.message); } diff --git a/index.d.ts b/index.d.ts index e748a5d1..7712e2c0 100644 --- a/index.d.ts +++ b/index.d.ts @@ -206,7 +206,6 @@ export class WorkerRegistry extends Registry { */ workerMetrics(): Promise; - addWorker(worker: Worker): void; /** * Sets the registry or registries to be aggregated. Call from workers to * use a registry/registries other than the default global registry. diff --git a/lib/cluster.js b/lib/cluster.js index 0c9a84a4..cdb7315f 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -39,7 +39,9 @@ const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes'; let registries = [Registry.globalRegistry]; let requestCtr = 0; // Concurrency control +let listenersAdded = false; const requests = new Map(); // Pending requests for workers' local metrics. +const workers = new Map(); class AggregatorRegistry extends Registry { /** @@ -49,54 +51,9 @@ class AggregatorRegistry extends Registry { constructor(regContentType = Registry.PROMETHEUS_CONTENT_TYPE) { super(regContentType); - if (cluster().isPrimary) { - this.workers = new Map(); - } - - addListeners(this); + addListeners(); } - addWorker(id) { - if (this.workers.has(id)) { - debug('duplicate worker announcement', id); - return; - } - - const worker = cluster().workers[id]; - - worker.on('disconnect', () => { - debug('worker disconnected', id); - this.workers.delete(id); - }); - - worker.on('message', message => { - if (message.type === GET_METRICS_RES) { - const request = requests.get(message.requestId); - - if (request === undefined) { - debug('unexpected results from worker', id); - return; - } - - const response = request.responseHandlers.get(id); - if (response === undefined) { - return; - } - request.responseHandlers.delete(id); - - if (message.error) { - response.reject(new Error(message.error)); - } else { - response.resolve({ - threadId: id, - metrics: message.metrics, - }); - } - } - }); - - this.workers.set(id, worker); - } /** * Gets aggregated metrics for all workers. The optional callback and * returned Promise resolve with the same value; either may be used. @@ -105,7 +62,7 @@ class AggregatorRegistry extends Registry { */ clusterMetrics() { const requestId = requestCtr++; - const orderedWorkers = [...this.workers.values()].sort( + const orderedWorkers = [...workers.values()].sort( (left, right) => left.id - right.id, ); @@ -216,46 +173,108 @@ class AggregatorRegistry extends Registry { * than once). * @returns {void} */ -function addListeners(registry) { +function addListeners() { + if (listenersAdded) { + return; + } + + listenersAdded = true; + if (cluster().isPrimary) { - // Listen for worker responses to requests for local metrics - cluster().on('message', (worker, message) => { - if (message.type === ANNOUNCEMENT) { - registry.addWorker(worker.id); - } - }); + replaceListener('message', cluster(), primaryListener); + replaceListener('disconnect', cluster(), disconnect); announce(); } else { - // Respond to master's requests for worker's local metrics. - process.on('message', async message => { - if (message.type === ANNOUNCEMENT) { - process.send({ type: ANNOUNCEMENT }); - } else if (message.type === GET_METRICS_REQ) { - const metrics = await Promise.all( - registries.map(r => r.getMetricsAsJSON()), - ); - - try { - process.send({ - type: GET_METRICS_RES, - requestId: message.requestId, - metrics, - }); - } catch (error) { - process.send({ - type: GET_METRICS_RES, - requestId: message.requestId, - error: error.message, - }); - } - } - }); + replaceListener('message', process, workerListener); + + if (typeof process.send !== 'function') { + debug('worker has no process.send()'); + } else if (!process.connected) { + debug('worker is not connected to parent process'); + } else { + process.send({ type: ANNOUNCEMENT }); + } + } +} +/** + * Watch for metrics events and aggregator announcements + * + * Whereas clusters are a top-level activity, multiple modules may start their + * own workers and require telemetry collection. + * @param message {MessageEvent} + */ +async function workerListener(message) { + if (message.type === ANNOUNCEMENT) { process.send({ type: ANNOUNCEMENT }); + } else if (message.type === GET_METRICS_REQ) { + try { + const metrics = await Promise.all( + registries.map(r => r.getMetricsAsJSON()), + ); + + process.send({ + type: GET_METRICS_RES, + requestId: message.requestId, + metrics, + }); + } catch (error) { + process.send({ + type: GET_METRICS_RES, + requestId: message.requestId, + error: error.message, + }); + } } } +/** + * Add workers to the aggregation list when they are announced. + * + * Whereas clusters are a top-level activity, multiple modules may start their + * own workers and require telemetry collection. + * @param event {MessageEvent} + */ + +async function primaryListener(worker, event) { + if (event.type === ANNOUNCEMENT) { + if (workers.has(worker.id)) { + debug('duplicate worker announcement', worker.id); + return; + } + + workers.set(worker.id, worker); + } else if (event.type === GET_METRICS_RES) { + const request = requests.get(event.requestId); + + if (request === undefined) { + debug('unexpected results from worker', worker.id); + return; + } + + const response = request.responseHandlers.get(worker.id); + if (response === undefined) { + return; + } + request.responseHandlers.delete(worker.id); + + if (event.error) { + response.reject(new Error(event.error)); + } else { + response.resolve({ + threadId: worker.id, + metrics: event.metrics, + }); + } + } +} + +function disconnect(event) { + debug('worker disconnected', event.id); + workers.delete(event.id); +} + function announce() { for (const worker of Object.values(cluster().workers)) { if (worker.isConnected()) { @@ -264,4 +283,27 @@ function announce() { } } +/** + * Replace any listeners with new ones. + * + * @param messageType + * @param emitter {EventEmitter} + * @param fn + */ +function replaceListener(messageType, emitter, fn) { + // Reloading a module creates a unique instance of each function, so the + // identity checks is cluster.off() will fail. + const functionString = fn.toString(); + + for (const listener of emitter.listeners(messageType)) { + // eslint-disable-next-line eqeqeq + if (functionString == listener) { + debug('removing duplicate listener', messageType); + emitter.off(messageType, listener); + } + } + + emitter.on(messageType, fn); +} + module.exports = AggregatorRegistry; diff --git a/lib/worker.js b/lib/worker.js index 081094b6..6b524d3a 100644 --- a/lib/worker.js +++ b/lib/worker.js @@ -36,8 +36,10 @@ const ANNOUNCEMENT_CHANNEL = new BroadcastChannel( ).unref(); let registries = [Registry.globalRegistry]; +let listenersAdded = false; let requestCtr = 0; // Concurrency control const requests = new Map(); // Pending requests for workers' local metrics. +const workers = new Map(); class WorkerRegistry extends Registry { /** @@ -54,59 +56,7 @@ class WorkerRegistry extends Registry { super(regContentType); this.primary = primary; - if (this.primary) { - this.channels = new Map(); - } - - addListeners(this); - } - - /** - * Add a worker to the aggregation list. - * Whereas clusters are a top-level activity, multiple modules may start their - * own workers and require telemetry collection. - * @param name {string} - */ - addWorker(name) { - if (this.channels.has(name)) { - debug('duplicate worker announcement', name); - return; - } - - const channel = new BroadcastChannel(name).unref(); - channel.addEventListener('close', () => { - this.channels.delete(name); - }); - - channel.addEventListener('message', event => { - const message = event.data; - - if (message.type === GET_METRICS_RES) { - const request = requests.get(message.requestId); - - if (request === undefined) { - debug('unexpected results from worker', name); - return; - } - - const response = request.responseHandlers.get(name); - if (response === undefined) { - return; - } - request.responseHandlers.delete(name); - - if (message.error) { - response.reject(new Error(message.error)); - } else { - response.resolve({ - threadId: message.threadId, - metrics: message.metrics, - }); - } - } - }); - - this.channels.set(name, channel); + addListeners(primary); } /** @@ -147,10 +97,15 @@ class WorkerRegistry extends Registry { }, 5_000), }; requests.set(requestId, request); - const responsePromises = [...this.channels.keys()].sort().map( - name => + + const orderedWorkers = [...workers.values()].sort( + (left, right) => left.threadId - right.threadId, + ); + + const responsePromises = orderedWorkers.map( + entry => new Promise((resolveResponse, rejectResponse) => { - responseHandlers.set(name, { + responseHandlers.set(entry.name, { resolve: resolveResponse, reject: rejectResponse, }); @@ -218,7 +173,17 @@ class WorkerRegistry extends Registry { /** * Watch for metrics collection events. */ -function addListeners(registry) { +function addListeners(primary) { + if (listenersAdded) { + return; + } + + listenersAdded = true; + + if (primary) { + ANNOUNCEMENT_CHANNEL.addEventListener('message', primaryListener); + } + const name = `@prometheus-io/client:worker:${threadId}`; const channel = new BroadcastChannel(name).unref(); @@ -226,9 +191,7 @@ function addListeners(registry) { const message = event.data; if (message.type === ANNOUNCEMENT) { - if (registry.primary) { - registry.addWorker(message.name); - } else if (message.primary) { + if (message.primary) { announce(name, false); } } else if (message.type === GET_METRICS_REQ) { @@ -253,7 +216,67 @@ function addListeners(registry) { } }); - announce(name, registry.primary); + announce(name, primary); +} + +/** + * Add workers to the aggregation list when they are announced. + * + * Whereas clusters are a top-level activity, multiple modules may start their + * own workers and require telemetry collection. + * @param event {MessageEvent} + */ + +async function primaryListener(event) { + const message = event.data; + + if (message.type === ANNOUNCEMENT) { + const workerName = message.name; + + if (workers.has(workerName)) { + debug('duplicate worker announcement', workerName); + return; + } + + const workerChannel = new BroadcastChannel(workerName, {}).unref(); + workers.set(workerName, { + name: workerName, + channel: workerChannel, + threadId: message.threadId, + }); + + workerChannel.addEventListener('close', () => { + workers.delete(workerName); + }); + + workerChannel.addEventListener('message', workerEvent => { + const workerMessage = workerEvent.data; + + if (workerMessage.type === GET_METRICS_RES) { + const request = requests.get(workerMessage.requestId); + + if (request === undefined) { + debug('unexpected results from worker', workerName); + return; + } + + const response = request.responseHandlers.get(workerName); + if (response === undefined) { + return; + } + request.responseHandlers.delete(workerName); + + if (workerMessage.error) { + response.reject(new Error(workerMessage.error)); + } else { + response.resolve({ + threadId: workerMessage.threadId, + metrics: workerMessage.metrics, + }); + } + } + }); + } } function announce(name, primary) { diff --git a/test/clusterTest.js b/test/clusterTest.js index fd532ec6..96c1424f 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -75,48 +75,56 @@ describe.each([ expect(metrics.trim()).toEqual(''); }); - it('aggregates worker responses in worker id order', async () => { - jest.resetModules(); + it("listeners don't accumulate", () => { + for (let i = 0; i < 30; i++) { + jest.resetModules(); + + const AggregatorRegistry = require('../lib/cluster'); + const ar = new AggregatorRegistry(regType); + } + }); - const registry = new Registry(regType); + it('aggregates worker responses in worker id order', async () => { const originalWorkers = cluster.workers; + jest.resetModules(); + const AggregatorRegistry = require('../lib/cluster'); + const registry = new AggregatorRegistry(regType); const workers = Object.fromEntries( [1, 2, 3].map(id => [ id, { id, isConnected: () => true, - on: jest.fn(), send: jest.fn(), }, ]), ); cluster.workers = workers; - Object.keys(workers).forEach(id => { - cluster.emit('message', workers[id], { type: ANNOUNCEMENT }); + Object.values(workers).forEach(worker => { + cluster.emit('message', worker, { type: ANNOUNCEMENT }); }); try { const result = registry.clusterMetrics(); - const calls = workers[1].send.mock.calls; - const requestId = calls[0][0].requestId; for (const [id, value] of [ [3, 0.3437699], [1, 0.5848208], [2, 0.5479198], ]) { - const listener = workers[id].on.mock.calls.at(-1)[1]; - listener({ + cluster.emit('message', workers[id], { type: GET_METRICS_RES, - requestId, + requestId: 0, metrics: [[metric(value)]], }); } await expect(result).resolves.toContain('test_metric 1.4765105'); } finally { + Object.values(workers).forEach(worker => { + cluster.emit('disconnect', worker); + }); cluster.workers = originalWorkers; } }, 6_000); diff --git a/test/workerTest.js b/test/workerTest.js index f61c85a4..9a3e420c 100644 --- a/test/workerTest.js +++ b/test/workerTest.js @@ -14,11 +14,11 @@ 'use strict'; -const { EventEmitter } = require('events'); const { setTimeout: delay } = require('timers/promises'); const { BroadcastChannel } = require('worker_threads'); const Registry = require('../lib/worker'); +const ANNOUNCEMENT = '@prometheus-io/client:announcement'; const GET_METRICS_REQ = '@prometheus-io/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes'; @@ -56,13 +56,19 @@ describe.each([ ); const responders = [1, 2, 3].map(threadId => { const name = `@prometheus-io/client:test-worker:${threadId}`; - registry.addWorker(name); - return { + const channel = new BroadcastChannel(name).unref(); + + announcementChannel.postMessage({ + type: ANNOUNCEMENT, + name, threadId, - channel: new BroadcastChannel(name), - }; + }); + + return { threadId, channel }; }); + await delay(5); // Let announcements arrive + let finishSendingResponses; const responsesSent = new Promise(resolve => { finishSendingResponses = resolve; @@ -93,21 +99,37 @@ describe.each([ } finally { announcementChannel.close(); for (const responder of responders) responder.channel.close(); - for (const channel of registry.channels.values()) channel.close(); } }); }); describe('message handling', () => { + it("listeners don't accumulate", () => { + for (let i = 0; i < 30; i++) { + jest.resetModules(); + + const AggregatorRegistry = require('../lib/worker'); + const ar = new AggregatorRegistry(regType); + } + }); + it('does not error out on unexpected (or late) responses', () => { jest.resetModules(); const WorkerRegistry = require('../lib/worker'); - const registry = new WorkerRegistry(regType); - const emitter = new EventEmitter(); - - registry.addWorker(emitter); + const announcementChannel = new BroadcastChannel( + '@prometheus-io/client:announce', + ); + const threadId = 20; + const name = `@prometheus-io/client:test-worker:${threadId}`; + const channel = new BroadcastChannel(name).unref(); + + announcementChannel.postMessage({ + type: ANNOUNCEMENT, + name, + threadId, + }); //Emulate a response that has been deleted from requests const unexpected = { @@ -117,9 +139,9 @@ describe.each([ }; try { - expect(() => emitter.emit('message', unexpected)).not.toThrow(); + expect(() => channel.postMessage(unexpected)).not.toThrow(); } finally { - for (const channel of registry.channels.values()) channel.close(); + channel.close(); } }); }); From 50b96f00000d8820b1560d20b40f336c97f462ab Mon Sep 17 00:00:00 2001 From: alencristen <299997878+alencristen@users.noreply.github.com> Date: Wed, 22 Jul 2026 07:00:18 -0400 Subject: [PATCH 5/6] fix(cluster): skip responses after IPC disconnect Signed-off-by: alencristen <299997878+alencristen@users.noreply.github.com> --- CHANGELOG.md | 1 + lib/cluster.js | 29 ++++++++++++++-------- test/clusterTest.js | 59 +++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 79 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index acbe5d74..577a9e7d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -48,6 +48,7 @@ This release marks our first release under the Prometheus umbrella. - Make cluster and worker-thread metric aggregation order deterministic - Export `MetricObject`, `MetricObjectWithValues`, `MetricValue` and `MetricValueWithName` from the TypeScript definitions - Improve cluster support to allow workers to opt out +- Abort cluster metric responses during process termination ### Added diff --git a/lib/cluster.js b/lib/cluster.js index cdb7315f..92159c2d 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -214,17 +214,26 @@ async function workerListener(message) { registries.map(r => r.getMetricsAsJSON()), ); - process.send({ - type: GET_METRICS_RES, - requestId: message.requestId, - metrics, - }); + if (!process.connected) { + debug('Connection to primary lost.'); + } else { + process.send({ + type: GET_METRICS_RES, + requestId: message.requestId, + metrics, + }); + } } catch (error) { - process.send({ - type: GET_METRICS_RES, - requestId: message.requestId, - error: error.message, - }); + debug('Error sending to primary', error); + if (!process.connected) { + debug('Connection to primary lost.'); + } else { + process.send({ + type: GET_METRICS_RES, + requestId: message.requestId, + error: error.message, + }); + } } } } diff --git a/test/clusterTest.js b/test/clusterTest.js index 96c1424f..323ff8dd 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -19,6 +19,7 @@ const process = require('process'); const Registry = require('../lib/cluster'); const ANNOUNCEMENT = '@prometheus-io/client:announcement'; +const GET_METRICS_REQ = '@prometheus-io/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes'; function metric(value) { @@ -168,3 +169,61 @@ describe.each([ }); }); }); + +describe('worker message handling', () => { + it('does not send metrics after the IPC channel disconnects', async () => { + jest.resetModules(); + jest.doMock('cluster', () => { + return { isPrimary: false }; + }); + + const messageListeners = new Set(process.listeners('message')); + const connectedDescriptor = Object.getOwnPropertyDescriptor( + process, + 'connected', + ); + const sendDescriptor = Object.getOwnPropertyDescriptor(process, 'send'); + const send = jest.fn(); + let listener; + + try { + Object.defineProperty(process, 'connected', { + configurable: true, + value: true, + writable: true, + }); + Object.defineProperty(process, 'send', { + configurable: true, + value: send, + }); + + const AggregatorRegistry = require('../lib/cluster'); + new AggregatorRegistry(); + + listener = process + .listeners('message') + .find(candidate => !messageListeners.has(candidate)); + expect(listener).toBeDefined(); + + listener({ type: GET_METRICS_REQ, requestId: 1 }); + process.connected = false; + await new Promise(resolve => setImmediate(resolve)); + + expect(send).toHaveBeenCalledTimes(1); // Announcement + } finally { + if (listener) process.removeListener('message', listener); + if (connectedDescriptor) { + Object.defineProperty(process, 'connected', connectedDescriptor); + } else { + delete process.connected; + } + if (sendDescriptor) { + Object.defineProperty(process, 'send', sendDescriptor); + } else { + delete process.send; + } + jest.dontMock('cluster'); + jest.resetModules(); + } + }); +}); From 1adc9d526562f370157c71cbda58d17a6912129f Mon Sep 17 00:00:00 2001 From: Jason Marshall Date: Tue, 28 Jul 2026 09:34:41 -0700 Subject: [PATCH 6/6] Simplify process.send sanity checks. 99.9% of the time process.send() is going to work. We don't need to guard it when it's already inside of a try block. Just guard the retry send. Also reduces the amount of excessive mocking going on in the tests by using jest more instead of creating our own mocks. Signed-off-by: Jason Marshall --- test/clusterTest.js | 42 +++++++++++++++++++----------------------- 1 file changed, 19 insertions(+), 23 deletions(-) diff --git a/test/clusterTest.js b/test/clusterTest.js index 323ff8dd..2b41c9a3 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -171,19 +171,24 @@ describe.each([ }); describe('worker message handling', () => { - it('does not send metrics after the IPC channel disconnects', async () => { + beforeEach(() => { jest.resetModules(); + }); + + it('does not send metrics after the IPC channel disconnects', async () => { jest.doMock('cluster', () => { return { isPrimary: false }; }); - const messageListeners = new Set(process.listeners('message')); const connectedDescriptor = Object.getOwnPropertyDescriptor( process, 'connected', ); - const sendDescriptor = Object.getOwnPropertyDescriptor(process, 'send'); - const send = jest.fn(); + + const AggregatorRegistry = require('../lib/cluster'); + new AggregatorRegistry(); + + const send = jest.spyOn(process, 'send'); let listener; try { @@ -192,38 +197,29 @@ describe('worker message handling', () => { value: true, writable: true, }); - Object.defineProperty(process, 'send', { - configurable: true, - value: send, - }); - const AggregatorRegistry = require('../lib/cluster'); - new AggregatorRegistry(); - - listener = process - .listeners('message') - .find(candidate => !messageListeners.has(candidate)); + listener = process.listeners('message').at(-1); expect(listener).toBeDefined(); - listener({ type: GET_METRICS_REQ, requestId: 1 }); + send.mockImplementationOnce(() => { + throw new Error('disconnected'); + }); + process.connected = false; + + listener({ type: GET_METRICS_REQ, requestId: 1 }); await new Promise(resolve => setImmediate(resolve)); - expect(send).toHaveBeenCalledTimes(1); // Announcement + expect(send).not.toHaveBeenCalled(); } finally { - if (listener) process.removeListener('message', listener); + process.removeListener('message', listener); if (connectedDescriptor) { Object.defineProperty(process, 'connected', connectedDescriptor); } else { delete process.connected; } - if (sendDescriptor) { - Object.defineProperty(process, 'send', sendDescriptor); - } else { - delete process.send; - } - jest.dontMock('cluster'); jest.resetModules(); + jest.clearAllMocks(); } }); });