use opentelemetry::trace::TracerProvider as _;
use opentelemetry_sdk::{
logs::SdkLoggerProvider, metrics::SdkMeterProvider, trace::SdkTracerProvider,
};
use tracing_subscriber::layer::SubscriberExt as _;
use tracing_subscriber::util::SubscriberInitExt as _;
use tracing_subscriber::{EnvFilter, Layer as _};
use crate::{LogFormat, TelemetryArgs};
#[derive(Debug, Clone)]
pub struct Telemetry {
logs: Option<SdkLoggerProvider>,
traces: Option<SdkTracerProvider>,
metrics: Option<SdkMeterProvider>,
drop_behavior: TelemetryDropBehavior,
}
#[derive(Debug, Clone, Copy, Default)]
pub enum TelemetryDropBehavior {
Flush,
#[default]
Shutdown,
}
impl Telemetry {
pub fn flush(&mut self) {
let Self {
logs,
traces,
metrics,
drop_behavior: _,
} = self;
if let Some(logs) = logs {
if let Err(err) = logs.force_flush() {
tracing::error!(%err, "failed to flush otel log provider");
}
}
if let Some(traces) = traces {
if let Err(err) = traces.force_flush() {
tracing::error!(%err, "failed to flush otel trace provider");
}
}
if let Some(metrics) = metrics {
if let Err(err) = metrics.force_flush() {
tracing::error!(%err, "failed to flush otel metric provider");
}
}
}
pub fn shutdown(&mut self) {
self.flush();
let Self {
logs,
traces,
metrics,
drop_behavior: _,
} = self;
if let Some(logs) = logs {
if let Err(err) = logs.shutdown() {
tracing::error!(%err, "failed to shutdown otel log provider");
}
}
if let Some(traces) = traces {
if let Err(err) = traces.shutdown() {
tracing::error!(%err, "failed to shutdown otel trace provider");
}
}
if let Some(metrics) = metrics {
if let Err(err) = metrics.shutdown() {
tracing::error!(%err, "failed to shutdown otel metric provider");
}
}
}
}
impl Drop for Telemetry {
fn drop(&mut self) {
match self.drop_behavior {
TelemetryDropBehavior::Flush => self.flush(),
TelemetryDropBehavior::Shutdown => self.shutdown(),
}
}
}
impl Telemetry {
#[must_use = "dropping this will flush and shutdown all telemetry systems"]
pub fn init(args: TelemetryArgs, drop_behavior: TelemetryDropBehavior) -> anyhow::Result<Self> {
let TelemetryArgs {
tracy_enabled,
enabled,
otel_enabled,
service_name,
attributes,
log_filter,
log_test_output,
log_format,
log_closed_spans,
log_otlp_enabled,
log_endpoint,
trace_filter,
trace_endpoint,
trace_sampler,
trace_sampler_args,
metric_endpoint,
metric_interval,
} = args;
if !enabled {
if tracy_enabled {
#[cfg(feature = "tracy")]
{
tracing_subscriber::registry()
.with(self::tracy::tracy_layer())
.try_init()?;
}
#[cfg(not(feature = "tracy"))]
{
anyhow::bail!(
"`TRACY_ENABLED=true` but the 'tracy' feature flag is not toggled"
);
}
}
return Ok(Self {
logs: None,
metrics: None,
traces: None,
drop_behavior,
});
}
#[expect(unsafe_code)]
unsafe {
std::env::set_var("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT", log_endpoint);
std::env::set_var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT", metric_endpoint);
std::env::set_var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT", trace_endpoint);
std::env::set_var("OTEL_METRIC_EXPORT_INTERVAL", metric_interval);
std::env::set_var("OTEL_RESOURCE_ATTRIBUTES", attributes);
std::env::set_var("OTEL_SERVICE_NAME", &service_name);
std::env::set_var("OTEL_TRACES_SAMPLER", trace_sampler);
std::env::set_var("OTEL_TRACES_SAMPLER_ARG", trace_sampler_args);
}
let create_filter = |base: &str, forced: &str| {
use crate::EnvFilterExt as _;
EnvFilter::new(base)
.add_directive_if_absent(base, "aws_smithy_runtime", forced)?
.add_directive_if_absent(base, "datafusion", forced)?
.add_directive_if_absent(base, "datafusion_optimizer", forced)?
.add_directive_if_absent(base, "h2", forced)?
.add_directive_if_absent(base, "hyper", forced)?
.add_directive_if_absent(base, "hyper_util", forced)?
.add_directive_if_absent(base, "lance-arrow", forced)?
.add_directive_if_absent(base, "lance-core", forced)?
.add_directive_if_absent(base, "lance-datafusion", forced)?
.add_directive_if_absent(base, "lance-encoding", forced)?
.add_directive_if_absent(base, "lance-file", forced)?
.add_directive_if_absent(base, "lance-index", forced)?
.add_directive_if_absent(base, "lance-io", forced)?
.add_directive_if_absent(base, "lance-linalg", forced)?
.add_directive_if_absent(base, "lance-table", forced)?
.add_directive_if_absent(base, "lance", forced)?
.add_directive_if_absent(base, "opentelemetry-otlp", forced)?
.add_directive_if_absent(base, "opentelemetry", forced)?
.add_directive_if_absent(base, "opentelemetry_sdk", forced)?
.add_directive_if_absent(base, "rustls", forced)?
.add_directive_if_absent(base, "sqlparser", forced)?
.add_directive_if_absent(base, "tonic", forced)?
.add_directive_if_absent(base, "tonic_web", forced)?
.add_directive_if_absent(base, "tower", forced)?
.add_directive_if_absent(base, "tower_http", forced)?
.add_directive_if_absent(base, "tower_web", forced)?
.add_directive_if_absent(base, "lance::index", "off")?
.add_directive_if_absent(base, "lance::dataset::scanner", "off")?
.add_directive_if_absent(base, "lance::dataset::builder", "off")?
.add_directive_if_absent(base, "lance_encoding", "off")
};
let layer_logs_and_traces_stdio = {
let layer = tracing_subscriber::fmt::layer()
.with_writer(std::io::stderr)
.with_file(true)
.with_line_number(true)
.with_target(false)
.with_thread_ids(true)
.with_thread_names(true)
.with_span_events(if log_closed_spans {
tracing_subscriber::fmt::format::FmtSpan::CLOSE
} else {
tracing_subscriber::fmt::format::FmtSpan::NONE
});
macro_rules! handle_format {
($format:ident) => {{
let layer = layer.event_format(tracing_subscriber::fmt::format().$format());
if log_test_output {
layer.with_test_writer().boxed()
} else {
layer.boxed()
}
}};
}
let layer = match log_format {
LogFormat::Pretty => handle_format!(pretty),
LogFormat::Compact => handle_format!(compact),
LogFormat::Json => handle_format!(json),
};
layer.with_filter(create_filter(&log_filter, "warn")?)
};
let (logger_provider, layer_logs_otlp) = if otel_enabled && log_otlp_enabled {
use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
let exporter = opentelemetry_otlp::LogExporter::builder()
.with_tonic() .build()?;
let provider = SdkLoggerProvider::builder()
.with_batch_exporter(exporter)
.build();
let layer = OpenTelemetryTracingBridge::new(&provider).boxed();
(
Some(provider),
Some(layer.with_filter(create_filter(&log_filter, "warn")?)),
)
} else {
(None, None)
};
let (tracer_provider, layer_traces_otlp) = if otel_enabled {
let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_tonic() .build()?;
let provider = SdkTracerProvider::builder()
.with_batch_exporter(exporter)
.build();
opentelemetry::global::set_text_map_propagator(
opentelemetry_sdk::propagation::TraceContextPropagator::new(),
);
opentelemetry::global::set_tracer_provider(provider.clone());
let layer = tracing_opentelemetry::layer()
.with_tracer(provider.tracer(service_name.clone()))
.with_filter(create_filter(&trace_filter, "info")?)
.boxed();
(Some(provider), Some(layer))
} else {
(None, None)
};
let metric_provider = if otel_enabled {
let exporter = opentelemetry_otlp::MetricExporter::builder()
.with_temporality(opentelemetry_sdk::metrics::Temporality::Cumulative)
.with_http() .build()?;
let provider = SdkMeterProvider::builder()
.with_periodic_exporter(exporter)
.build();
opentelemetry::global::set_meter_provider(provider.clone());
Some(provider)
} else {
None
};
if tracy_enabled {
#[cfg(feature = "tracy")]
{
tracing::warn!(
"using tracy in addition to standard telemetry stack, consider `TELEMETRY_ENABLED=false`"
);
tracing_subscriber::registry()
.with(layer_logs_otlp)
.with(layer_logs_and_traces_stdio)
.with(layer_traces_otlp)
.with(self::tracy::tracy_layer())
.try_init()?;
}
#[cfg(not(feature = "tracy"))]
{
anyhow::bail!("`TRACY_ENABLED=true` but the 'tracy' feature flag is not toggled");
}
} else {
tracing_subscriber::registry()
.with(layer_logs_otlp)
.with(layer_logs_and_traces_stdio)
.with(layer_traces_otlp)
.try_init()?;
}
Ok(Self {
drop_behavior,
logs: logger_provider,
traces: tracer_provider,
metrics: metric_provider,
})
}
}
#[cfg(feature = "tracy")]
mod tracy {
#[derive(Default)]
pub struct TracyConfig(tracing_subscriber::fmt::format::DefaultFields);
impl tracing_tracy::Config for TracyConfig {
type Formatter = tracing_subscriber::fmt::format::DefaultFields;
fn formatter(&self) -> &Self::Formatter {
&self.0
}
fn format_fields_in_zone_name(&self) -> bool {
false
}
}
pub fn tracy_layer() -> tracing_tracy::TracyLayer<TracyConfig> {
tracing_tracy::TracyLayer::new(TracyConfig::default())
}
}