#[cfg(feature = "signals")]
use std::sync::Arc;
#[cfg(feature = "signals")]
use anyhow::Result;
#[cfg(feature = "signals")]
use tako_rs_core::plugins::TakoPlugin;
#[cfg(feature = "signals")]
use tako_rs_core::router::Router;
#[cfg(feature = "signals")]
use tako_rs_core::signals::Signal;
#[cfg(feature = "signals")]
use tako_rs_core::signals::app_events;
#[cfg(feature = "signals")]
use tako_rs_core::signals::ids;
#[cfg(feature = "signals")]
pub trait MetricsBackend: Send + Sync + 'static {
fn on_request_completed(&self, signal: &Signal);
fn on_route_request_completed(&self, signal: &Signal);
fn on_connection_opened(&self, signal: &Signal);
fn on_connection_closed(&self, signal: &Signal);
}
pub const DEFAULT_LATENCY_BUCKETS_SEC: &[f64] = &[
0.001, 0.0025, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0,
];
#[cfg(feature = "signals")]
pub struct MetricsPlugin<B: MetricsBackend> {
backend: Arc<B>,
}
#[cfg(feature = "signals")]
impl<B: MetricsBackend> Clone for MetricsPlugin<B> {
fn clone(&self) -> Self {
Self {
backend: Arc::clone(&self.backend),
}
}
}
#[cfg(feature = "signals")]
impl<B: MetricsBackend> MetricsPlugin<B> {
pub fn new(backend: B) -> Self {
Self {
backend: Arc::new(backend),
}
}
}
#[cfg(feature = "signals")]
impl<B: MetricsBackend> TakoPlugin for MetricsPlugin<B> {
fn name(&self) -> &'static str {
"MetricsPlugin"
}
#[cfg(feature = "signals")]
fn setup(&self, _router: &Router) -> Result<()> {
let backend = self.backend.clone();
let app_arbiter = app_events();
app_arbiter.on(ids::REQUEST_COMPLETED, move |signal: Signal| {
let backend = backend.clone();
async move {
backend.on_request_completed(&signal);
}
});
let backend_conn = self.backend.clone();
app_arbiter.on(ids::CONNECTION_OPENED, move |signal: Signal| {
let backend = backend_conn.clone();
async move {
backend.on_connection_opened(&signal);
}
});
let backend_close = self.backend.clone();
app_arbiter.on(ids::CONNECTION_CLOSED, move |signal: Signal| {
let backend = backend_close.clone();
async move {
backend.on_connection_closed(&signal);
}
});
let backend_route = self.backend.clone();
let mut rx = app_arbiter.subscribe_prefix("route.request.");
#[cfg(not(feature = "compio"))]
tokio::spawn(async move {
while let Ok(signal) = rx.recv().await {
backend_route.on_route_request_completed(&signal);
}
});
#[cfg(feature = "compio")]
compio::runtime::spawn(async move {
while let Ok(signal) = rx.recv().await {
backend_route.on_route_request_completed(&signal);
}
})
.detach();
Ok(())
}
#[cfg(not(feature = "signals"))]
fn setup(&self, _router: &Router) -> Result<()> {
Ok(())
}
}