#[allow(dead_code)]
pub mod catalog;
#[cfg(test)]
mod assets;
mod exporter;
pub mod http;
pub mod metrics;
mod spans;
#[cfg(test)]
pub(crate) mod testing;
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,
CONVERGENCE_PRICING_BOUNDARY, 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::logs::LoggerProvider as _;
use opentelemetry::trace::TracerProvider as _;
use opentelemetry::{KeyValue, global};
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";
pub const INSTANCE_ID_ENV: &str = "AXOND_INSTANCE_ID";
const MAX_INSTANCE_ID_BYTES: usize = 128;
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>,
instance_id: Option<String>,
}
impl TelemetryConfig {
pub fn from_env() -> Result<Self, TelemetryError> {
Self::from_values_with_instance(
std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok().as_deref(),
std::env::var("OTEL_EXPORTER_OTLP_PROTOCOL").ok().as_deref(),
std::env::var(INSTANCE_ID_ENV).ok().as_deref(),
)
}
fn from_values_with_instance(
endpoint: Option<&str>,
protocol: Option<&str>,
instance_id: Option<&str>,
) -> Result<Self, TelemetryError> {
let Some(endpoint) = non_empty(endpoint) else {
return Ok(Self {
endpoint: None,
instance_id: validate_instance_id(instance_id)?,
});
};
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),
instance_id: validate_instance_id(instance_id)?,
})
}
}
fn non_empty(value: Option<&str>) -> Option<String> {
value
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_owned)
}
fn validate_instance_id(value: Option<&str>) -> Result<Option<String>, TelemetryError> {
let Some(value) = non_empty(value) else {
return Ok(None);
};
if value.len() > MAX_INSTANCE_ID_BYTES {
return Err(TelemetryError(format!(
"{INSTANCE_ID_ENV} must be at most {MAX_INSTANCE_ID_BYTES} bytes"
)));
}
if !value
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-'))
{
return Err(TelemetryError(format!(
"{INSTANCE_ID_ENV} may contain only ASCII letters, digits, `.`, `_`, and `-`"
)));
}
Ok(Some(value))
}
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(instance_id: Option<&str>) -> Resource {
let mut builder = Resource::builder().with_service_name(SERVICE_NAME);
if let Some(instance_id) = instance_id {
builder =
builder.with_attributes([KeyValue::new("service.instance.id", instance_id.to_owned())]);
}
builder.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(config.instance_id.as_deref()))
.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(config.instance_id.as_deref()))
.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(config.instance_id.as_deref()))
.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_with_instance(None, None, None).expect("default config");
assert!(config.endpoint.is_none());
let config = TelemetryConfig::from_values_with_instance(Some(" "), None, None)
.expect("blank is off");
assert!(config.endpoint.is_none());
}
#[test]
fn rejects_unsupported_protocol_and_scheme() {
assert!(
TelemetryConfig::from_values_with_instance(
Some("http://collector:4318"),
Some("grpc"),
None
)
.is_err()
);
assert!(
TelemetryConfig::from_values_with_instance(Some("collector:4318"), None, None).is_err()
);
}
#[test]
fn instance_identity_is_optional_bounded_and_validated() {
let without =
TelemetryConfig::from_values_with_instance(Some("http://collector:4318"), None, None)
.expect("an instance id is optional");
assert_eq!(without.instance_id, None);
let with = TelemetryConfig::from_values_with_instance(
Some("http://collector:4318"),
None,
Some("gateway-a_1.example"),
)
.expect("the documented identity alphabet is accepted");
assert_eq!(with.instance_id.as_deref(), Some("gateway-a_1.example"));
for invalid in ["gateway/a", "gateway a", "gateway:a"] {
assert!(
TelemetryConfig::from_values_with_instance(
Some("http://collector:4318"),
None,
Some(invalid),
)
.is_err(),
"{invalid} must not become a resource identity"
);
}
let too_long = "x".repeat(MAX_INSTANCE_ID_BYTES + 1);
assert!(
TelemetryConfig::from_values_with_instance(
Some("http://collector:4318"),
None,
Some(&too_long),
)
.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(None);
assert_eq!(
resource
.get(&opentelemetry::Key::from_static_str("service.name"))
.map(|v| v.as_str().to_string()),
Some(SERVICE_NAME.to_string())
);
}
#[test]
fn resource_carries_the_optional_instance_identity() {
let resource = resource(Some("gateway-a"));
assert_eq!(
resource
.get(&opentelemetry::Key::from_static_str("service.instance.id"))
.map(|value| value.as_str().to_string()),
Some("gateway-a".to_owned())
);
}
}