use crate::config::TelemetryConfig;
use crate::errors::TelemetryError;
use crate::sampling::Signal;
pub(crate) mod adopt;
#[cfg(feature = "otel")]
mod async_runtime;
#[cfg(feature = "otel")]
mod bounded;
#[cfg(feature = "otel")]
mod endpoint;
#[cfg(feature = "otel")]
mod flush;
#[cfg(feature = "otel-grpc")]
mod grpc;
#[cfg(feature = "otel")]
pub(crate) mod logs;
#[cfg(feature = "otel")]
pub(crate) mod metrics;
#[cfg(feature = "otel")]
pub(crate) mod resilient;
#[cfg(feature = "otel")]
mod resource;
#[cfg(feature = "otel")]
pub(crate) mod traces;
#[cfg(feature = "otel")]
pub(crate) fn setup_otel(config: &TelemetryConfig) -> Result<(), TelemetryError> {
let resource = resource::build_resource(config);
traces::install_tracer_provider(config, resource.clone())?;
metrics::install_meter_provider(config, resource.clone())?;
logs::install_logger_provider(config, resource)?;
Ok(())
}
#[cfg(feature = "otel")]
fn map_exporter_build<T, E: std::fmt::Display>(
result: Result<T, E>,
signal: &str,
) -> Result<T, TelemetryError> {
result.map_err(|err| TelemetryError::new(format!("OTLP {signal} exporter build failed: {err}")))
}
#[cfg(not(feature = "otel"))]
pub(crate) fn setup_otel(_config: &TelemetryConfig) -> Result<(), TelemetryError> {
Ok(())
}
#[cfg(feature = "otel")]
pub(crate) use bounded::_reset_abandoned_workers_for_tests;
#[cfg(all(test, feature = "otel"))]
pub(crate) use bounded::drain_deadline;
#[cfg(feature = "otel")]
pub(crate) use bounded::{bounded_flush, bounded_teardown};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum DrainOutcome {
Drained,
Failed,
TimedOut,
}
pub(crate) fn flush_otel(timeout_seconds: Option<f64>) -> DrainOutcome {
flush_otel_by_signal(timeout_seconds).worst()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct SignalDrainOutcomes {
pub logs: DrainOutcome,
pub traces: DrainOutcome,
pub metrics: DrainOutcome,
}
impl SignalDrainOutcomes {
pub(crate) fn worst(self) -> DrainOutcome {
let signals = [self.logs, self.traces, self.metrics];
if signals.contains(&DrainOutcome::TimedOut) {
return DrainOutcome::TimedOut;
}
if signals.contains(&DrainOutcome::Failed) {
return DrainOutcome::Failed;
}
DrainOutcome::Drained
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct SignalDrains {
pub logs: bool,
pub traces: bool,
pub metrics: bool,
}
fn note_blocking_drain() {
if tokio::runtime::Handle::try_current().is_err() {
return;
}
let owned = owned_signals();
if owned.logs {
crate::health::increment_async_blocking_risk(Signal::Logs);
}
if owned.traces {
crate::health::increment_async_blocking_risk(Signal::Traces);
}
if owned.metrics {
crate::health::increment_async_blocking_risk(Signal::Metrics);
}
}
pub(crate) fn flush_otel_by_signal(timeout_seconds: Option<f64>) -> SignalDrainOutcomes {
note_blocking_drain();
#[cfg(feature = "otel")]
{
let drain_logs = || flush::flush_logger_provider(timeout_seconds);
let drain_traces = || flush::flush_tracer_provider(timeout_seconds);
std::thread::scope(|scope| {
let logs = std::thread::Builder::new()
.name("provide-logs-drain".to_string())
.spawn_scoped(scope, drain_logs);
let traces = std::thread::Builder::new()
.name("provide-traces-drain".to_string())
.spawn_scoped(scope, drain_traces);
let metrics = flush::flush_meter_provider(timeout_seconds);
SignalDrainOutcomes {
logs: bounded::join_or_inline(logs, drain_logs),
traces: bounded::join_or_inline(traces, drain_traces),
metrics,
}
})
}
#[cfg(not(feature = "otel"))]
{
let _ = timeout_seconds;
SignalDrainOutcomes {
logs: DrainOutcome::Drained,
traces: DrainOutcome::Drained,
metrics: DrainOutcome::Drained,
}
}
}
pub(crate) fn owned_signals() -> SignalDrains {
#[cfg(feature = "otel")]
{
SignalDrains {
logs: logs::logger_provider_installed(),
traces: traces::tracer_provider_installed(),
metrics: metrics::meter_provider_installed(),
}
}
#[cfg(not(feature = "otel"))]
{
SignalDrains {
logs: false,
traces: false,
metrics: false,
}
}
}
#[cfg(feature = "otel")]
pub(crate) fn traces_provider_effective() -> bool {
crate::runtime::tracing_enabled_by_loaded_config()
&& (traces::tracer_provider_installed() || adopt::traces_adopted())
}
#[cfg(feature = "otel")]
pub(crate) fn metrics_provider_effective() -> bool {
crate::runtime::metrics_enabled_by_loaded_config()
&& (metrics::meter_provider_installed() || adopt::metrics_adopted())
}
pub(crate) fn shutdown_otel(timeout_seconds: Option<f64>) {
note_blocking_drain();
adopt::release_adopted_providers();
#[cfg(feature = "otel")]
{
let teardown_logs = || logs::shutdown_logger_provider(timeout_seconds);
let teardown_metrics = || metrics::shutdown_meter_provider(timeout_seconds);
std::thread::scope(|scope| {
let logs = std::thread::Builder::new()
.name("provide-logs-teardown".to_string())
.spawn_scoped(scope, teardown_logs);
let metrics = std::thread::Builder::new()
.name("provide-metrics-teardown".to_string())
.spawn_scoped(scope, teardown_metrics);
traces::shutdown_tracer_provider(timeout_seconds);
bounded::join_or_inline(logs, teardown_logs);
bounded::join_or_inline(metrics, teardown_metrics);
});
}
#[cfg(not(feature = "otel"))]
{
let _ = timeout_seconds;
}
}
pub(crate) fn otel_installed() -> bool {
#[cfg(feature = "otel")]
{
traces::tracer_provider_installed()
|| metrics::meter_provider_installed()
|| logs::logger_provider_installed()
}
#[cfg(not(feature = "otel"))]
{
false
}
}
pub fn otel_installed_for_tests() -> bool {
otel_installed()
}
pub fn _reset_otel_for_tests() {
shutdown_otel(None);
}
#[cfg(test)]
#[path = "mod_tests.rs"]
mod mod_tests;
#[cfg(all(test, feature = "otel"))]
#[path = "bounded_flush_tests.rs"]
mod bounded_flush_tests;
#[cfg(all(test, feature = "otel"))]
#[path = "async_blocking_risk_tests.rs"]
mod async_blocking_risk_tests;