Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ project adheres to [Semantic Versioning](http://semver.org/).

### Added

- Record cluster and worker thread scrape failures in internal histograms,
including timeouts, worker-reported errors, and failures with no known worker errors.

## [0.16.0] - 2026-08-24

This release marks our first release as a Prometheus subproject.
Expand Down
17 changes: 17 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -630,6 +630,23 @@ aggregation method, set the `aggregator` property in the metric config to one of
'sum', 'first', 'min', 'max', 'average' or 'omit'. (See `lib/metrics/version.js`
for an example.)

Failed cluster collections are recorded in the
`prom_client_cluster_worker_scrape_failures` histogram. Worker thread collections
use `prom_client_worker_scrape_failures`. Each failed collection records one
observation: the number of outstanding worker responses on a timeout, or the
number of worker-reported errors received before the collection rejects. Other
collection errors record zero when no worker failures are known. Successful
collections do not add observations.

Failed collections reject without returning partial metrics. These internal
histograms are registered in the global registry when their corresponding
`ClusterRegistry` or `WorkerRegistry` is constructed. Read the coordinating
process or thread's failure observations through `register.metrics()`.
`clusterMetrics()` also includes the coordinator's selected registries, whereas
`workerMetrics()` retains its existing worker-only collection behavior.
Custom registries selected with `setRegistries()` do not automatically include
the internal histogram; register it explicitly if needed.

If you need to expose metrics about an individual worker, you can include a
value that is unique to the worker (such as the worker ID or process ID) in a
label. (See `example/server.js` for an example using
Expand Down
24 changes: 23 additions & 1 deletion lib/cluster.js
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
*/

const { debuglog } = require('node:util');
const Histogram = require('./histogram');
const Registry = require('./registry');
const { waitFor } = require('./util');

Expand All @@ -41,6 +42,14 @@ const GET_METRICS_REQ = '@prometheus-io/client:getMetricsReq';
const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes';
const GOODBYE = '@prometheus-io/client:goodbye';

const clusterWorkerScrapeFailures = new Histogram({
name: 'prom_client_cluster_worker_scrape_failures',
help: 'Number of workers that failed to return metrics during a failed cluster scrape.',
buckets: [0, 1, 2, 4, 8, 16, 32],
// Register on construction, so importing the module does not add metrics.
registers: [],
});

let registries = [Registry.globalRegistry];
let listenersAdded = false;
let requestCtr = 0; // Concurrency control
Expand All @@ -55,6 +64,7 @@ class AggregatorRegistry extends Registry {
*/
constructor(regContentType = Registry.PROMETHEUS_CONTENT_TYPE) {
super(regContentType);
Registry.globalRegistry.registerMetric(clusterWorkerScrapeFailures);

addListeners();
}
Expand Down Expand Up @@ -90,6 +100,7 @@ class AggregatorRegistry extends Registry {
const responsePromises = [this.#selfMetrics(), ...workerMetrics];
const request = {
responseHandlers,
workerFailures: 0,
promise: waitFor(
this.#gather(requestId, metricSnapshot, responsePromises),
5_000,
Expand All @@ -105,7 +116,17 @@ class AggregatorRegistry extends Registry {

return await request.promise;
} catch (err) {
if (err.message === 'Timeout') {
const timedOut =
err.message === 'Timeout' && request.workerFailures === 0;
let failedWorkers = request.workerFailures;

if (timedOut) {
failedWorkers += request.responseHandlers.size;
}

clusterWorkerScrapeFailures.observe(failedWorkers);

if (timedOut) {
throw new Error(
`Operation timed out. ${request.responseHandlers.size} outstanding responses.`,
);
Expand Down Expand Up @@ -348,6 +369,7 @@ async function primaryListener(worker, event) {
request.responseHandlers.delete(worker.id);

if (event.error) {
request.workerFailures++;
response.reject(new Error(event.error));
} else {
response.resolve({
Expand Down
32 changes: 27 additions & 5 deletions lib/worker.js
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

const { debuglog } = require('node:util');
const worker = require('node:worker_threads');
const Histogram = require('./histogram');
const Registry = require('./registry');
const { waitFor } = require('./util');

Expand All @@ -36,6 +37,14 @@ const GET_METRICS_REQ = '@prometheus-io/client:getMetricsReq';
const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes';
const GOODBYE = '@prometheus-io/client:goodbye';

const workerScrapeFailures = new Histogram({
name: 'prom_client_worker_scrape_failures',
help: 'Number of workers that failed to return metrics during a failed worker scrape.',
buckets: [0, 1, 2, 4, 8, 16, 32],
// Register on construction, so importing the module does not add metrics.
registers: [],
});

const ANNOUNCEMENT_CHANNEL = new BroadcastChannel(
'@prometheus-io/client:announce',
).unref();
Expand All @@ -60,6 +69,7 @@ class WorkerRegistry extends Registry {
primary = isMainThread,
) {
super(regContentType);
Registry.globalRegistry.registerMetric(workerScrapeFailures);
this.primary = primary;

addListeners(primary);
Expand Down Expand Up @@ -100,6 +110,7 @@ class WorkerRegistry extends Registry {

const request = {
responseHandlers,
workerFailures: 0,
promise: waitFor(
this.#gather(requestId, metricSnapshot, responsePromises),
5_000,
Expand All @@ -117,7 +128,17 @@ class WorkerRegistry extends Registry {

return await request.promise;
} catch (err) {
if (err.message === 'Timeout') {
const timedOut =
err.message === 'Timeout' && request.workerFailures === 0;
let failedWorkers = request.workerFailures;

if (timedOut) {
failedWorkers += request.responseHandlers.size;
}

workerScrapeFailures.observe(failedWorkers);

if (timedOut) {
throw new Error(
`Operation timed out. ${request.responseHandlers.size} outstanding responses.`,
);
Expand Down Expand Up @@ -277,11 +298,11 @@ function addListeners(primary) {
announce(name, false);
}
} else if (message.type === GET_METRICS_REQ) {
const metrics = await Promise.all(
registries.map(r => r.getMetricsAsJSON()),
);

try {
const metrics = await Promise.all(
registries.map(r => r.getMetricsAsJSON()),
);

channel.postMessage({
type: GET_METRICS_RES,
requestId: message.requestId,
Expand Down Expand Up @@ -375,6 +396,7 @@ async function primaryListener(event) {
request.responseHandlers.delete(workerName);

if (workerMessage.error) {
request.workerFailures++;
response.reject(new Error(workerMessage.error));
} else {
response.resolve({
Expand Down
Loading
Loading