use std::time::Duration;
use opentelemetry::global;
use opentelemetry::trace::TracerProvider as _;
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::propagation::TraceContextPropagator;
use opentelemetry_sdk::runtime;
use opentelemetry_sdk::trace::span_processor_with_async_runtime::BatchSpanProcessor;
use opentelemetry_sdk::trace::{BatchConfigBuilder, SdkTracerProvider};
use tracing_subscriber::EnvFilter;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
pub struct Telemetry {
provider: Option<SdkTracerProvider>,
}
impl Drop for Telemetry {
fn drop(&mut self) {
if let Some(provider) = self.provider.take()
&& let Err(err) = provider.shutdown()
{
eprintln!("telemetry shutdown: {err}");
}
}
}
pub fn init(service: &'static str) -> Telemetry {
global::set_text_map_propagator(TraceContextPropagator::new());
let provider = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")
.is_ok()
.then(|| {
let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_http()
.build()
.expect("OTLP exporter builds from its environment");
let processor = BatchSpanProcessor::builder(exporter, runtime::Tokio)
.with_batch_config(
BatchConfigBuilder::default()
.with_scheduled_delay(Duration::from_millis(500))
.with_max_queue_size(16_384)
.with_max_export_batch_size(2_048)
.build(),
)
.build();
let provider = SdkTracerProvider::builder()
.with_span_processor(processor)
.with_resource(Resource::builder().with_service_name(service).build())
.build();
global::set_tracer_provider(provider.clone());
provider
});
let export = provider
.as_ref()
.map(|provider| tracing_opentelemetry::layer().with_tracer(provider.tracer(service)));
let filter = EnvFilter::try_from_default_env()
.unwrap_or_else(|_| EnvFilter::new("info"))
.add_directive("opentelemetry=warn".parse().expect("static directive"))
.add_directive("opentelemetry_sdk=warn".parse().expect("static directive"))
.add_directive("opentelemetry-otlp=warn".parse().expect("static directive"))
.add_directive("hyper_util=warn".parse().expect("static directive"))
.add_directive("reqwest=warn".parse().expect("static directive"));
tracing_subscriber::registry()
.with(filter)
.with(tracing_subscriber::fmt::layer().compact())
.with(export)
.init();
Telemetry { provider }
}