#[cfg(not(feature = "std"))]
use alloc::sync::Arc;
use core::sync::atomic::{AtomicU32, AtomicU64, Ordering};
#[cfg(feature = "std")]
use std::sync::Arc;
use crate::utils::BasisPoints;
#[derive(Debug, Default)]
pub struct ServletMetrics {
queue_depth: AtomicU32,
}
impl ServletMetrics {
pub fn new() -> Self {
Self { queue_depth: AtomicU32::new(0) }
}
pub fn enqueue(&self) {
self.queue_depth.fetch_add(1, Ordering::Relaxed);
}
pub fn dequeue(&self) {
self.queue_depth.fetch_sub(1, Ordering::Relaxed);
}
pub fn queue_depth(&self) -> u32 {
self.queue_depth.load(Ordering::Relaxed)
}
}
#[derive(Debug)]
pub struct LatencyTracker {
ema_us: AtomicU64,
alpha_bps: u16,
target_us: u64,
}
impl LatencyTracker {
pub const fn new(alpha_bps: u16, target_us: u64) -> Self {
Self { ema_us: AtomicU64::new(0), alpha_bps, target_us }
}
pub fn record(&self, latency_us: u64) {
let alpha = self.alpha_bps as u64;
let one_minus_alpha = 10000u64.saturating_sub(alpha);
loop {
let current = self.ema_us.load(Ordering::Relaxed);
let new_ema = (alpha.saturating_mul(latency_us) + one_minus_alpha.saturating_mul(current)) / 10000;
let success = Ordering::Release;
let failure = Ordering::Relaxed;
match self.ema_us.compare_exchange_weak(current, new_ema, success, failure) {
Ok(_) => break,
Err(_) => continue, }
}
}
pub fn utilization(&self) -> BasisPoints {
let ema = self.ema_us.load(Ordering::Relaxed);
if self.target_us == 0 {
return BasisPoints::MAX;
}
let ratio_bps = ema.saturating_mul(10000) / self.target_us;
BasisPoints::new_saturating(ratio_bps.min(10000) as u16)
}
pub fn ema_microseconds(&self) -> u64 {
self.ema_us.load(Ordering::Relaxed)
}
pub fn reset(&self) {
self.ema_us.store(0, Ordering::Relaxed);
}
}
impl Default for LatencyTracker {
fn default() -> Self {
Self::new(2000, 100_000)
}
}
const DEFAULT_LATENCY_WEIGHT_BPS: u16 = 7000;
#[derive(Debug)]
pub struct UtilizationReporter {
tracker: LatencyTracker,
queue: Option<Arc<ServletMetrics>>,
queue_capacity: u32,
latency_weight_bps: u16,
}
impl UtilizationReporter {
pub const fn new(alpha_bps: u16, target_us: u64) -> Self {
Self {
tracker: LatencyTracker::new(alpha_bps, target_us),
queue: None,
queue_capacity: 100,
latency_weight_bps: DEFAULT_LATENCY_WEIGHT_BPS,
}
}
pub fn with_queue(
alpha_bps: u16,
target_us: u64,
queue: Arc<ServletMetrics>,
capacity: u32,
latency_weight_bps: u16,
) -> Self {
Self {
tracker: LatencyTracker::new(alpha_bps, target_us),
queue: Some(queue),
queue_capacity: capacity,
latency_weight_bps,
}
}
pub fn record_latency(&self, duration: core::time::Duration) {
self.tracker.record(duration.as_micros() as u64);
}
pub fn record_microseconds(&self, latency_us: u64) {
self.tracker.record(latency_us);
}
pub fn utilization(&self) -> BasisPoints {
let latency_util = self.tracker.utilization().get() as u32;
let Some(ref queue) = self.queue else {
return self.tracker.utilization();
};
let queue_util = queue
.queue_depth()
.saturating_mul(10000)
.checked_div(self.queue_capacity)
.unwrap_or(0);
let latency_weight = self.latency_weight_bps as u32;
let queue_weight = 10000u32.saturating_sub(latency_weight);
let combined = (latency_weight.saturating_mul(latency_util) + queue_weight.saturating_mul(queue_util)) / 10000;
BasisPoints::new_saturating(combined.min(10000) as u16)
}
pub fn ema_microseconds(&self) -> u64 {
self.tracker.ema_microseconds()
}
pub fn queue_depth(&self) -> Option<u32> {
self.queue.as_ref().map(|q| q.queue_depth())
}
pub fn reset(&self) {
self.tracker.reset();
}
}
impl Default for UtilizationReporter {
fn default() -> Self {
Self::new(2000, 100_000)
}
}