From f2d04e82810f7cff117c9c398dcb181b760a2e08 Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Mon, 31 Aug 2026 15:48:01 -0400 Subject: [PATCH 1/2] fix(otlp): emit non-overlapping delta metric windows --- .../metrics/periodic_metric_reader.js | 24 +++++++++----- packages/dd-trace/src/span_stats.js | 11 ++++--- .../test/opentelemetry/metrics.spec.js | 31 +++++++++++++++++++ packages/dd-trace/test/span_stats.spec.js | 26 ++++++++++++++++ 4 files changed, 80 insertions(+), 12 deletions(-) diff --git a/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js b/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js index 3bdf1408f80..72806606b7f 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js +++ b/packages/dd-trace/src/opentelemetry/metrics/periodic_metric_reader.js @@ -28,15 +28,18 @@ const { nowUnixNano } = require('./time') * * @typedef {SumCumulativeState | HistogramCumulativeState} CumulativeStateValue * + * @typedef {{ value: number, timeUnixNano: number }} SumLastExportedState + * * @typedef {{ * count: number, * sum: number, * min?: number, * max?: number, - * bucketCounts: number[] + * bucketCounts: number[], + * timeUnixNano: number * }} HistogramLastExportedState * - * @typedef {number | HistogramLastExportedState} LastExportedStateValue + * @typedef {SumLastExportedState | HistogramLastExportedState} LastExportedStateValue */ /** @@ -431,7 +434,7 @@ class MetricAggregator { } } - this.#applyDeltaTemporality(metricsMap.values(), lastExportedState) + this.#applyDeltaTemporality(metricsMap.values(), lastExportedState, nowUnixNano()) return metricsMap } @@ -475,26 +478,31 @@ class MetricAggregator { * * @param {Iterable} metrics - The metrics to apply delta temporality to * @param {Map} lastExportedState - The last exported state of the metrics + * @param {number} collectionTime - The collection timestamp in nanoseconds * @returns {void} */ - #applyDeltaTemporality (metrics, lastExportedState) { + #applyDeltaTemporality (metrics, lastExportedState, collectionTime) { for (const metric of metrics) { if (metric.temporality === TEMPORALITY.DELTA && this.#isDeltaType(metric.type)) { const scopeKey = this.#getScopeKey(metric.instrumentationScope) for (const dataPoint of metric.dataPointMap.values()) { const stateKey = this.#getStateKey(scopeKey, metric.name, metric.type, dataPoint.attrKey) + dataPoint.timeUnixNano = collectionTime if (metric.type === METRIC_TYPES.COUNTER || metric.type === METRIC_TYPES.OBSERVABLECOUNTER) { - const lastValue = lastExportedState.get(stateKey) || 0 + const lastState = lastExportedState.get(stateKey) const currentValue = dataPoint.value - dataPoint.value = currentValue - lastValue - lastExportedState.set(stateKey, currentValue) + dataPoint.startTimeUnixNano = lastState?.timeUnixNano ?? + dataPoint.startTimeUnixNano ?? dataPoint.timeUnixNano + dataPoint.value = currentValue - (lastState?.value ?? 0) + lastExportedState.set(stateKey, { value: currentValue, timeUnixNano: dataPoint.timeUnixNano }) } else if (metric.type === METRIC_TYPES.HISTOGRAM) { const lastState = lastExportedState.get(stateKey) || { count: 0, sum: 0, bucketCounts: new Array(dataPoint.bucketCounts.length).fill(0), + timeUnixNano: dataPoint.startTimeUnixNano, } const currentState = { count: dataPoint.count, @@ -502,7 +510,9 @@ class MetricAggregator { min: dataPoint.min, max: dataPoint.max, bucketCounts: [...dataPoint.bucketCounts], + timeUnixNano: dataPoint.timeUnixNano, } + dataPoint.startTimeUnixNano = lastState.timeUnixNano dataPoint.count = currentState.count - lastState.count dataPoint.sum = currentState.sum - lastState.sum dataPoint.bucketCounts = currentState.bucketCounts.map( diff --git a/packages/dd-trace/src/span_stats.js b/packages/dd-trace/src/span_stats.js index 1045ce967c3..18d25648d7d 100644 --- a/packages/dd-trace/src/span_stats.js +++ b/packages/dd-trace/src/span_stats.js @@ -251,11 +251,11 @@ class SpanStatsProcessor { * @param {Function} [done] */ forceFlush (done) { - this.#flush(done) + this.#flush(done, true) } - #flush (done) { - const drained = this.#drainBuckets() + #flush (done, force = false) { + const drained = this.#drainBuckets(this.otlpExporter && !force ? Date.now() * 1e6 : Infinity) if (this.enabled && !this.otlpExporter) { this.exporter.export({ @@ -310,12 +310,13 @@ class SpanStatsProcessor { .record(span) } - #drainBuckets () { + #drainBuckets (cutoff) { const drained = [] for (const [timeNs, bucket] of this.buckets.entries()) { + if (timeNs + this.bucketSizeNs > cutoff) continue drained.push({ timeNs, bucket }) + this.buckets.delete(timeNs) } - this.buckets.clear() return drained } diff --git a/packages/dd-trace/test/opentelemetry/metrics.spec.js b/packages/dd-trace/test/opentelemetry/metrics.spec.js index fd11cc31618..ad8226c3eba 100644 --- a/packages/dd-trace/test/opentelemetry/metrics.spec.js +++ b/packages/dd-trace/test/opentelemetry/metrics.spec.js @@ -522,6 +522,37 @@ describe('OpenTelemetry Meter Provider', () => { }, 120) }) + it('advances DELTA start times between exports', () => { + const clock = sinon.useFakeTimers() + const exported = [] + mockOtlpExport((decoded) => { + exported.push(decoded.resourceMetrics[0].scopeMetrics[0].metrics) + }) + + setupMetrics({ OTEL_EXPORTER_OTLP_METRICS_TEMPORALITY_PREFERENCE: 'DELTA' }) + const meter = metrics.getMeter('app') + const counter = meter.createCounter('requests') + const histogram = meter.createHistogram('latency') + + counter.add(5) + histogram.record(10) + clock.tick(100) + counter.add(3) + histogram.record(20) + clock.tick(100) + + assert.strictEqual(exported.length, 2) + for (const name of ['requests', 'latency']) { + const metric = exported.map(metrics => metrics.find(metric => metric.name === name)) + const points = metric.map(metric => (metric.sum || metric.histogram).dataPoints[0]) + assert(points[0].timeUnixNano > points[0].startTimeUnixNano) + assert.strictEqual(points[1].startTimeUnixNano, points[0].timeUnixNano) + assert(points[1].timeUnixNano > points[1].startTimeUnixNano) + } + assert.strictEqual(exported[1].find(metric => metric.name === 'requests').sum.dataPoints[0].asInt, 3) + assert.strictEqual(exported[1].find(metric => metric.name === 'latency').histogram.dataPoints[0].count, 1) + }) + it('LOWMEMORY uses DELTA for sync counters', (done) => { const validator = mockOtlpExport((decoded) => { const counter = decoded.resourceMetrics[0].scopeMetrics[0].metrics[0] diff --git a/packages/dd-trace/test/span_stats.spec.js b/packages/dd-trace/test/span_stats.spec.js index 052f88ff2f5..d9cc33a764e 100644 --- a/packages/dd-trace/test/span_stats.spec.js +++ b/packages/dd-trace/test/span_stats.spec.js @@ -634,6 +634,32 @@ describe('SpanStatsProcessor', () => { assert.strictEqual(bucketSizeNs, p.bucketSizeNs) }) + it('should keep an open OTLP bucket until it is complete', () => { + const clock = sinon.useFakeTimers({ now: 12_345_000 }) + try { + otlpExporter.export.resetHistory() + const p = new SpanStatsProcessor(config, otlpExporter) + clearTimeout(p.timer) + p.onSpanFinished(topLevelSpan) + + p.onInterval() + + assert.ok(otlpExporter.export.notCalled) + assert.strictEqual(p.buckets.size, 1) + + p.onSpanFinished(topLevelSpan) + clock.setSystemTime(12_350_000) + p.onInterval() + + assert.ok(otlpExporter.export.calledOnce) + assert.strictEqual(p.buckets.size, 0) + const [drained] = otlpExporter.export.firstCall.args + assert.strictEqual(drained[0].bucket.values().next().value.topLevelOkDistribution.count, 2) + } finally { + clock.restore() + } + }) + it('should split OTLP trace roots when their attribute is exported', () => { const childSpan = { ...topLevelSpan, parent_id: { equals: () => false } } const processor = new SpanStatsProcessor(config, otlpExporter) From b19bdae70f97b51bdb2beda067a0c40e81eb1993 Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Thu, 17 Sep 2026 15:46:06 -0400 Subject: [PATCH 2/2] fix(otlp): use continuous span metric windows --- .../metrics/otlp_span_stats_transformer.js | 6 ++-- packages/dd-trace/src/span_stats.js | 33 ++++++++++++++----- .../otlp_span_stats_transformer.spec.js | 8 +++-- packages/dd-trace/test/span_stats.spec.js | 33 ++++++++++--------- 4 files changed, 50 insertions(+), 30 deletions(-) diff --git a/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_transformer.js b/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_transformer.js index 4f1aaf980cf..45712a0ac2d 100644 --- a/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_transformer.js +++ b/packages/dd-trace/src/opentelemetry/metrics/otlp_span_stats_transformer.js @@ -68,7 +68,7 @@ class OtlpStatsTransformer extends OtlpTransformerBase { } /** - * @param {Array<{timeNs: number, bucket: import('../../span_stats').SpanBuckets}>} drained + * @param {Array<{timeNs: number, durationNs?: number, bucket: import('../../span_stats').SpanBuckets}>} drained * @param {number} bucketSizeNs */ transform (drained, bucketSizeNs) { @@ -89,9 +89,9 @@ class OtlpStatsTransformer extends OtlpTransformerBase { const dataPoints = [] - for (const { timeNs, bucket } of drained) { + for (const { timeNs, durationNs = bucketSizeNs, bucket } of drained) { const distributions = new Map() - const endTimeNs = timeNs + bucketSizeNs + const endTimeNs = timeNs + durationNs const startNano = isJson ? String(timeNs) : timeNs const endNano = isJson ? String(endTimeNs) : endTimeNs diff --git a/packages/dd-trace/src/span_stats.js b/packages/dd-trace/src/span_stats.js index 18d25648d7d..d45b881c4a3 100644 --- a/packages/dd-trace/src/span_stats.js +++ b/packages/dd-trace/src/span_stats.js @@ -16,6 +16,7 @@ const { const { ORIGIN_KEY, TOP_LEVEL_KEY, SVC_SRC_KEY, GRPC_STATUS_NAMES } = require('./constants') const id = require('./id') const log = require('./log') +const { nowUnixNano } = require('./opentelemetry/metrics/time') const GRPC_STATUS_CODE_MAP = Object.fromEntries(GRPC_STATUS_NAMES.map((name, i) => [name, String(i)])) const ZERO_ID = id('0') @@ -202,6 +203,7 @@ class TimeBuckets extends Map { class SpanStatsProcessor { #config + #otlpStartTimeNs /** * @param {import('./config/config-base')} config @@ -231,6 +233,7 @@ class SpanStatsProcessor { this.hostname = os.hostname() this.enabled = enabled this.otlpExporter = otlpExporter || null + this.#otlpStartTimeNs = otlpExporter ? nowUnixNano() : undefined this.env = env this.#config = config this.sequence = 0 @@ -251,11 +254,12 @@ class SpanStatsProcessor { * @param {Function} [done] */ forceFlush (done) { - this.#flush(done, true) + this.#flush(done) + if (this.otlpExporter) this.timer.refresh() } - #flush (done, force = false) { - const drained = this.#drainBuckets(this.otlpExporter && !force ? Date.now() * 1e6 : Infinity) + #flush (done) { + const drained = this.otlpExporter ? this.#drainOtlpBucket() : this.#drainBuckets() if (this.enabled && !this.otlpExporter) { this.exporter.export({ @@ -293,7 +297,7 @@ class SpanStatsProcessor { this.otlpExporter.export(drained, this.bucketSizeNs, done) } } else if (this.otlpExporter) { - if (typeof this.otlpExporter.flush === 'function') this.otlpExporter.flush(done) + if (done && typeof this.otlpExporter.flush === 'function') this.otlpExporter.flush(done) else done?.() } else done?.() } @@ -302,24 +306,35 @@ class SpanStatsProcessor { if (!this.enabled && !this.otlpExporter) return if (!span.metrics[TOP_LEVEL_KEY] && !span.metrics[MEASURED]) return - const spanEndNs = span.start + span.duration - const bucketTime = spanEndNs - (spanEndNs % this.bucketSizeNs) + let bucketTime = this.#otlpStartTimeNs + if (!this.otlpExporter) { + const spanEndNs = span.start + span.duration + bucketTime = spanEndNs - (spanEndNs % this.bucketSizeNs) + } this.buckets.forTime(bucketTime) .forSpan(span) .record(span) } - #drainBuckets (cutoff) { + #drainBuckets () { const drained = [] for (const [timeNs, bucket] of this.buckets.entries()) { - if (timeNs + this.bucketSizeNs > cutoff) continue drained.push({ timeNs, bucket }) - this.buckets.delete(timeNs) } + this.buckets.clear() return drained } + #drainOtlpBucket () { + const startTimeNs = this.#otlpStartTimeNs + const endTimeNs = Math.max(nowUnixNano(), startTimeNs + 1e3) + const bucket = this.buckets.get(startTimeNs) + this.buckets = new TimeBuckets(true) + this.#otlpStartTimeNs = endTimeNs + return bucket ? [{ timeNs: startTimeNs, durationNs: endTimeNs - startTimeNs, bucket }] : [] + } + #toV06Payload (drained) { const { bucketSizeNs } = this return drained.map(({ timeNs, bucket }) => ({ diff --git a/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_transformer.spec.js b/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_transformer.spec.js index cb914c9c8fa..c3acd5a1c09 100644 --- a/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_transformer.spec.js +++ b/packages/dd-trace/test/opentelemetry/metrics/otlp_span_stats_transformer.spec.js @@ -341,13 +341,15 @@ describe('OtlpStatsTransformer', () => { assert.strictEqual(serviceByResource['GET /bar'], 'svc-other') }) - it('sets timestamps from the bucket time and size', () => { + it('sets timestamps from the collection window', () => { const timeNs = 12340000000000 - const dp = dataPointsOf(JSON.parse(transformer.transform(makeDrained(timeNs, [makeSpan()]), BUCKET_SIZE_NS)))[0] + const drained = makeDrained(timeNs, [makeSpan()]) + drained[0].durationNs = 123456 + const dp = dataPointsOf(JSON.parse(transformer.transform(drained, BUCKET_SIZE_NS)))[0] assert.deepStrictEqual( { start: dp.startTimeUnixNano, end: dp.timeUnixNano }, - { start: String(timeNs), end: String(timeNs + BUCKET_SIZE_NS) } + { start: String(timeNs), end: String(timeNs + drained[0].durationNs) } ) }) diff --git a/packages/dd-trace/test/span_stats.spec.js b/packages/dd-trace/test/span_stats.spec.js index d9cc33a764e..9549e2260cf 100644 --- a/packages/dd-trace/test/span_stats.spec.js +++ b/packages/dd-trace/test/span_stats.spec.js @@ -634,27 +634,30 @@ describe('SpanStatsProcessor', () => { assert.strictEqual(bucketSizeNs, p.bucketSizeNs) }) - it('should keep an open OTLP bucket until it is complete', () => { + it('uses continuous OTLP windows and restarts the interval after force flush', () => { const clock = sinon.useFakeTimers({ now: 12_345_000 }) try { - otlpExporter.export.resetHistory() - const p = new SpanStatsProcessor(config, otlpExporter) - clearTimeout(p.timer) - p.onSpanFinished(topLevelSpan) + const localExporter = { + export: sinon.stub().callsFake((_drained, _bucketSizeNs, done) => done?.()), + flush: sinon.stub().callsFake(done => done?.()), + } + const p = new SpanStatsProcessor(config, localExporter) - p.onInterval() + p.onSpanFinished(topLevelSpan) + clock.tick(3_000) + p.forceFlush(() => {}) + p.onSpanFinished(topLevelSpan) - assert.ok(otlpExporter.export.notCalled) - assert.strictEqual(p.buckets.size, 1) + clock.tick(7_000) + assert.ok(localExporter.export.calledOnce) - p.onSpanFinished(topLevelSpan) - clock.setSystemTime(12_350_000) - p.onInterval() + clock.tick(3_000) + assert.ok(localExporter.export.calledTwice) + const [first] = localExporter.export.firstCall.args[0] + const [second] = localExporter.export.secondCall.args[0] - assert.ok(otlpExporter.export.calledOnce) - assert.strictEqual(p.buckets.size, 0) - const [drained] = otlpExporter.export.firstCall.args - assert.strictEqual(drained[0].bucket.values().next().value.topLevelOkDistribution.count, 2) + assert.strictEqual(second.timeNs, first.timeNs + first.durationNs) + assert.strictEqual(second.durationNs, 10_000 * 1e6) } finally { clock.restore() }