Skip to main content

rskit_logging/
otlp.rs

1//! OpenTelemetry Logs bridge with OTLP export.
2//!
3//! Bridges [`tracing`] events to the OpenTelemetry Logs SDK
4//! so they can be exported via OTLP (gRPC or HTTP) to a collector.
5//!
6//! # Feature gate
7//!
8//! This module is only available when the `otlp` cargo feature is enabled.
9//!
10//! # Example
11//!
12//! ```rust,ignore
13//! use rskit_logging::otlp::{OtlpConfig, OtlpProvider};
14//!
15//! let cfg = OtlpConfig { enabled: true, ..Default::default() };
16//! let provider = OtlpProvider::new(&cfg, "my-svc", "production", "1.0.0")?;
17//! // provider is Some when enabled — add its layer to the subscriber stack.
18//! ```
19
20use std::collections::HashMap;
21
22use opentelemetry::KeyValue;
23use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
24use opentelemetry_otlp::{WithExportConfig, WithHttpConfig};
25use opentelemetry_sdk::Resource;
26use opentelemetry_sdk::logs::{SdkLogger, SdkLoggerProvider};
27use tracing::Subscriber;
28use tracing_subscriber::registry::LookupSpan;
29
30use crate::error::{self, LoggingResult};
31
32// ── Config ──────────────────────────────────────────────────────────────────
33
34/// Configuration for OTLP log export.
35#[derive(Debug, Clone)]
36pub struct OtlpConfig {
37    /// Master switch — when `false`, [`OtlpProvider::new`] returns `Ok(None)`.
38    pub enabled: bool,
39    /// Collector endpoint (default: `"http://localhost:4317"`).
40    pub endpoint: String,
41    /// Protocol: `"grpc"` or `"http"` (default: `"grpc"`).
42    pub protocol: String,
43    /// Additional HTTP headers (e.g. auth tokens).
44    ///
45    /// Supported only when [`OtlpConfig::protocol`] is `"http"` because the underlying HTTP
46    /// and gRPC OTLP exporters expose different header types.
47    pub headers: HashMap<String, String>,
48}
49
50impl Default for OtlpConfig {
51    fn default() -> Self {
52        Self {
53            enabled: false,
54            endpoint: "http://localhost:4317".to_string(),
55            protocol: "grpc".to_string(),
56            headers: HashMap::new(),
57        }
58    }
59}
60
61// ── Provider ────────────────────────────────────────────────────────────────
62
63/// Manages the OpenTelemetry [`SdkLoggerProvider`] for OTLP export.
64///
65/// Create via [`OtlpProvider::new`],
66/// then call [`OtlpProvider::layer`] to obtain a [`tracing_subscriber::Layer`] that can be composed into the subscriber stack.
67///
68/// The provider **must** be shut down gracefully via [`OtlpProvider::shutdown`] (or by dropping the [`crate::LoggingGuard`]) to flush pending log records.
69pub struct OtlpProvider {
70    provider: SdkLoggerProvider,
71}
72
73impl OtlpProvider {
74    /// Create a new OTLP provider.
75    ///
76    /// Returns `Ok(None)` when `cfg.enabled` is `false`.
77    pub fn new(
78        cfg: &OtlpConfig,
79        service_name: &str,
80        environment: &str,
81        version: &str,
82    ) -> LoggingResult<Option<Self>> {
83        if !cfg.enabled {
84            return Ok(None);
85        }
86
87        let resource = Resource::builder_empty()
88            .with_attributes([
89                KeyValue::new("service.name", service_name.to_string()),
90                KeyValue::new("deployment.environment", environment.to_string()),
91                KeyValue::new("service.version", version.to_string()),
92            ])
93            .build();
94
95        let exporter = build_exporter(cfg)?;
96
97        let provider = SdkLoggerProvider::builder()
98            .with_resource(resource)
99            .with_batch_exporter(exporter)
100            .build();
101
102        Ok(Some(Self { provider }))
103    }
104
105    /// Get the OpenTelemetry tracing layer for use with [`tracing_subscriber`].
106    ///
107    /// The returned layer converts every [`tracing`] event into an OpenTelemetry log record
108    /// and forwards it to the OTLP exporter.
109    pub fn layer<S>(&self) -> OpenTelemetryTracingBridge<SdkLoggerProvider, SdkLogger>
110    where
111        S: Subscriber + for<'a> LookupSpan<'a>,
112    {
113        OpenTelemetryTracingBridge::new(&self.provider)
114    }
115
116    /// Gracefully shut down the provider, flushing all pending log records.
117    pub fn shutdown(self) -> LoggingResult<()> {
118        self.provider.shutdown().map_err(error::otlp_shutdown)
119    }
120}
121
122// ── Internal helpers ────────────────────────────────────────────────────────
123
124/// Build an OTLP [`LogExporter`](opentelemetry_otlp::LogExporter) from config.
125fn build_exporter(cfg: &OtlpConfig) -> LoggingResult<opentelemetry_otlp::LogExporter> {
126    match cfg.protocol.as_str() {
127        "http" => {
128            let exporter = opentelemetry_otlp::LogExporter::builder()
129                .with_http()
130                .with_endpoint(&cfg.endpoint)
131                .with_headers(cfg.headers.clone())
132                .build()
133                .map_err(error::otlp_exporter)?;
134            Ok(exporter)
135        }
136        "grpc" => {
137            if !cfg.headers.is_empty() {
138                return Err(error::grpc_headers_not_supported());
139            }
140            let exporter = opentelemetry_otlp::LogExporter::builder()
141                .with_tonic()
142                .with_endpoint(&cfg.endpoint)
143                .build()
144                .map_err(error::otlp_exporter)?;
145            Ok(exporter)
146        }
147        other => Err(error::invalid_protocol(other)),
148    }
149}
150
151// ── Tests ───────────────────────────────────────────────────────────────────
152
153#[cfg(test)]
154mod tests {
155    use super::*;
156
157    #[test]
158    fn default_config_is_disabled() {
159        let cfg = OtlpConfig::default();
160        assert!(!cfg.enabled);
161        assert_eq!(cfg.endpoint, "http://localhost:4317");
162        assert_eq!(cfg.protocol, "grpc");
163        assert!(cfg.headers.is_empty());
164    }
165
166    #[test]
167    fn disabled_config_returns_none() {
168        let cfg = OtlpConfig::default();
169        let result = OtlpProvider::new(&cfg, "test-svc", "test", "0.1.0");
170        assert!(result.is_ok());
171        assert!(result.unwrap().is_none());
172    }
173
174    #[test]
175    fn config_clone_preserves_values() {
176        let mut cfg = OtlpConfig {
177            enabled: true,
178            endpoint: "http://collector:4317".to_string(),
179            ..Default::default()
180        };
181        cfg.headers
182            .insert("x-api-key".to_string(), "secret".to_string());
183
184        let cloned = cfg.clone();
185        assert!(cloned.enabled);
186        assert_eq!(cloned.endpoint, "http://collector:4317");
187        assert_eq!(cloned.headers.get("x-api-key").unwrap(), "secret");
188    }
189
190    #[test]
191    fn config_debug_format() {
192        let cfg = OtlpConfig::default();
193        let debug = format!("{cfg:?}");
194        assert!(debug.contains("OtlpConfig"));
195        assert!(debug.contains("enabled"));
196    }
197
198    #[test]
199    fn invalid_protocol_returns_typed_error() {
200        let cfg = OtlpConfig {
201            enabled: true,
202            protocol: "udp".to_string(),
203            ..Default::default()
204        };
205        let err = match OtlpProvider::new(&cfg, "test-svc", "test", "0.1.0") {
206            Ok(_) => panic!("unsupported protocol must fail"),
207            Err(err) => err,
208        };
209        assert_eq!(err.code(), rskit_errors::ErrorCode::InvalidInput);
210    }
211
212    #[test]
213    fn grpc_headers_return_typed_error() {
214        let mut cfg = OtlpConfig {
215            enabled: true,
216            ..Default::default()
217        };
218        cfg.headers
219            .insert("x-api-key".to_string(), "secret".to_string());
220        let err = match OtlpProvider::new(&cfg, "test-svc", "test", "0.1.0") {
221            Ok(_) => panic!("grpc headers are unsupported"),
222            Err(err) => err,
223        };
224        assert_eq!(err.code(), rskit_errors::ErrorCode::InvalidInput);
225    }
226}