diff --git a/queries/cdmq/README.md b/queries/cdmq/README.md index 6feaeba7..dcf2879e 100644 --- a/queries/cdmq/README.md +++ b/queries/cdmq/README.md @@ -12,6 +12,10 @@ npm install The contents of this directory contain a collection of scripts in Javascript intended to be executed with [node.js](https://nodejs.org). These scripts get data from an OpenSearch instance. The data must be in Common Data Format, whose index mapping definitions are documented in [cdm.js](./cdm.js)'s `indexDefs` object. The scripts here are meant to help inspect, compare, and export data from benchmarks and performance & resource-utilization tools, in order to report and investigate performance. +Native distribution statistics use OpenSearch Point-in-Time search with `search_after` for stable deep pagination. The Crucible controller currently ships OpenSearch 3.6.0, which is the tested compatibility target for this feature. Deployments without PIT support receive an explicit `NATIVE_STATS_PIT_UNSUPPORTED` error. + +Resource limits are configurable on the CDM server with positive-integer environment variables. Defaults are `CDM_NATIVE_STATS_MAX_DOCUMENTS=250000`, `CDM_NATIVE_STATS_MAX_INTERVALS=500000`, `CDM_NATIVE_STATS_MAX_RUNTIME_MS=300000`, and `CDM_NATIVE_STATS_PAGE_SIZE=1000`. Exceeding a limit returns an error; the server never silently falls back to statistics over resolution buckets. + In order to generate this data, you must run a benchmark via automation framework which uses the Common Data Format and index that data into OpenSearch. One of those automation frameworks is the [crucible](https://github.com/perftool-incubator/crucible) project. A subproject of crucible, [crucible-examples](https://github.com/perftool-incubator/crucible-examples), includes scenarios to run some of these benchmarks. ## Terms @@ -487,6 +491,15 @@ node ./get-metric-data.js --period --source iostat --type kB-sec --breako So far all of the metrics have been represented as a single value for a specific time period. When `--period` is used, the script finds the begin and end times for this period, which in most cases, has a duration equal to the measurement time in the benchmark itself (around 90 seconds in these examples). One can also specify `--run`, `--begin`, and `--end` instead of `--period`, should they need to focus on a different period of time. However, for benchmark metrics (such as uperf), it is important to limit the begin and end to within the actual measurement period for that sample. Conversely, tool metrics can use a begin and end spanning any time period within the run, as the tool collection tends to run continuously for any particular run. Whatever time period is used, one can also use `--resolution` to divide this time period into multiple data-samples, in order to generate things like line graphs: +The optional `--distribution-stats` option returns duration-weighted statistics over the native reconstructed timeline, independently of `--resolution`. For example: + +```bash +node ./get-metric-data.js --period --source uperf --type Gbps \ + --distribution-stats min,max,mean,median,stddev,p95 +``` + +This describes variation in the underlying metric over time rather than variation in caller-selected display buckets. The JSON response includes these statistics in `distributionStats`, keyed by breakout label, while the existing `values` response remains unchanged. Native statistics are opt-in because they stream all matching metric documents and are subject to resource limits. Statistics use population standard deviation, duration-weighted nearest-rank percentiles, and `median` is equivalent to `p50`; a single native interval has `stddev=0`. + # node ./get-metric-data.js --period 4F1014D6-AD33-11EC-94E3-ADE96E3275F7 --source sar-net --type L2-Gbps --breakout csid=1,cstype=worker,type=physical,direction=tx,dev --filter gt:0.01 --resolution 10 Checking for httpd...appears to be running Checking for OpenSearch...appears to be running diff --git a/queries/cdmq/cdm.js b/queries/cdmq/cdm.js index 66d5c787..f6f86707 100644 --- a/queries/cdmq/cdm.js +++ b/queries/cdmq/cdm.js @@ -1,7 +1,20 @@ //# vim: autoindent tabstop=2 shiftwidth=2 expandtab softtabstop=2 filetype=javascript var request = require('sync-request'); var thenRequest = require('then-request'); +const { calculateDistributionStats, reconstructTimeline, resampleTimeline, validateRequestedStats } = require('./native-stats'); var bigQuerySize = 262144; + +function nativeStatsLimit(name, fallback) { + const value = process.env[name]; + if (typeof value === 'undefined') return fallback; + const limit = Number(value); + if (!Number.isSafeInteger(limit) || limit <= 0) { + const error = new Error(name + ' must be a positive integer'); + error.code = 'NATIVE_STATS_CONFIG'; + throw error; + } + return limit; +} const docTypes = { v7dev: ['run', 'tag', 'iteration', 'param', 'sample', 'period', 'metric_desc', 'metric_data'], v8dev: ['run', 'tag', 'iteration', 'param', 'sample', 'period', 'metric_desc', 'metric_data'], @@ -3699,6 +3712,44 @@ getMetricDataFromIdsSets = async function (instance, sets, metricGroupIdsByLabel // Each template has prefix/suffix pairs for the 4 query types, // with __IDS__ as placeholder for the metric UUID list. var defaultAggregation = sets[idx].defaultAggregation || 'sum'; + + // When native statistics are requested, use the reconstructed timeline for + // both the output series and the statistics. This avoids running the legacy + // resolution query and then fetching the same documents again for stats. + if (sets[idx].distributionStats) { + valueSets[idx] = { distributionStats: {} }; + const nativeIndexName = getIndexName('metric_data', instance, yearDotMonth); + const sortedNativeLabels = Object.keys(metricGroupIdsByLabelSets[idx]).sort(); + const nativeStatsConcurrency = 4; + for (let nativeStart = 0; nativeStart < sortedNativeLabels.length; nativeStart += nativeStatsConcurrency) { + const nativeLabels = sortedNativeLabels.slice(nativeStart, nativeStart + nativeStatsConcurrency); + const nativeResults = await Promise.all( + nativeLabels.map(async (label) => { + const native = await getNativeMetricStats( + instance, + metricGroupIdsByLabelSets[idx][label], + begin, + end, + defaultAggregation, + sets[idx].distributionStats, + yearDotMonth, + { includeTimeline: true, indexName: nativeIndexName } + ); + return { + label: label, + values: resampleTimeline(native.timeline, begin, end, resolution, defaultAggregation), + stats: native.stats + }; + }) + ); + nativeResults.forEach((result) => { + valueSets[idx][result.label] = result.values; + valueSets[idx].distributionStats[result.label] = result.stats; + }); + } + continue; + } + var timeRangeTemplates = []; var thisBegin = begin; var thisEnd = begin + duration; @@ -3861,6 +3912,157 @@ getMetricDataFromIdsSets = async function (instance, sets, metricGroupIdsByLabel exports.getMetricDataFromIdsSets = getMetricDataFromIdsSets; +// Stream the documents needed to reconstruct one native aggregate timeline. +// This is deliberately separate from getMetricDataFromIdsSets(): resolution +// bucketing can use aggregations, while native statistics must see every +// boundary where any selected metric ID changes value. +async function getNativeMetricStats( + instance, + metricIds, + begin, + end, + aggregation, + requestedStats, + yearDotMonth, + options = {} +) { + validateRequestedStats(requestedStats); + const maxDocuments = options.maxDocuments || nativeStatsLimit('CDM_NATIVE_STATS_MAX_DOCUMENTS', 250000); + const maxIntervals = options.maxIntervals || nativeStatsLimit('CDM_NATIVE_STATS_MAX_INTERVALS', 500000); + const maxRuntimeMs = options.maxRuntimeMs || nativeStatsLimit('CDM_NATIVE_STATS_MAX_RUNTIME_MS', 300000); + const pageSize = options.pageSize || nativeStatsLimit('CDM_NATIVE_STATS_PAGE_SIZE', 1000); + const indexName = options.indexName || getIndexName('metric_data', instance, yearDotMonth); + const baseUrl = 'http://' + instance.host; + const headers = instance.header || { 'Content-Type': 'application/json' }; + const fetchImpl = options.fetch || fetch; + let pitId; + const deadline = Date.now() + maxRuntimeMs; + + async function send(method, url, body) { + const response = await fetchImpl(url, { + method: method, + headers: headers, + body: body === undefined ? undefined : JSON.stringify(body) + }); + if (!response.ok) { + throw new Error('OpenSearch request failed with HTTP status ' + response.status); + } + return response.json(); + } + + try { + let pit; + try { + pit = await send('POST', baseUrl + '/' + indexName + '/_search/point_in_time?keep_alive=1m', {}); + } catch (error) { + error.code = 'NATIVE_STATS_PIT_UNSUPPORTED'; + error.message = 'native distribution statistics require OpenSearch Point-in-Time search support: ' + error.message; + throw error; + } + pitId = pit.pit_id || pit.id; + if (!pitId) { + const error = new Error('OpenSearch did not return a point-in-time ID'); + error.code = 'NATIVE_STATS_PIT_UNSUPPORTED'; + throw error; + } + + const documentsById = {}; + let searchAfter; + let documentCount = 0; + let firstPage = true; + while (true) { + if (Date.now() > deadline) { + const error = new Error('native distribution statistics processing time limit exceeded'); + error.code = 'NATIVE_STATS_LIMIT'; + throw error; + } + const query = { + size: pageSize, + track_total_hits: firstPage ? maxDocuments + 1 : false, + pit: { id: pitId, keep_alive: '1m' }, + sort: [ + { 'metric_data.begin': 'asc' }, + { 'metric_data.end': 'asc' }, + { 'metric_desc.metric_desc-uuid': 'asc' } + ], + docvalue_fields: [ + { field: 'metric_desc.metric_desc-uuid' }, + { field: 'metric_data.begin', format: 'epoch_millis' }, + { field: 'metric_data.end', format: 'epoch_millis' }, + { field: 'metric_data.value' } + ], + _source: false, + query: { + bool: { + filter: [ + { range: { 'metric_data.end': { gte: begin } } }, + { range: { 'metric_data.begin': { lte: end } } }, + { terms: { 'metric_desc.metric_desc-uuid': metricIds } } + ] + } + } + }; + if (searchAfter) query.search_after = searchAfter; + + const response = await send('POST', baseUrl + '/_search', query); + const total = response.hits && response.hits.total; + if (firstPage && total && total.value > maxDocuments) { + const error = new Error('native distribution statistics document limit exceeded'); + error.code = 'NATIVE_STATS_LIMIT'; + throw error; + } + + const hits = (response.hits && response.hits.hits) || []; + documentCount += hits.length; + if (documentCount > maxDocuments) { + const error = new Error('native distribution statistics document limit exceeded'); + error.code = 'NATIVE_STATS_LIMIT'; + throw error; + } + hits.forEach((hit) => { + const fields = hit.fields || {}; + const valueOf = (name) => { + const values = fields[name]; + return Array.isArray(values) ? values[0] : values; + }; + const metricId = valueOf('metric_desc.metric_desc-uuid'); + if (!documentsById[metricId]) documentsById[metricId] = []; + documentsById[metricId].push({ + begin: Number(valueOf('metric_data.begin')), + end: Number(valueOf('metric_data.end')), + value: Number(valueOf('metric_data.value')) + }); + }); + + if (hits.length < pageSize) break; + searchAfter = hits[hits.length - 1].sort; + if (!searchAfter) throw new Error('OpenSearch response did not include sort values for search_after'); + firstPage = false; + } + + metricIds.forEach((metricId) => { + if (!documentsById[metricId]) { + const error = new Error('native distribution statistics missing metric ID ' + metricId); + error.code = 'NATIVE_STATS_DATA_QUALITY'; + throw error; + } + }); + const timeline = reconstructTimeline(documentsById, Number(begin), Number(end), aggregation, { maxIntervals }); + const stats = calculateDistributionStats(timeline, requestedStats); + return options.includeTimeline ? { timeline: timeline, stats: stats } : stats; + } finally { + if (pitId) { + try { + await send('DELETE', baseUrl + '/_search/point_in_time', { pit_id: pitId }); + } catch (error) { + console.error('Failed to close OpenSearch point-in-time: ' + error.message); + } + } + } +} + +exports.getNativeMetricStats = getNativeMetricStats; + // -------------------------------------------------------------------------------------------------------------- // Generates 1 or more values for 1 or more groups for a metric of a particular source // (tool or benchmark) and type (iops, l2-Gbps, ints/sec, etc). @@ -4111,13 +4313,15 @@ getMetricDataSets = async function (instance, sets, yearDotMonth) { for (var i = 0; i < sets.length; i++) { // Rearrange the actual data into 'values' section - Object.keys(dataSets[i]).forEach((label) => { + Object.keys(dataSets[i]) + .filter((label) => label !== 'distributionStats') + .forEach((label) => { if (isUndefined(dataSets[i].values)) { dataSets[i].values = {}; } dataSets[i].values[label] = dataSets[i][label]; delete dataSets[i][label]; - }); + }); // Build the label-decoder and the remaining breakouts dataSets[i].usedBreakouts = sets[i].breakout; dataSets[i].valueSeriesLabelDecoder = ''; @@ -4148,6 +4352,7 @@ getMetricDataSets = async function (instance, sets, yearDotMonth) { ) ) { delete dataSets[i].values[metric]; + if (dataSets[i].distributionStats) delete dataSets[i].distributionStats[metric]; } }); } diff --git a/queries/cdmq/get-metric-data.js b/queries/cdmq/get-metric-data.js index 95790329..086f0144 100644 --- a/queries/cdmq/get-metric-data.js +++ b/queries/cdmq/get-metric-data.js @@ -163,6 +163,11 @@ async function main() { '[optional] Filter out (do not output) metrics which do not pass the conditional. gt=greater-than, ge=greater-than-or-equal, lt=less-than, le=less-than-or-equal' ) .option('--aggregation ', '[optional] Override the default aggregation method for this query') + .option( + '--distribution-stats ', + '[optional] Return duration-weighted native-timeline statistics (min,max,mean,median,stddev,pNN)', + (value) => value.split(',').map((stat) => stat.trim()).filter(Boolean) + ) .option( '--allow-incompatible-aggregation', '[optional] Allow an aggregation explicitly disallowed by the metric definition' @@ -212,6 +217,7 @@ async function main() { breakout: program.breakout, // Send as array to preserve complex breakout syntax filter: program.filter, aggregation: program.aggregation, + 'distribution-stats': program.distributionStats, 'allow-incompatible-aggregation': program.allowIncompatibleAggregation, instances: program.instances.length > 0 ? program.instances : undefined }; @@ -424,6 +430,19 @@ async function main() { } console.log(line); } + + if (metric_data.distributionStats && program.outputContent != 'headers') { + console.log('\nDistribution statistics (native timeline):'); + Object.keys(metric_data.distributionStats) + .sort((a, b) => a.localeCompare(b, undefined, { numeric: true, sensitivity: 'base' })) + .forEach((label) => { + const stats = metric_data.distributionStats[label]; + const formatted = Object.keys(stats) + .map((stat) => stat + '=' + Number(stats[stat]).toFixed(program.decimalPlaces)) + .join(' '); + console.log(' ' + (label || '') + ': ' + formatted); + }); + } } main(); diff --git a/queries/cdmq/native-stats.js b/queries/cdmq/native-stats.js new file mode 100644 index 00000000..f39fde42 --- /dev/null +++ b/queries/cdmq/native-stats.js @@ -0,0 +1,231 @@ +'use strict'; + +const VALID_AGGREGATIONS = new Set(['sum', 'avg', 'min', 'max']); + +function fail(message, code = 'NATIVE_STATS_DATA_QUALITY') { + const error = new Error('Invalid metric timeline: ' + message); + if (code) error.code = code; + throw error; +} + +function normalizeDocuments(documents, metricId, begin, end) { + if (!Array.isArray(documents) || documents.length === 0) { + fail('metric ID ' + metricId + ' has no documents'); + } + + const clipped = documents + .map((document) => { + if (!Number.isFinite(document.begin) || !Number.isFinite(document.end) || !Number.isFinite(document.value)) { + fail('metric ID ' + metricId + ' contains a non-numeric document'); + } + if (document.begin > document.end) { + fail('metric ID ' + metricId + ' contains a reversed interval'); + } + const clippedBegin = Math.max(begin, document.begin); + const clippedEnd = Math.min(end, document.end); + return clippedBegin <= clippedEnd ? { begin: clippedBegin, end: clippedEnd, value: document.value } : null; + }) + .filter(Boolean) + .sort((left, right) => left.begin - right.begin || left.end - right.end); + + if (clipped.length === 0) { + fail('metric ID ' + metricId + ' does not cover the query range'); + } + + let nextBegin = begin; + clipped.forEach((document) => { + if (document.begin > nextBegin) { + fail('metric ID ' + metricId + ' has a gap'); + } + if (document.begin < nextBegin) { + fail('metric ID ' + metricId + ' has overlapping intervals'); + } + nextBegin = document.end + 1; + }); + if (nextBegin <= end) { + fail('metric ID ' + metricId + ' has a gap'); + } + + return clipped; +} + +function aggregateValues(values, aggregation, metricCount) { + if (aggregation === 'sum') return values.reduce((sum, value) => sum + value, 0); + if (aggregation === 'avg') return values.reduce((sum, value) => sum + value, 0) / metricCount; + if (aggregation === 'min') return Math.min(...values); + return Math.max(...values); +} + +function validateRequestedStats(requestedStats) { + if (!Array.isArray(requestedStats) || requestedStats.length === 0) { + fail('at least one distribution statistic must be requested'); + } + requestedStats.forEach((stat) => { + const percentile = typeof stat === 'string' && /^p(\d{1,3})$/.exec(stat); + if (!['min', 'max', 'mean', 'median', 'stddev'].includes(stat) && (!percentile || Number(percentile[1]) > 100)) { + fail('invalid distribution statistic ' + stat); + } + }); + return requestedStats; +} + +/** + * Reconstruct the piecewise-constant aggregate timeline for selected metric IDs. + * Input documents use CDM's inclusive millisecond intervals. + */ +function reconstructTimeline(documentsById, begin, end, aggregation = 'sum', options = {}) { + if (!Number.isInteger(begin) || !Number.isInteger(end) || begin > end) { + fail('invalid query range'); + } + if (!VALID_AGGREGATIONS.has(aggregation)) { + fail('unsupported aggregation ' + aggregation); + } + + const metricIds = Object.keys(documentsById); + if (metricIds.length === 0) fail('no metric IDs were selected'); + const maxIntervals = options.maxIntervals || Infinity; + + const events = new Map(); + const addEvent = (time, type, metricId, value) => { + if (!events.has(time)) events.set(time, { end: [], start: [] }); + events.get(time)[type].push({ metricId, value }); + }; + + metricIds.forEach((metricId) => { + normalizeDocuments(documentsById[metricId], metricId, begin, end).forEach((document) => { + addEvent(document.begin, 'start', metricId, document.value); + addEvent(document.end + 1, 'end', metricId, document.value); + }); + }); + + const active = new Map(); + const boundaries = [...events.keys()].sort((left, right) => left - right); + const timeline = []; + let previous = begin; + + boundaries.forEach((boundary) => { + if (boundary > end + 1) return; + if (boundary > previous) { + if (active.size !== metricIds.length) { + fail('aggregate timeline has a missing metric value'); + } + if (timeline.length >= maxIntervals) { + fail('native interval limit exceeded', 'NATIVE_STATS_LIMIT'); + } + const values = [...active.values()]; + timeline.push({ + begin: previous, + end: boundary - 1, + duration: boundary - previous, + value: aggregateValues(values, aggregation, metricIds.length) + }); + } + + const event = events.get(boundary); + event.end.forEach(({ metricId }) => active.delete(metricId)); + event.start.forEach(({ metricId, value }) => active.set(metricId, value)); + previous = boundary; + }); + + if (previous <= end) fail('aggregate timeline does not cover the query range'); + return timeline; +} + +function percentile(intervals, quantile, totalDuration) { + const target = quantile * totalDuration; + let cumulative = 0; + const ordered = [...intervals].sort((left, right) => left.value - right.value); + for (const interval of ordered) { + cumulative += interval.duration; + if (cumulative >= target) return interval.value; + } + return ordered[ordered.length - 1].value; +} + +/** Calculate requested duration-weighted statistics over a reconstructed timeline. */ +function calculateDistributionStats(intervals, requestedStats) { + if (!Array.isArray(intervals) || intervals.length === 0) fail('cannot calculate statistics for an empty timeline'); + validateRequestedStats(requestedStats); + + let totalDuration = 0; + let mean = 0; + let m2 = 0; + let minimum = Infinity; + let maximum = -Infinity; + intervals.forEach((interval) => { + if (!Number.isFinite(interval.value) || !Number.isFinite(interval.duration) || interval.duration <= 0) { + fail('timeline contains an invalid interval'); + } + const weight = interval.duration; + const newTotal = totalDuration + weight; + const delta = interval.value - mean; + mean += (weight / newTotal) * delta; + m2 += weight * delta * (interval.value - mean); + totalDuration = newTotal; + minimum = Math.min(minimum, interval.value); + maximum = Math.max(maximum, interval.value); + }); + + const result = {}; + const standardDeviation = Math.sqrt(Math.max(0, m2 / totalDuration)); + requestedStats.forEach((stat) => { + if (stat === 'min') result.min = minimum; + else if (stat === 'max') result.max = maximum; + else if (stat === 'mean') result.mean = mean; + else if (stat === 'stddev') result.stddev = standardDeviation; + else if (stat === 'median') result.median = percentile(intervals, 0.5, totalDuration); + else { + const match = /^p(\d{1,3})$/.exec(stat); + if (!match || Number(match[1]) > 100) fail('invalid statistic ' + stat); + result[stat] = percentile(intervals, Number(match[1]) / 100, totalDuration); + } + }); + return result; +} + +function resampleTimeline(intervals, begin, end, resolution, aggregation) { + if (!Number.isInteger(begin) || !Number.isInteger(end) || begin > end) fail('invalid query range'); + if (!Number.isInteger(resolution) || resolution <= 0) fail('resolution must be a positive integer'); + if (!VALID_AGGREGATIONS.has(aggregation)) fail('unsupported aggregation ' + aggregation); + + const windowDuration = Math.floor((end - begin) / resolution); + if (windowDuration <= 0) fail('resolution is greater than the query duration'); + + const values = []; + let intervalIndex = 0; + let windowBegin = begin; + let windowEnd = begin + windowDuration; + while (windowBegin <= end) { + if (windowEnd > end) windowEnd = end; + let weightedValue = 0; + let totalWeight = 0; + let extreme = aggregation === 'min' ? Infinity : -Infinity; + while (intervalIndex < intervals.length && intervals[intervalIndex].end < windowBegin) intervalIndex++; + for (let index = intervalIndex; index < intervals.length && intervals[index].begin <= windowEnd; index++) { + const interval = intervals[index]; + const overlapBegin = Math.max(windowBegin, interval.begin); + const overlapEnd = Math.min(windowEnd, interval.end); + if (overlapBegin > overlapEnd) continue; + const duration = overlapEnd - overlapBegin + 1; + if (aggregation === 'min') extreme = Math.min(extreme, interval.value); + else if (aggregation === 'max') extreme = Math.max(extreme, interval.value); + else { + weightedValue += interval.value * duration; + totalWeight += duration; + } + } + const value = aggregation === 'min' || aggregation === 'max' ? extreme : weightedValue / totalWeight; + if (!Number.isFinite(value)) fail('native timeline does not cover a resolution window'); + values.push({ begin: windowBegin, end: windowEnd, value: value }); + windowBegin = windowEnd + 1; + windowEnd += windowDuration + 1; + } + return values; +} + +module.exports = { + calculateDistributionStats, + reconstructTimeline, + resampleTimeline, + validateRequestedStats +}; diff --git a/queries/cdmq/package.json b/queries/cdmq/package.json index 06a80194..563c9f36 100644 --- a/queries/cdmq/package.json +++ b/queries/cdmq/package.json @@ -4,7 +4,10 @@ "description": "Query Utils for Common Data Model", "main": "index.js", "scripts": { - "test": "cdmq" + "test": "cdmq", + "test:native-stats": "node --test test/native-stats.test.js", + "test:native-stats:integration": "node --test test/native-stats-opensearch.test.js", + "test:rest": "node --test test/rest-api.test.js" }, "author": "Andrew Theurer", "license": "ISC", diff --git a/queries/cdmq/server.js b/queries/cdmq/server.js index 717ed635..63772c58 100755 --- a/queries/cdmq/server.js +++ b/queries/cdmq/server.js @@ -4,6 +4,7 @@ const fs = require('fs'); const path = require('path'); const app = express(); const cdm = require('./cdm'); +const { validateRequestedStats } = require('./native-stats'); const PORT = process.env.PORT || 3000; const { Command } = require('commander'); const program = new Command(); @@ -1659,6 +1660,7 @@ app.post('/api/v1/metric-data', async (req, res) => { 'breakout', 'filter', 'aggregation', + 'distribution-stats', 'allow-incompatible-aggregation', 'instances' ]; @@ -1677,6 +1679,7 @@ app.post('/api/v1/metric-data', async (req, res) => { breakout, filter, aggregation, + 'distribution-stats': distributionStats, 'allow-incompatible-aggregation': allowIncompatibleAggregation, instances: reqInstances } = req.body; @@ -1689,6 +1692,32 @@ app.post('/api/v1/metric-data', async (req, res) => { }); } + if (typeof distributionStats === 'string') { + distributionStats = distributionStats + .split(',') + .map((stat) => stat.trim()) + .filter(Boolean); + } + if (typeof distributionStats !== 'undefined' && !Array.isArray(distributionStats)) { + return res.status(400).json({ + code: 'INVALID_DISTRIBUTION_STATS', + error: 'distribution-stats must be an array or comma-separated string' + }); + } + if (Array.isArray(distributionStats) && distributionStats.length === 0) { + return res.status(400).json({ + code: 'INVALID_DISTRIBUTION_STATS', + error: 'distribution-stats must contain at least one statistic' + }); + } + if (Array.isArray(distributionStats)) { + try { + validateRequestedStats(distributionStats); + } catch (error) { + return res.status(400).json({ code: 'INVALID_DISTRIBUTION_STATS', error: error.message }); + } + } + var reqStart = Date.now(); var breakoutStr = Array.isArray(breakout) ? breakout @@ -1784,6 +1813,7 @@ app.post('/api/v1/metric-data', async (req, res) => { breakout: breakout, filter: filter, aggregation: aggregation, + distributionStats: distributionStats, allowIncompatibleAggregation: allowIncompatibleAggregation === true }; var resp = await cdm.getMetricDataSets(instance, [set], yearDotMonth); @@ -1806,9 +1836,25 @@ app.post('/api/v1/metric-data', async (req, res) => { res.json(metric_data); } catch (error) { serverError('Error in /api/v1/metric-data:', error); - res.status(500).json({ - code: 'INTERNAL_ERROR', - error: 'Internal server error while fetching metric data', + const status = + error.code === 'NATIVE_STATS_LIMIT' + ? 413 + : error.code === 'NATIVE_STATS_DATA_QUALITY' + ? 422 + : error.code === 'NATIVE_STATS_PIT_UNSUPPORTED' + ? 501 + : 500; + res.status(status).json({ + code: error.code || 'INTERNAL_ERROR', + error: + [ + 'NATIVE_STATS_CONFIG', + 'NATIVE_STATS_DATA_QUALITY', + 'NATIVE_STATS_LIMIT', + 'NATIVE_STATS_PIT_UNSUPPORTED' + ].includes(error.code) + ? error.message + : 'Internal server error while fetching metric data', details: error.message }); } diff --git a/queries/cdmq/test/native-stats-opensearch.test.js b/queries/cdmq/test/native-stats-opensearch.test.js new file mode 100644 index 00000000..e83a0b70 --- /dev/null +++ b/queries/cdmq/test/native-stats-opensearch.test.js @@ -0,0 +1,175 @@ +'use strict'; + +const assert = require('node:assert/strict'); +const test = require('node:test'); +const { getNativeMetricStats } = require('../cdm'); + +function response(body) { + return { ok: true, status: 200, json: async () => body }; +} + +test('streams PIT pages and closes the PIT', async () => { + const requests = []; + const bodies = [ + { pit_id: 'pit-1' }, + { + hits: { + total: { value: 2, relation: 'eq' }, + hits: [ + { + sort: [0, 4, 'a'], + fields: { + 'metric_desc.metric_desc-uuid': ['a'], + 'metric_data.begin': [0], + 'metric_data.end': [4], + 'metric_data.value': [10] + } + }, + { + sort: [0, 4, 'b'], + fields: { + 'metric_desc.metric_desc-uuid': ['b'], + 'metric_data.begin': [0], + 'metric_data.end': [4], + 'metric_data.value': [2] + } + } + ] + } + }, + { hits: { total: { value: 0, relation: 'eq' }, hits: [] } }, + { succeeded: true } + ]; + const fetchImpl = async (url, request) => { + requests.push({ method: request.method, url, request: JSON.parse(request.body || '{}') }); + return response(bodies.shift()); + }; + + const stats = await getNativeMetricStats( + { host: 'opensearch.example', header: { Authorization: 'Basic test' } }, + ['a', 'b'], + 0, + 4, + 'avg', + ['mean', 'stddev'], + '@2026.09', + { fetch: fetchImpl, indexName: 'cdm-v10dev-metric_data@2026.09', pageSize: 2 } + ); + + assert.deepEqual(stats, { mean: 6, stddev: 0 }); + assert.match(requests[0].url, /point_in_time/); + assert.equal(requests[1].request.pit.id, 'pit-1'); + assert.deepEqual(requests[1].request.search_after, undefined); + assert.deepEqual(requests[2].request.search_after, [0, 4, 'b']); + assert.equal(requests.at(-1).method, 'DELETE'); +}); + +test('stops pagination after a short final page', async () => { + const requests = []; + const bodies = [ + { pit_id: 'pit-short' }, + { + hits: { + total: { value: 1, relation: 'eq' }, + hits: [ + { + sort: [0, 4, 'a'], + fields: { + 'metric_desc.metric_desc-uuid': ['a'], + 'metric_data.begin': [0], + 'metric_data.end': [4], + 'metric_data.value': [10] + } + } + ] + } + }, + { succeeded: true } + ]; + const fetchImpl = async (url, request) => { + requests.push({ method: request.method, url, request: JSON.parse(request.body || '{}') }); + return response(bodies.shift()); + }; + + const stats = await getNativeMetricStats({ host: 'opensearch.example' }, ['a'], 0, 4, 'sum', ['mean'], '@2026.09', { + fetch: fetchImpl, + indexName: 'metric_data', + pageSize: 2 + }); + + assert.deepEqual(stats, { mean: 10 }); + assert.deepEqual( + requests.map((request) => request.method), + ['POST', 'POST', 'DELETE'] + ); +}); + +test('closes the PIT when the document limit is exceeded', async () => { + const methods = []; + const fetchImpl = async (url, request) => { + methods.push(request.method); + return response( + methods.length === 1 ? { pit_id: 'pit-2' } : { hits: { total: { value: 2, relation: 'eq' }, hits: [] } } + ); + }; + + await assert.rejects( + getNativeMetricStats({ host: 'opensearch.example' }, ['a'], 0, 4, 'sum', ['mean'], '@2026.09', { + fetch: fetchImpl, + indexName: 'metric_data', + maxDocuments: 1 + }), + { code: 'NATIVE_STATS_LIMIT' } + ); + assert.deepEqual(methods, ['POST', 'POST', 'DELETE']); +}); + +test('rejects invalid statistics before opening a PIT', async () => { + let requestCount = 0; + const fetchImpl = async () => { + requestCount++; + return response({ pit_id: 'unexpected' }); + }; + + await assert.rejects( + getNativeMetricStats({ host: 'opensearch.example' }, ['a'], 0, 4, 'sum', ['p101'], '@2026.09', { + fetch: fetchImpl, + indexName: 'metric_data' + }), + /invalid distribution statistic/ + ); + assert.equal(requestCount, 0); +}); + +test('reports unsupported PIT search clearly', async () => { + const fetchImpl = async () => ({ + ok: false, + status: 404, + json: async () => ({ error: 'missing endpoint' }) + }); + + await assert.rejects( + getNativeMetricStats({ host: 'opensearch.example' }, ['a'], 0, 4, 'sum', ['mean'], '@2026.09', { + fetch: fetchImpl, + indexName: 'metric_data' + }), + (error) => error.code === 'NATIVE_STATS_PIT_UNSUPPORTED' && /Point-in-Time/.test(error.message) + ); +}); + +test('rejects invalid configured resource limits', async () => { + const previous = process.env.CDM_NATIVE_STATS_MAX_DOCUMENTS; + process.env.CDM_NATIVE_STATS_MAX_DOCUMENTS = 'not-a-number'; + try { + await assert.rejects( + getNativeMetricStats({ host: 'opensearch.example' }, ['a'], 0, 4, 'sum', ['mean'], '@2026.09', { + fetch: async () => response({ pit_id: 'unexpected' }), + indexName: 'metric_data' + }), + { code: 'NATIVE_STATS_CONFIG' } + ); + } finally { + if (typeof previous === 'undefined') delete process.env.CDM_NATIVE_STATS_MAX_DOCUMENTS; + else process.env.CDM_NATIVE_STATS_MAX_DOCUMENTS = previous; + } +}); diff --git a/queries/cdmq/test/native-stats.test.js b/queries/cdmq/test/native-stats.test.js new file mode 100644 index 00000000..cb644350 --- /dev/null +++ b/queries/cdmq/test/native-stats.test.js @@ -0,0 +1,183 @@ +'use strict'; + +const assert = require('node:assert/strict'); +const test = require('node:test'); +const { + calculateDistributionStats, + reconstructTimeline, + resampleTimeline, + validateRequestedStats +} = require('../native-stats'); + +test('reconstructs aligned metric intervals and calculates weighted statistics', () => { + const timeline = reconstructTimeline( + { + a: [ + { begin: 0, end: 4, value: 10 }, + { begin: 5, end: 9, value: 20 } + ], + b: [ + { begin: 0, end: 4, value: 2 }, + { begin: 5, end: 9, value: 4 } + ] + }, + 0, + 9, + 'avg' + ); + + assert.deepEqual(timeline, [ + { begin: 0, end: 4, duration: 5, value: 6 }, + { begin: 5, end: 9, duration: 5, value: 12 } + ]); + assert.deepEqual(calculateDistributionStats(timeline, ['min', 'max', 'mean', 'median', 'stddev', 'p95']), { + min: 6, + max: 12, + mean: 9, + median: 6, + stddev: 3, + p95: 12 + }); +}); + +test('handles misaligned intervals and clips to the query range', () => { + const timeline = reconstructTimeline( + { + a: [ + { begin: -2, end: 2, value: 1 }, + { begin: 3, end: 8, value: 3 } + ], + b: [ + { begin: 0, end: 4, value: 10 }, + { begin: 5, end: 8, value: 20 } + ] + }, + 0, + 8, + 'sum' + ); + + assert.deepEqual(timeline, [ + { begin: 0, end: 2, duration: 3, value: 11 }, + { begin: 3, end: 4, duration: 2, value: 13 }, + { begin: 5, end: 8, duration: 4, value: 23 } + ]); + assert.equal(calculateDistributionStats(timeline, ['mean']).mean, (3 * 11 + 2 * 13 + 4 * 23) / 9); +}); + +test('rejects gaps and overlaps', () => { + assert.throws( + () => + reconstructTimeline( + { + a: [ + { begin: 0, end: 1, value: 1 }, + { begin: 3, end: 4, value: 1 } + ] + }, + 0, + 4 + ), + (error) => error.code === 'NATIVE_STATS_DATA_QUALITY' && /has a gap/.test(error.message) + ); + assert.throws( + () => + reconstructTimeline( + { + a: [ + { begin: 0, end: 2, value: 1 }, + { begin: 2, end: 4, value: 1 } + ] + }, + 0, + 4 + ), + (error) => error.code === 'NATIVE_STATS_DATA_QUALITY' && /overlapping intervals/.test(error.message) + ); +}); + +test('returns zero standard deviation for a single interval', () => { + const timeline = reconstructTimeline({ a: [{ begin: 0, end: 99, value: 7 }] }, 0, 99); + assert.deepEqual(calculateDistributionStats(timeline, ['stddev', 'median']), { stddev: 0, median: 7 }); +}); + +test('supports pointwise min and max aggregation', () => { + const documents = { + a: [{ begin: 0, end: 9, value: 10 }], + b: [{ begin: 0, end: 9, value: 3 }] + }; + assert.equal(reconstructTimeline(documents, 0, 9, 'min')[0].value, 3); + assert.equal(reconstructTimeline(documents, 0, 9, 'max')[0].value, 10); +}); + +test('supports all aggregation modes across multiple metric IDs', () => { + const documents = { + a: [ + { begin: 0, end: 4, value: 10 }, + { begin: 5, end: 9, value: 20 } + ], + b: [ + { begin: 0, end: 4, value: 2 }, + { begin: 5, end: 9, value: 4 } + ] + }; + assert.deepEqual( + reconstructTimeline(documents, 0, 9, 'sum').map((interval) => interval.value), + [12, 24] + ); + assert.deepEqual( + reconstructTimeline(documents, 0, 9, 'avg').map((interval) => interval.value), + [6, 12] + ); + assert.deepEqual( + reconstructTimeline(documents, 0, 9, 'min').map((interval) => interval.value), + [2, 4] + ); + assert.deepEqual( + reconstructTimeline(documents, 0, 9, 'max').map((interval) => interval.value), + [10, 20] + ); +}); + +test('uses duration for weighted percentiles', () => { + const intervals = [ + { begin: 0, end: 0, duration: 1, value: 1 }, + { begin: 1, end: 9, duration: 9, value: 10 } + ]; + assert.deepEqual(calculateDistributionStats(intervals, ['mean', 'median', 'p10']), { + mean: 9.1, + median: 10, + p10: 1 + }); +}); + +test('validates requested statistic names and percentile ranges', () => { + assert.deepEqual(validateRequestedStats(['mean', 'p95']), ['mean', 'p95']); + assert.throws(() => validateRequestedStats([]), /at least one/); + assert.throws(() => validateRequestedStats(['p101']), /invalid distribution statistic/); + assert.throws(() => validateRequestedStats(['variance']), /invalid distribution statistic/); +}); + +test('resamples the native timeline for values and preserves aggregation semantics', () => { + const intervals = [ + { begin: 0, end: 4, duration: 5, value: 10 }, + { begin: 5, end: 9, duration: 5, value: 20 } + ]; + assert.deepEqual(resampleTimeline(intervals, 0, 9, 2, 'sum'), [ + { begin: 0, end: 4, value: 10 }, + { begin: 5, end: 9, value: 20 } + ]); + assert.deepEqual(resampleTimeline(intervals, 0, 9, 2, 'min'), [ + { begin: 0, end: 4, value: 10 }, + { begin: 5, end: 9, value: 20 } + ]); + + const misaligned = [ + { begin: 0, end: 2, duration: 3, value: 10 }, + { begin: 3, end: 9, duration: 7, value: 20 } + ]; + assert.deepEqual(resampleTimeline(misaligned, 0, 9, 2, 'avg'), [ + { begin: 0, end: 4, value: 14 }, + { begin: 5, end: 9, value: 20 } + ]); +}); diff --git a/queries/cdmq/test/rest-api.test.js b/queries/cdmq/test/rest-api.test.js new file mode 100644 index 00000000..dff7378a --- /dev/null +++ b/queries/cdmq/test/rest-api.test.js @@ -0,0 +1,63 @@ +'use strict'; + +const assert = require('node:assert/strict'); +const test = require('node:test'); + +const baseUrl = process.env.CDM_TEST_SERVER_URL || 'http://localhost:3000'; + +async function postMetricData(body) { + const response = await fetch(baseUrl + '/api/v1/metric-data', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify(body) + }); + return { status: response.status, body: await response.json() }; +} + +test('CDM server health endpoint responds', async () => { + const response = await fetch(baseUrl + '/health'); + assert.equal(response.status, 200); + assert.equal((await response.json()).status, 'OK'); +}); + +test('metric-data rejects invalid distribution statistics', async () => { + const result = await postMetricData({ 'distribution-stats': ['p101'] }); + assert.equal(result.status, 400); + assert.equal(result.body.code, 'INVALID_DISTRIBUTION_STATS'); +}); + +const requiredQueryVariables = ['CDM_TEST_RUN', 'CDM_TEST_SOURCE', 'CDM_TEST_TYPE', 'CDM_TEST_BEGIN', 'CDM_TEST_END']; +const hasQueryFixture = requiredQueryVariables.every((name) => process.env[name]); + +test('metric-data statistics are independent of output resolution', { skip: !hasQueryFixture }, async () => { + const common = { + run: process.env.CDM_TEST_RUN, + source: process.env.CDM_TEST_SOURCE, + type: process.env.CDM_TEST_TYPE, + begin: Number(process.env.CDM_TEST_BEGIN), + end: Number(process.env.CDM_TEST_END), + 'distribution-stats': ['min', 'max', 'mean', 'median', 'stddev', 'p95'] + }; + const resolutionOne = await postMetricData({ ...common, resolution: 1 }); + const resolutionTen = await postMetricData({ ...common, resolution: 10 }); + assert.equal(resolutionOne.status, 200); + assert.equal(resolutionTen.status, 200); + assert.deepEqual(resolutionOne.body.distributionStats, resolutionTen.body.distributionStats); +}); + +test('metric-data accepts all aggregation modes', { skip: !hasQueryFixture }, async () => { + const common = { + run: process.env.CDM_TEST_RUN, + source: process.env.CDM_TEST_SOURCE, + type: process.env.CDM_TEST_TYPE, + begin: Number(process.env.CDM_TEST_BEGIN), + end: Number(process.env.CDM_TEST_END), + resolution: 1, + 'distribution-stats': ['min', 'max', 'mean', 'stddev'] + }; + for (const aggregation of ['sum', 'avg', 'min', 'max']) { + const result = await postMetricData({ ...common, aggregation }); + assert.equal(result.status, 200, aggregation + ' aggregation failed: ' + JSON.stringify(result.body)); + assert.ok(result.body.distributionStats); + } +}); diff --git a/queries/cdmq/web-ui-requirements.md b/queries/cdmq/web-ui-requirements.md index 20c9c4be..ce7b8dcc 100644 --- a/queries/cdmq/web-ui-requirements.md +++ b/queries/cdmq/web-ui-requirements.md @@ -60,7 +60,8 @@ The UI communicates exclusively with the CDM API server. No direct OpenSearch ac "end": 1770958601744, "resolution": 1, "breakout": ["hostname", "num", "type"], - "filter": "gt:0.01" + "filter": "gt:0.01", + "distribution-stats": ["min", "max", "mean", "median", "stddev", "p95"] } ``` @@ -73,12 +74,25 @@ The UI communicates exclusively with the CDM API server. No direct OpenSearch ac }, "usedBreakouts": ["hostname", "num"], "remainingBreakouts": ["type", "core", "package"], - "valueSeriesLabelDecoder": "--" + "valueSeriesLabelDecoder": "--", + "distributionStats": { + "": { + "min": 0.42, + "max": 0.91, + "mean": 0.68, + "median": 0.67, + "stddev": 0.12, + "p95": 0.89 + } + } } ``` - `resolution=1` returns a single averaged value per label - `resolution=N` returns N time-series datapoints per label +- `distribution-stats` is optional and returns duration-weighted statistics over the native reconstructed timeline, independent of `resolution` +- Supported statistics are `min`, `max`, `mean`, `median`, `stddev`, and exact duration-weighted `pNN` percentiles such as `p95` +- `stddev` is population standard deviation, percentiles use duration-weighted nearest rank, `median` is `p50`, and a single native interval has `stddev=0` - `breakout` controls grouping dimensions; `remainingBreakouts` shows what's still available - Labels encode the breakout values: `-` for breakouts `hostname,num` - Values are numeric (e.g., mpstat values are 0.0–1.0 where 1.0 = 100% busy) diff --git a/queries/cdmq/web-ui/ARCHITECTURE.md b/queries/cdmq/web-ui/ARCHITECTURE.md index b0002fd3..203393c9 100644 --- a/queries/cdmq/web-ui/ARCHITECTURE.md +++ b/queries/cdmq/web-ui/ARCHITECTURE.md @@ -187,7 +187,7 @@ enabling autocomplete dropdowns in the search UI. All accept optional | Method | Endpoint | Body | Returns | |--------|----------|------|---------| -| POST | `/api/v1/metric-data` | `{ run, period, source, type, resolution, breakout, filter }` | `{ values, usedBreakouts, remainingBreakouts }` | +| POST | `/api/v1/metric-data` | `{ run, period, source, type, resolution, breakout, filter, distribution-stats }` | `{ values, usedBreakouts, remainingBreakouts, distributionStats? }` | ### CDM Library Functions Added (`cdm.js`) diff --git a/queries/cdmq/web-ui/DESIGN.md b/queries/cdmq/web-ui/DESIGN.md index 1700ce5b..b4dad9be 100644 --- a/queries/cdmq/web-ui/DESIGN.md +++ b/queries/cdmq/web-ui/DESIGN.md @@ -643,7 +643,7 @@ These endpoints use a `resolveRun` middleware that finds the OpenSearch instance - POST `/api/v1/run/:id/iterations/params`, `/primary-metric`, `/samples`, `/primary-period-name` - POST `/api/v1/run/:id/samples/statuses`, `/primary-period-id` - POST `/api/v1/run/:id/periods/range`, `/metric-types` -- POST `/api/v1/metric-data` +- POST `/api/v1/metric-data` (optionally include `distribution-stats` to obtain duration-weighted native-timeline statistics alongside the resolution series) ---