#[cfg(not(any(test, feature = "test-utils")))]
use std::io::IsTerminal;
use std::path::{Path, PathBuf};
use thiserror::Error;
use tracing::level_filters::LevelFilter;
use tracing_appender::non_blocking::WorkerGuard;
use tracing_subscriber::fmt::format::FmtSpan;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::{SubscriberInitExt, TryInitError};
use tracing_subscriber::{EnvFilter, Layer, Registry};
#[derive(Debug, Error)]
pub enum TelemetryError {
#[error("OXEN_LOG_DIR set but is empty, cannot enable JSON file logging.")]
EmptyLogDir,
#[error("Requested JSON file logging cannot be enabled because OXEN_LOG_DIR is a file: {0}")]
LogDirIsFile(PathBuf),
#[error("Failed to create log directory ({0}): {1}")]
CreateLogDir(PathBuf, std::io::Error),
#[error("Failed to initialize tracing: {0}")]
InitFail(#[from] TryInitError),
#[cfg(feature = "otel")]
#[error("Unknown {var} value: {value} (expected grpc, http, http/protobuf, or http/json)")]
UnknownProtocol { var: &'static str, value: String },
}
pub type BoxedLayer = Box<dyn Layer<Registry> + Send + Sync>;
pub struct TracingGuard {
_file_guard: Option<WorkerGuard>,
#[cfg(feature = "otel")]
_tracer_provider: Option<opentelemetry_sdk::trace::SdkTracerProvider>,
}
impl TracingGuard {
pub async fn shutdown(&self) {
#[cfg(feature = "otel")]
if let Some(provider) = self._tracer_provider.clone() {
match tokio::task::spawn_blocking(move || shutdown_provider(&provider)).await {
Ok(()) => {}
Err(e) => eprintln!("warning: OTel tracer provider shutdown task failed: {e}"),
}
}
}
}
#[cfg(feature = "otel")]
fn shutdown_provider(provider: &opentelemetry_sdk::trace::SdkTracerProvider) {
match provider.shutdown() {
Ok(()) | Err(opentelemetry_sdk::error::OTelSdkError::AlreadyShutdown) => {}
Err(e) => eprintln!("warning: OTel tracer provider shutdown failed: {e}"),
}
}
impl Drop for TracingGuard {
fn drop(&mut self) {
#[cfg(feature = "otel")]
if let Some(ref provider) = self._tracer_provider {
shutdown_provider(provider);
}
}
}
#[cfg(feature = "otel")]
mod atexit_flush {
use std::sync::OnceLock;
static PROVIDER: OnceLock<opentelemetry_sdk::trace::SdkTracerProvider> = OnceLock::new();
extern "C" fn on_exit() {
if let Some(provider) = PROVIDER.get() {
let _ = provider.shutdown();
}
}
pub(super) fn register(provider: opentelemetry_sdk::trace::SdkTracerProvider) -> bool {
if PROVIDER.set(provider).is_err() {
return false;
}
unsafe extern "C" {
safe fn atexit(f: extern "C" fn()) -> core::ffi::c_int;
}
let registered = atexit(on_exit) == 0;
if !registered {
eprintln!("warning: failed to register OTel atexit flush handler");
}
registered
}
}
pub fn init_tracing(app_name: &str, default: LevelFilter) -> Result<TracingGuard, TelemetryError> {
init_tracing_with_layer(app_name, default, None)
}
pub fn init_tracing_with_layer(
app_name: &str,
default: LevelFilter,
extra: Option<BoxedLayer>,
) -> Result<TracingGuard, TelemetryError> {
let log_directives = log_filter_directives(default);
let span_events = std::env::var("OXEN_FMT_SPAN")
.ok()
.map(|v| parse_fmt_span(&v))
.unwrap_or(FmtSpan::NONE);
#[cfg(any(test, feature = "test-utils"))]
let stderr_layer = tracing_subscriber::fmt::layer()
.with_writer(tracing_subscriber::fmt::TestWriter::default())
.with_target(true)
.with_ansi(false)
.with_span_events(span_events);
#[cfg(not(any(test, feature = "test-utils")))]
let stderr_layer = tracing_subscriber::fmt::layer()
.with_writer(std::io::stderr)
.with_target(true)
.with_ansi(std::io::stderr().is_terminal())
.with_span_events(span_events);
let maybe_log_dir = match std::env::var("OXEN_LOG_DIR").ok() {
Some(log_dir) => {
let created_log_dir = create_log_dir(&log_dir)?;
Some(created_log_dir)
}
None => None,
};
let (m_json_layer, m_worker_guard) = if let Some(ref log_dir) = maybe_log_dir {
let (jl, wg) = json_file_logging(app_name, log_dir);
(Some(jl), Some(wg))
} else {
(None, None)
};
let log_filter = env_filter(&log_directives, default);
let registry = tracing_subscriber::registry()
.with(extra.map(|layer| layer.with_filter(log_filter.clone())))
.with(m_json_layer.map(|layer| layer.with_filter(log_filter.clone())))
.with(stderr_layer.with_filter(log_filter));
#[cfg(feature = "otel")]
let (m_otel_layer, m_tracer_provider, m_endpoint_p) = match otel_endpoint() {
Some((var, endpoint)) => {
let protocol = otel_protocol()?;
match build_otel_layer(app_name, &protocol) {
(Some(layer), Some(provider)) => {
atexit_flush::register(provider.clone());
(
Some(layer),
Some(provider),
Some(format!("{protocol} (protobuf) -> {var}={endpoint}")),
)
}
_ => (None, None, None),
}
}
None => (None, None, None),
};
#[cfg(feature = "otel")]
{
let otel_directives = otel_filter_directives();
let otel_layer = m_otel_layer
.map(|layer| layer.with_filter(env_filter(&otel_directives, OTEL_DEFAULT_FILTER)));
registry.with(otel_layer).try_init()?;
if let Some(protocol_and_endpoint) = m_endpoint_p {
log::info!(
"OpenTelemetry tracing enabled (endpoint: {protocol_and_endpoint}, span filter: {otel_directives})"
);
}
}
#[cfg(not(feature = "otel"))]
{
registry.try_init()?;
if non_empty_env("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT").is_some()
|| non_empty_env("OTEL_EXPORTER_OTLP_ENDPOINT").is_some()
{
log::error!(
"An OTLP endpoint is configured but the otel feature is not enabled! (Ignoring)"
)
}
}
if let Some(ref log_dir) = maybe_log_dir {
log::info!(
"JSON file logging enabled (log directory: {})",
log_dir.display()
);
}
Ok(TracingGuard {
_file_guard: m_worker_guard,
#[cfg(feature = "otel")]
_tracer_provider: m_tracer_provider,
})
}
#[cfg(feature = "otel")]
const OTEL_FILTER_ENV: &str = "OXEN_OTEL_FILTER";
#[cfg(feature = "otel")]
const OTEL_DEFAULT_FILTER: LevelFilter = LevelFilter::INFO;
fn log_filter_directives(default: LevelFilter) -> String {
non_empty_env(EnvFilter::DEFAULT_ENV).unwrap_or_else(|| default.to_string())
}
#[cfg(feature = "otel")]
fn otel_filter_directives() -> String {
non_empty_env(OTEL_FILTER_ENV).unwrap_or_else(|| OTEL_DEFAULT_FILTER.to_string())
}
#[cfg(feature = "otel")]
fn env_names_a_service(
service_name_var: Option<&str>,
attributes_service_name: Option<&str>,
) -> bool {
[service_name_var, attributes_service_name]
.into_iter()
.flatten()
.any(|name| !name.trim().is_empty())
}
#[cfg(feature = "otel")]
fn otel_endpoint() -> Option<(&'static str, String)> {
[
"OTEL_EXPORTER_OTLP_TRACES_ENDPOINT",
"OTEL_EXPORTER_OTLP_ENDPOINT",
]
.into_iter()
.find_map(|var| non_empty_env(var).map(|value| (var, value)))
}
fn non_empty_env(name: &str) -> Option<String> {
std::env::var(name)
.ok()
.filter(|value| !value.trim().is_empty())
}
fn env_filter(directives: &str, fallback: LevelFilter) -> EnvFilter {
EnvFilter::builder()
.with_default_directive(fallback.into())
.parse_lossy(directives)
}
fn create_log_dir(oxen_log_dir: &str) -> Result<PathBuf, TelemetryError> {
let oxen_log_dir = oxen_log_dir.trim();
if oxen_log_dir.is_empty() {
Err(TelemetryError::EmptyLogDir)
} else {
let log_dir = PathBuf::from(oxen_log_dir);
if log_dir.is_file() {
Err(TelemetryError::LogDirIsFile(log_dir))
} else {
match std::fs::create_dir_all(&log_dir) {
Ok(()) => Ok(log_dir),
Err(e) => Err(TelemetryError::CreateLogDir(log_dir, e)),
}
}
}
}
fn json_file_logging<S>(app_name: &str, log_dir: &Path) -> (impl Layer<S>, WorkerGuard)
where
S: tracing::Subscriber + for<'span> tracing_subscriber::registry::LookupSpan<'span>,
{
let file_appender = tracing_appender::rolling::daily(log_dir, app_name);
let (non_blocking, guard) = tracing_appender::non_blocking(file_appender);
let layer = tracing_subscriber::fmt::layer()
.json()
.with_writer(non_blocking)
.with_target(true)
.with_thread_ids(true)
.with_file(true)
.with_line_number(true);
(layer, guard)
}
fn parse_fmt_span(value: &str) -> FmtSpan {
let upper = value.to_uppercase();
if !upper.contains('|') {
return parse_fmt_span_token(&upper);
}
let mut span = FmtSpan::NONE;
for part in upper.split('|') {
span |= parse_fmt_span_token(part.trim());
}
span
}
fn parse_fmt_span_token(token: &str) -> FmtSpan {
match token {
"1" | "TRUE" | "CLOSE" => FmtSpan::CLOSE,
"NEW" => FmtSpan::NEW,
"ENTER" => FmtSpan::ENTER,
"EXIT" => FmtSpan::EXIT,
"ACTIVE" => FmtSpan::ACTIVE,
"FULL" => FmtSpan::FULL,
"NONE" => FmtSpan::NONE,
other => {
eprintln!("warning: unknown OXEN_FMT_SPAN component: {other:?}, ignoring");
FmtSpan::NONE
}
}
}
#[cfg(feature = "otel")]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Protocol {
Grpc,
Http,
}
#[cfg(feature = "otel")]
impl std::fmt::Display for Protocol {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Protocol::Grpc => write!(f, "grpc"),
Protocol::Http => write!(f, "http"),
}
}
}
#[cfg(feature = "otel")]
fn otel_protocol() -> Result<Protocol, TelemetryError> {
let (var, configured) = [
"OTEL_EXPORTER_OTLP_TRACES_PROTOCOL",
"OTEL_EXPORTER_OTLP_PROTOCOL",
]
.into_iter()
.find_map(|var| non_empty_env(var).map(|value| (var, Some(value))))
.unwrap_or(("OTEL_EXPORTER_OTLP_PROTOCOL", None));
parse_otel_protocol(var, configured.as_deref())
}
#[cfg(feature = "otel")]
fn parse_otel_protocol(
var: &'static str,
configured: Option<&str>,
) -> Result<Protocol, TelemetryError> {
match configured
.map(|value| value.trim().to_lowercase())
.as_deref()
{
None | Some("http" | "http/protobuf" | "http/json") => Ok(Protocol::Http),
Some("grpc") => Ok(Protocol::Grpc),
Some(unknown) => Err(TelemetryError::UnknownProtocol {
var,
value: unknown.to_string(),
}),
}
}
#[cfg(feature = "otel")]
fn build_otel_layer<S>(
app_name: &str,
protocol: &Protocol,
) -> (
Option<tracing_opentelemetry::OpenTelemetryLayer<S, opentelemetry_sdk::trace::SdkTracer>>,
Option<opentelemetry_sdk::trace::SdkTracerProvider>,
)
where
S: tracing::Subscriber + for<'span> tracing_subscriber::registry::LookupSpan<'span>,
{
use opentelemetry::trace::TracerProvider;
use opentelemetry::{Key, KeyValue};
use opentelemetry_otlp::WithExportConfig;
use opentelemetry_otlp::WithTonicConfig;
use opentelemetry_otlp::tonic_types::transport::ClientTlsConfig;
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::propagation::TraceContextPropagator;
use opentelemetry_sdk::resource::{EnvResourceDetector, ResourceDetector};
use opentelemetry_sdk::trace::{BatchConfigBuilder, BatchSpanProcessor, SdkTracerProvider};
let exporter = match protocol {
Protocol::Http => {
match opentelemetry_otlp::SpanExporter::builder()
.with_http()
.with_protocol(opentelemetry_otlp::Protocol::HttpBinary)
.build()
{
Ok(e) => e,
Err(err) => {
eprintln!("[ERROR] failed to build OTel HTTP exporter: {err}");
return (None, None);
}
}
}
Protocol::Grpc => {
match opentelemetry_otlp::SpanExporter::builder()
.with_tonic()
.with_tls_config(ClientTlsConfig::new().with_native_roots())
.build()
{
Ok(e) => e,
Err(err) => {
eprintln!("[ERROR] failed to build OTel gRPC exporter: {err}");
return (None, None);
}
}
}
};
let attributes_service_name = EnvResourceDetector::new()
.detect()
.get(&Key::from_static_str("service.name"))
.map(|name| name.to_string());
let mut builder = Resource::builder().with_attributes([KeyValue::new(
"service.version",
crate::constants::OXEN_VERSION,
)]);
if !env_names_a_service(
std::env::var("OTEL_SERVICE_NAME").ok().as_deref(),
attributes_service_name.as_deref(),
) {
builder = builder.with_service_name(app_name.to_string());
}
let resource = builder.build();
let mut batch_config = BatchConfigBuilder::default();
if std::env::var_os("OTEL_BSP_MAX_QUEUE_SIZE").is_none() {
batch_config = batch_config.with_max_queue_size(4096);
}
if std::env::var_os("OTEL_BSP_SCHEDULE_DELAY").is_none() {
batch_config = batch_config.with_scheduled_delay(std::time::Duration::from_secs(2));
}
let processor = BatchSpanProcessor::builder(exporter)
.with_batch_config(batch_config.build())
.build();
let provider = SdkTracerProvider::builder()
.with_span_processor(processor)
.with_resource(resource)
.build();
opentelemetry::global::set_text_map_propagator(TraceContextPropagator::new());
let tracer = provider.tracer("oxen");
let layer = tracing_opentelemetry::layer().with_tracer(tracer);
(Some(layer), Some(provider))
}
#[cfg(test)]
mod tests {
use super::*;
use tracing_subscriber::fmt::format::FmtSpan;
#[test]
fn a_filter_with_no_usable_directive_falls_back_to_the_given_level() {
for directives in ["liboxen=verbose", "warn=oops=bad", "="] {
assert_eq!(
env_filter(directives, LevelFilter::WARN).max_level_hint(),
Some(LevelFilter::WARN),
"{directives} should fall back to the given level"
);
}
}
#[test]
fn usable_directives_win_over_the_fallback() {
assert_eq!(
env_filter("debug", LevelFilter::WARN).max_level_hint(),
Some(LevelFilter::DEBUG)
);
assert_eq!(
env_filter("off", LevelFilter::WARN).max_level_hint(),
Some(LevelFilter::OFF)
);
assert_eq!(
env_filter("liboxen=verbose,oxen_server=debug", LevelFilter::WARN).max_level_hint(),
Some(LevelFilter::DEBUG)
);
}
#[test]
fn token_close() {
assert_eq!(parse_fmt_span_token("CLOSE"), FmtSpan::CLOSE);
assert_eq!(parse_fmt_span_token("1"), FmtSpan::CLOSE);
assert_eq!(parse_fmt_span_token("TRUE"), FmtSpan::CLOSE);
assert_eq!(parse_fmt_span("cLosE"), FmtSpan::CLOSE);
assert_eq!(parse_fmt_span("tRuE"), FmtSpan::CLOSE);
}
#[test]
fn token_new() {
assert_eq!(parse_fmt_span_token("NEW"), FmtSpan::NEW);
assert_eq!(parse_fmt_span("NeW"), FmtSpan::NEW);
}
#[test]
fn token_enter() {
assert_eq!(parse_fmt_span_token("ENTER"), FmtSpan::ENTER);
assert_eq!(parse_fmt_span("eNteR"), FmtSpan::ENTER);
}
#[test]
fn token_exit() {
assert_eq!(parse_fmt_span_token("EXIT"), FmtSpan::EXIT);
assert_eq!(parse_fmt_span("exIT"), FmtSpan::EXIT);
}
#[test]
fn token_active() {
assert_eq!(parse_fmt_span_token("ACTIVE"), FmtSpan::ACTIVE);
assert_eq!(parse_fmt_span("aCtIvE"), FmtSpan::ACTIVE);
}
#[test]
fn token_full() {
assert_eq!(parse_fmt_span_token("FULL"), FmtSpan::FULL);
assert_eq!(parse_fmt_span("FUll"), FmtSpan::FULL);
}
#[test]
fn token_none() {
assert_eq!(parse_fmt_span_token("NONE"), FmtSpan::NONE);
assert_eq!(parse_fmt_span("NonE"), FmtSpan::NONE);
}
#[test]
fn token_unknown_returns_none() {
assert_eq!(parse_fmt_span_token("BOGUS"), FmtSpan::NONE);
assert_eq!(parse_fmt_span("bogus"), FmtSpan::NONE);
}
#[test]
fn combined_new_close() {
assert_eq!(parse_fmt_span("NEW|CLOSE"), FmtSpan::NEW | FmtSpan::CLOSE);
assert_eq!(parse_fmt_span("new|close"), FmtSpan::NEW | FmtSpan::CLOSE);
assert_eq!(parse_fmt_span("NEW | CLOSE"), FmtSpan::NEW | FmtSpan::CLOSE);
}
#[test]
fn combined_active_close() {
assert_eq!(
parse_fmt_span("ACTIVE|CLOSE"),
FmtSpan::ACTIVE | FmtSpan::CLOSE
);
}
#[test]
fn combined_full_new() {
assert_eq!(parse_fmt_span("FULL|NEW"), FmtSpan::FULL | FmtSpan::NEW);
}
#[test]
fn combined_unknown_component_ignored() {
assert_eq!(parse_fmt_span("NEW|BOGUS"), FmtSpan::NEW | FmtSpan::NONE);
}
#[test]
fn combined_all_four_lifecycle() {
assert_eq!(
parse_fmt_span("NEW|ENTER|EXIT|CLOSE"),
FmtSpan::NEW | FmtSpan::ENTER | FmtSpan::EXIT | FmtSpan::CLOSE
);
}
#[cfg(feature = "otel")]
mod otel_tests {
use super::super::{
Protocol, TelemetryError, env_names_a_service, otel_filter_directives,
parse_otel_protocol,
};
#[test]
fn no_configured_protocol_is_http() {
assert_eq!(
parse_otel_protocol("OTEL_EXPORTER_OTLP_PROTOCOL", None).unwrap(),
Protocol::Http
);
}
#[test]
fn standard_protocol_values_select_a_transport() {
let var = "OTEL_EXPORTER_OTLP_PROTOCOL";
assert_eq!(
parse_otel_protocol(var, Some("grpc")).unwrap(),
Protocol::Grpc
);
for value in [
"http",
"http/protobuf",
"http/json",
"HTTP/protobuf",
" http ",
] {
assert_eq!(
parse_otel_protocol(var, Some(value)).unwrap(),
Protocol::Http,
"{value} should select HTTP"
);
}
}
#[test]
fn rejects_an_unknown_protocol() {
for var in [
"OTEL_EXPORTER_OTLP_TRACES_PROTOCOL",
"OTEL_EXPORTER_OTLP_PROTOCOL",
] {
for value in ["https", "http/proto", "tcp"] {
let result = parse_otel_protocol(var, Some(value));
match result {
Err(TelemetryError::UnknownProtocol {
var: named,
value: v,
}) => {
assert_eq!(
named, var,
"the rejected value should name its own variable"
);
assert_eq!(v, value);
}
other => panic!("{var}={value} should be rejected, got {other:?}"),
}
}
}
}
#[test]
fn the_environment_names_the_service_when_it_says_so() {
assert!(env_names_a_service(Some("checkout"), None));
assert!(env_names_a_service(None, Some("checkout")));
assert!(env_names_a_service(Some("checkout"), Some("billing")));
}
#[test]
fn a_blank_service_name_names_nothing() {
for blank in [None, Some(""), Some(" ")] {
assert!(
!env_names_a_service(blank, None),
"OTEL_SERVICE_NAME {blank:?} should not count as naming a service"
);
assert!(
!env_names_a_service(None, blank),
"service.name {blank:?} should not count as naming a service"
);
assert!(
!env_names_a_service(blank, blank),
"both blank at {blank:?} should not count as naming a service"
);
}
}
#[test]
fn span_export_defaults_to_info() {
if std::env::var_os("OXEN_OTEL_FILTER").is_none() {
assert_eq!(otel_filter_directives(), "info");
}
}
}
}