diff --git a/lib/cluster.js b/lib/cluster.js index 84118c36..f9ce4382 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -10,20 +10,31 @@ const cluster = require('cluster'); const Registry = require('./registry'); +const Gauge = require('./gauge'); const util = require('./util'); const aggregators = require('./metricAggregators').aggregators; const GET_METRICS_REQ = 'prom-client:getMetricsReq'; const GET_METRICS_RES = 'prom-client:getMetricsRes'; +const REG_METRICS_WORKER = 'prom-client:registerMetricsWorker'; let registries = [Registry.globalRegistry]; let requestCtr = 0; // Concurrency control let listenersAdded = false; +const coordinatedWorkers = new Set(); const requests = new Map(); // Pending requests for workers' local metrics. class AggregatorRegistry extends Registry { - constructor() { + /** + Create an AggregatorRegistry instance. Accepts an optional `options` object. + + Options are: + coordinated If false (default), request metrics from all cluster workers. If true, request metrics only from workers that have required prom-client. + * @param {object?} options object + */ + constructor({ coordinated } = {}) { super(); + this.coordinated = coordinated || false; addListeners(); } @@ -32,13 +43,15 @@ class AggregatorRegistry extends Registry { * returned Promise resolve with the same value; either may be used. * @param {Function?} callback (err, metrics) => any * @return {Promise} Promise that resolves with the aggregated - * metrics. + * metrics. */ clusterMetrics(callback) { const requestId = requestCtr++; return new Promise((resolve, reject) => { - const nWorkers = Object.keys(cluster.workers).length; + const nWorkers = this.coordinated + ? coordinatedWorkers.size + : Object.keys(cluster.workers).length; function done(err, result) { // Don't resolve/reject the promise if a callback is provided @@ -56,6 +69,7 @@ class AggregatorRegistry extends Registry { const request = { responses: [], + workerCount: nWorkers, pending: nWorkers, done, errorTimeout: setTimeout(() => { @@ -71,7 +85,19 @@ class AggregatorRegistry extends Registry { type: GET_METRICS_REQ, requestId }; - for (const id in cluster.workers) cluster.workers[id].send(message); + const workers = this.coordinated + ? coordinatedWorkers + : Object.keys(cluster.workers); + for (const id of workers) { + const worker = cluster.workers[id]; + if (worker === undefined || worker.isDead() || !worker.isConnected()) { + request.pending--; + request.workerCount--; + coordinatedWorkers.delete(id); // Set is safe to mutate while iterating + } else { + worker.send(message); + } + } }); } @@ -81,7 +107,7 @@ class AggregatorRegistry extends Registry { * the method specified by their `aggregator` property, or by summation if * `aggregator` is undefined. * @param {Array} metricsArr Array of metrics, each of which created by - * `registry.getMetricsAsJSON()`. + * `registry.getMetricsAsJSON()`. * @return {Registry} aggregated registry. */ static aggregate(metricsArr) { @@ -122,7 +148,7 @@ class AggregatorRegistry extends Registry { * Sets the registry or registries to be aggregated. Call from workers to * use a registry/registries other than the default global registry. * @param {Array|Registry} regs Registry or registries to be - * aggregated. + * aggregated. * @return {void} */ static setRegistries(regs) { @@ -167,9 +193,19 @@ function addListeners() { if (request.failed) return; // Callback already run with Error. const registry = AggregatorRegistry.aggregate(request.responses); + const g = new Gauge({ + name: 'nodejs_prom_client_cluster_workers', + help: 'Number of connected cluster workers reporting to prometheus', + registers: [registry] + }); + g.set(request.workerCount); const promString = registry.metrics(); request.done(null, promString); } + } else if (message.type === REG_METRICS_WORKER) { + //setup coordinated workers + const workerId = message.workerId; + coordinatedWorkers.add(workerId); } }); } @@ -186,6 +222,11 @@ if (cluster.isWorker) { }); } }); + + process.send({ + type: REG_METRICS_WORKER, + workerId: cluster.worker.id + }); } module.exports = AggregatorRegistry;