use anyhow::Result;
use tako_rs_core::router::Router;
use crate::plugins::metrics::MetricsPlugin;
#[cfg(feature = "metrics-opentelemetry")]
pub mod opentelemetry_backend {
use opentelemetry::KeyValue;
use opentelemetry::metrics::Counter;
use opentelemetry::metrics::Meter;
use tako_rs_core::signals::Signal;
use crate::plugins::metrics::MetricsBackend;
pub struct OtelMetricsBackend {
http_requests_total: Counter<u64>,
http_route_requests_total: Counter<u64>,
connections_opened_total: Counter<u64>,
connections_closed_total: Counter<u64>,
}
impl OtelMetricsBackend {
pub fn new(meter: Meter) -> Self {
let http_requests_total = meter.u64_counter("tako_http_requests_total").build();
let http_route_requests_total = meter.u64_counter("tako_route_requests_total").build();
let connections_opened_total = meter.u64_counter("tako_connections_opened_total").build();
let connections_closed_total = meter.u64_counter("tako_connections_closed_total").build();
Self {
http_requests_total,
http_route_requests_total,
connections_opened_total,
connections_closed_total,
}
}
}
fn transport_label(signal: &Signal) -> &'static str {
if signal.metadata.get("protocol").map(String::as_str) == Some("h3") {
"h3"
} else if signal.metadata.get("tls").map(String::as_str) == Some("true") {
"tls"
} else if signal.metadata.contains_key("unix_path") {
"unix"
} else {
"tcp"
}
}
fn route_label(signal: &Signal) -> String {
signal
.metadata
.get("route")
.cloned()
.unwrap_or_else(|| "unmatched".to_string())
}
impl MetricsBackend for OtelMetricsBackend {
fn on_request_completed(&self, signal: &Signal) {
let method = signal.metadata.get("method").cloned().unwrap_or_default();
let route = route_label(signal);
let status = signal.metadata.get("status").cloned().unwrap_or_default();
self.http_requests_total.add(
1,
&[
KeyValue::new("method", method),
KeyValue::new("route", route),
KeyValue::new("status", status),
],
);
}
fn on_route_request_completed(&self, signal: &Signal) {
let method = signal.metadata.get("method").cloned().unwrap_or_default();
let route = route_label(signal);
let status = signal.metadata.get("status").cloned().unwrap_or_default();
self.http_route_requests_total.add(
1,
&[
KeyValue::new("method", method),
KeyValue::new("route", route),
KeyValue::new("status", status),
],
);
}
fn on_connection_opened(&self, signal: &Signal) {
self
.connections_opened_total
.add(1, &[KeyValue::new("transport", transport_label(signal))]);
}
fn on_connection_closed(&self, signal: &Signal) {
self
.connections_closed_total
.add(1, &[KeyValue::new("transport", transport_label(signal))]);
}
}
}
#[cfg(feature = "metrics-opentelemetry")]
#[derive(Clone)]
pub struct OtelMetricsConfig {
pub meter_name: &'static str,
pub endpoint: String,
}
#[cfg(feature = "metrics-opentelemetry")]
impl Default for OtelMetricsConfig {
fn default() -> Self {
Self {
meter_name: "tako",
endpoint: "http://localhost:4318/v1/metrics".to_string(),
}
}
}
#[cfg(feature = "metrics-opentelemetry")]
impl OtelMetricsConfig {
pub fn with_endpoint(mut self, endpoint: impl Into<String>) -> Self {
self.endpoint = endpoint.into();
self
}
pub fn with_meter_name(mut self, name: &'static str) -> Self {
self.meter_name = name;
self
}
pub fn install(
self,
router: &mut Router,
) -> Result<opentelemetry_sdk::metrics::SdkMeterProvider> {
use opentelemetry::global;
use opentelemetry_otlp::WithExportConfig;
let exporter = opentelemetry_otlp::MetricExporter::builder()
.with_http()
.with_endpoint(&self.endpoint)
.build()
.map_err(|e| anyhow::anyhow!("failed to create OTLP metric exporter: {e}"))?;
let meter_provider = opentelemetry_sdk::metrics::SdkMeterProvider::builder()
.with_periodic_exporter(exporter)
.build();
global::set_meter_provider(meter_provider.clone());
let meter = global::meter(self.meter_name);
let backend = opentelemetry_backend::OtelMetricsBackend::new(meter);
let plugin = MetricsPlugin::new(backend);
router.plugin(plugin);
Ok(meter_provider)
}
}