newton-core 0.7.2

newton protocol core sdk
//! Shared telemetry initialization for Newton services.

use super::log::{logger_layer, LoggerConfig};
use opentelemetry::{global, trace::TracerProvider as _, KeyValue};
use opentelemetry_sdk::{trace::SdkTracerProvider, Resource};
use tracing_subscriber::{
    filter::filter_fn, layer::SubscriberExt, util::SubscriberInitExt, EnvFilter, Layer, Registry,
};

const DEFAULT_TRACE_FILTER: &str = "off";

/// Application-specific values used to initialize telemetry.
#[derive(Debug, Clone, Copy)]
pub struct TelemetryConfig {
    service_name: &'static str,
    service_version: &'static str,
    required_trace_directive: Option<&'static str>,
}

impl TelemetryConfig {
    /// Creates telemetry configuration for a service.
    pub const fn new(service_name: &'static str, service_version: &'static str) -> Self {
        Self {
            service_name,
            service_version,
            required_trace_directive: None,
        }
    }

    /// Requires an exact filter directive before OTLP export is enabled.
    ///
    /// Use this when a service exports only a deliberately enabled trace target.
    pub const fn requiring_trace_directive(mut self, directive: &'static str) -> Self {
        self.required_trace_directive = Some(directive);
        self
    }

    fn trace_export_enabled(self, filter: &str) -> bool {
        match self.required_trace_directive {
            Some(required) => filter.split(',').map(str::trim).any(|directive| directive == required),
            None => filter.trim() != DEFAULT_TRACE_FILTER,
        }
    }
}

/// Flushes completed OpenTelemetry spans when the process exits normally.
#[derive(Debug, Default)]
pub struct TelemetryGuard {
    provider: Option<SdkTracerProvider>,
}

impl Drop for TelemetryGuard {
    fn drop(&mut self) {
        if let Some(provider) = self.provider.take() {
            if let Err(error) = provider.shutdown() {
                eprintln!("failed to flush OpenTelemetry traces during shutdown: {error}");
            }
        }
    }
}

/// Initializes Newton's standard logger and optional OTLP trace export.
///
/// The exporter uses the SDK batch processor, keeping network export outside
/// instrumented application futures. The returned guard must live until process
/// shutdown so buffered spans can be flushed.
pub fn init_telemetry(logger: LoggerConfig, config: TelemetryConfig) -> eyre::Result<TelemetryGuard> {
    init_telemetry_with_layers(logger, config, Vec::new())
}

/// Same as [`init_telemetry`], plus caller-supplied layers.
///
/// `newton-metric` depends on `newton-core`, so core cannot reference the
/// metrics crate directly. Binaries pass its pipeline-stage metrics layer
/// (`newton_metric::stages::stage_metrics_layer()`) through here instead.
///
/// Extra layers are installed on BOTH paths -- with and without OTLP export --
/// so metrics derived from spans do not silently depend on whether tracing
/// export happens to be configured.
pub fn init_telemetry_with_layers(
    logger: LoggerConfig,
    config: TelemetryConfig,
    extra_layers: Vec<Box<dyn Layer<Registry> + Send + Sync>>,
) -> eyre::Result<TelemetryGuard> {
    let endpoint_configured = std::env::var_os("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT").is_some();
    let trace_filter = std::env::var("OTEL_TRACES_FILTER").unwrap_or_else(|_| DEFAULT_TRACE_FILTER.to_string());

    if !endpoint_configured || !config.trace_export_enabled(&trace_filter) {
        let mut layers: Vec<Box<dyn Layer<Registry> + Send + Sync>> = vec![logger_layer(logger)];
        layers.extend(extra_layers);
        tracing_subscriber::registry()
            .with(layers)
            .try_init()
            .map_err(|error| eyre::eyre!("failed to initialize tracing subscriber: {error}"))?;
        return Ok(TelemetryGuard::default());
    }

    let protocol = std::env::var("OTEL_EXPORTER_OTLP_TRACES_PROTOCOL").unwrap_or_else(|_| "http/protobuf".to_string());
    if protocol != "http/protobuf" {
        eyre::bail!(
            "{} OTLP export supports http/protobuf, but configured protocol is {protocol}",
            config.service_name
        );
    }

    let exporter = opentelemetry_otlp::SpanExporter::builder()
        .with_http()
        .build()
        .map_err(|error| eyre::eyre!("failed to build OTLP trace exporter: {error}"))?;
    let service_name = std::env::var("OTEL_SERVICE_NAME").unwrap_or_else(|_| config.service_name.to_string());
    let instance_id = std::env::var("HOSTNAME").unwrap_or_else(|_| format!("pid-{}", std::process::id()));
    let resource = Resource::builder()
        .with_service_name(service_name)
        .with_attributes([
            KeyValue::new("service.version", config.service_version),
            KeyValue::new("service.instance.id", instance_id),
        ])
        .build();

    // The SDK builder reads the standard sampler, span-limit, and batch-export
    // environment variables.
    let provider = SdkTracerProvider::builder()
        .with_resource(resource)
        .with_batch_exporter(exporter)
        .build();
    let tracer = provider.tracer(config.service_name);
    let otel_layer = tracing_opentelemetry::layer()
        .with_tracer(tracer)
        .with_filter(EnvFilter::new(trace_filter))
        // Logs keep their normal formatting and destination; only spans are
        // exported through OTLP.
        .with_filter(filter_fn(|metadata| metadata.is_span()))
        .boxed();
    let mut layers: Vec<Box<dyn Layer<Registry> + Send + Sync>> = vec![logger_layer(logger), otel_layer];
    layers.extend(extra_layers);

    tracing_subscriber::registry()
        .with(layers)
        .try_init()
        .map_err(|error| eyre::eyre!("failed to initialize tracing subscriber: {error}"))?;
    global::set_tracer_provider(provider.clone());

    Ok(TelemetryGuard {
        provider: Some(provider),
    })
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn default_gate_exports_any_explicitly_enabled_filter() {
        let config = TelemetryConfig::new("test-service", "test-version");

        assert!(!config.trace_export_enabled("off"));
        assert!(config.trace_export_enabled("info"));
        assert!(config.trace_export_enabled("off,test_target=debug"));
    }

    #[test]
    fn directive_gate_requires_an_exact_directive() {
        let config = TelemetryConfig::new("test-service", "test-version")
            .requiring_trace_directive("newton::task_evaluation=debug");

        assert!(config.trace_export_enabled("off,newton::task_evaluation=debug"));
        assert!(!config.trace_export_enabled("debug"));
        assert!(!config.trace_export_enabled("off,newton::task_evaluation=info"));
    }
}