diff --git a/CHANGELOG.md b/CHANGELOG.md index de63ceba..a1f31d0a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/README.md b/README.md index fdff3b7d..08df59ae 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/lib/cluster.js b/lib/cluster.js index b3b04004..21a688d0 100644 --- a/lib/cluster.js +++ b/lib/cluster.js @@ -23,6 +23,7 @@ */ const { debuglog } = require('node:util'); +const Histogram = require('./histogram'); const Registry = require('./registry'); const { waitFor } = require('./util'); @@ -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 @@ -55,6 +64,7 @@ class AggregatorRegistry extends Registry { */ constructor(regContentType = Registry.PROMETHEUS_CONTENT_TYPE) { super(regContentType); + Registry.globalRegistry.registerMetric(clusterWorkerScrapeFailures); addListeners(); } @@ -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, @@ -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.`, ); @@ -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({ diff --git a/lib/worker.js b/lib/worker.js index aa5c3c64..4d727256 100644 --- a/lib/worker.js +++ b/lib/worker.js @@ -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'); @@ -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(); @@ -60,6 +69,7 @@ class WorkerRegistry extends Registry { primary = isMainThread, ) { super(regContentType); + Registry.globalRegistry.registerMetric(workerScrapeFailures); this.primary = primary; addListeners(primary); @@ -100,6 +110,7 @@ class WorkerRegistry extends Registry { const request = { responseHandlers, + workerFailures: 0, promise: waitFor( this.#gather(requestId, metricSnapshot, responsePromises), 5_000, @@ -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.`, ); @@ -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, @@ -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({ diff --git a/test/clusterTest.js b/test/clusterTest.js index c475266c..e1c2ca71 100644 --- a/test/clusterTest.js +++ b/test/clusterTest.js @@ -24,6 +24,8 @@ const ANNOUNCEMENT = '@prometheus-io/client:announcement'; const GET_METRICS_REQ = '@prometheus-io/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes'; const GOODBYE = '@prometheus-io/client:goodbye'; +const CLUSTER_WORKER_SCRAPE_FAILURES = + 'prom_client_cluster_worker_scrape_failures'; function metric(value) { return { @@ -91,6 +93,7 @@ describe.each([ afterEach(() => { cluster.off('message', listener); + jest.useRealTimers(); jest.restoreAllMocks(); }); @@ -98,11 +101,7 @@ describe.each([ const ar = new AggregatorRegistry(regType); const metrics = await ar.clusterMetrics(); - if (regType === Registry.OPENMETRICS_CONTENT_TYPE) { - expect(metrics).toContain('# EOF\n'); - } else { - expect(metrics.trim()).toEqual(''); - } + expect(metrics).toContain(`${CLUSTER_WORKER_SCRAPE_FAILURES}_count 0`); }); it('formats in the correct content type', async () => { @@ -112,7 +111,7 @@ describe.each([ if (regType === Registry.OPENMETRICS_CONTENT_TYPE) { expect(metrics).toContain('# EOF\n'); } else { - expect(metrics.trim()).toEqual(''); + expect(metrics).not.toContain('# EOF\n'); } }); @@ -188,6 +187,203 @@ describe.each([ } }); + it('records the number of workers that time out', async () => { + jest.useFakeTimers(); + + const originalWorkers = cluster.workers; + const registry = new AggregatorRegistry(regType); + const workers = Object.fromEntries( + [1, 2, 3].map(id => [ + id, + { + id, + isConnected: () => true, + send: jest.fn(), + }, + ]), + ); + cluster.workers = workers; + + Object.values(workers).forEach(worker => { + cluster.emit('message', worker, { type: ANNOUNCEMENT }); + }); + + try { + await discovery; + + const failedResult = registry.clusterMetrics(); + const rejection = expect(failedResult).rejects.toThrow( + 'Operation timed out. 2 outstanding responses.', + ); + + cluster.emit('message', workers[1], { + type: GET_METRICS_RES, + requestId: 0, + metrics: [[]], + }); + + await jest.advanceTimersByTimeAsync(5_000); + await rejection; + + const recoveredResult = registry.clusterMetrics(); + Object.values(workers).forEach(worker => { + cluster.emit('message', worker, { + type: GET_METRICS_RES, + requestId: 1, + metrics: [[]], + }); + }); + + await expect(recoveredResult).resolves.toContain( + `${CLUSTER_WORKER_SCRAPE_FAILURES}_sum 2`, + ); + await expect(recoveredResult).resolves.toContain( + `${CLUSTER_WORKER_SCRAPE_FAILURES}_count 1`, + ); + } finally { + Object.values(workers).forEach(worker => { + cluster.emit('disconnect', worker); + }); + cluster.workers = originalWorkers; + } + }); + + it('records worker-reported scrape errors', async () => { + const originalWorkers = cluster.workers; + const registry = new AggregatorRegistry(regType); + const worker = { + id: 1, + isConnected: () => true, + send: jest.fn(), + }; + cluster.workers = [worker]; + cluster.emit('message', worker, { type: ANNOUNCEMENT }); + + try { + await discovery; + + const failedResult = registry.clusterMetrics(); + const rejection = expect(failedResult).rejects.toThrow( + 'worker collection failed', + ); + + cluster.emit('message', worker, { + type: GET_METRICS_RES, + requestId: 0, + error: 'worker collection failed', + }); + await rejection; + + const recoveredResult = registry.clusterMetrics(); + cluster.emit('message', worker, { + type: GET_METRICS_RES, + requestId: 1, + metrics: [[]], + }); + + await expect(recoveredResult).resolves.toContain( + `${CLUSTER_WORKER_SCRAPE_FAILURES}_sum 1`, + ); + await expect(recoveredResult).resolves.toContain( + `${CLUSTER_WORKER_SCRAPE_FAILURES}_count 1`, + ); + } finally { + cluster.emit('disconnect', worker); + cluster.workers = originalWorkers; + } + }); + + it.each(['default', 'custom', 'custom with internal histogram'])( + 'records zero for a %s registry failure and respects selected registries', + async registryType => { + const MetricRegistry = require('../lib/registry'); + const source = + registryType === 'default' + ? MetricRegistry.globalRegistry + : new MetricRegistry(); + AggregatorRegistry.setRegistries(source); + const registry = new AggregatorRegistry(regType); + const histogram = MetricRegistry.globalRegistry.getSingleMetric( + CLUSTER_WORKER_SCRAPE_FAILURES, + ); + expect(histogram).toBeDefined(); + if (registryType === 'custom with internal histogram') { + source.registerMetric(histogram); + } + const error = new TypeError('local collection failed'); + jest.spyOn(source, 'getMetricsAsJSON').mockRejectedValueOnce(error); + + await expect(registry.clusterMetrics()).rejects.toBe(error); + + for (let scrape = 0; scrape < 2; scrape++) { + const metrics = await registry.clusterMetrics(); + if (registryType === 'custom') { + expect(metrics).not.toContain(CLUSTER_WORKER_SCRAPE_FAILURES); + continue; + } + expect( + metrics + .split('\n') + .filter(line => + line.startsWith(`${CLUSTER_WORKER_SCRAPE_FAILURES}_count `), + ), + ).toEqual([`${CLUSTER_WORKER_SCRAPE_FAILURES}_count 1`]); + expect(metrics).toContain( + `${CLUSTER_WORKER_SCRAPE_FAILURES}_sum 0\n`, + ); + expect(metrics).toContain( + `${CLUSTER_WORKER_SCRAPE_FAILURES}_bucket{le="0"} 1\n`, + ); + } + + const observed = await histogram.get(); + expect(observed.values).toContainEqual({ + metricName: `${CLUSTER_WORKER_SCRAPE_FAILURES}_count`, + labels: {}, + value: 1, + }); + MetricRegistry.globalRegistry.resetMetrics(); + expect((await histogram.get()).values).toEqual([]); + source.getMetricsAsJSON.mockRejectedValueOnce(error); + await expect(registry.clusterMetrics()).rejects.toBe(error); + await registry.clusterMetrics(); + await expect( + MetricRegistry.globalRegistry.metrics(), + ).resolves.toContain(`${CLUSTER_WORKER_SCRAPE_FAILURES}_count 1\n`); + }, + ); + + it('does not count waiting workers when a worker reports an error named Timeout', async () => { + jest.useFakeTimers(); + const registry = new AggregatorRegistry(regType); + const workers = [1, 2].map(id => { + return { id, isConnected: () => true, send: jest.fn() }; + }); + for (const worker of workers) { + cluster.emit('message', worker, { type: ANNOUNCEMENT }); + } + + try { + const rejected = expect(registry.clusterMetrics()).rejects.toThrow( + /^Timeout$/, + ); + cluster.emit('message', workers[0], { + type: GET_METRICS_RES, + requestId: 0, + error: 'Timeout', + }); + await rejected; + await jest.advanceTimersByTimeAsync(0); + expect(jest.getTimerCount()).toBe(0); + } finally { + for (const worker of workers) cluster.emit('disconnect', worker); + } + + const metrics = await registry.clusterMetrics(); + expect(metrics).toContain(`${CLUSTER_WORKER_SCRAPE_FAILURES}_count 1\n`); + expect(metrics).toContain(`${CLUSTER_WORKER_SCRAPE_FAILURES}_sum 1\n`); + }); + it('accumulate stats from terminated workers', async () => { const originalWorkers = cluster.workers; const registry = new AggregatorRegistry(regType); @@ -338,7 +534,12 @@ describe.each([ ], }; - await expect(metrics).resolves.toEqual([[expected]]); + const histogram = AggregatorRegistry.globalRegistry.getSingleMetric( + CLUSTER_WORKER_SCRAPE_FAILURES, + ); + await expect(metrics).resolves.toEqual([ + [await histogram.get(), expected], + ]); } finally { jest.dontMock('cluster'); gauge.remove(); diff --git a/test/workerTest.js b/test/workerTest.js index 3bdf6c0b..5a3015ae 100644 --- a/test/workerTest.js +++ b/test/workerTest.js @@ -23,6 +23,7 @@ const ANNOUNCEMENT = '@prometheus-io/client:announcement'; const GET_METRICS_REQ = '@prometheus-io/client:getMetricsReq'; const GET_METRICS_RES = '@prometheus-io/client:getMetricsRes'; const GOODBYE = '@prometheus-io/client:goodbye'; +const WORKER_SCRAPE_FAILURES = 'prom_client_worker_scrape_failures'; function metric(value) { return { @@ -81,7 +82,7 @@ describe.each([ it('works properly if there are no workers', async () => { const metrics = await registry.workerMetrics(); - expect(metrics).toEqual(''); + expect(metrics).toBe(''); }); it('formats in the correct content type', async () => { @@ -285,7 +286,12 @@ describe.each([ ], }; - await expect(metrics).resolves.toEqual([[expected]]); + const histogram = AggregatorRegistry.globalRegistry.getSingleMetric( + WORKER_SCRAPE_FAILURES, + ); + await expect(metrics).resolves.toEqual([ + [await histogram.get(), expected], + ]); } finally { channel.close(); } @@ -377,3 +383,227 @@ describe.each([ }); }); }); + +describe.each([ + ['Prometheus', Registry.PROMETHEUS_CONTENT_TYPE], + ['OpenMetrics', Registry.OPENMETRICS_CONTENT_TYPE], +])('%s worker scrape failures', (tag, regType) => { + let WorkerRegistry; + let registry; + let channels; + let announcementChannel; + + beforeEach(() => { + jest.resetModules(); + jest.useFakeTimers(); + channels = []; + + // Deliver messages explicitly so timeout and late-response tests do not + // depend on BroadcastChannel scheduling or wall-clock delays. + jest.doMock('node:worker_threads', () => { + return { + isMainThread: true, + threadId: 0, + BroadcastChannel: class { + constructor(name) { + this.name = name; + this.listeners = []; + this.postMessage = jest.fn(); + channels.push(this); + } + + unref() { + return this; + } + + close() {} + + addEventListener(type, listener) { + if (type === 'message') this.listeners.push(listener); + } + + receive(data) { + return Promise.all( + this.listeners.map(listener => listener({ data })), + ); + } + }, + }; + }); + + WorkerRegistry = require('../lib/worker'); + registry = new WorkerRegistry(regType); + WorkerRegistry.globalRegistry.setContentType(regType); + announcementChannel = channels[0]; + }); + + afterEach(() => { + jest.useRealTimers(); + jest.restoreAllMocks(); + jest.dontMock('node:worker_threads'); + }); + + async function announceWorker(threadId) { + const name = `@prometheus-io/client:worker:${threadId}`; + await announcementChannel.receive({ type: ANNOUNCEMENT, name, threadId }); + return channels.find(channel => channel.name === name); + } + + function expectFailures(metrics, count, sum, zeros = 0) { + expect( + metrics + .split('\n') + .filter(line => line.startsWith(`${WORKER_SCRAPE_FAILURES}_count `)), + ).toEqual([`${WORKER_SCRAPE_FAILURES}_count ${count}`]); + expect(metrics).toContain(`${WORKER_SCRAPE_FAILURES}_sum ${sum}\n`); + expect(metrics).toContain( + `${WORKER_SCRAPE_FAILURES}_bucket{le="0"} ${zeros}\n`, + ); + expect(metrics.endsWith('# EOF\n')).toBe( + regType === Registry.OPENMETRICS_CONTENT_TYPE, + ); + } + + it('retains consecutive timeouts after the last worker leaves and ignores late errors', async () => { + const workers = []; + for (const id of [1, 2, 3]) workers.push(await announceWorker(id)); + + for (const requestId of [0, 1]) { + const rejection = expect(registry.workerMetrics()).rejects.toThrow( + `Operation timed out. ${2 - requestId} outstanding responses.`, + ); + for (let index = 0; index <= requestId; index++) { + await workers[index].receive({ + type: GET_METRICS_RES, + requestId, + metrics: [[]], + }); + } + await jest.advanceTimersByTimeAsync(5_000); + await rejection; + + await workers[2].receive({ + type: GET_METRICS_RES, + requestId, + error: 'late failure', + }); + } + + for (const worker of workers) { + await worker.receive({ type: GOODBYE, metrics: [[]] }); + } + expect(WorkerRegistry.workerCount()).toBe(0); + await expect(registry.workerMetrics()).resolves.toBe(''); + const recovered = await WorkerRegistry.globalRegistry.metrics(); + expectFailures(recovered, 2, 3); + expect(recovered).toContain(`${WORKER_SCRAPE_FAILURES}_bucket{le="1"} 1\n`); + expect(recovered).toContain(`${WORKER_SCRAPE_FAILURES}_bucket{le="2"} 2\n`); + expectFailures(await WorkerRegistry.globalRegistry.metrics(), 2, 3); + }); + + it.each(['worker collection failed', 'Timeout'])( + 'records the worker error %s once without counting a waiting worker', + async errorMessage => { + const worker = await announceWorker(1); + const waitingWorker = await announceWorker(2); + const rejection = expect(registry.workerMetrics()).rejects.toThrow( + new Error(errorMessage), + ); + const error = { + type: GET_METRICS_RES, + requestId: 0, + error: errorMessage, + }; + await worker.receive(error); + await rejection; + await worker.receive(error); + + const recovered = registry.workerMetrics(); + await worker.receive({ + type: GET_METRICS_RES, + requestId: 1, + metrics: [[metric(7)]], + }); + await waitingWorker.receive({ + type: GET_METRICS_RES, + requestId: 1, + metrics: [[]], + }); + const metrics = await recovered; + expect(metrics).not.toContain(WORKER_SCRAPE_FAILURES); + expectFailures(await WorkerRegistry.globalRegistry.metrics(), 1, 1); + expect(metrics).toContain('test_metric 7\n'); + }, + ); + + it('records zero for aggregation errors without adding local metrics to worker responses', async () => { + const worker = await announceWorker(1); + const gauge = new (require('../lib/gauge'))({ + name: 'coordinator_value', + help: 'test', + }); + gauge.set(17); + const rejected = expect(registry.workerMetrics()).rejects.toThrow( + "'invalid' is not a defined aggregator.", + ); + await worker.receive({ + type: GET_METRICS_RES, + requestId: 0, + metrics: [[{ ...metric(7), aggregator: 'invalid' }]], + }); + await rejected; + expectFailures(await WorkerRegistry.globalRegistry.metrics(), 1, 0, 1); + + const recovered = registry.workerMetrics(); + await worker.receive({ + type: GET_METRICS_RES, + requestId: 1, + metrics: [[metric(7)]], + }); + const metrics = await recovered; + expect(metrics).toContain('test_metric 7\n'); + expect(metrics).not.toContain('coordinator_value'); + expect(metrics).not.toContain(WORKER_SCRAPE_FAILURES); + expectFailures(await WorkerRegistry.globalRegistry.metrics(), 1, 0, 1); + }); + + it.each(['default', 'custom'])( + 'reports %s registry collection errors to the parent', + async registryType => { + const MetricRegistry = require('../lib/registry'); + const source = + registryType === 'default' + ? MetricRegistry.globalRegistry + : new MetricRegistry(); + WorkerRegistry.setRegistries(source); + jest + .spyOn(source, 'getMetricsAsJSON') + .mockRejectedValueOnce(new Error('worker collection failed')); + const workerChannel = channels[1]; + + await announcementChannel.receive({ + type: GET_METRICS_REQ, + requestId: 42, + }); + expect(workerChannel.postMessage).toHaveBeenLastCalledWith({ + type: GET_METRICS_RES, + requestId: 42, + error: 'worker collection failed', + }); + + await announcementChannel.receive({ + type: GET_METRICS_REQ, + requestId: 43, + }); + expect(workerChannel.postMessage).toHaveBeenLastCalledWith({ + type: GET_METRICS_RES, + requestId: 43, + threadId: 0, + metrics: [await source.getMetricsAsJSON()], + }); + expect( + MetricRegistry.globalRegistry.getSingleMetric(WORKER_SCRAPE_FAILURES), + ).toBeDefined(); + }, + ); +});