1use 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#[derive(Debug, Clone)]
36pub struct OtlpConfig {
37 pub enabled: bool,
39 pub endpoint: String,
41 pub protocol: String,
43 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
61pub struct OtlpProvider {
70 provider: SdkLoggerProvider,
71}
72
73impl OtlpProvider {
74 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 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 pub fn shutdown(self) -> LoggingResult<()> {
118 self.provider.shutdown().map_err(error::otlp_shutdown)
119 }
120}
121
122fn 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#[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}