use super::MetricsSink;
use dashmap::DashMap;
use opentelemetry_api::metrics::{Counter, Gauge, Histogram, Meter};
use opentelemetry_api::{KeyValue, StringValue};
use std::sync::Arc;
#[derive(Clone, Debug)]
pub struct OtelMetricsSink {
meter: Meter,
counters: Arc<DashMap<&'static str, Counter<u64>>>,
histograms: Arc<DashMap<&'static str, Histogram<f64>>>,
gauges: Arc<DashMap<&'static str, Gauge<f64>>>,
}
impl OtelMetricsSink {
pub fn new(meter: Meter) -> Self {
Self {
meter,
counters: Arc::new(DashMap::new()),
histograms: Arc::new(DashMap::new()),
gauges: Arc::new(DashMap::new()),
}
}
}
#[inline]
fn to_kv(labels: &[(&'static str, &str)]) -> Vec<KeyValue> {
labels
.iter()
.map(|(k, v)| KeyValue::new(*k, StringValue::from(v.to_string())))
.collect()
}
impl MetricsSink for OtelMetricsSink {
fn increment_counter(&self, name: &'static str, value: u64, labels: &[(&'static str, &str)]) {
let counter = self
.counters
.entry(name)
.or_insert_with(|| self.meter.u64_counter(name).build())
.clone();
counter.add(value, &to_kv(labels));
}
fn record_histogram(&self, name: &'static str, value: f64, labels: &[(&'static str, &str)]) {
let hist = self
.histograms
.entry(name)
.or_insert_with(|| self.meter.f64_histogram(name).build())
.clone();
hist.record(value, &to_kv(labels));
}
fn set_gauge(&self, name: &'static str, value: f64, labels: &[(&'static str, &str)]) {
let gauge = self
.gauges
.entry(name)
.or_insert_with(|| self.meter.f64_gauge(name).build())
.clone();
gauge.record(value, &to_kv(labels));
}
}
#[cfg(test)]
mod tests {
use super::*;
use opentelemetry_api::metrics::MeterProvider as _;
use opentelemetry_sdk::metrics::SdkMeterProvider;
fn make_sink() -> OtelMetricsSink {
let provider = SdkMeterProvider::default();
let meter = provider.meter("asx_test");
OtelMetricsSink::new(meter)
}
#[test]
fn counter_does_not_panic() {
let sink = make_sink();
sink.increment_counter("asx_test_messages_total", 1, &[("partner_id", "partner-a")]);
sink.increment_counter("asx_test_messages_total", 5, &[("partner_id", "partner-a")]);
}
#[test]
fn histogram_does_not_panic() {
let sink = make_sink();
sink.record_histogram(
"asx_test_processing_duration_seconds",
0.042,
&[("action", "send")],
);
}
#[test]
fn gauge_does_not_panic() {
let sink = make_sink();
sink.set_gauge("asx_test_queue_depth", 7.0, &[("partner_id", "p2")]);
}
#[test]
fn different_names_use_distinct_instruments() {
let sink = make_sink();
sink.increment_counter("asx_counter_a_total", 1, &[]);
sink.increment_counter("asx_counter_b_total", 1, &[]);
assert_eq!(sink.counters.len(), 2);
}
#[test]
fn same_name_reuses_instrument() {
let sink = make_sink();
for _ in 0..5 {
sink.increment_counter("asx_reuse_test_total", 1, &[]);
}
assert_eq!(sink.counters.len(), 1);
}
}