use crate::memory_limiter::MemoryPressureState;
use cpu_time::ProcessTime;
use otel_arrow_dfe_telemetry::instrument::{Gauge, ObserveUpDownCounter};
use otel_arrow_dfe_telemetry::metrics::MetricSet;
use otel_arrow_dfe_telemetry::registry::{EntityKey, TelemetryRegistryHandle};
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use otel_arrow_dfe_telemetry_macros::metric_set;
use std::time::Instant;
#[metric_set(name = "engine")]
#[derive(Debug, Default, Clone)]
pub struct EngineMetrics {
#[metric(unit = "{By}")]
pub memory_rss: ObserveUpDownCounter<u64>,
#[metric(unit = "{1}")]
pub cpu_utilization: Gauge<f64>,
#[metric(unit = "{state}")]
pub memory_pressure_state: Gauge<u64>,
#[metric(unit = "{By}")]
pub process_memory_usage_bytes: Gauge<u64>,
#[metric(unit = "{By}")]
pub process_memory_soft_limit_bytes: Gauge<u64>,
#[metric(unit = "{By}")]
pub process_memory_hard_limit_bytes: Gauge<u64>,
}
pub struct EngineMetricsMonitor {
metrics: MetricSet<EngineMetrics>,
reporter: MetricsReporter,
registry: TelemetryRegistryHandle,
wall_start: Instant,
cpu_start: ProcessTime,
num_cores: usize,
memory_pressure_state: MemoryPressureState,
}
impl EngineMetricsMonitor {
#[must_use]
pub fn new(
registry: TelemetryRegistryHandle,
entity_key: EntityKey,
reporter: MetricsReporter,
memory_pressure_state: MemoryPressureState,
) -> Self {
let metrics = registry.register_metric_set_for_entity::<EngineMetrics>(entity_key);
let num_cores = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1);
Self {
metrics,
reporter,
registry,
wall_start: Instant::now(),
cpu_start: ProcessTime::now(),
num_cores,
memory_pressure_state,
}
}
pub fn update(&mut self) {
self.metrics.memory_rss.observe(get_rss_bytes());
let now_wall = Instant::now();
let now_cpu = ProcessTime::now();
let wall_delta = now_wall.duration_since(self.wall_start);
let cpu_delta = now_cpu.duration_since(self.cpu_start);
let wall_secs = wall_delta.as_secs_f64();
if wall_secs > 0.0 {
let utilization =
(cpu_delta.as_secs_f64() / (wall_secs * self.num_cores as f64)).clamp(0.0, 1.0);
self.metrics.cpu_utilization.set(utilization);
} else {
self.metrics.cpu_utilization.set(0.0);
}
self.metrics
.memory_pressure_state
.set(self.memory_pressure_state.level() as u64);
self.metrics
.process_memory_usage_bytes
.set(self.memory_pressure_state.usage_bytes());
self.metrics
.process_memory_soft_limit_bytes
.set(self.memory_pressure_state.soft_limit_bytes());
self.metrics
.process_memory_hard_limit_bytes
.set(self.memory_pressure_state.hard_limit_bytes());
self.wall_start = now_wall;
self.cpu_start = now_cpu;
}
pub fn report(&mut self) -> Result<(), otel_arrow_dfe_telemetry::error::Error> {
self.reporter.report(&mut self.metrics)
}
pub async fn finish_reporting_until(
&mut self,
deadline: Instant,
) -> Result<(), otel_arrow_dfe_telemetry::error::Error> {
self.update();
let _ = self
.reporter
.report_reliably_until(&mut self.metrics, deadline)
.await?;
self.reporter.flush_until(deadline).await
}
}
fn get_rss_bytes() -> u64 {
memory_stats::memory_stats()
.map(|stats| stats.physical_mem as u64)
.unwrap_or(0)
}
impl Drop for EngineMetricsMonitor {
fn drop(&mut self) {
let _ = self
.registry
.unregister_metric_set(self.metrics.metric_set_key());
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::context::ControllerContext;
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
#[test]
fn engine_metrics_reports_nonzero_rss() {
let registry = TelemetryRegistryHandle::new();
let controller = ControllerContext::new(registry.clone());
let entity_key = controller.register_engine_entity();
let (_rx, reporter) = MetricsReporter::create_new_and_receiver(16);
let mut monitor = EngineMetricsMonitor::new(
registry,
entity_key,
reporter,
controller.memory_pressure_state(),
);
monitor.update();
assert!(
monitor.metrics.memory_rss.get() > 0,
"memory_rss should report non-zero process RSS"
);
}
#[test]
fn engine_metrics_report_succeeds() {
let registry = TelemetryRegistryHandle::new();
let controller = ControllerContext::new(registry.clone());
let entity_key = controller.register_engine_entity();
let (_rx, reporter) = MetricsReporter::create_new_and_receiver(16);
let mut monitor = EngineMetricsMonitor::new(
registry,
entity_key,
reporter,
controller.memory_pressure_state(),
);
monitor.update();
assert!(monitor.report().is_ok());
}
#[test]
fn engine_metrics_cpu_utilization_in_range() {
let registry = TelemetryRegistryHandle::new();
let controller = ControllerContext::new(registry.clone());
let entity_key = controller.register_engine_entity();
let (_rx, reporter) = MetricsReporter::create_new_and_receiver(16);
let mut monitor = EngineMetricsMonitor::new(
registry,
entity_key,
reporter,
controller.memory_pressure_state(),
);
let start = Instant::now();
while start.elapsed() < std::time::Duration::from_millis(10) {
let _ = std::hint::black_box(0u64.wrapping_add(1));
}
monitor.update();
let util = monitor.metrics.cpu_utilization.get();
assert!(
(0.0..=1.0).contains(&util),
"cpu_utilization should be in [0, 1], got {util}"
);
}
#[test]
fn engine_metrics_expose_process_memory_limiter_usage_and_limits() {
let registry = TelemetryRegistryHandle::new();
let controller = ControllerContext::new(registry.clone());
let state = controller.memory_pressure_state();
state.configure(crate::memory_limiter::MemoryPressureBehaviorConfig {
retry_after_secs: 1,
fail_readiness_on_hard: true,
mode: otel_arrow_dfe_config::policy::MemoryLimiterMode::Enforce,
});
state.set_sample_for_tests(
crate::memory_limiter::MemoryPressureLevel::Soft,
95,
90,
100,
);
let entity_key = controller.register_engine_entity();
let (_rx, reporter) = MetricsReporter::create_new_and_receiver(16);
let mut monitor = EngineMetricsMonitor::new(registry, entity_key, reporter, state);
monitor.update();
assert_eq!(monitor.metrics.memory_pressure_state.get(), 1);
assert_eq!(monitor.metrics.process_memory_usage_bytes.get(), 95);
assert_eq!(monitor.metrics.process_memory_soft_limit_bytes.get(), 90);
assert_eq!(monitor.metrics.process_memory_hard_limit_bytes.get(), 100);
}
}