From 8a6941ae210c9e11e717c8e126cf8ab325e42091 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. --- 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 4e6f08b8..3e0d77ac 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/client:announcement'; const GET_METRICS_REQ = '@prometheus/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus/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 8f4e3b3a..b5e5975b 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/client:announcement'; const GET_METRICS_REQ = '@prometheus/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus/client:getMetricsRes'; const ANNOUNCEMENT_CHANNEL = new BroadcastChannel( '@prometheus/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/client:worker:${threadId}`; const channel = new BroadcastChannel(name); diff --git a/test/clusterTest.js b/test/clusterTest.js index 70dcd081..25b7c0e5 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/client:announcement'; const GET_METRICS_RES = '@prometheus/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 0bd94ae42dbe0225a2d73e769864b109b78661c6 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. --- lib/worker.js | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/lib/worker.js b/lib/worker.js index b5e5975b..0ace811f 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/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 3c366f7fabf05ce3efd9e53b14a12983feba5fad 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 --- example/cluster.js | 4 ++++ lib/cluster.js | 13 ++++++++++--- test/clusterTest.js | 20 +++++++++++++++++++- 3 files changed, 33 insertions(+), 4 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 3e0d77ac..0f95dd51 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,11 +152,18 @@ class AggregatorRegistry extends Registry { }), ); - if (responsePromises.length === 0) { + const myMetrics = Promise.all( + registries.map(r => r.getMetricsAsJSON()), + ).then(metrics => { + return { metrics }; + }); + + const allMetrics = [myMetrics, ...workerMetrics]; + if (allMetrics.length === 0) { debug('No workers found for requestId', requestId); process.nextTick(() => done(undefined, '')); } else { - Promise.all(responsePromises) + 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 25b7c0e5..7dbbcc29 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,24 @@ 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' }); + 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'); + + gauge.remove(); + }); }); describe('message handling', () => { From c0b814032e77424fe42f56e80c292c79183d5f88 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 | 1 + example/cluster.js | 1 + lib/cluster.js | 164 ++++++++++++++++++++------------------ lib/worker.js | 188 +++++++++++++++++++++++++------------------- test/clusterTest.js | 23 ++++-- test/workerTest.js | 47 ++++++++--- 6 files changed, 244 insertions(+), 180 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 11be4a31..f9dc096c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -41,6 +41,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 primary to report metrics and workers to opt out ### Added diff --git a/example/cluster.js b/example/cluster.js index 3472cdb3..acd35258 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.log(ex); res.statusCode = 500; res.send(ex.message); } diff --git a/lib/cluster.js b/lib/cluster.js index 0f95dd51..3ea20d1f 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -39,7 +39,9 @@ const GET_METRICS_RES = '@prometheus/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,43 +173,96 @@ 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); - } + cluster().on('message', primaryListener); + cluster().on('disconnect', deadWorker => { + debug('worker disconnected', deadWorker.id); + workers.delete(deadWorker.id); }); 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, - }); - } - } - }); + process.on('message', workerListener); + 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) { + 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, + }); + } + } +} + +/** + * 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, + }); + } } } diff --git a/lib/worker.js b/lib/worker.js index 0ace811f..f296b446 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 channels = 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,7 +97,7 @@ class WorkerRegistry extends Registry { }, 5_000), }; requests.set(requestId, request); - const responsePromises = [...this.channels.keys()].sort().map( + const responsePromises = [...channels.keys()].sort().map( name => new Promise((resolveResponse, rejectResponse) => { responseHandlers.set(name, { @@ -218,42 +168,114 @@ class WorkerRegistry extends Registry { /** * Watch for metrics collection events. */ -function addListeners(registry) { +function addListeners(primary) { + if (listenersAdded) { + return; + } + + listenersAdded = true; + + const name = `@prometheus/client:worker:${threadId}`; + if (primary) { + ANNOUNCEMENT_CHANNEL.addEventListener('message', primaryListener); + } + + ANNOUNCEMENT_CHANNEL.addEventListener('message', workerListener); + announce(name, primary); +} + +/** + * 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 event {MessageEvent} + */ +async function workerListener(event) { const name = `@prometheus/client:worker:${threadId}`; const channel = new BroadcastChannel(name).unref(); + const message = event.data; - ANNOUNCEMENT_CHANNEL.addEventListener('message', async event => { - const message = event.data; + if (message.type === ANNOUNCEMENT) { + if (message.primary) { + announce(name, false); + } + } else if (message.type === GET_METRICS_REQ) { + const metrics = await Promise.all( + registries.map(r => r.getMetricsAsJSON()), + ); + + try { + channel.postMessage({ + type: GET_METRICS_RES, + requestId: message.requestId, + threadId, + metrics, + }); + } catch (error) { + channel.postMessage({ + type: GET_METRICS_RES, + requestId: message.requestId, + error: error.message, + }); + } + } +} - if (message.type === ANNOUNCEMENT) { - if (registry.primary) { - registry.addWorker(message.name); - } else if (message.primary) { - announce(name, false); - } - } else if (message.type === GET_METRICS_REQ) { - const metrics = await Promise.all( - registries.map(r => r.getMetricsAsJSON()), - ); +/** + * 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} + */ - try { - channel.postMessage({ - type: GET_METRICS_RES, - requestId: message.requestId, - threadId, - metrics, - }); - } catch (error) { - channel.postMessage({ - type: GET_METRICS_RES, - requestId: message.requestId, - error: error.message, - }); - } +async function primaryListener(event) { + const message = event.data; + + if (message.type === ANNOUNCEMENT) { + const name = message.name; + + if (channels.has(name)) { + debug('duplicate worker announcement', name); + return; } - }); - announce(name, registry.primary); + const channel = new BroadcastChannel(name, {}).unref(); + channels.set(name, channel); + + channel.addEventListener('close', () => { + 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, + }); + } + } + }); + } } function announce(name, primary) { diff --git a/test/clusterTest.js b/test/clusterTest.js index 7dbbcc29..962c385e 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -75,26 +75,32 @@ describe.each([ expect(metrics.trim()).toEqual(''); }); + it("listeners don't accumulate", () => { + const AggregatorRegistry = require('../lib/cluster'); + for (let i = 0; i < 30; i++) { + jest.resetModules(); + + const ar = new AggregatorRegistry(regType); + } + }); + 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 => [ 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 { @@ -107,8 +113,7 @@ describe.each([ [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, metrics: [[metric(value)]], @@ -117,7 +122,9 @@ describe.each([ await expect(result).resolves.toContain('test_metric 1.4765105'); } finally { - cluster.workers = originalWorkers; + Object.values(workers).forEach(worker => { + cluster.emit('disconnect', worker); + }); } }, 6_000); diff --git a/test/workerTest.js b/test/workerTest.js index 76c3ad45..9ae46ce3 100644 --- a/test/workerTest.js +++ b/test/workerTest.js @@ -19,6 +19,7 @@ const { setTimeout: delay } = require('timers/promises'); const { BroadcastChannel } = require('worker_threads'); const Registry = require('../lib/worker'); +const ANNOUNCEMENT = '@prometheus/client:announcement'; const GET_METRICS_REQ = '@prometheus/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus/client:getMetricsRes'; @@ -53,16 +54,22 @@ describe.each([ const registry = new Registry(regType); const announcementChannel = new BroadcastChannel( '@prometheus/client:announce', - ); + ).unref(); const responders = [1, 2, 3].map(threadId => { const name = `@prometheus/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,7 +100,6 @@ describe.each([ } finally { announcementChannel.close(); for (const responder of responders) responder.channel.close(); - for (const channel of registry.channels.values()) channel.close(); } }); }); @@ -103,11 +109,19 @@ describe.each([ jest.resetModules(); const WorkerRegistry = require('../lib/worker'); - const registry = new WorkerRegistry(regType); - const emitter = new EventEmitter(); - - registry.addWorker(emitter); + const announcementChannel = new BroadcastChannel( + '@prometheus/client:announce', + ).unref(); + const threadId = 20; + const name = `@prometheus/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 +131,18 @@ 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(); + } + }); + + it("listeners don't accumulate", () => { + const AggregatorRegistry = require('../lib/worker'); + for (let i = 0; i < 30; i++) { + jest.resetModules(); + + const ar = new AggregatorRegistry(regType); } }); }); From efa76912140ff824795dd4b7230355f486c4669e 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 f9dc096c..c51136e6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -42,6 +42,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 primary to report metrics and workers to opt out +- Abort cluster metric responses during process termination ### Added diff --git a/lib/cluster.js b/lib/cluster.js index 3ea20d1f..eadcbb18 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -210,17 +210,26 @@ async function workerListener(message) { ); try { - 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 962c385e..73e1471e 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/client:announcement'; +const GET_METRICS_REQ = '@prometheus/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus/client:getMetricsRes'; function metric(value) { @@ -164,3 +165,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 3bc15bb061bc36de16dcc9d5d9ceb096afcef788 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 --- lib/cluster.js | 14 +++++--------- test/clusterTest.js | 37 +++++++++++++++++-------------------- 2 files changed, 22 insertions(+), 29 deletions(-) diff --git a/lib/cluster.js b/lib/cluster.js index eadcbb18..4a5f023d 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -210,15 +210,11 @@ async function workerListener(message) { ); try { - if (!process.connected) { - debug('Connection to primary lost.'); - } else { - process.send({ - type: GET_METRICS_RES, - requestId: message.requestId, - metrics, - }); - } + process.send({ + type: GET_METRICS_RES, + requestId: message.requestId, + metrics, + }); } catch (error) { debug('Error sending to primary', error); if (!process.connected) { diff --git a/test/clusterTest.js b/test/clusterTest.js index 73e1471e..d5890e43 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -167,19 +167,20 @@ 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 send = jest.spyOn(process, 'send'); let listener; try { @@ -188,38 +189,34 @@ 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 }); + expect(send).toHaveBeenCalledTimes(1); + send.mockReset(); + + 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).toHaveBeenCalledTimes(1); } 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(); } }); });