use std::sync::{
Arc, Once, OnceLock,
atomic::{AtomicBool, Ordering},
};
use tracing_subscriber::{
EnvFilter,
filter::filter_fn,
layer::{Layer, SubscriberExt},
util::SubscriberInitExt,
};
use crate::span_exporter::{AdkSpanExporter, AdkSpanLayer, is_runtime_span};
pub(crate) static INIT: Once = Once::new();
static ADK_EXPORTER: OnceLock<Arc<AdkSpanExporter>> = OnceLock::new();
static ADK_EXPORTER_INSTALLED: AtomicBool = AtomicBool::new(false);
#[derive(Debug, thiserror::Error)]
pub enum TelemetryError {
#[error("telemetry init failed: {0}")]
Init(String),
}
pub fn init_telemetry(service_name: &str) -> Result<(), TelemetryError> {
INIT.call_once(|| {
let filter = EnvFilter::try_from_default_env()
.or_else(|_| EnvFilter::try_new("info"))
.unwrap_or_else(|_| EnvFilter::new("info"));
tracing_subscriber::registry()
.with(filter)
.with(
tracing_subscriber::fmt::layer()
.with_target(true)
.with_thread_ids(true)
.with_line_number(true),
)
.init();
tracing::info!(service.name = service_name, "telemetry initialized");
});
Ok(())
}
#[cfg(feature = "otlp")]
pub fn init_with_otlp(service_name: &str, endpoint: &str) -> Result<(), TelemetryError> {
use opentelemetry::trace::TracerProvider;
use opentelemetry_otlp::WithExportConfig;
use tracing_opentelemetry::OpenTelemetryLayer;
let endpoint = endpoint.to_string();
let service_name = service_name.to_string();
let init_error: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
INIT.call_once(|| {
let resource = opentelemetry_sdk::Resource::builder_empty()
.with_attributes([opentelemetry::KeyValue::new("service.name", service_name.clone())])
.build();
let span_exporter = match opentelemetry_otlp::SpanExporter::builder()
.with_tonic()
.with_endpoint(&endpoint)
.build()
{
Ok(e) => e,
Err(e) => {
*init_error.lock().unwrap_or_else(|p| p.into_inner()) =
Some(format!("failed to build OTLP span exporter: {e}"));
return;
}
};
let tracer_provider = opentelemetry_sdk::trace::SdkTracerProvider::builder()
.with_batch_exporter(span_exporter)
.with_resource(resource.clone())
.build();
let tracer = tracer_provider.tracer("adk-telemetry");
opentelemetry::global::set_tracer_provider(tracer_provider);
let metric_exporter = match opentelemetry_otlp::MetricExporter::builder()
.with_tonic()
.with_endpoint(&endpoint)
.build()
{
Ok(e) => e,
Err(e) => {
*init_error.lock().unwrap_or_else(|p| p.into_inner()) =
Some(format!("failed to build OTLP metric exporter: {e}"));
return;
}
};
let meter_provider = opentelemetry_sdk::metrics::SdkMeterProvider::builder()
.with_periodic_exporter(metric_exporter)
.with_resource(resource)
.build();
opentelemetry::global::set_meter_provider(meter_provider);
let telemetry_layer = OpenTelemetryLayer::new(tracer);
let filter = EnvFilter::try_from_default_env()
.or_else(|_| EnvFilter::try_new("info"))
.unwrap_or_else(|_| EnvFilter::new("info"));
tracing_subscriber::registry()
.with(
tracing_subscriber::fmt::layer()
.with_target(true)
.with_thread_ids(true)
.with_line_number(true)
.with_filter(filter),
)
.with(telemetry_layer)
.init();
tracing::info!(
service.name = service_name,
otlp.endpoint = %endpoint,
"telemetry initialized with OpenTelemetry"
);
});
if let Some(err) = init_error.lock().unwrap_or_else(|p| p.into_inner()).take() {
return Err(TelemetryError::Init(err));
}
Ok(())
}
#[cfg(feature = "otlp")]
pub(crate) mod otlp_pipeline {
use super::TelemetryError;
use opentelemetry::trace::TracerProvider;
use opentelemetry_otlp::{WithExportConfig, WithTonicConfig};
pub(crate) trait ExporterHook {
fn configure<B: WithTonicConfig>(&self, builder: B) -> B;
}
pub(crate) struct NoopHook;
impl ExporterHook for NoopHook {
fn configure<B: WithTonicConfig>(&self, builder: B) -> B {
builder
}
}
pub(crate) fn build_tracer<H: ExporterHook>(
resource: opentelemetry_sdk::Resource,
endpoint: &str,
hook: &H,
) -> Result<opentelemetry_sdk::trace::SdkTracer, TelemetryError> {
let span_exporter = hook
.configure(
opentelemetry_otlp::SpanExporter::builder().with_tonic().with_endpoint(endpoint),
)
.build()
.map_err(|e| {
TelemetryError::Init(format!("failed to build OTLP span exporter: {e}"))
})?;
let tracer_provider = opentelemetry_sdk::trace::SdkTracerProvider::builder()
.with_batch_exporter(span_exporter)
.with_resource(resource)
.build();
let tracer = tracer_provider.tracer("adk-telemetry");
opentelemetry::global::set_tracer_provider(tracer_provider);
Ok(tracer)
}
}
#[cfg(feature = "otlp")]
pub fn build_otlp_layer<S>(
service_name: &str,
endpoint: &str,
) -> Result<Box<dyn tracing_subscriber::Layer<S> + Send + Sync>, TelemetryError>
where
S: tracing::Subscriber
+ for<'span> tracing_subscriber::registry::LookupSpan<'span>
+ Send
+ Sync,
{
use opentelemetry_otlp::WithExportConfig;
use tracing_opentelemetry::OpenTelemetryLayer;
let resource = opentelemetry_sdk::Resource::builder_empty()
.with_attributes([opentelemetry::KeyValue::new("service.name", service_name.to_string())])
.build();
let tracer = otlp_pipeline::build_tracer(resource.clone(), endpoint, &otlp_pipeline::NoopHook)?;
let metric_exporter = opentelemetry_otlp::MetricExporter::builder()
.with_tonic()
.with_endpoint(endpoint)
.build()
.map_err(|e| TelemetryError::Init(format!("failed to build OTLP metric exporter: {e}")))?;
let meter_provider = opentelemetry_sdk::metrics::SdkMeterProvider::builder()
.with_periodic_exporter(metric_exporter)
.with_resource(resource)
.build();
opentelemetry::global::set_meter_provider(meter_provider);
Ok(Box::new(OpenTelemetryLayer::new(tracer)))
}
pub fn shutdown_telemetry() {
#[cfg(feature = "otlp")]
{
opentelemetry::global::set_tracer_provider(
opentelemetry::trace::noop::NoopTracerProvider::new(),
);
}
}
pub fn init_with_adk_exporter(service_name: &str) -> Result<Arc<AdkSpanExporter>, TelemetryError> {
initialize_adk_exporter_with(&INIT, &ADK_EXPORTER, &ADK_EXPORTER_INSTALLED, |exporter| {
let filter = EnvFilter::try_from_default_env()
.or_else(|_| EnvFilter::try_new("info"))
.unwrap_or_else(|_| EnvFilter::new("info"));
let adk_layer = AdkSpanLayer::new(exporter).with_filter(filter_fn(|metadata| {
metadata.is_span() && is_runtime_span(metadata.name())
}));
tracing_subscriber::registry()
.with(
tracing_subscriber::fmt::layer()
.with_target(true)
.with_thread_ids(true)
.with_line_number(true)
.with_filter(filter),
)
.with(adk_layer)
.init();
tracing::info!(service.name = service_name, "telemetry initialized with ADK span exporter");
})
}
fn initialize_adk_exporter_with<F>(
init: &Once,
exporter_cell: &OnceLock<Arc<AdkSpanExporter>>,
installed: &AtomicBool,
install: F,
) -> Result<Arc<AdkSpanExporter>, TelemetryError>
where
F: FnOnce(Arc<AdkSpanExporter>),
{
let exporter = exporter_cell.get_or_init(|| Arc::new(AdkSpanExporter::new())).clone();
init.call_once(|| {
install(exporter.clone());
installed.store(true, Ordering::Release);
});
if installed.load(Ordering::Acquire) {
Ok(exporter)
} else {
Err(TelemetryError::Init(
"global telemetry was already initialized without the ADK in-process exporter; \
initialize the ADK exporter before other global telemetry modes"
.to_string(),
))
}
}
#[cfg(feature = "sqlite")]
pub fn init_with_sqlite(
service_name: &str,
db_path: impl AsRef<std::path::Path>,
) -> Result<Arc<crate::sqlite::SqliteSpanExporter>, TelemetryError> {
let exporter = Arc::new(crate::sqlite::SqliteSpanExporter::new(db_path)?);
let exporter_clone = exporter.clone();
INIT.call_once(|| {
let filter = EnvFilter::try_from_default_env()
.or_else(|_| EnvFilter::try_new("info"))
.unwrap_or_else(|_| EnvFilter::new("info"));
let adk_layer = AdkSpanLayer::new(exporter_clone).with_filter(filter_fn(|metadata| {
metadata.is_span() && is_runtime_span(metadata.name())
}));
tracing_subscriber::registry()
.with(
tracing_subscriber::fmt::layer()
.with_target(true)
.with_thread_ids(true)
.with_line_number(true)
.with_filter(filter),
)
.with(adk_layer)
.init();
tracing::info!(
service.name = service_name,
"telemetry initialized with SQLite span exporter"
);
});
Ok(exporter)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicUsize;
#[test]
fn repeated_adk_initialization_reuses_the_registered_exporter() {
let init = Once::new();
let exporter = OnceLock::new();
let installed = AtomicBool::new(false);
let installations = AtomicUsize::new(0);
let first = initialize_adk_exporter_with(&init, &exporter, &installed, |_| {
installations.fetch_add(1, Ordering::Relaxed);
})
.unwrap();
let second = initialize_adk_exporter_with(&init, &exporter, &installed, |_| {
installations.fetch_add(1, Ordering::Relaxed);
})
.unwrap();
assert!(Arc::ptr_eq(&first, &second));
assert_eq!(installations.load(Ordering::Relaxed), 1);
}
#[test]
fn adk_initialization_rejects_an_incompatible_existing_global_mode() {
let init = Once::new();
init.call_once(|| {});
let exporter = OnceLock::new();
let installed = AtomicBool::new(false);
let error = initialize_adk_exporter_with(&init, &exporter, &installed, |_| {})
.expect_err("an earlier telemetry mode must not return a disconnected exporter");
assert!(error.to_string().contains("already initialized without the ADK"));
}
}