diff --git a/Cargo.lock b/Cargo.lock index 7ebd6bdb3db..ba3f508486e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -13681,7 +13681,6 @@ dependencies = [ "tower 0.5.2", "tower-http 0.6.11", "tracing", - "tracing-opentelemetry", "urlencoding", "uuid", "vrl", diff --git a/src/common/telemetry/src/logging.rs b/src/common/telemetry/src/logging.rs index 83ceb7aef9d..a9b50930786 100644 --- a/src/common/telemetry/src/logging.rs +++ b/src/common/telemetry/src/logging.rs @@ -18,7 +18,8 @@ mod file_retention; use std::collections::HashMap; use std::env; use std::io::IsTerminal; -use std::sync::{Arc, Mutex, Once, RwLock}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex, Once}; use std::time::Duration; use common_base::readable_size::ReadableSize; @@ -60,35 +61,55 @@ pub static LOG_RELOAD_HANDLE: OnceCell, Registry>; type OtelTraceLayer = tracing_opentelemetry::OpenTelemetryLayer; +struct TraceLayerState { + enabled: AtomicBool, + layer: OnceCell, +} + #[derive(Clone)] pub struct TraceReloadHandle { - inner: Arc>>, + inner: Arc, } impl TraceReloadHandle { - fn new(inner: Arc>>) -> Self { + fn new(inner: Arc) -> Self { Self { inner } } - pub fn reload(&self, new_layer: Option) { - let mut guard = self.inner.write().unwrap(); - *guard = new_layer; - drop(guard); + /// Enables or disables OTLP data collection, initializing it on first enable. + /// Disabling stops new spans, events, fields and links. Existing spans still + /// finish their lifecycle and export on close. + pub fn set_enabled(&self, enabled: bool) -> Result<(), &'static str> { + self.set_enabled_with(enabled, || { + get_or_init_tracer().map(|tracer| tracing_opentelemetry::layer().with_tracer(tracer)) + }) + } + fn set_enabled_with( + &self, + enabled: bool, + init: impl FnOnce() -> Result, + ) -> Result<(), &'static str> { + if enabled { + self.inner.layer.get_or_try_init(init)?; + } + self.inner.enabled.store(enabled, Ordering::Release); callsite::rebuild_interest_cache(); + Ok(()) } } -/// A tracing layer that can be dynamically reloaded. -/// -/// Mostly copied from [`tracing_subscriber::reload::Layer`]. +/// An OTLP layer with a runtime switch and a stable address for downcasts. struct TraceLayer { - inner: Arc>>, + inner: Arc, } impl TraceLayer { fn new(initial: Option) -> (Self, TraceReloadHandle) { - let inner = Arc::new(RwLock::new(initial)); + let inner = Arc::new(TraceLayerState { + enabled: AtomicBool::new(initial.is_some()), + layer: initial.map(OnceCell::with_value).unwrap_or_default(), + }); ( Self { inner: inner.clone(), @@ -98,17 +119,11 @@ impl TraceLayer { } fn with_layer(&self, f: impl FnOnce(&OtelTraceLayer) -> R) -> Option { - self.inner - .read() - .ok() - .and_then(|guard| guard.as_ref().map(f)) + self.inner.layer.get().map(f) } - fn with_layer_mut(&self, f: impl FnOnce(&mut OtelTraceLayer) -> R) -> Option { - self.inner - .write() - .ok() - .and_then(|mut guard| guard.as_mut().map(f)) + fn is_enabled(&self) -> bool { + self.inner.enabled.load(Ordering::Acquire) } } @@ -117,10 +132,6 @@ impl tracing_subscriber::Layer for TraceLayer { let _ = self.with_layer(|layer| layer.on_register_dispatch(subscriber)); } - fn on_layer(&mut self, subscriber: &mut DynSubscriber) { - let _ = self.with_layer_mut(|layer| layer.on_layer(subscriber)); - } - fn register_callsite( &self, metadata: &'static tracing::Metadata<'static>, @@ -144,7 +155,9 @@ impl tracing_subscriber::Layer for TraceLayer { id: &tracing::span::Id, ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>, ) { - let _ = self.with_layer(|layer| layer.on_new_span(attrs, id, ctx)); + if self.is_enabled() { + let _ = self.with_layer(|layer| layer.on_new_span(attrs, id, ctx)); + } } fn max_level_hint(&self) -> Option { @@ -157,7 +170,9 @@ impl tracing_subscriber::Layer for TraceLayer { values: &tracing::span::Record<'_>, ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>, ) { - let _ = self.with_layer(|layer| layer.on_record(span, values, ctx)); + if self.is_enabled() { + let _ = self.with_layer(|layer| layer.on_record(span, values, ctx)); + } } fn on_follows_from( @@ -166,7 +181,9 @@ impl tracing_subscriber::Layer for TraceLayer { follows: &tracing::span::Id, ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>, ) { - let _ = self.with_layer(|layer| layer.on_follows_from(span, follows, ctx)); + if self.is_enabled() { + let _ = self.with_layer(|layer| layer.on_follows_from(span, follows, ctx)); + } } fn event_enabled( @@ -183,7 +200,9 @@ impl tracing_subscriber::Layer for TraceLayer { event: &tracing::Event<'_>, ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>, ) { - let _ = self.with_layer(|layer| layer.on_event(event, ctx)); + if self.is_enabled() { + let _ = self.with_layer(|layer| layer.on_event(event, ctx)); + } } fn on_enter( @@ -220,11 +239,13 @@ impl tracing_subscriber::Layer for TraceLayer { } unsafe fn downcast_raw(&self, id: std::any::TypeId) -> Option<*const ()> { - self.inner.read().ok().and_then(|guard| { - guard - .as_ref() - .and_then(|layer| unsafe { layer.downcast_raw(id) }) - }) + // Keep downcasts available while disabled: an in-flight WithContext + // callback may still need the layer. OnceCell keeps both addresses valid + // for the subscriber's lifetime, even across concurrent toggles. + self.inner + .layer + .get() + .and_then(|layer| unsafe { layer.downcast_raw(id) }) } } @@ -764,8 +785,234 @@ where #[cfg(test)] mod tests { + use std::sync::Barrier; + use std::sync::atomic::AtomicUsize; + + use opentelemetry::trace::TraceContextExt; + use opentelemetry_sdk::trace::SdkTracerProvider; + use tracing_opentelemetry::OpenTelemetrySpanExt; + use super::*; + #[test] + fn test_trace_switch_preserves_active_spans() { + let provider = SdkTracerProvider::builder() + .with_sampler(Sampler::AlwaysOn) + .build(); + for initially_enabled in [false, true] { + let new_layer = + || tracing_opentelemetry::layer().with_tracer(provider.tracer("switch")); + let (filter, _) = tracing_subscriber::reload::Layer::new( + Targets::new().with_default(tracing::Level::INFO), + ); + let (layer, handle) = TraceLayer::new(initially_enabled.then(new_layer)); + let dispatch = tracing::Dispatch::new(Registry::default().with(filter).with(layer)); + tracing::dispatcher::with_default(&dispatch, || { + let initial = tracing::info_span!("initial"); + assert!(!initial.is_disabled()); + assert_eq!( + initial.context().span().span_context().is_valid(), + initially_enabled + ); + handle.set_enabled_with(true, || Ok(new_layer())).unwrap(); + let parent = tracing::info_span!("parent"); + let parent_context = parent.context(); + assert!(parent_context.span().span_context().is_valid()); + let previous_context = opentelemetry::Context::current(); + let entered = parent.enter(); + handle.set_enabled(false).unwrap(); + let disabled = tracing::info_span!("disabled"); + assert!(!disabled.is_disabled()); + assert!(!disabled.context().span().span_context().is_valid()); + assert_eq!( + parent.context().span().span_context(), + parent_context.span().span_context() + ); + assert!(dispatch.downcast_ref::().is_some()); + drop(entered); + assert_eq!( + opentelemetry::Context::current().span().span_context(), + previous_context.span().span_context() + ); + let disabled_entered = disabled.enter(); + handle + .set_enabled_with(true, || panic!("layer must not be replaced")) + .unwrap(); + drop(disabled_entered); + let child = tracing::info_span!(parent: &parent, "child"); + assert_eq!( + child.context().span().span_context().trace_id(), + parent_context.span().span_context().trace_id() + ); + assert!(!disabled.context().span().span_context().is_valid()); + }); + } + } + + #[test] + fn test_trace_switch_retries_initialization() { + let (filter, _) = tracing_subscriber::reload::Layer::new( + Targets::new().with_default(tracing::Level::INFO), + ); + let (layer, handle) = TraceLayer::new(None); + let subscriber = Registry::default().with(filter).with(layer); + tracing::subscriber::with_default(subscriber, || { + handle + .set_enabled_with(false, || panic!("disabling must not initialize OTLP")) + .unwrap(); + assert_eq!( + handle.set_enabled_with(true, || Err("initialization failed")), + Err("initialization failed") + ); + let failed = tracing::info_span!("after_failed_enable"); + assert!(!failed.context().span().span_context().is_valid()); + let provider = SdkTracerProvider::builder() + .with_sampler(Sampler::AlwaysOn) + .build(); + handle + .set_enabled_with(true, || { + Ok(tracing_opentelemetry::layer().with_tracer(provider.tracer("retry"))) + }) + .unwrap(); + let enabled = tracing::info_span!("after_successful_enable"); + assert!(enabled.context().span().span_context().is_valid()); + }); + } + + #[test] + fn test_trace_switch_stops_collecting_fields_when_disabled() { + struct CountFormatting(AtomicUsize); + + impl std::fmt::Debug for CountFormatting { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + self.0.fetch_add(1, Ordering::Relaxed); + f.write_str("value") + } + } + + let value = CountFormatting(AtomicUsize::new(0)); + let provider = SdkTracerProvider::builder() + .with_sampler(Sampler::AlwaysOn) + .build(); + let (filter, _) = tracing_subscriber::reload::Layer::new( + Targets::new().with_default(tracing::Level::INFO), + ); + let (layer, handle) = TraceLayer::new(Some( + tracing_opentelemetry::layer().with_tracer(provider.tracer("fields")), + )); + tracing::subscriber::with_default(Registry::default().with(filter).with(layer), || { + let admitted = tracing::info_span!("admitted", field = tracing::field::Empty); + admitted.record("field", tracing::field::debug(&value)); + admitted.in_scope(|| tracing::info!(value = ?value)); + assert_eq!(value.0.load(Ordering::Relaxed), 2); + + handle.set_enabled(false).unwrap(); + let unadmitted = tracing::info_span!("unadmitted", field = tracing::field::Empty); + for span in [&admitted, &unadmitted] { + span.record("field", tracing::field::debug(&value)); + span.in_scope(|| tracing::info!(value = ?value)); + } + assert_eq!(value.0.load(Ordering::Relaxed), 2); + + handle + .set_enabled_with(true, || panic!("layer must not be replaced")) + .unwrap(); + admitted.record("field", tracing::field::debug(&value)); + admitted.in_scope(|| tracing::info!(value = ?value)); + assert_eq!(value.0.load(Ordering::Relaxed), 4); + }); + } + + #[test] + fn test_trace_switch_finishes_admitted_spans() { + #[derive(Debug)] + struct Exporter(std::sync::mpsc::Sender); + + impl opentelemetry_sdk::trace::SpanExporter for Exporter { + async fn export( + &self, + batch: Vec, + ) -> opentelemetry_sdk::error::OTelSdkResult { + for span in batch { + self.0.send(span.name.into_owned()).unwrap(); + } + Ok(()) + } + } + + let (sender, receiver) = std::sync::mpsc::channel(); + let provider = SdkTracerProvider::builder() + .with_sampler(Sampler::AlwaysOn) + .with_simple_exporter(Exporter(sender)) + .build(); + let (filter, _) = tracing_subscriber::reload::Layer::new( + Targets::new().with_default(tracing::Level::INFO), + ); + let (layer, handle) = TraceLayer::new(Some( + tracing_opentelemetry::layer().with_tracer(provider.tracer("export")), + )); + tracing::subscriber::with_default(Registry::default().with(filter).with(layer), || { + let active = tracing::info_span!("admitted"); + handle.set_enabled(false).unwrap(); + let disabled = tracing::info_span!("not_admitted"); + drop(active); + drop(disabled); + }); + provider.force_flush().unwrap(); + assert_eq!(receiver.try_iter().collect::>(), ["admitted"]); + } + + #[test] + fn test_trace_switch_concurrent_context_access() { + let provider = SdkTracerProvider::builder() + .with_sampler(Sampler::AlwaysOn) + .build(); + let (filter, _) = tracing_subscriber::reload::Layer::new( + Targets::new().with_default(tracing::Level::INFO), + ); + let (layer, handle) = TraceLayer::new(Some( + tracing_opentelemetry::layer().with_tracer(provider.tracer("concurrent")), + )); + let dispatch = tracing::Dispatch::new(Registry::default().with(filter).with(layer)); + let parent = tracing::dispatcher::with_default(&dispatch, || { + tracing::info_span!("concurrent_parent") + }); + let context = parent.context(); + let barrier = Barrier::new(4); + std::thread::scope(|scope| { + for _ in 0..3 { + scope.spawn(|| { + tracing::dispatcher::with_default(&dispatch, || { + let _entered = parent.enter(); + barrier.wait(); + for _ in 0..500 { + assert_eq!( + parent.context().span().span_context(), + context.span().span_context() + ); + { + let child = tracing::info_span!("concurrent_child"); + let _entered = child.enter(); + let _ = child.context(); + tracing::info!("concurrent event"); + } + assert_eq!( + opentelemetry::Context::current().span().span_context(), + context.span().span_context() + ); + } + }); + }); + } + barrier.wait(); + for i in 0..1000 { + handle + .set_enabled_with(i % 2 == 0, || panic!("layer must not be replaced")) + .unwrap(); + } + }); + } + #[test] fn test_logging_options_deserialization_default() { let json = r#"{}"#; diff --git a/src/servers/Cargo.toml b/src/servers/Cargo.toml index 5fa10313e04..8ffd97e2c0f 100644 --- a/src/servers/Cargo.toml +++ b/src/servers/Cargo.toml @@ -140,7 +140,6 @@ tonic-web = "0.14" tower = { workspace = true, features = ["full"] } tower-http = { version = "0.6", features = ["full"] } tracing.workspace = true -tracing-opentelemetry.workspace = true urlencoding = "2.1" uuid.workspace = true vrl.workspace = true diff --git a/src/servers/src/http/dyn_trace.rs b/src/servers/src/http/dyn_trace.rs index dcdb74c56ab..7fe1880cf6b 100644 --- a/src/servers/src/http/dyn_trace.rs +++ b/src/servers/src/http/dyn_trace.rs @@ -14,7 +14,7 @@ use axum::http::StatusCode; use axum::response::IntoResponse; -use common_telemetry::{TRACE_RELOAD_HANDLE, get_or_init_tracer, info}; +use common_telemetry::{TRACE_RELOAD_HANDLE, info}; use crate::error::{InvalidParameterSnafu, Result}; @@ -34,20 +34,14 @@ pub async fn dyn_trace_handler(enable_str: String) -> Result )); }; - if enable { - let tracer = match get_or_init_tracer() { - Ok(tracer) => tracer, - Err(reason) => { - return Ok((StatusCode::SERVICE_UNAVAILABLE, reason.to_string())); - } - }; + if let Err(reason) = trace_reload_handle.set_enabled(enable) { + return Ok((StatusCode::SERVICE_UNAVAILABLE, reason.to_string())); + } - let trace_layer = tracing_opentelemetry::layer().with_tracer(tracer); - trace_reload_handle.reload(Some(trace_layer)); + if enable { info!("trace enabled"); Ok((StatusCode::OK, "trace enabled".to_string())) } else { - trace_reload_handle.reload(None); info!("trace disabled"); Ok((StatusCode::OK, "trace disabled".to_string())) }