use opentelemetry_sdk::metrics::SdkMeterProvider;
pub fn init_otlp_pipeline(
service_name: &str,
) -> Result<SdkMeterProvider, Box<dyn std::error::Error>> {
use opentelemetry::KeyValue as Kv;
use opentelemetry_otlp::MetricExporter;
use opentelemetry_sdk::metrics::PeriodicReader;
use opentelemetry_sdk::Resource;
let exporter = MetricExporter::builder().with_tonic().build()?;
let reader = PeriodicReader::builder(exporter).build();
let resource = Resource::builder()
.with_attributes([Kv::new("service.name", service_name.to_owned())])
.build();
let provider = SdkMeterProvider::builder()
.with_reader(reader)
.with_resource(resource)
.build();
opentelemetry::global::set_meter_provider(provider.clone());
Ok(provider)
}
#[cfg(test)]
mod tests {
use super::init_otlp_pipeline;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use tokio::io::AsyncReadExt as _;
#[tokio::test]
async fn pipeline_exports_to_the_configured_endpoint() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind a stand-in collector");
let endpoint = format!("http://{}", listener.local_addr().expect("addr"));
let saw_preface = Arc::new(AtomicBool::new(false));
let total_bytes = Arc::new(AtomicUsize::new(0));
{
let saw_preface = Arc::clone(&saw_preface);
let total_bytes = Arc::clone(&total_bytes);
tokio::spawn(async move {
while let Ok((mut stream, _)) = listener.accept().await {
let saw_preface = Arc::clone(&saw_preface);
let total_bytes = Arc::clone(&total_bytes);
tokio::spawn(async move {
const PREFACE: &[u8] = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n";
let mut seen = Vec::new();
let mut buf = [0_u8; 4096];
while let Ok(n) = stream.read(&mut buf).await {
if n == 0 {
break;
}
seen.extend_from_slice(&buf[..n]);
total_bytes.store(seen.len(), Ordering::SeqCst);
if seen.starts_with(PREFACE) {
saw_preface.store(true, Ordering::SeqCst);
}
}
});
}
});
}
std::env::set_var("OTEL_EXPORTER_OTLP_ENDPOINT", &endpoint);
let provider = init_otlp_pipeline("a2a-pipeline-test").expect("pipeline builds");
let meter = opentelemetry::global::meter("a2a-pipeline-test");
let counter = meter.u64_counter("a2a_pipeline_probe").build();
counter.add(1, &[]);
let flusher = tokio::task::spawn_blocking(move || {
let _ = provider.force_flush();
});
for _ in 0..100 {
if saw_preface.load(Ordering::SeqCst) && total_bytes.load(Ordering::SeqCst) > 24 {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
flusher.abort();
let bytes = total_bytes.load(Ordering::SeqCst);
assert!(
saw_preface.load(Ordering::SeqCst),
"nothing spoke HTTP/2 to {endpoint} — the pipeline built no \
exporter, or never became the global meter provider. {bytes} \
bytes arrived"
);
assert!(
bytes > 24,
"only the HTTP/2 preface reached {endpoint}: the exporter \
connected but pushed no payload"
);
}
}