use super::log::{logger_layer, LoggerConfig};
use opentelemetry::{global, trace::TracerProvider as _, KeyValue};
use opentelemetry_sdk::{trace::SdkTracerProvider, Resource};
use tracing_subscriber::{
filter::filter_fn, layer::SubscriberExt, util::SubscriberInitExt, EnvFilter, Layer, Registry,
};
const DEFAULT_TRACE_FILTER: &str = "off";
#[derive(Debug, Clone, Copy)]
pub struct TelemetryConfig {
service_name: &'static str,
service_version: &'static str,
required_trace_directive: Option<&'static str>,
}
impl TelemetryConfig {
pub const fn new(service_name: &'static str, service_version: &'static str) -> Self {
Self {
service_name,
service_version,
required_trace_directive: None,
}
}
pub const fn requiring_trace_directive(mut self, directive: &'static str) -> Self {
self.required_trace_directive = Some(directive);
self
}
fn trace_export_enabled(self, filter: &str) -> bool {
match self.required_trace_directive {
Some(required) => filter.split(',').map(str::trim).any(|directive| directive == required),
None => filter.trim() != DEFAULT_TRACE_FILTER,
}
}
}
#[derive(Debug, Default)]
pub struct TelemetryGuard {
provider: Option<SdkTracerProvider>,
}
impl Drop for TelemetryGuard {
fn drop(&mut self) {
if let Some(provider) = self.provider.take() {
if let Err(error) = provider.shutdown() {
eprintln!("failed to flush OpenTelemetry traces during shutdown: {error}");
}
}
}
}
pub fn init_telemetry(logger: LoggerConfig, config: TelemetryConfig) -> eyre::Result<TelemetryGuard> {
let endpoint_configured = std::env::var_os("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT").is_some();
let trace_filter = std::env::var("OTEL_TRACES_FILTER").unwrap_or_else(|_| DEFAULT_TRACE_FILTER.to_string());
if !endpoint_configured || !config.trace_export_enabled(&trace_filter) {
tracing_subscriber::registry()
.with(logger_layer(logger))
.try_init()
.map_err(|error| eyre::eyre!("failed to initialize tracing subscriber: {error}"))?;
return Ok(TelemetryGuard::default());
}
let protocol = std::env::var("OTEL_EXPORTER_OTLP_TRACES_PROTOCOL").unwrap_or_else(|_| "http/protobuf".to_string());
if protocol != "http/protobuf" {
eyre::bail!(
"{} OTLP export supports http/protobuf, but configured protocol is {protocol}",
config.service_name
);
}
let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_http()
.build()
.map_err(|error| eyre::eyre!("failed to build OTLP trace exporter: {error}"))?;
let service_name = std::env::var("OTEL_SERVICE_NAME").unwrap_or_else(|_| config.service_name.to_string());
let instance_id = std::env::var("HOSTNAME").unwrap_or_else(|_| format!("pid-{}", std::process::id()));
let resource = Resource::builder()
.with_service_name(service_name)
.with_attributes([
KeyValue::new("service.version", config.service_version),
KeyValue::new("service.instance.id", instance_id),
])
.build();
let provider = SdkTracerProvider::builder()
.with_resource(resource)
.with_batch_exporter(exporter)
.build();
let tracer = provider.tracer(config.service_name);
let otel_layer = tracing_opentelemetry::layer()
.with_tracer(tracer)
.with_filter(EnvFilter::new(trace_filter))
.with_filter(filter_fn(|metadata| metadata.is_span()))
.boxed();
let layers: Vec<Box<dyn Layer<Registry> + Send + Sync>> = vec![logger_layer(logger), otel_layer];
tracing_subscriber::registry()
.with(layers)
.try_init()
.map_err(|error| eyre::eyre!("failed to initialize tracing subscriber: {error}"))?;
global::set_tracer_provider(provider.clone());
Ok(TelemetryGuard {
provider: Some(provider),
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_gate_exports_any_explicitly_enabled_filter() {
let config = TelemetryConfig::new("test-service", "test-version");
assert!(!config.trace_export_enabled("off"));
assert!(config.trace_export_enabled("info"));
assert!(config.trace_export_enabled("off,test_target=debug"));
}
#[test]
fn directive_gate_requires_an_exact_directive() {
let config = TelemetryConfig::new("test-service", "test-version")
.requiring_trace_directive("newton::task_evaluation=debug");
assert!(config.trace_export_enabled("off,newton::task_evaluation=debug"));
assert!(!config.trace_export_enabled("debug"));
assert!(!config.trace_export_enabled("off,newton::task_evaluation=info"));
}
}