diff --git a/CHANGELOG.md b/CHANGELOG.md index d46db5b6eb2..66ef8e85cf0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -34,6 +34,7 @@ For notes on migrating to 2.x / 0.200.x see [the upgrade guide](doc/upgrade-to-2 * perf(sdk-metrics): defer allocation of HrTime to accumulation creation [#6839](https://github.com/open-telemetry/opentelemetry-js/pull/6839) @legendecas * chore(\*): migrate use of sdk-trace-base and sdk-trace-node to sdk-trace [#6851](https://github.com/open-telemetry/opentelemetry-js/pull/6851) @trentm +* perf(sdk-metrics): optionally capture active context for sync instruments [#6848](https://github.com/open-telemetry/opentelemetry-js/pull/6848) @legendecas ## 2.8.0 diff --git a/packages/sdk-metrics/src/Instruments.ts b/packages/sdk-metrics/src/Instruments.ts index f61ead4f812..62711ad7acc 100644 --- a/packages/sdk-metrics/src/Instruments.ts +++ b/packages/sdk-metrics/src/Instruments.ts @@ -16,7 +16,7 @@ import type { ObservableGauge, ObservableUpDownCounter, } from '@opentelemetry/api'; -import { context as contextApi, diag, ValueType } from '@opentelemetry/api'; +import { diag, ValueType } from '@opentelemetry/api'; import type { InstrumentDescriptor } from './InstrumentDescriptor'; import type { ObservableRegistry } from './state/ObservableRegistry'; import type { @@ -39,7 +39,7 @@ export class SyncInstrument { protected _record( value: number, attributes: Attributes = {}, - context: Context = contextApi.active() + context?: Context ) { if (typeof value !== 'number') { diag.warn( diff --git a/packages/sdk-metrics/src/state/AsyncMetricStorage.ts b/packages/sdk-metrics/src/state/AsyncMetricStorage.ts index bd5051c4783..6c3101088c0 100644 --- a/packages/sdk-metrics/src/state/AsyncMetricStorage.ts +++ b/packages/sdk-metrics/src/state/AsyncMetricStorage.ts @@ -28,12 +28,12 @@ export class AsyncMetricStorage> private _aggregationCardinalityLimit?: number; private _deltaMetricStorage: DeltaMetricProcessor; private _temporalMetricStorage: TemporalMetricProcessor; - private _attributesProcessor: IAttributesProcessor; + private _attributesProcessor?: IAttributesProcessor; constructor( _instrumentDescriptor: InstrumentDescriptor, aggregator: Aggregator, - attributesProcessor: IAttributesProcessor, + attributesProcessor: IAttributesProcessor | undefined, collectorHandles: MetricCollectorHandle[], aggregationCardinalityLimit?: number ) { @@ -51,6 +51,10 @@ export class AsyncMetricStorage> } record(measurements: AttributeHashMap, observationTime: HrTime) { + if (this._attributesProcessor === undefined) { + this._deltaMetricStorage.batchCumulate(measurements, observationTime); + return; + } const processed = new AttributeHashMap(); for (const [attributes, value] of measurements.entries()) { processed.set(this._attributesProcessor.process(attributes), value); diff --git a/packages/sdk-metrics/src/state/DeltaMetricProcessor.ts b/packages/sdk-metrics/src/state/DeltaMetricProcessor.ts index 1156a86a449..47079a865d2 100644 --- a/packages/sdk-metrics/src/state/DeltaMetricProcessor.ts +++ b/packages/sdk-metrics/src/state/DeltaMetricProcessor.ts @@ -3,7 +3,7 @@ * SPDX-License-Identifier: Apache-2.0 */ -import type { Context, HrTime, Attributes } from '@opentelemetry/api'; +import type { HrTime, Attributes } from '@opentelemetry/api'; import { millisToHrTime } from '@opentelemetry/core'; import type { Maybe } from '../utils'; import { hashAttributes } from '../utils'; @@ -33,12 +33,7 @@ export class DeltaMetricProcessor> { this._overflowHashCode = hashAttributes(this._overflowAttributes); } - record( - value: number, - attributes: Attributes, - _context: Context, - collectionTime: number - ) { + record(value: number, attributes: Attributes, collectionTime: number) { let accumulation = this._activeCollectionStorage.get(attributes); if (!accumulation) { diff --git a/packages/sdk-metrics/src/state/MeterSharedState.ts b/packages/sdk-metrics/src/state/MeterSharedState.ts index 898cc5ad9a3..4f226ea7a79 100644 --- a/packages/sdk-metrics/src/state/MeterSharedState.ts +++ b/packages/sdk-metrics/src/state/MeterSharedState.ts @@ -20,7 +20,6 @@ import { ObservableRegistry } from './ObservableRegistry'; import { SyncMetricStorage } from './SyncMetricStorage'; import type { Accumulation, Aggregator } from '../aggregator/types'; import type { IAttributesProcessor } from '../view/AttributesProcessor'; -import { createNoopAttributesProcessor } from '../view/AttributesProcessor'; import type { MetricStorage } from './MetricStorage'; /** @@ -166,7 +165,7 @@ export class MeterSharedState { const storage = new MetricStorageType( descriptor, aggregator, - createNoopAttributesProcessor(), + undefined, [collector], cardinalityLimit ) as R; @@ -190,7 +189,7 @@ interface MetricStorageConstructor { new ( instrumentDescriptor: InstrumentDescriptor, aggregator: Aggregator>, - attributesProcessor: IAttributesProcessor, + attributesProcessor: IAttributesProcessor | undefined, collectors: MetricCollectorHandle[], aggregationCardinalityLimit?: number ): MetricStorage; diff --git a/packages/sdk-metrics/src/state/MultiWritableMetricStorage.ts b/packages/sdk-metrics/src/state/MultiWritableMetricStorage.ts index 31323a347da..fae06909830 100644 --- a/packages/sdk-metrics/src/state/MultiWritableMetricStorage.ts +++ b/packages/sdk-metrics/src/state/MultiWritableMetricStorage.ts @@ -4,6 +4,7 @@ */ import type { Context, Attributes } from '@opentelemetry/api'; +import { context as contextApi } from '@opentelemetry/api'; import type { WritableMetricStorage } from './WritableMetricStorage'; /** @@ -11,16 +12,24 @@ import type { WritableMetricStorage } from './WritableMetricStorage'; */ export class MultiMetricStorage implements WritableMetricStorage { private readonly _backingStorages: WritableMetricStorage[]; + readonly hasAttributeProcessor: boolean; + constructor(backingStorages: WritableMetricStorage[]) { this._backingStorages = backingStorages; + this.hasAttributeProcessor = backingStorages.some( + s => s.hasAttributeProcessor + ); } record( value: number, attributes: Attributes, - context: Context, + context: Context | undefined, recordTime: number ) { + if (this.hasAttributeProcessor && context === undefined) { + context = contextApi.active(); + } const storages = this._backingStorages; for (let i = 0; i < storages.length; i++) { storages[i].record(value, attributes, context, recordTime); diff --git a/packages/sdk-metrics/src/state/SyncMetricStorage.ts b/packages/sdk-metrics/src/state/SyncMetricStorage.ts index 98e3568e87f..5a3d443572a 100644 --- a/packages/sdk-metrics/src/state/SyncMetricStorage.ts +++ b/packages/sdk-metrics/src/state/SyncMetricStorage.ts @@ -4,6 +4,7 @@ */ import type { Context, HrTime, Attributes } from '@opentelemetry/api'; +import { context as contextApi } from '@opentelemetry/api'; import type { WritableMetricStorage } from './WritableMetricStorage'; import type { Accumulation, Aggregator } from '../aggregator/types'; import type { InstrumentDescriptor } from '../InstrumentDescriptor'; @@ -27,12 +28,12 @@ export class SyncMetricStorage> private _aggregationCardinalityLimit?: number; private _deltaMetricStorage: DeltaMetricProcessor; private _temporalMetricStorage: TemporalMetricProcessor; - private _attributesProcessor: IAttributesProcessor; + private _attributesProcessor?: IAttributesProcessor; constructor( instrumentDescriptor: InstrumentDescriptor, aggregator: Aggregator, - attributesProcessor: IAttributesProcessor, + attributesProcessor: IAttributesProcessor | undefined, collectorHandles: MetricCollectorHandle[], aggregationCardinalityLimit?: number ) { @@ -47,16 +48,24 @@ export class SyncMetricStorage> collectorHandles ); this._attributesProcessor = attributesProcessor; + this.hasAttributeProcessor = attributesProcessor !== undefined; } + readonly hasAttributeProcessor: boolean; + record( value: number, attributes: Attributes, - context: Context, + context: Context | undefined, recordTime: number ) { - attributes = this._attributesProcessor.process(attributes, context); - this._deltaMetricStorage.record(value, attributes, context, recordTime); + if (this._attributesProcessor !== undefined) { + attributes = this._attributesProcessor.process( + attributes, + context ?? contextApi.active() + ); + } + this._deltaMetricStorage.record(value, attributes, recordTime); } /** diff --git a/packages/sdk-metrics/src/state/WritableMetricStorage.ts b/packages/sdk-metrics/src/state/WritableMetricStorage.ts index af8892752d0..b08719d49bf 100644 --- a/packages/sdk-metrics/src/state/WritableMetricStorage.ts +++ b/packages/sdk-metrics/src/state/WritableMetricStorage.ts @@ -13,11 +13,14 @@ import type { AttributeHashMap } from './HashMap'; * An interface representing SyncMetricStorage with type parameters removed. */ export interface WritableMetricStorage { + /** Whether this storage has an attribute processor that needs context. */ + hasAttributeProcessor: boolean; + /** Records a measurement. */ record( value: number, attributes: Attributes, - context: Context, + context: Context | undefined, recordTime: number ): void; } diff --git a/packages/sdk-metrics/test/performance/benchmark/counter-add.js b/packages/sdk-metrics/test/performance/benchmark/counter-add.js index f75b8cd09fe..0f58fe7d147 100644 --- a/packages/sdk-metrics/test/performance/benchmark/counter-add.js +++ b/packages/sdk-metrics/test/performance/benchmark/counter-add.js @@ -4,12 +4,28 @@ */ const Benchmark = require('benchmark'); -const { MeterProvider } = require('../../../build/src'); +const { + MeterProvider, + createAllowListAttributesProcessor, +} = require('../../../build/src'); const provider = new MeterProvider(); const meter = provider.getMeter('bench'); const counter = meter.createCounter('bench.counter'); +const viewProvider = new MeterProvider({ + views: [ + { + instrumentName: 'bench.counter.view', + attributesProcessors: [ + createAllowListAttributesProcessor(['service', 'route']), + ], + }, + ], +}); +const viewMeter = viewProvider.getMeter('bench'); +const viewCounter = viewMeter.createCounter('bench.counter.view'); + const fixedAttributes = { service: 'checkout', route: '/api/cart', @@ -42,4 +58,8 @@ suite.add('Counter.add (varied attributes, 100 combos)', () => { counter.add(1, variedAttributes[variedIdx++ % 100]); }); +suite.add('Counter.add (with view attribute processor)', () => { + viewCounter.add(1, fixedAttributes); +}); + suite.run(); diff --git a/packages/sdk-metrics/test/state/DeltaMetricProcessor.test.ts b/packages/sdk-metrics/test/state/DeltaMetricProcessor.test.ts index 1c5d68fb992..58a922609cc 100644 --- a/packages/sdk-metrics/test/state/DeltaMetricProcessor.test.ts +++ b/packages/sdk-metrics/test/state/DeltaMetricProcessor.test.ts @@ -3,7 +3,6 @@ * SPDX-License-Identifier: Apache-2.0 */ -import * as api from '@opentelemetry/api'; import * as assert from 'assert'; import { DropAggregator, SumAggregator } from '../../src/aggregator'; import { DeltaMetricProcessor } from '../../src/state/DeltaMetricProcessor'; @@ -17,7 +16,7 @@ describe('DeltaMetricProcessor', () => { for (const value of commonValues) { for (const attributes of commonAttributes) { - metricProcessor.record(value, attributes, api.context.active(), 0); + metricProcessor.record(value, attributes, 0); } } }); @@ -27,7 +26,7 @@ describe('DeltaMetricProcessor', () => { for (const value of commonValues) { for (const attributes of commonAttributes) { - metricProcessor.record(value, attributes, api.context.active(), 0); + metricProcessor.record(value, attributes, 0); } } }); @@ -134,9 +133,9 @@ describe('DeltaMetricProcessor', () => { it('should export', () => { const metricProcessor = new DeltaMetricProcessor(new SumAggregator(true)); - metricProcessor.record(1, { attribute: '1' }, api.ROOT_CONTEXT, 0); - metricProcessor.record(2, { attribute: '1' }, api.ROOT_CONTEXT, 1000); - metricProcessor.record(1, { attribute: '2' }, api.ROOT_CONTEXT, 2000); + metricProcessor.record(1, { attribute: '1' }, 0); + metricProcessor.record(2, { attribute: '1' }, 1000); + metricProcessor.record(1, { attribute: '2' }, 2000); let accumulations = metricProcessor.collect(); assert.strictEqual(accumulations.size, 2); diff --git a/packages/sdk-metrics/test/state/MultiWritableMetricStorage.test.ts b/packages/sdk-metrics/test/state/MultiWritableMetricStorage.test.ts index f22c3074f2e..43d9327ed9c 100644 --- a/packages/sdk-metrics/test/state/MultiWritableMetricStorage.test.ts +++ b/packages/sdk-metrics/test/state/MultiWritableMetricStorage.test.ts @@ -4,15 +4,21 @@ */ import * as api from '@opentelemetry/api'; -import type { Attributes } from '@opentelemetry/api'; +import type { Attributes, Context } from '@opentelemetry/api'; import * as assert from 'assert'; +import { SumAggregator } from '../../src/aggregator'; +import { AggregationTemporality } from '../../src/export/AggregationTemporality'; +import type { MetricCollectorHandle } from '../../src/state/MetricCollector'; import { MultiMetricStorage } from '../../src/state/MultiWritableMetricStorage'; +import { SyncMetricStorage } from '../../src/state/SyncMetricStorage'; import type { WritableMetricStorage } from '../../src/state/WritableMetricStorage'; +import type { IAttributesProcessor } from '../../src/view/AttributesProcessor'; import type { Measurement } from '../util'; import { assertMeasurementEqual, commonAttributes, commonValues, + defaultInstrumentDescriptor, } from '../util'; describe('MultiMetricStorage', () => { @@ -22,13 +28,14 @@ describe('MultiMetricStorage', () => { for (const value of commonValues) { for (const attribute of commonAttributes) { - metricStorage.record(value, attribute, api.context.active(), 0); + metricStorage.record(value, attribute, undefined, 0); } } }); it('record with multiple backing storages', () => { class TestWritableMetricStorage implements WritableMetricStorage { + hasAttributeProcessor = false; records: Measurement[] = []; record( value: number, @@ -68,5 +75,69 @@ describe('MultiMetricStorage', () => { assertMeasurementEqual(backingStorage2.records[idx], expected); } }); + + it('should resolve active context for attribute processors', () => { + const deltaCollector: MetricCollectorHandle = { + selectAggregationTemporality: () => AggregationTemporality.DELTA, + selectCardinalityLimit: () => 2000, + }; + + const processor: IAttributesProcessor = { + process(incoming: Attributes, context?: Context) { + assert.strictEqual(context, api.context.active()); + return incoming; + }, + }; + + const storage1 = new SyncMetricStorage( + defaultInstrumentDescriptor, + new SumAggregator(true), + processor, + [deltaCollector] + ); + const storage2 = new SyncMetricStorage( + defaultInstrumentDescriptor, + new SumAggregator(true), + processor, + [deltaCollector] + ); + const multi = new MultiMetricStorage([storage1, storage2]); + + multi.record(1, {}, undefined, 0); + }); + + it('should pass provided context to attribute processor', () => { + const deltaCollector: MetricCollectorHandle = { + selectAggregationTemporality: () => AggregationTemporality.DELTA, + selectCardinalityLimit: () => 2000, + }; + + const expectedContext = api.ROOT_CONTEXT.setValue( + api.createContextKey('test'), + 'value' + ); + const processor: IAttributesProcessor = { + process(incoming: Attributes, context?: Context) { + assert.strictEqual(context, expectedContext); + return incoming; + }, + }; + + const storage1 = new SyncMetricStorage( + defaultInstrumentDescriptor, + new SumAggregator(true), + processor, + [deltaCollector] + ); + const storage2 = new SyncMetricStorage( + defaultInstrumentDescriptor, + new SumAggregator(true), + processor, + [deltaCollector] + ); + const multi = new MultiMetricStorage([storage1, storage2]); + + multi.record(1, {}, expectedContext, 0); + }); }); }); diff --git a/packages/sdk-metrics/test/state/SyncMetricStorage.test.ts b/packages/sdk-metrics/test/state/SyncMetricStorage.test.ts index d880f295811..7a17bcaea24 100644 --- a/packages/sdk-metrics/test/state/SyncMetricStorage.test.ts +++ b/packages/sdk-metrics/test/state/SyncMetricStorage.test.ts @@ -4,6 +4,7 @@ */ import * as api from '@opentelemetry/api'; +import type { Attributes, Context } from '@opentelemetry/api'; import * as assert from 'assert'; import { SumAggregator } from '../../src/aggregator'; @@ -126,4 +127,44 @@ describe('SyncMetricStorage', () => { }); }); }); + + describe('attribute processor receives context', () => { + it('should pass provided context to attribute processor', () => { + const expectedContext = api.ROOT_CONTEXT.setValue( + api.createContextKey('test'), + 'value' + ); + const attributeProcessor = { + process(incoming: Attributes, context?: Context) { + assert.strictEqual(context, expectedContext); + return incoming; + }, + }; + const metricStorage = new SyncMetricStorage( + defaultInstrumentDescriptor, + new SumAggregator(true), + attributeProcessor, + [deltaCollector] + ); + + metricStorage.record(1, {}, expectedContext, 0); + }); + + it('should resolve active context when context is undefined', () => { + const attributeProcessor = { + process(incoming: Attributes, context?: Context) { + assert.strictEqual(context, api.context.active()); + return incoming; + }, + }; + const metricStorage = new SyncMetricStorage( + defaultInstrumentDescriptor, + new SumAggregator(true), + attributeProcessor, + [deltaCollector] + ); + + metricStorage.record(1, {}, undefined, 0); + }); + }); }); diff --git a/packages/sdk-metrics/test/state/TemporalMetricProcessor.test.ts b/packages/sdk-metrics/test/state/TemporalMetricProcessor.test.ts index 4291179bc12..3388bb1ccdc 100644 --- a/packages/sdk-metrics/test/state/TemporalMetricProcessor.test.ts +++ b/packages/sdk-metrics/test/state/TemporalMetricProcessor.test.ts @@ -3,7 +3,6 @@ * SPDX-License-Identifier: Apache-2.0 */ -import * as api from '@opentelemetry/api'; import * as assert from 'assert'; import * as sinon from 'sinon'; import { SumAggregator } from '../../src/aggregator'; @@ -48,7 +47,7 @@ describe('TemporalMetricProcessor', () => { const temporalMetricStorage = new TemporalMetricProcessor(aggregator, [ deltaCollector1, ]); - deltaMetricStorage.record(1, {}, api.context.active(), 1000); + deltaMetricStorage.record(1, {}, 1000); { const metric = temporalMetricStorage.buildMetrics( deltaCollector1, @@ -67,7 +66,7 @@ describe('TemporalMetricProcessor', () => { assertDataPoint(metric.dataPoints[0], {}, 1, [1, 0], [2, 2]); } - deltaMetricStorage.record(2, {}, api.context.active(), 3000); + deltaMetricStorage.record(2, {}, 3000); { const metric = temporalMetricStorage.buildMetrics( deltaCollector1, @@ -113,7 +112,7 @@ describe('TemporalMetricProcessor', () => { deltaCollector2, ]); - deltaMetricStorage.record(1, {}, api.context.active(), 1000); + deltaMetricStorage.record(1, {}, 1000); { const metric = temporalMetricStorage.buildMetrics( deltaCollector1, @@ -165,7 +164,7 @@ describe('TemporalMetricProcessor', () => { cumulativeCollector1, ]); - deltaMetricStorage.record(1, {}, api.context.active(), 1000); + deltaMetricStorage.record(1, {}, 1000); { const metric = temporalMetricStorage.buildMetrics( cumulativeCollector1, @@ -184,7 +183,7 @@ describe('TemporalMetricProcessor', () => { assertDataPoint(metric.dataPoints[0], {}, 1, [1, 0], [2, 2]); } - deltaMetricStorage.record(2, {}, api.context.active(), 3000); + deltaMetricStorage.record(2, {}, 3000); { const metric = temporalMetricStorage.buildMetrics( cumulativeCollector1, @@ -217,7 +216,7 @@ describe('TemporalMetricProcessor', () => { deltaCollector1, ]); - deltaMetricStorage.record(1, {}, api.context.active(), 1000); + deltaMetricStorage.record(1, {}, 1000); { const metric = temporalMetricStorage.buildMetrics( cumulativeCollector1, @@ -236,7 +235,7 @@ describe('TemporalMetricProcessor', () => { assertDataPoint(metric.dataPoints[0], {}, 1, [1, 0], [2, 2]); } - deltaMetricStorage.record(2, {}, api.context.active(), 3000); + deltaMetricStorage.record(2, {}, 3000); { const metric = temporalMetricStorage.buildMetrics( deltaCollector1,