#[allow(dead_code)]
pub mod catalog;
mod exporter;
pub mod http;
pub mod metrics;
mod spans;
pub use http::TelemetryLayer;
pub use metrics::{record_last_known_good, record_revision_rejection};
#[allow(unused_imports)]
pub use spans::{
ATTEMPT_ERROR, ATTEMPT_OK, CONVERGENCE_BOOT, CONVERGENCE_NOTIFIED, CONVERGENCE_POLLED,
LEASE_ERROR, LEASE_PARKED, LEASE_RATE_LIMITED, LEASE_SERVED, RELOAD_APPLIED, RELOAD_REJECTED,
config_reload_span, credential_lease_span, finish_config_reload, finish_credential_lease,
finish_revision_convergence, finish_upstream_attempt, record_attempt_timeout, record_request,
record_routing, record_streamed, revision_convergence_span, trace_id, upstream_attempt_span,
};
use std::sync::OnceLock;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, Instant};
use opentelemetry::global;
use opentelemetry::logs::LoggerProvider as _;
use opentelemetry::trace::TracerProvider as _;
use opentelemetry_otlp::{Protocol, WithExportConfig, WithHttpConfig};
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::logs::{SdkLogger, SdkLoggerProvider};
use opentelemetry_sdk::metrics::SdkMeterProvider;
use opentelemetry_sdk::propagation::TraceContextPropagator;
use opentelemetry_sdk::trace::SdkTracerProvider;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
use tracing_subscriber::{EnvFilter, Layer};
pub const SERVICE_NAME: &str = "axond";
const FLUSH_TIMEOUT: Duration = Duration::from_secs(5);
const SIGNALS: [&str; 3] = ["traces", "metrics", "logs"];
static EXPORTING: AtomicBool = AtomicBool::new(false);
static USAGE_LOGGER: OnceLock<SdkLogger> = OnceLock::new();
pub const USAGE_SCOPE: &str = "axond.usage";
pub fn usage_logger() -> Option<SdkLogger> {
USAGE_LOGGER.get().cloned()
}
pub fn is_exporting() -> bool {
EXPORTING.load(Ordering::Relaxed)
}
#[derive(Debug, thiserror::Error)]
#[error("{0}")]
pub struct TelemetryError(String);
#[derive(Debug, Clone, Default)]
pub struct TelemetryConfig {
pub endpoint: Option<String>,
}
impl TelemetryConfig {
pub fn from_env() -> Result<Self, TelemetryError> {
Self::from_values(
std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok().as_deref(),
std::env::var("OTEL_EXPORTER_OTLP_PROTOCOL").ok().as_deref(),
)
}
fn from_values(endpoint: Option<&str>, protocol: Option<&str>) -> Result<Self, TelemetryError> {
let Some(endpoint) = non_empty(endpoint) else {
return Ok(Self::default());
};
if !(endpoint.starts_with("http://") || endpoint.starts_with("https://")) {
return Err(TelemetryError(
"OTEL_EXPORTER_OTLP_ENDPOINT must be an http:// or https:// URL".to_owned(),
));
}
match non_empty(protocol).as_deref() {
None | Some("http/protobuf") => {}
Some(other) => {
return Err(TelemetryError(format!(
"OTEL_EXPORTER_OTLP_PROTOCOL=`{other}` is unsupported: axond exports OTLP/HTTP, so point the endpoint at the collector's HTTP receiver"
)));
}
}
Ok(Self {
endpoint: Some(endpoint),
})
}
}
fn non_empty(value: Option<&str>) -> Option<String> {
value
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_owned)
}
pub struct TelemetryGuard {
tracer: Option<SdkTracerProvider>,
meter: Option<SdkMeterProvider>,
logger: Option<SdkLoggerProvider>,
}
impl TelemetryGuard {
pub fn shutdown(&mut self, deadline: Instant) -> Vec<(&'static str, String)> {
let remaining = || deadline.saturating_duration_since(Instant::now());
let mut failures = Vec::new();
if let Some(provider) = self.tracer.take()
&& let Err(error) = provider.shutdown_with_timeout(remaining())
{
failures.push(("traces", error.to_string()));
}
if let Some(provider) = self.meter.take()
&& let Err(error) = provider.shutdown_with_timeout(remaining())
{
failures.push(("metrics", error.to_string()));
}
if let Some(provider) = self.logger.take()
&& let Err(error) = provider.shutdown_with_timeout(remaining())
{
failures.push(("logs", error.to_string()));
}
for (signal, error) in &failures {
tracing::error!(
signal,
error = %error,
"telemetry exporter did not drain within the shutdown bound"
);
}
failures
}
}
impl Drop for TelemetryGuard {
fn drop(&mut self) {
let _ = self.shutdown(Instant::now() + FLUSH_TIMEOUT);
}
}
fn resource() -> Resource {
Resource::builder().with_service_name(SERVICE_NAME).build()
}
pub fn init() -> Result<TelemetryGuard, TelemetryError> {
init_with(TelemetryConfig::from_env()?)
}
type Filtered = tracing_subscriber::layer::Layered<EnvFilter, tracing_subscriber::Registry>;
type OtelLayer = Box<dyn Layer<Filtered> + Send + Sync>;
fn install(filter: EnvFilter, otel: Option<OtelLayer>) -> Result<(), TelemetryError> {
tracing_subscriber::registry()
.with(filter)
.with(otel)
.with(tracing_subscriber::fmt::layer().json())
.try_init()
.map_err(|e| TelemetryError(format!("subscriber initialization failed: {e}")))
}
fn init_with(config: TelemetryConfig) -> Result<TelemetryGuard, TelemetryError> {
let filter =
EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info,axond=info"));
let Some(endpoint) = config.endpoint else {
install(filter, None)?;
return Ok(TelemetryGuard {
tracer: None,
meter: None,
logger: None,
});
};
let client = exporter::ExportClient::new()?;
let span_exporter = opentelemetry_otlp::SpanExporter::builder()
.with_http()
.with_protocol(Protocol::HttpBinary)
.with_http_client(client.clone())
.with_endpoint(signal_endpoint(&endpoint, "traces"))
.build()
.map_err(|e| TelemetryError(format!("OTLP span exporter configuration failed: {e}")))?;
let tracer_provider = SdkTracerProvider::builder()
.with_batch_exporter(span_exporter)
.with_resource(resource())
.build();
let metric_exporter = opentelemetry_otlp::MetricExporter::builder()
.with_http()
.with_protocol(Protocol::HttpBinary)
.with_http_client(client.clone())
.with_endpoint(signal_endpoint(&endpoint, "metrics"))
.build()
.map_err(|e| TelemetryError(format!("OTLP metric exporter configuration failed: {e}")))?;
let meter_provider = SdkMeterProvider::builder()
.with_periodic_exporter(metric_exporter)
.with_resource(resource())
.build();
let log_exporter = opentelemetry_otlp::LogExporter::builder()
.with_http()
.with_protocol(Protocol::HttpBinary)
.with_http_client(client)
.with_endpoint(signal_endpoint(&endpoint, "logs"))
.build()
.map_err(|e| TelemetryError(format!("OTLP log exporter configuration failed: {e}")))?;
let logger_provider = SdkLoggerProvider::builder()
.with_batch_exporter(log_exporter)
.with_resource(resource())
.build();
let _ = USAGE_LOGGER.set(logger_provider.logger(USAGE_SCOPE));
let otel: OtelLayer =
Box::new(tracing_opentelemetry::layer().with_tracer(tracer_provider.tracer(SERVICE_NAME)));
global::set_text_map_propagator(TraceContextPropagator::new());
global::set_tracer_provider(tracer_provider.clone());
global::set_meter_provider(meter_provider.clone());
install(filter, Some(otel))?;
metrics::init();
EXPORTING.store(true, Ordering::Relaxed);
Ok(TelemetryGuard {
tracer: Some(tracer_provider),
meter: Some(meter_provider),
logger: Some(logger_provider),
})
}
fn signal_endpoint(endpoint: &str, signal: &str) -> String {
let base = SIGNALS
.iter()
.fold(endpoint.trim_end_matches('/'), |base, s| {
base.trim_end_matches(&format!("/v1/{s}"))
});
format!("{}/v1/{signal}", base.trim_end_matches('/'))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn no_endpoint_means_telemetry_is_off() {
let config = TelemetryConfig::from_values(None, None).expect("default config");
assert!(config.endpoint.is_none());
let config = TelemetryConfig::from_values(Some(" "), None).expect("blank is off");
assert!(config.endpoint.is_none());
}
#[test]
fn rejects_unsupported_protocol_and_scheme() {
assert!(TelemetryConfig::from_values(Some("http://collector:4318"), Some("grpc")).is_err());
assert!(TelemetryConfig::from_values(Some("collector:4318"), None).is_err());
}
#[test]
fn signal_paths_are_appended_once() {
assert_eq!(
signal_endpoint("http://collector:4318", "traces"),
"http://collector:4318/v1/traces"
);
assert_eq!(
signal_endpoint("http://collector:4318/v1/metrics", "metrics"),
"http://collector:4318/v1/metrics"
);
assert_eq!(
signal_endpoint("http://collector:4318/v1/traces", "metrics"),
"http://collector:4318/v1/metrics"
);
assert_eq!(
signal_endpoint("http://collector:4318/otlp/", "traces"),
"http://collector:4318/otlp/v1/traces"
);
}
#[test]
fn resource_carries_the_service_name() {
let resource = resource();
assert_eq!(
resource
.get(&opentelemetry::Key::from_static_str("service.name"))
.map(|v| v.as_str().to_string()),
Some(SERVICE_NAME.to_string())
);
}
}