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..454e405d 100644 --- a/README.md +++ b/README.md @@ -630,6 +630,24 @@ 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. Their observations +are exposed by subsequent successful calls to `clusterMetrics()` or +`workerMetrics()`, respectively, even if there are no workers left. These internal +histograms are registered in the global registry when their corresponding +`ClusterRegistry` or `WorkerRegistry` is constructed. Aggregation also includes +the coordinating process or thread's metrics. If custom registries selected with +`setRegistries()` do not contain the internal histogram, it is included separately +so failure observations remain available without duplicating it in the default +registry path. + 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..29c25c15 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(); } @@ -88,11 +98,14 @@ class AggregatorRegistry extends Registry { ); const responsePromises = [this.#selfMetrics(), ...workerMetrics]; + const timeoutError = new Error('Timeout'); const request = { responseHandlers, + workerFailures: 0, promise: waitFor( this.#gather(requestId, metricSnapshot, responsePromises), 5_000, + timeoutError, ), }; @@ -105,7 +118,22 @@ class AggregatorRegistry extends Registry { return await request.promise; } catch (err) { - if (err.message === 'Timeout') { + const timedOut = err === timeoutError; + let failedWorkers = request.workerFailures; + + if (timedOut) { + failedWorkers += request.responseHandlers.size; + } + + clusterWorkerScrapeFailures.observe(failedWorkers); + + // Sending can throw before request.promise is awaited. + for (const response of request.responseHandlers.values()) { + response.reject(err); + } + request.promise.catch(() => {}); + + if (timedOut) { throw new Error( `Operation timed out. ${request.responseHandlers.size} outstanding responses.`, ); @@ -118,8 +146,22 @@ class AggregatorRegistry extends Registry { } async #selfMetrics() { + const metrics = await Promise.all( + registries.map(r => r.getMetricsAsJSON()), + ); + // Custom registries may not contain the coordinator's internal metric. + if ( + !registries.some( + r => + r.getSingleMetric(clusterWorkerScrapeFailures.name) === + clusterWorkerScrapeFailures, + ) + ) { + metrics.push([await clusterWorkerScrapeFailures.get()]); + } + return { - metrics: await Promise.all(registries.map(r => r.getMetricsAsJSON())), + metrics, }; } @@ -348,6 +390,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/util.js b/lib/util.js index 69f12d52..0113f4da 100644 --- a/lib/util.js +++ b/lib/util.js @@ -173,9 +173,10 @@ exports.nowTimestamp = function nowTimestamp() { * Async functions with a timeout. * @param promise {Promise} * @param limit {number} + * @param {Error} [timeoutError] - Error used when this deadline expires. * @returns {Promise} */ -exports.waitFor = async function waitFor(promise, limit = 5_000) { +exports.waitFor = async function waitFor(promise, limit = 5_000, timeoutError) { let resolve, reject; const resultPromise = new Promise((res, rej) => { @@ -184,7 +185,7 @@ exports.waitFor = async function waitFor(promise, limit = 5_000) { }); const timeout = setTimeout(() => { - reject(new Error('Timeout')); + reject(timeoutError ?? new Error('Timeout')); }, limit); try { diff --git a/lib/worker.js b/lib/worker.js index aa5c3c64..3851aafe 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); @@ -78,17 +88,12 @@ class WorkerRegistry extends Registry { ); if (orderedWorkers.length === 0) { - if (historicMetrics.length === 0) { - debug('No data found for requestId', requestId); - return ''; - } else { - debug('No workers found for requestId', requestId); - } + debug('No workers found for requestId', requestId); } const metricSnapshot = historicMetrics; const responseHandlers = new Map(); - const responsePromises = orderedWorkers.map( + const workerMetrics = orderedWorkers.map( entry => new Promise((resolveResponse, rejectResponse) => { responseHandlers.set(entry.name, { @@ -98,11 +103,15 @@ class WorkerRegistry extends Registry { }), ); + const responsePromises = [this.#selfMetrics(), ...workerMetrics]; + const timeoutError = new Error('Timeout'); const request = { responseHandlers, + workerFailures: 0, promise: waitFor( this.#gather(requestId, metricSnapshot, responsePromises), 5_000, + timeoutError, ), }; @@ -117,7 +126,22 @@ class WorkerRegistry extends Registry { return await request.promise; } catch (err) { - if (err.message === 'Timeout') { + const timedOut = err === timeoutError; + let failedWorkers = request.workerFailures; + + if (timedOut) { + failedWorkers += request.responseHandlers.size; + } + + workerScrapeFailures.observe(failedWorkers); + + // Sending can throw before request.promise is awaited. + for (const response of request.responseHandlers.values()) { + response.reject(err); + } + request.promise.catch(() => {}); + + if (timedOut) { throw new Error( `Operation timed out. ${request.responseHandlers.size} outstanding responses.`, ); @@ -129,6 +153,23 @@ class WorkerRegistry extends Registry { } } + async #selfMetrics() { + const metrics = await Promise.all( + registries.map(r => r.getMetricsAsJSON()), + ); + // Custom registries may not contain the coordinator's internal metric. + if ( + !registries.some( + r => + r.getSingleMetric(workerScrapeFailures.name) === workerScrapeFailures, + ) + ) { + metrics.push([await workerScrapeFailures.get()]); + } + + return { metrics }; + } + /** * Collect the data for a metrics request. * @param requestId {number} @@ -277,11 +318,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 +416,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..0f83c3ac 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,219 @@ 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 without workers and exports it once', + 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('Timeout'); + 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(); + 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`, + ); + } + + MetricRegistry.globalRegistry.resetMetrics(); + expect((await histogram.get()).values).toEqual([]); + source.getMetricsAsJSON.mockRejectedValueOnce(error); + await expect(registry.clusterMetrics()).rejects.toBe(error); + await expect(registry.clusterMetrics()).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('records synchronous send failures and settles the pending request', async () => { + jest.useFakeTimers(); + const registry = new AggregatorRegistry(regType); + const error = new TypeError('worker.send failed'); + const worker = { + id: 1, + isConnected: () => true, + send: jest.fn(() => { + throw error; + }), + }; + cluster.emit('message', worker, { type: ANNOUNCEMENT }); + + try { + await expect(registry.clusterMetrics()).rejects.toBe(error); + await jest.advanceTimersByTimeAsync(0); + expect(jest.getTimerCount()).toBe(0); + await expect(registry.shutdown()).resolves.toBeUndefined(); + } finally { + 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 0\n`); + }); + it('accumulate stats from terminated workers', async () => { const originalWorkers = cluster.workers; const registry = new AggregatorRegistry(regType); @@ -338,7 +550,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/utilTest.js b/test/utilTest.js index 11f5ceef..e1e2ba19 100644 --- a/test/utilTest.js +++ b/test/utilTest.js @@ -56,6 +56,21 @@ describe('utils', () => { await expect(waitFor(Promise.resolve('foo'))).resolves.toEqual('foo'); }); + it('rejects with the supplied error when the deadline expires', async () => { + const timeoutError = new Error('Timeout'); + await expect( + waitFor(new Promise(() => {}), 1, timeoutError), + ).rejects.toBe(timeoutError); + }); + + it('preserves a promise rejection with the same message as the deadline', async () => { + const timeoutError = new Error('Timeout'); + const collectorError = new Error('Timeout'); + await expect( + waitFor(Promise.reject(collectorError), 5_000, timeoutError), + ).rejects.toBe(collectorError); + }); + it('Rejects on a promise rejection', async () => { const promise = waitFor(Promise.reject(new Error('nope'))); await expect(promise).rejects.toThrow('nope'); diff --git a/test/workerTest.js b/test/workerTest.js index 3bdf6c0b..ec938d90 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).toContain(`${WORKER_SCRAPE_FAILURES}_count 0`); }); 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,258 @@ 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); + 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); + const recovered = await registry.workerMetrics(); + 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 registry.workerMetrics(), 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; + expectFailures(metrics, 1, 1); + expect(metrics).toContain('test_metric 7\n'); + }, + ); + + it.each([0, 1])( + 'records zero for a broadcast failure with %i workers and settles the request', + async workerCount => { + const worker = workerCount ? await announceWorker(1) : undefined; + const error = new TypeError('postMessage failed'); + announcementChannel.postMessage.mockImplementationOnce(() => { + throw error; + }); + + await expect(registry.workerMetrics()).rejects.toBe(error); + await jest.advanceTimersByTimeAsync(0); + expect(jest.getTimerCount()).toBe(0); + await expect(registry.shutdown()).resolves.toBeUndefined(); + + const recovered = registry.workerMetrics(); + if (worker) { + await worker.receive({ + type: GET_METRICS_RES, + requestId: 1, + metrics: [[]], + }); + } + expectFailures(await recovered, 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(); + }, + ); + + it.each(['default', 'custom', 'custom with internal histogram'])( + 'exports a %s registry failure once and preserves a collector error named Timeout', + async registryType => { + const MetricRegistry = require('../lib/registry'); + const Gauge = require('../lib/gauge'); + const source = + registryType === 'default' + ? MetricRegistry.globalRegistry + : new MetricRegistry(); + WorkerRegistry.setRegistries(source); + const histogram = MetricRegistry.globalRegistry.getSingleMetric( + WORKER_SCRAPE_FAILURES, + ); + expect(histogram).toBeDefined(); + if (registryType === 'custom with internal histogram') + source.registerMetric(histogram); + const gauge = new Gauge({ + name: 'coordinator_value', + help: 'test', + registers: [source], + }); + gauge.set(17); + const error = new TypeError('Timeout'); + jest.spyOn(source, 'getMetricsAsJSON').mockRejectedValueOnce(error); + await expect(registry.workerMetrics()).rejects.toBe(error); + + for (let scrape = 0; scrape < 2; scrape++) { + const metrics = await registry.workerMetrics(); + expectFailures(metrics, 1, 0, 1); + expect(metrics).toContain('coordinator_value 17\n'); + } + MetricRegistry.globalRegistry.resetMetrics(); + expect((await histogram.get()).values).toEqual([]); + source.getMetricsAsJSON.mockRejectedValueOnce(error); + await expect(registry.workerMetrics()).rejects.toBe(error); + expectFailures(await registry.workerMetrics(), 1, 0, 1); + }, + ); +});