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(format!("{}/v1/traces", endpoint.trim_end_matches('/')))
.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();
}
const MALLOC_TRIM_INTERVAL_ENV: &str = "MYKO_MALLOC_TRIM_INTERVAL_SECS";
pub fn start_malloc_trim_probe() {
let Some(interval_secs) = std::env::var(MALLOC_TRIM_INTERVAL_ENV)
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&s| s > 0)
else {
return;
};
#[cfg(all(target_os = "linux", target_env = "gnu"))]
{
use std::sync::OnceLock;
static STARTED: OnceLock<()> = OnceLock::new();
if STARTED.set(()).is_err() {
return;
}
if jemalloc_linked() {
tracing::warn!(
target: "myko_server::mem_probe",
"{MALLOC_TRIM_INTERVAL_ENV} is set but this binary links jemalloc \
(_rjem_mallctl resolved) — malloc_trim only trims glibc arenas, which \
jemalloc bypasses. Probe disabled; use jemalloc's allocated-vs-resident \
stats instead."
);
return;
}
let _ = std::thread::Builder::new()
.name("myko-malloc-trim".to_string())
.spawn(move || run_malloc_trim_loop(interval_secs))
.map_err(|e| {
tracing::warn!(
target: "myko_server::mem_probe",
"Failed to spawn malloc_trim probe thread: {e}"
)
});
}
#[cfg(not(all(target_os = "linux", target_env = "gnu")))]
{
let _ = interval_secs;
tracing::warn!(
target: "myko_server::mem_probe",
"{MALLOC_TRIM_INTERVAL_ENV} is set but malloc_trim is glibc-only; probe disabled"
);
}
}
#[cfg(all(target_os = "linux", target_env = "gnu"))]
fn run_malloc_trim_loop(interval_secs: u64) {
unsafe extern "C" {
fn malloc_trim(pad: usize) -> i32;
}
let mb = |bytes: u64| bytes as f64 / (1024.0 * 1024.0);
let mut warned_wrong_allocator = false;
loop {
std::thread::sleep(Duration::from_secs(interval_secs));
let before = rss_bytes();
let released = unsafe { malloc_trim(0) } == 1;
let after = rss_bytes();
if let (Some(before), Some(after)) = (before, after) {
let released_bytes = before.saturating_sub(after);
tracing::info!(
target: "myko_server::mem_probe",
"[malloc_trim] rss_before={:.2}MB rss_after={:.2}MB released={:.2}MB ({})",
mb(before),
mb(after),
mb(released_bytes),
if released { "pages returned" } else { "no-op" },
);
if !warned_wrong_allocator
&& released_bytes < 16 * 1024 * 1024
&& after > 1024 * 1024 * 1024
{
warned_wrong_allocator = true;
tracing::warn!(
target: "myko_server::mem_probe",
"[malloc_trim] trim released almost nothing against {:.0}MB RSS. Either \
this heap is genuinely live, or this binary sets a non-glibc \
#[global_allocator] (e.g. jemalloc) that malloc_trim cannot touch — \
check the host's main.rs before drawing conclusions; under jemalloc \
use its allocated-vs-resident stats instead of this probe.",
mb(after),
);
}
} else {
tracing::warn!(
target: "myko_server::mem_probe",
"[malloc_trim] VmRSS unavailable in /proc/self/status; trim ran unmeasured"
);
}
}
}
#[cfg(all(target_os = "linux", target_env = "gnu"))]
fn jemalloc_linked() -> bool {
use core::ffi::{c_char, c_void};
unsafe extern "C" {
fn dlsym(handle: *mut c_void, symbol: *const c_char) -> *mut c_void;
}
unsafe { !dlsym(std::ptr::null_mut(), c"_rjem_mallctl".as_ptr()).is_null() }
}
#[cfg(all(target_os = "linux", target_env = "gnu"))]
fn rss_bytes() -> Option<u64> {
let status = std::fs::read_to_string("/proc/self/status").ok()?;
let line = status.lines().find(|l| l.starts_with("VmRSS:"))?;
let kb = line.split_whitespace().nth(1)?.parse::<u64>().ok()?;
Some(kb * 1024)
}
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(format!("{}/v1/metrics", endpoint.trim_end_matches('/')))
.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()
}