use std::collections::HashMap;
use opentelemetry::KeyValue;
use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
use opentelemetry_otlp::{WithExportConfig, WithHttpConfig};
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::logs::{SdkLogger, SdkLoggerProvider};
use tracing::Subscriber;
use tracing_subscriber::registry::LookupSpan;
use crate::error::{self, LoggingResult};
#[derive(Debug, Clone)]
pub struct OtlpConfig {
pub enabled: bool,
pub endpoint: String,
pub protocol: String,
pub headers: HashMap<String, String>,
}
impl Default for OtlpConfig {
fn default() -> Self {
Self {
enabled: false,
endpoint: "http://localhost:4317".to_string(),
protocol: "grpc".to_string(),
headers: HashMap::new(),
}
}
}
pub struct OtlpProvider {
provider: SdkLoggerProvider,
}
impl OtlpProvider {
pub fn new(
cfg: &OtlpConfig,
service_name: &str,
environment: &str,
version: &str,
) -> LoggingResult<Option<Self>> {
if !cfg.enabled {
return Ok(None);
}
let resource = Resource::builder_empty()
.with_attributes([
KeyValue::new("service.name", service_name.to_string()),
KeyValue::new("deployment.environment", environment.to_string()),
KeyValue::new("service.version", version.to_string()),
])
.build();
let exporter = build_exporter(cfg)?;
let provider = SdkLoggerProvider::builder()
.with_resource(resource)
.with_batch_exporter(exporter)
.build();
Ok(Some(Self { provider }))
}
pub fn layer<S>(&self) -> OpenTelemetryTracingBridge<SdkLoggerProvider, SdkLogger>
where
S: Subscriber + for<'a> LookupSpan<'a>,
{
OpenTelemetryTracingBridge::new(&self.provider)
}
pub fn shutdown(self) -> LoggingResult<()> {
self.provider.shutdown().map_err(error::otlp_shutdown)
}
}
fn build_exporter(cfg: &OtlpConfig) -> LoggingResult<opentelemetry_otlp::LogExporter> {
match cfg.protocol.as_str() {
"http" => {
let exporter = opentelemetry_otlp::LogExporter::builder()
.with_http()
.with_endpoint(&cfg.endpoint)
.with_headers(cfg.headers.clone())
.build()
.map_err(error::otlp_exporter)?;
Ok(exporter)
}
"grpc" => {
if !cfg.headers.is_empty() {
return Err(error::grpc_headers_not_supported());
}
let exporter = opentelemetry_otlp::LogExporter::builder()
.with_tonic()
.with_endpoint(&cfg.endpoint)
.build()
.map_err(error::otlp_exporter)?;
Ok(exporter)
}
other => Err(error::invalid_protocol(other)),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_config_is_disabled() {
let cfg = OtlpConfig::default();
assert!(!cfg.enabled);
assert_eq!(cfg.endpoint, "http://localhost:4317");
assert_eq!(cfg.protocol, "grpc");
assert!(cfg.headers.is_empty());
}
#[test]
fn disabled_config_returns_none() {
let cfg = OtlpConfig::default();
let result = OtlpProvider::new(&cfg, "test-svc", "test", "0.1.0");
assert!(result.is_ok());
assert!(result.unwrap().is_none());
}
#[test]
fn config_clone_preserves_values() {
let mut cfg = OtlpConfig {
enabled: true,
endpoint: "http://collector:4317".to_string(),
..Default::default()
};
cfg.headers
.insert("x-api-key".to_string(), "secret".to_string());
let cloned = cfg.clone();
assert!(cloned.enabled);
assert_eq!(cloned.endpoint, "http://collector:4317");
assert_eq!(cloned.headers.get("x-api-key").unwrap(), "secret");
}
#[test]
fn config_debug_format() {
let cfg = OtlpConfig::default();
let debug = format!("{cfg:?}");
assert!(debug.contains("OtlpConfig"));
assert!(debug.contains("enabled"));
}
#[test]
fn invalid_protocol_returns_typed_error() {
let cfg = OtlpConfig {
enabled: true,
protocol: "udp".to_string(),
..Default::default()
};
let err = match OtlpProvider::new(&cfg, "test-svc", "test", "0.1.0") {
Ok(_) => panic!("unsupported protocol must fail"),
Err(err) => err,
};
assert_eq!(err.code(), rskit_errors::ErrorCode::InvalidInput);
}
#[test]
fn grpc_headers_return_typed_error() {
let mut cfg = OtlpConfig {
enabled: true,
..Default::default()
};
cfg.headers
.insert("x-api-key".to_string(), "secret".to_string());
let err = match OtlpProvider::new(&cfg, "test-svc", "test", "0.1.0") {
Ok(_) => panic!("grpc headers are unsupported"),
Err(err) => err,
};
assert_eq!(err.code(), rskit_errors::ErrorCode::InvalidInput);
}
}