use core::fmt;
use std::env;
use opentelemetry::trace::TracerProvider as _;
use opentelemetry_otlp::{MetricExporter, SpanExporter, WithExportConfig as _};
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::metrics::SdkMeterProvider;
use opentelemetry_sdk::trace::SdkTracerProvider;
use tracing_subscriber::Registry;
use tracing_subscriber::layer::SubscriberExt as _;
pub use crate::metrics::bind_global_meter;
const DEFAULT_SERVICE_NAME: &str = "keel";
const GATE_VAR: &str = "KEEL_OTEL";
const ENDPOINT_VAR: &str = "OTEL_EXPORTER_OTLP_ENDPOINT";
#[derive(Debug)]
pub struct OtelGuard {
tracer_provider: SdkTracerProvider,
meter_provider: SdkMeterProvider,
}
impl Drop for OtelGuard {
fn drop(&mut self) {
if let Err(error) = self.tracer_provider.shutdown() {
tracing::warn!(%error, "otel span exporter shutdown failed; buffered spans may be lost");
}
if let Err(error) = self.meter_provider.shutdown() {
tracing::warn!(%error, "otel metric exporter shutdown failed; the final collection may be lost");
}
}
}
#[derive(Debug)]
#[non_exhaustive]
pub enum OtelInitError {
Exporter(opentelemetry_otlp::ExporterBuildError),
Subscriber(tracing::subscriber::SetGlobalDefaultError),
}
impl fmt::Display for OtelInitError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Exporter(error) => write!(f, "failed to build an OTLP exporter: {error}"),
Self::Subscriber(error) => {
write!(
f,
"a global tracing subscriber is already installed: {error}"
)
}
}
}
}
impl std::error::Error for OtelInitError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
Self::Exporter(error) => Some(error),
Self::Subscriber(error) => Some(error),
}
}
}
impl From<opentelemetry_otlp::ExporterBuildError> for OtelInitError {
fn from(error: opentelemetry_otlp::ExporterBuildError) -> Self {
Self::Exporter(error)
}
}
impl From<tracing::subscriber::SetGlobalDefaultError> for OtelInitError {
fn from(error: tracing::subscriber::SetGlobalDefaultError) -> Self {
Self::Subscriber(error)
}
}
#[must_use]
pub fn env_opt_in() -> Option<bool> {
opt_in_from(env::var(GATE_VAR).ok().as_deref())
}
#[must_use]
pub fn export_enabled(policy_endpoint: Option<&str>) -> bool {
env_opt_in().unwrap_or_else(|| policy_endpoint.is_some_and(|e| !e.trim().is_empty()))
}
#[must_use]
pub fn resolve_endpoint(policy_endpoint: Option<&str>) -> Option<String> {
endpoint_from(env::var(ENDPOINT_VAR).ok().as_deref(), policy_endpoint)
}
fn opt_in_from(raw: Option<&str>) -> Option<bool> {
let value = raw?.trim().to_ascii_lowercase();
if value.is_empty() {
return None;
}
Some(!matches!(value.as_str(), "0" | "false" | "no" | "off"))
}
fn endpoint_from(env_endpoint: Option<&str>, policy_endpoint: Option<&str>) -> Option<String> {
if env_endpoint.is_some_and(|e| !e.trim().is_empty()) {
return None; }
policy_endpoint
.map(str::trim)
.filter(|e| !e.is_empty())
.map(str::to_owned)
}
pub fn init_otlp(endpoint: Option<&str>) -> Result<OtelGuard, OtelInitError> {
let mut span_builder = SpanExporter::builder().with_tonic();
let mut metric_builder = MetricExporter::builder().with_tonic();
if let Some(endpoint) = endpoint {
span_builder = span_builder.with_endpoint(endpoint);
metric_builder = metric_builder.with_endpoint(endpoint);
}
let span_exporter = span_builder.build()?;
let metric_exporter = metric_builder.build()?;
let resource = Resource::builder()
.with_service_name(DEFAULT_SERVICE_NAME)
.build();
let tracer_provider = SdkTracerProvider::builder()
.with_resource(resource.clone())
.with_batch_exporter(span_exporter)
.build();
let tracer = tracer_provider.tracer("keel-core");
let layer = tracing_opentelemetry::layer().with_tracer(tracer);
let subscriber = Registry::default().with(layer);
tracing::subscriber::set_global_default(subscriber)?;
let meter_provider = SdkMeterProvider::builder()
.with_resource(resource)
.with_periodic_exporter(metric_exporter)
.build();
opentelemetry::global::set_meter_provider(meter_provider.clone());
bind_global_meter();
Ok(OtelGuard {
tracer_provider,
meter_provider,
})
}
#[cfg(test)]
mod tests {
use super::{endpoint_from, opt_in_from};
#[test]
fn opt_in_tristate() {
assert_eq!(opt_in_from(None), None);
assert_eq!(opt_in_from(Some("")), None);
assert_eq!(opt_in_from(Some(" ")), None);
assert_eq!(opt_in_from(Some("0")), Some(false));
assert_eq!(opt_in_from(Some("false")), Some(false));
assert_eq!(opt_in_from(Some(" No ")), Some(false));
assert_eq!(opt_in_from(Some("OFF")), Some(false));
assert_eq!(opt_in_from(Some("1")), Some(true));
assert_eq!(opt_in_from(Some("true")), Some(true));
assert_eq!(opt_in_from(Some("yes")), Some(true));
assert_eq!(opt_in_from(Some("collector-a")), Some(true));
}
#[test]
fn endpoint_env_wins_over_policy() {
assert_eq!(
endpoint_from(Some("http://env:4317"), Some("http://policy:4317")),
None
);
assert_eq!(
endpoint_from(None, Some("http://policy:4317")),
Some("http://policy:4317".to_owned())
);
assert_eq!(
endpoint_from(Some(" "), Some(" http://policy:4317 ")),
Some("http://policy:4317".to_owned())
);
assert_eq!(endpoint_from(None, None), None);
assert_eq!(endpoint_from(None, Some(" ")), None);
}
#[test]
fn export_enabled_matrix() {
use super::export_enabled;
if std::env::var_os(super::GATE_VAR).is_none() {
assert!(export_enabled(Some("http://collector:4317")));
assert!(!export_enabled(Some(" ")));
assert!(!export_enabled(None));
}
}
}