diff --git a/packages/core-shared/src/instrumentation/otel/otel-metrics.test.ts b/packages/core-shared/src/instrumentation/otel/otel-metrics.test.ts new file mode 100644 index 0000000..a143803 --- /dev/null +++ b/packages/core-shared/src/instrumentation/otel/otel-metrics.test.ts @@ -0,0 +1,120 @@ +import { describe, it, expect, beforeEach, afterEach } from "vitest"; +import { metrics } from "@opentelemetry/api"; +import { + InMemoryMetricExporter, + MeterProvider, + PeriodicExportingMetricReader, + AggregationTemporality, +} from "@opentelemetry/sdk-metrics"; +import { OtelMetrics } from "./otel-metrics"; + +function setupMeterProvider(): { + exporter: InMemoryMetricExporter; + provider: MeterProvider; +} { + const exporter = new InMemoryMetricExporter(AggregationTemporality.CUMULATIVE); + const reader = new PeriodicExportingMetricReader({ + exporter, + exportIntervalMillis: 100, + }); + const provider = new MeterProvider({ readers: [reader] }); + metrics.setGlobalMeterProvider(provider); + return { exporter, provider }; +} + +describe("OtelMetrics", () => { + let exporter: InMemoryMetricExporter; + let provider: MeterProvider; + + beforeEach(() => { + ({ exporter, provider } = setupMeterProvider()); + }); + + afterEach(async () => { + await provider.shutdown(); + metrics.disable(); + }); + + it("counter() records a counter measurement", async () => { + const otelMetrics = new OtelMetrics(); + otelMetrics.counter("http.requests", 1, { method: "GET" }); + + await provider.forceFlush(); + const metrics_ = exporter.getMetrics(); + const counterMetric = metrics_ + .flatMap((rm) => rm.scopeMetrics) + .flatMap((sm) => sm.metrics) + .find((m) => m.descriptor.name === "http.requests"); + + expect(counterMetric).toBeDefined(); + expect(counterMetric!.dataPoints).toHaveLength(1); + expect(counterMetric!.dataPoints[0]!.value).toBe(1); + }); + + it("counter() defaults value to 1 when omitted", async () => { + const otelMetrics = new OtelMetrics(); + otelMetrics.counter("events.processed"); + + await provider.forceFlush(); + const metrics_ = exporter.getMetrics(); + const counterMetric = metrics_ + .flatMap((rm) => rm.scopeMetrics) + .flatMap((sm) => sm.metrics) + .find((m) => m.descriptor.name === "events.processed"); + + expect(counterMetric).toBeDefined(); + expect(counterMetric!.dataPoints[0]!.value).toBe(1); + }); + + it("histogram() records a histogram measurement", async () => { + const otelMetrics = new OtelMetrics(); + otelMetrics.histogram("http.duration", 250, { route: "/api/me" }); + + await provider.forceFlush(); + const metrics_ = exporter.getMetrics(); + const histogramMetric = metrics_ + .flatMap((rm) => rm.scopeMetrics) + .flatMap((sm) => sm.metrics) + .find((m) => m.descriptor.name === "http.duration"); + + expect(histogramMetric).toBeDefined(); + expect(histogramMetric!.dataPoints).toHaveLength(1); + // Histogram data points carry a Histogram aggregate value with sum/count/buckets. + const dp = histogramMetric!.dataPoints[0] as { + value: { sum?: number; count: number }; + }; + expect(dp.value.sum).toBe(250); + expect(dp.value.count).toBe(1); + }); + + it("gauge() records via UpDownCounter", async () => { + const otelMetrics = new OtelMetrics(); + otelMetrics.gauge("queue.depth", 5, { queue: "emails" }); + + await provider.forceFlush(); + const metrics_ = exporter.getMetrics(); + const gaugeMetric = metrics_ + .flatMap((rm) => rm.scopeMetrics) + .flatMap((sm) => sm.metrics) + .find((m) => m.descriptor.name === "queue.depth"); + + expect(gaugeMetric).toBeDefined(); + expect(gaugeMetric!.dataPoints).toHaveLength(1); + }); + + it("lazily caches instrument instances — same counter object reused across calls", async () => { + const otelMetrics = new OtelMetrics(); + otelMetrics.counter("reuse.test", 1); + otelMetrics.counter("reuse.test", 2); + + await provider.forceFlush(); + const metrics_ = exporter.getMetrics(); + const counterMetric = metrics_ + .flatMap((rm) => rm.scopeMetrics) + .flatMap((sm) => sm.metrics) + .find((m) => m.descriptor.name === "reuse.test"); + + // Cumulative: should accumulate both adds (1+2=3) + expect(counterMetric!.dataPoints[0]!.value).toBe(3); + }); +}); diff --git a/packages/core-shared/src/instrumentation/otel/otel-metrics.ts b/packages/core-shared/src/instrumentation/otel/otel-metrics.ts new file mode 100644 index 0000000..5cb5f2d --- /dev/null +++ b/packages/core-shared/src/instrumentation/otel/otel-metrics.ts @@ -0,0 +1,70 @@ +import { metrics } from "@opentelemetry/api"; +import type { Counter, Histogram, UpDownCounter } from "@opentelemetry/api"; +import type { IMetrics, MetricAttributeValue } from "../metrics.interface"; + +/** + * OTel-backed IMetrics implementation. + * + * - counter → OTel Counter (monotonic, add-only) + * - histogram → OTel Histogram + * - gauge → OTel UpDownCounter (synchronous emit). Known limitation: this + * accumulates deltas, not point-in-time values. True "set to + * absolute" semantics require ObservableGauge with a callback; + * deferred to a v2 interface. + * + * Instrument instances are lazily created and cached per name so repeated + * calls to the same metric name reuse the same OTel instrument. + */ +export class OtelMetrics implements IMetrics { + private readonly meter = metrics.getMeter("@repo/core-shared", "1.0.0"); + private readonly counters = new Map(); + private readonly histograms = new Map(); + private readonly gauges = new Map(); + + counter( + name: string, + value = 1, + attributes?: Record, + ): void { + let counter = this.counters.get(name); + if (!counter) { + counter = this.meter.createCounter(name); + this.counters.set(name, counter); + } + counter.add(value, attributes); + } + + histogram( + name: string, + value: number, + attributes?: Record, + ): void { + let histogram = this.histograms.get(name); + if (!histogram) { + histogram = this.meter.createHistogram(name); + this.histograms.set(name, histogram); + } + histogram.record(value, attributes); + } + + /** + * Records a gauge value via UpDownCounter. + * + * Note: UpDownCounter accumulates a running delta — each call adds to the + * previous value rather than replacing it. This is a synchronous approximation + * of gauge semantics. For true "set to absolute value" behaviour, use an + * ObservableGauge with a periodic callback instead (a future v2 addition). + */ + gauge( + name: string, + value: number, + attributes?: Record, + ): void { + let gauge = this.gauges.get(name); + if (!gauge) { + gauge = this.meter.createUpDownCounter(name); + this.gauges.set(name, gauge); + } + gauge.add(value, attributes); + } +}