diff --git a/CHANGELOG.md b/CHANGELOG.md index 11be4a31..c51136e6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -41,6 +41,8 @@ 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 +- Abort cluster metric responses during process termination ### Added diff --git a/example/cluster.js b/example/cluster.js index 2686a27e..acd35258 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 }); } @@ -32,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 4e6f08b8..4a5f023d 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,6 +32,8 @@ 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'; @@ -38,10 +41,16 @@ 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 { + /** + * Create a Registry. + * @param regContentType + */ constructor(regContentType = Registry.PROMETHEUS_CONTENT_TYPE) { super(regContentType); + addListeners(); } @@ -53,9 +62,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 = [...workers.values()].sort( + (left, right) => left.id - right.id, + ); return new Promise((resolve, reject) => { let settled = false; @@ -78,37 +87,44 @@ 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 workerMetrics = 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); + 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(allMetrics) + .then(responses => responses.flatMap(response => response.metrics)) + .then(metrics => Registry.aggregate(metrics).metrics()) + .then(result => done(undefined, result), done); + } }); } @@ -158,53 +174,108 @@ class AggregatorRegistry extends Registry { * @returns {void} */ function addListeners() { - if (listenersAdded) return; + 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 === GET_METRICS_RES) { - const request = requests.get(message.requestId); + cluster().on('message', primaryListener); + cluster().on('disconnect', deadWorker => { + debug('worker disconnected', deadWorker.id); + workers.delete(deadWorker.id); + }); - if (request === undefined) { - return; - } + announce(); + } else { + process.on('message', workerListener); + process.send({ type: ANNOUNCEMENT }); + } +} - const response = request.responseHandlers.get(worker.id); - if (response === undefined) { - return; - } - request.responseHandlers.delete(worker.id); +/** + * 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()), + ); - if (message.error) { - response.reject(new Error(message.error)); - } else { - response.resolve(message.metrics); - } - } - }); - } 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, - }); - }); + try { + process.send({ + type: GET_METRICS_RES, + requestId: message.requestId, + metrics, + }); + } catch (error) { + 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, + }); } - }); + } + } +} + +/** + * 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 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..f296b446 100644 --- a/lib/worker.js +++ b/lib/worker.js @@ -22,23 +22,24 @@ * 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; +let requestCtr = 0; // Concurrency control const requests = new Map(); // Pending requests for workers' local metrics. +const channels = new Map(); class WorkerRegistry extends Registry { /** @@ -55,57 +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)) { - return; - } - - const channel = new BroadcastChannel(name); - 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) { - 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); } /** @@ -115,8 +66,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 +97,7 @@ class WorkerRegistry extends Registry { }, 5_000), }; requests.set(requestId, request); - const responsePromises = [...this.channels.keys()].map( + const responsePromises = [...channels.keys()].sort().map( name => new Promise((resolveResponse, rejectResponse) => { responseHandlers.set(name, { @@ -163,19 +114,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); }); } @@ -222,7 +168,7 @@ class WorkerRegistry extends Registry { /** * Watch for metrics collection events. */ -function addListeners(registry) { +function addListeners(primary) { if (listenersAdded) { return; } @@ -230,42 +176,106 @@ function addListeners(registry) { listenersAdded = true; const name = `@prometheus/client:worker:${threadId}`; - const channel = new BroadcastChannel(name); + 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; - channel.unref(); + 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, + }); + } + } +} - ANNOUNCEMENT_CHANNEL.addEventListener('message', async event => { - const message = event.data; +/** + * 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} + */ - 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()), - ); +async function primaryListener(event) { + const message = event.data; - 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) { + 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 70dcd081..d5890e43 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -18,6 +18,8 @@ const cluster = require('cluster'); 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) { @@ -71,11 +73,21 @@ 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("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 () => { - const originalWorkers = cluster.workers; + jest.resetModules(); + const registry = new Registry(regType); const workers = Object.fromEntries( [1, 2, 3].map(id => [ id, @@ -88,10 +100,14 @@ describe.each([ ); cluster.workers = workers; + Object.values(workers).forEach(worker => { + cluster.emit('message', worker, { 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], @@ -107,8 +123,28 @@ 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); + + 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(); }); }); @@ -129,3 +165,58 @@ describe.each([ }); }); }); + +describe('worker message handling', () => { + beforeEach(() => { + jest.resetModules(); + }); + + it('does not send metrics after the IPC channel disconnects', async () => { + jest.doMock('cluster', () => { + return { isPrimary: false }; + }); + + const connectedDescriptor = Object.getOwnPropertyDescriptor( + process, + 'connected', + ); + const send = jest.spyOn(process, 'send'); + let listener; + + try { + Object.defineProperty(process, 'connected', { + configurable: true, + value: true, + writable: true, + }); + + const AggregatorRegistry = require('../lib/cluster'); + new AggregatorRegistry(); + + listener = process.listeners('message').at(-1); + expect(listener).toBeDefined(); + + 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); + } finally { + process.removeListener('message', listener); + if (connectedDescriptor) { + Object.defineProperty(process, 'connected', connectedDescriptor); + } else { + delete process.connected; + } + jest.resetModules(); + jest.clearAllMocks(); + } + }); +}); 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); } }); });