use std::time::Duration;
use crate::error::OtelError;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub enum OtelProtocol {
#[default]
Grpc,
Http,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct OtelConfig {
pub endpoint: String,
pub protocol: OtelProtocol,
pub service_name: String,
pub service_version: Option<String>,
pub headers: Vec<(String, String)>,
pub export_interval: Duration,
pub export_timeout: Duration,
pub shutdown_timeout: Duration,
pub with_fmt_layer: bool,
pub env_filter: Option<String>,
}
impl OtelConfig {
pub fn localhost(service_name: impl Into<String>) -> Self {
Self {
endpoint: "http://localhost:4317".into(),
protocol: OtelProtocol::Grpc,
service_name: service_name.into(),
service_version: None,
headers: Vec::new(),
export_interval: Duration::from_secs(5),
export_timeout: Duration::from_secs(10),
shutdown_timeout: Duration::from_secs(3),
with_fmt_layer: false,
env_filter: None,
}
}
pub fn from_env() -> Result<Self, OtelError> {
let mut cfg =
Self::localhost(std::env::var("OTEL_SERVICE_NAME").unwrap_or_else(|_| "id_effect".into()));
if let Ok(endpoint) = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT") {
if endpoint.trim().is_empty() {
return Err(OtelError::Config(
"OTEL_EXPORTER_OTLP_ENDPOINT must not be empty".into(),
));
}
cfg.endpoint = endpoint;
}
if let Ok(protocol) = std::env::var("OTEL_EXPORTER_OTLP_PROTOCOL") {
cfg.protocol = parse_protocol(&protocol)?;
}
if let Ok(version) = std::env::var("OTEL_SERVICE_VERSION")
&& !version.trim().is_empty()
{
cfg.service_version = Some(version);
}
if let Ok(headers) = std::env::var("OTEL_EXPORTER_OTLP_HEADERS") {
cfg.headers = parse_header_list(&headers)?;
}
if let Ok(ms) = std::env::var("OTEL_METRIC_EXPORT_INTERVAL") {
let millis: u64 = ms
.parse()
.map_err(|_| OtelError::Config(format!("invalid OTEL_METRIC_EXPORT_INTERVAL: {ms}")))?;
cfg.export_interval = Duration::from_millis(millis);
}
if let Ok(ms) = std::env::var("OTEL_EXPORTER_OTLP_TIMEOUT") {
let millis: u64 = ms
.parse()
.map_err(|_| OtelError::Config(format!("invalid OTEL_EXPORTER_OTLP_TIMEOUT: {ms}")))?;
cfg.export_timeout = Duration::from_millis(millis);
}
if let Ok(filter) = std::env::var("RUST_LOG")
&& !filter.trim().is_empty()
{
cfg.env_filter = Some(filter);
}
Ok(cfg)
}
#[inline]
pub fn with_fmt_layer(mut self, enabled: bool) -> Self {
self.with_fmt_layer = enabled;
self
}
#[inline]
pub fn with_endpoint(mut self, endpoint: impl Into<String>) -> Self {
self.endpoint = endpoint.into();
self
}
#[inline]
pub fn with_protocol(mut self, protocol: OtelProtocol) -> Self {
self.protocol = protocol;
self
}
}
pub mod config_keys {
pub const ENDPOINT: &str = "otel.endpoint";
pub const PROTOCOL: &str = "otel.protocol";
pub const SERVICE_NAME: &str = "otel.service_name";
pub const SERVICE_VERSION: &str = "otel.service_version";
pub const HEADERS: &str = "otel.headers";
}
#[cfg(feature = "config")]
pub fn load_otel_config(
provider: &dyn id_effect_config::ConfigProvider,
) -> Result<OtelConfig, OtelError> {
use id_effect_config::config;
let mut cfg = OtelConfig::from_env()?;
if let Ok(endpoint) = config::string(provider, config_keys::ENDPOINT) {
if endpoint.trim().is_empty() {
return Err(OtelError::Config(format!(
"{} must not be empty",
config_keys::ENDPOINT
)));
}
cfg.endpoint = endpoint;
}
if let Ok(protocol) = config::string(provider, config_keys::PROTOCOL) {
cfg.protocol = parse_protocol(&protocol)?;
}
if let Ok(name) = config::string(provider, config_keys::SERVICE_NAME) {
if name.trim().is_empty() {
return Err(OtelError::Config(format!(
"{} must not be empty",
config_keys::SERVICE_NAME
)));
}
cfg.service_name = name;
}
if let Ok(version) = config::string(provider, config_keys::SERVICE_VERSION)
&& !version.trim().is_empty()
{
cfg.service_version = Some(version);
}
if let Ok(headers) = config::string(provider, config_keys::HEADERS) {
cfg.headers = parse_header_list(&headers)?;
}
Ok(cfg)
}
fn parse_protocol(raw: &str) -> Result<OtelProtocol, OtelError> {
match raw.trim().to_ascii_lowercase().as_str() {
"grpc" | "grpc/protobuf" => Ok(OtelProtocol::Grpc),
"http" | "http/protobuf" | "http-proto" => Ok(OtelProtocol::Http),
other => Err(OtelError::Config(format!(
"unsupported OTLP protocol {other:?} (expected grpc or http)"
))),
}
}
fn parse_header_list(raw: &str) -> Result<Vec<(String, String)>, OtelError> {
if raw.trim().is_empty() {
return Ok(Vec::new());
}
raw
.split(',')
.map(|pair| {
let (key, value) = pair.split_once('=').ok_or_else(|| {
OtelError::Config(format!(
"invalid OTLP header pair {pair:?} (expected key=value)"
))
})?;
Ok((key.trim().to_string(), value.trim().to_string()))
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_headers_accepts_comma_separated_pairs() {
let headers = parse_header_list("authorization=Bearer x, x-tenant=acme").unwrap();
assert_eq!(
headers,
vec![
("authorization".into(), "Bearer x".into()),
("x-tenant".into(), "acme".into()),
]
);
}
#[test]
fn parse_protocol_rejects_unknown_values() {
assert!(parse_protocol("nats").is_err());
}
}