use std::sync::OnceLock;
pub fn install() -> metrics_exporter_prometheus::PrometheusHandle {
static HANDLE: OnceLock<metrics_exporter_prometheus::PrometheusHandle> = OnceLock::new();
HANDLE
.get_or_init(|| {
#[cfg(feature = "otlp")]
{
match ot_endpoint() {
Ok(Some(endpoint)) => {
return install_with_otlp(endpoint)
.expect("installing the OTLP metric fanout");
}
Err(error) => {
panic!("OTLP metric export is configured but invalid: {error}");
}
Ok(None) => {}
}
}
#[cfg(not(feature = "otlp"))]
{
if std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")
.ok()
.is_some_and(|value| !value.trim().is_empty())
{
panic!(
"OTEL_EXPORTER_OTLP_ENDPOINT is set but this binary was built without \
the `otlp` feature — refusing to silently skip the configured export; \
rebuild with --features otlp or unset the variable"
);
}
}
metrics_exporter_prometheus::PrometheusBuilder::new()
.install_recorder()
.expect("install prometheus recorder")
})
.clone()
}
#[cfg(feature = "otlp")]
fn ot_endpoint() -> anyhow::Result<Option<String>> {
let endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT")
.or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT"))
.ok()
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty());
match endpoint {
Some(endpoint) if !endpoint.contains("://") => Err(anyhow::anyhow!(
"the OTLP endpoint must carry its scheme (http:// or https://): {endpoint}"
)),
other => Ok(other),
}
}
#[cfg(feature = "otlp")]
fn install_with_otlp(
endpoint: String,
) -> anyhow::Result<metrics_exporter_prometheus::PrometheusHandle> {
use opentelemetry_otlp::WithExportConfig;
let exporter = opentelemetry_otlp::MetricExporter::builder()
.with_tonic()
.with_endpoint(endpoint)
.build()
.map_err(|error| anyhow::anyhow!("building the OTLP metric exporter: {error:?}"))?;
let reader = opentelemetry_sdk::metrics::PeriodicReader::builder(exporter).build();
let (_provider, otel_recorder) =
metrics_exporter_opentelemetry::Recorder::builder("exocortex-node")
.with_meter_provider(|builder| builder.with_reader(reader))
.build();
let prometheus = metrics_exporter_prometheus::PrometheusBuilder::new().build_recorder();
let handle = prometheus.handle();
let fanout = metrics_util::layers::FanoutBuilder::default()
.add_recorder(prometheus)
.add_recorder(otel_recorder)
.build();
metrics::set_global_recorder(fanout)
.map_err(|_| anyhow::anyhow!("a global metrics recorder is already installed"))?;
Ok(handle)
}
#[cfg(all(test, feature = "otlp"))]
mod tests {
use opentelemetry_sdk::metrics::data::ResourceMetrics;
use opentelemetry_sdk::metrics::exporter::PushMetricExporter;
use opentelemetry_sdk::metrics::PeriodicReader;
#[derive(Clone, Default)]
struct CapturingExporter {
seen: std::sync::Arc<std::sync::Mutex<Vec<String>>>,
}
impl PushMetricExporter for CapturingExporter {
fn export(
&self,
metrics: &ResourceMetrics,
) -> impl std::future::Future<Output = opentelemetry_sdk::error::OTelSdkResult> + Send
{
let mut names = Vec::new();
for scope in metrics.scope_metrics() {
for metric in scope.metrics() {
names.push(metric.name().to_string());
}
}
self.seen.lock().unwrap().extend(names);
std::future::ready(Ok(()))
}
fn force_flush(&self) -> opentelemetry_sdk::error::OTelSdkResult {
Ok(())
}
fn shutdown_with_timeout(
&self,
_timeout: std::time::Duration,
) -> opentelemetry_sdk::error::OTelSdkResult {
Ok(())
}
fn temporality(&self) -> opentelemetry_sdk::metrics::Temporality {
opentelemetry_sdk::metrics::Temporality::Cumulative
}
}
#[test]
fn the_fanout_delivers_instruments_to_the_otel_leg() {
let capturing = CapturingExporter::default();
let reader = PeriodicReader::builder(capturing.clone()).build();
let (provider, otel_recorder) =
metrics_exporter_opentelemetry::Recorder::builder("exocortex-test")
.with_meter_provider(|builder| builder.with_reader(reader))
.build();
let fanout = metrics_util::layers::FanoutBuilder::default()
.add_recorder(otel_recorder)
.build();
metrics::with_local_recorder(&fanout, || {
metrics::counter!("exocortex_test_otlp_fanout_total").increment(7);
});
provider.force_flush().expect("flush");
assert!(
capturing
.seen
.lock()
.unwrap()
.iter()
.any(|name| name.contains("exocortex_test_otlp_fanout")),
"the counter reached the OTel exporter leg: {:?}",
capturing.seen.lock().unwrap()
);
}
}