use opentelemetry::trace::TracerProvider as _;
use opentelemetry_otlp::WithExportConfig;
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::trace::SdkTracerProvider;
use serde::{Deserialize, Serialize};
use tracing_opentelemetry::OpenTelemetryLayer;
use tracing_subscriber::Layer as _;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
pub enum OtelTracingProtocol {
#[default]
Grpc,
Http,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(default)]
pub struct OtelTracingConfig {
pub enabled: bool,
pub endpoint: String,
pub protocol: OtelTracingProtocol,
pub service_name: String,
pub sample_ratio: f64,
pub batch_scheduled_delay_ms: u64,
pub batch_max_queue_size: usize,
pub batch_max_export_batch_size: usize,
pub export_timeout_ms: u64,
}
impl Default for OtelTracingConfig {
fn default() -> Self {
Self {
enabled: true,
endpoint: "http://localhost:4317".into(),
protocol: OtelTracingProtocol::Grpc,
service_name: String::new(),
sample_ratio: 0.05,
batch_scheduled_delay_ms: 5_000,
batch_max_queue_size: 2_048,
batch_max_export_batch_size: 512,
export_timeout_ms: 10_000,
}
}
}
impl OtelTracingConfig {
#[must_use]
pub fn from_cascade() -> Self {
#[cfg(feature = "config")]
{
if let Some(cfg) = crate::config::try_get()
&& let Ok(settings) = cfg.unmarshal_key_registered::<Self>("otel_tracing")
{
return settings;
}
}
Self::default()
}
#[must_use]
pub fn is_active(&self) -> bool {
self.enabled && !resolve(self).endpoint.is_empty()
}
}
#[derive(Debug, thiserror::Error)]
pub enum OtelTracingError {
#[error("OTLP {protocol:?} span exporter: {source}")]
ExporterBuild {
protocol: OtelTracingProtocol,
source: opentelemetry_otlp::ExporterBuildError,
},
}
fn resolve(config: &OtelTracingConfig) -> OtelTracingConfig {
let mut resolved = config.clone();
if let Ok(v) = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT") {
resolved.endpoint = v;
}
resolved.endpoint = resolved.endpoint.trim().to_string();
if let Ok(v) = std::env::var("OTEL_EXPORTER_OTLP_PROTOCOL") {
resolved.protocol = match v.as_str() {
"http/protobuf" | "http" => OtelTracingProtocol::Http,
_ => OtelTracingProtocol::Grpc,
};
}
if let Ok(v) = std::env::var("OTEL_SERVICE_NAME") {
resolved.service_name = v;
}
if let Ok(v) = std::env::var("OTEL_TRACES_SAMPLER_ARG")
&& let Ok(ratio) = v.parse::<f64>()
{
resolved.sample_ratio = ratio;
}
resolved
}
fn build_span_exporter(
config: &OtelTracingConfig,
) -> Result<opentelemetry_otlp::SpanExporter, OtelTracingError> {
let timeout = std::time::Duration::from_millis(config.export_timeout_ms);
let result = match config.protocol {
OtelTracingProtocol::Grpc => opentelemetry_otlp::SpanExporter::builder()
.with_tonic()
.with_endpoint(&config.endpoint)
.with_timeout(timeout)
.build(),
OtelTracingProtocol::Http => opentelemetry_otlp::SpanExporter::builder()
.with_http()
.with_endpoint(&config.endpoint)
.with_timeout(timeout)
.build(),
};
result.map_err(|source| OtelTracingError::ExporterBuild {
protocol: config.protocol,
source,
})
}
pub fn build_tracer_layer<S>(
config: &OtelTracingConfig,
) -> Result<
(
OpenTelemetryLayer<S, opentelemetry_sdk::trace::Tracer>,
SdkTracerProvider,
),
OtelTracingError,
>
where
S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
{
let resolved = resolve(config);
let exporter = build_span_exporter(&resolved)?;
let resource = Resource::builder()
.with_service_name(resolved.service_name.clone())
.build();
let batch_config = opentelemetry_sdk::trace::BatchConfigBuilder::default()
.with_max_queue_size(resolved.batch_max_queue_size)
.with_max_export_batch_size(resolved.batch_max_export_batch_size)
.with_scheduled_delay(std::time::Duration::from_millis(
resolved.batch_scheduled_delay_ms,
))
.build();
let exporter = crate::otel_backoff::GatedSpanExporter::new(
exporter,
std::time::Duration::from_millis(resolved.batch_scheduled_delay_ms),
);
let processor = opentelemetry_sdk::trace::BatchSpanProcessor::builder(exporter)
.with_batch_config(batch_config)
.build();
let sampler = opentelemetry_sdk::trace::Sampler::ParentBased(Box::new(
opentelemetry_sdk::trace::Sampler::TraceIdRatioBased(resolved.sample_ratio),
));
let provider = SdkTracerProvider::builder()
.with_span_processor(processor)
.with_sampler(sampler)
.with_resource(resource)
.build();
let tracer = provider.tracer("scalo");
opentelemetry::global::set_tracer_provider(provider.clone());
let layer = tracing_opentelemetry::layer().with_tracer(tracer);
Ok((layer, provider))
}
const SELF_TELEMETRY_TARGETS: [&str; 8] = [
"opentelemetry",
"opentelemetry_sdk",
"h2",
"hyper",
"hyper_util",
"tonic",
"tower",
"reqwest",
];
fn is_self_telemetry(target: &str) -> bool {
SELF_TELEMETRY_TARGETS.iter().any(|crate_name| {
target
.strip_prefix(crate_name)
.is_some_and(|rest| rest.is_empty() || rest.starts_with("::"))
})
}
type SelfTelemetryFilter = tracing_subscriber::filter::FilterFn<fn(&tracing::Metadata<'_>) -> bool>;
fn keep_out_of_export(meta: &tracing::Metadata<'_>) -> bool {
!is_self_telemetry(meta.target())
}
fn self_telemetry_filter() -> SelfTelemetryFilter {
tracing_subscriber::filter::FilterFn::new(
keep_out_of_export as fn(&tracing::Metadata<'_>) -> bool,
)
}
static TRACER_PROVIDER: std::sync::OnceLock<SdkTracerProvider> = std::sync::OnceLock::new();
static INIT_STATUS: std::sync::OnceLock<String> = std::sync::OnceLock::new();
pub fn layer_if_active<S>(
config: &OtelTracingConfig,
) -> Option<
tracing_subscriber::filter::Filtered<
OpenTelemetryLayer<S, opentelemetry_sdk::trace::Tracer>,
SelfTelemetryFilter,
S,
>,
>
where
S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
{
if !config.is_active() {
let _ = INIT_STATUS.set(
"OTLP span export disabled (otel_tracing.enabled=false or blank endpoint)".to_string(),
);
return None;
}
match build_tracer_layer(config) {
Ok((layer, provider)) => {
let _ = TRACER_PROVIDER.set(provider);
let resolved = resolve(config);
let _ = INIT_STATUS.set(format!(
"OTLP span export enabled -> {} (sample_ratio {}, max_queue {})",
resolved.endpoint, resolved.sample_ratio, resolved.batch_max_queue_size
));
Some(layer.with_filter(self_telemetry_filter()))
}
Err(e) => {
let _ = INIT_STATUS.set(format!(
"OTLP span export unavailable, continuing without it: {e}"
));
None
}
}
}
pub fn log_init_status() {
if let Some(status) = INIT_STATUS.get() {
tracing::info!("{status}");
}
}
pub fn shutdown() {
if let Some(provider) = TRACER_PROVIDER.get()
&& let Err(e) = provider.shutdown_with_timeout(SHUTDOWN_TIMEOUT)
{
tracing::debug!(error = %e, "OTel tracer provider shutdown");
}
}
const SHUTDOWN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn config_default_round_trip() {
let cfg = OtelTracingConfig::default();
assert_eq!(cfg.protocol, OtelTracingProtocol::Grpc);
assert!(!cfg.endpoint.is_empty());
assert!(cfg.enabled, "span export is on by default");
assert!(cfg.service_name.is_empty());
}
#[test]
fn export_is_off_when_disabled_or_unaddressed() {
assert!(OtelTracingConfig::default().is_active());
assert!(
!OtelTracingConfig {
enabled: false,
..OtelTracingConfig::default()
}
.is_active(),
"enabled=false must switch it off"
);
assert!(
!OtelTracingConfig {
endpoint: " ".to_string(),
..OtelTracingConfig::default()
}
.is_active(),
"a blank endpoint must switch it off"
);
}
#[test]
fn batch_defaults_bound_what_can_be_held() {
let cfg = OtelTracingConfig::default();
assert!(
cfg.batch_max_queue_size > 0,
"an unbounded queue would let a dead collector grow memory"
);
assert!(
cfg.batch_max_export_batch_size <= cfg.batch_max_queue_size,
"the SDK rejects a batch size above the queue size"
);
assert!(
cfg.sample_ratio > 0.0 && cfg.sample_ratio <= 1.0,
"sample ratio {} is outside 0..1",
cfg.sample_ratio
);
}
#[test]
fn the_export_path_is_kept_out_of_the_export() {
for target in [
"hyper",
"hyper::client::conn",
"h2::codec",
"tonic::transport::channel",
"opentelemetry_sdk::trace::span_processor",
"tower::buffer",
] {
assert!(
is_self_telemetry(target),
"{target} is on the export path and must not be exported"
);
}
}
#[test]
fn app_targets_that_merely_share_a_prefix_are_still_exported() {
for target in [
"hyperion",
"hyperi_thing::worker",
"towerbridge",
"h2o",
"reqwest_middleware_of_ours",
"dfe_receiver::ingest",
] {
assert!(
!is_self_telemetry(target),
"{target} is application telemetry and must still be exported"
);
}
}
#[test]
fn the_sampler_arg_env_var_overrides_the_ratio() {
temp_env::with_var("OTEL_TRACES_SAMPLER_ARG", Some("1.0"), || {
let r = resolve(&OtelTracingConfig::default());
assert!(
(r.sample_ratio - 1.0).abs() < f64::EPSILON,
"OTEL_TRACES_SAMPLER_ARG must win over the configured ratio"
);
});
}
#[test]
fn resolve_picks_up_env_overrides() {
temp_env::with_vars(
[
(
"OTEL_EXPORTER_OTLP_ENDPOINT",
Some("http://my-collector:4317"),
),
("OTEL_EXPORTER_OTLP_PROTOCOL", Some("http/protobuf")),
("OTEL_SERVICE_NAME", Some("test-service")),
],
|| {
let r = resolve(&OtelTracingConfig::default());
assert_eq!(r.endpoint, "http://my-collector:4317");
assert_eq!(r.protocol, OtelTracingProtocol::Http);
assert_eq!(r.service_name, "test-service");
},
);
}
}