use std::{sync::Arc, time::Duration};
use hyphae::Gettable;
use myko::store::StoreRegistry;
use opentelemetry::{KeyValue, global, trace::TracerProvider};
use opentelemetry_otlp::WithExportConfig;
use opentelemetry_sdk::{Resource, metrics::SdkMeterProvider, trace::SdkTracerProvider};
use tracing_subscriber::{
EnvFilter, Layer, layer::SubscriberExt, registry::LookupSpan, util::SubscriberInitExt,
};
const DEFAULT_METRICS_INTERVAL_SECS: u64 = 60;
pub struct TelemetryGuard {
tracer_provider: Option<SdkTracerProvider>,
meter_provider: Option<SdkMeterProvider>,
}
impl Drop for TelemetryGuard {
fn drop(&mut self) {
if let Some(provider) = self.tracer_provider.take()
&& let Err(e) = provider.shutdown()
{
eprintln!("myko telemetry: tracer provider shutdown error: {e}");
}
if let Some(provider) = self.meter_provider.take()
&& let Err(e) = provider.shutdown()
{
eprintln!("myko telemetry: meter provider shutdown error: {e}");
}
}
}
pub fn init_from_env() -> TelemetryGuard {
let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info"));
let fmt_layer = tracing_subscriber::fmt::layer();
let (otel_layer, guard) = otel_layer_from_env();
tracing_subscriber::registry()
.with(filter)
.with(fmt_layer)
.with(otel_layer)
.init();
guard.unwrap_or(TelemetryGuard {
tracer_provider: None,
meter_provider: None,
})
}
pub fn otel_layer_from_env<S>() -> (Option<impl Layer<S> + Send + Sync>, Option<TelemetryGuard>)
where
S: tracing::Subscriber + for<'a> LookupSpan<'a> + Send + Sync,
{
let Ok(endpoint) = std::env::var("MYKO_TRACING_ENDPOINT") else {
return (None, None);
};
let resource = Resource::builder().with_service_name("myko-server").build();
let tracer_provider = build_tracer_provider(&endpoint, resource.clone());
let meter_provider = build_meter_provider(&endpoint, resource);
global::set_meter_provider(meter_provider.clone());
let tracer = tracer_provider.tracer("myko-server");
let otel_layer = tracing_opentelemetry::layer().with_tracer(tracer);
(
Some(otel_layer),
Some(TelemetryGuard {
tracer_provider: Some(tracer_provider),
meter_provider: Some(meter_provider),
}),
)
}
fn build_tracer_provider(endpoint: &str, resource: Resource) -> SdkTracerProvider {
let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_http()
.with_endpoint(endpoint)
.build()
.expect("failed to build OTLP/HTTP trace exporter");
SdkTracerProvider::builder()
.with_batch_exporter(exporter)
.with_resource(resource)
.build()
}
pub fn register_item_count_gauge(registry: Arc<StoreRegistry>) {
let meter = global::meter("myko-server");
let _gauge = meter
.u64_observable_gauge("myko.store.item_count")
.with_description("Live entity count per store, sampled on each metrics export")
.with_callback(move |observer| {
for entity_type in registry.entity_types() {
let count = registry.get_or_create(&entity_type).len().get() as u64;
observer.observe(
count,
&[KeyValue::new("entity_type", entity_type.to_string())],
);
}
})
.build();
}
fn build_meter_provider(endpoint: &str, resource: Resource) -> SdkMeterProvider {
let interval_secs = std::env::var("MYKO_MEM_PROFILE_INTERVAL_SECS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.unwrap_or(DEFAULT_METRICS_INTERVAL_SECS);
let exporter = opentelemetry_otlp::MetricExporter::builder()
.with_http()
.with_endpoint(endpoint)
.build()
.expect("failed to build OTLP/HTTP metrics exporter");
let reader = opentelemetry_sdk::metrics::PeriodicReader::builder(exporter)
.with_interval(Duration::from_secs(interval_secs))
.build();
SdkMeterProvider::builder()
.with_reader(reader)
.with_resource(resource)
.build()
}