use crate::context::PipelineContext;
use cpu_time::ThreadTime;
use otel_arrow_dfe_telemetry::instrument::{Counter, Gauge, ObserveCounter, ObserveUpDownCounter};
use otel_arrow_dfe_telemetry::metrics::MetricSet;
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
use otel_arrow_dfe_telemetry_macros::metric_set;
use std::time::Instant;
#[cfg(all(not(windows), feature = "jemalloc"))]
use tikv_jemalloc_ctl::{thread, thread::ThreadLocal};
#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "openbsd"))]
use nix::sys::resource::{UsageWho, getrusage};
#[metric_set(name = "pipeline")]
#[derive(Debug, Default, Clone)]
pub struct PipelineMetrics {
#[metric(unit = "{s}")]
pub uptime: Gauge<f64>,
#[metric(unit = "{By}")]
pub memory_usage: ObserveUpDownCounter<u64>,
#[metric(unit = "{By}")]
pub memory_allocated: ObserveCounter<u64>,
#[metric(unit = "{By}")]
pub memory_freed: ObserveCounter<u64>,
#[metric(unit = "{By}")]
pub memory_allocated_delta: Counter<u64>,
#[metric(unit = "{By}")]
pub memory_freed_delta: Counter<u64>,
#[metric(unit = "{s}")]
pub cpu_time: Counter<f64>,
#[metric(unit = "{1}")]
pub cpu_utilization: Gauge<f64>,
#[metric(unit = "{1}")]
pub context_switches_voluntary: ObserveCounter<u64>,
#[metric(unit = "{1}")]
pub context_switches_involuntary: ObserveCounter<u64>,
#[metric(unit = "{1}")]
pub page_faults_minor: ObserveCounter<u64>,
#[metric(unit = "{1}")]
pub page_faults_major: ObserveCounter<u64>,
}
#[metric_set(name = "tokio.runtime")]
#[derive(Debug, Default, Clone)]
pub struct TokioRuntimeMetrics {
#[metric(unit = "{thread}")]
pub worker_count: ObserveUpDownCounter<u64>,
#[metric(unit = "{task}")]
pub task_active_count: ObserveUpDownCounter<u64>,
#[metric(unit = "{task}")]
pub global_task_queue_size: ObserveUpDownCounter<u64>,
#[cfg(target_has_atomic = "64")]
#[metric(unit = "{s}")]
pub worker_busy_time: ObserveCounter<f64>,
#[cfg(target_has_atomic = "64")]
#[metric(unit = "{park}")]
pub worker_park_count: ObserveCounter<u64>,
#[cfg(target_has_atomic = "64")]
#[metric(unit = "{unpark}")]
pub worker_park_unpark_count: ObserveCounter<u64>,
#[cfg(tokio_unstable)]
#[metric(unit = "{task}")]
pub blocking_task_queue_size: ObserveUpDownCounter<u64>,
#[cfg(tokio_unstable)]
#[metric(unit = "{thread}")]
pub blocking_thread_count: ObserveUpDownCounter<u64>,
#[cfg(tokio_unstable)]
#[metric(unit = "{thread}")]
pub blocking_thread_idle_count: ObserveUpDownCounter<u64>,
#[cfg(tokio_unstable)]
#[metric(unit = "{task}")]
pub worker_local_queue_size: ObserveUpDownCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{spawn}")]
pub spawned_tasks_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{remote}")]
pub remote_schedule_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{yield}")]
pub budget_forced_yield_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{noop}")]
pub worker_noop_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{steal}")]
pub worker_steal_success_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{steal}")]
pub worker_steal_attempt_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{poll}")]
pub worker_poll_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{schedule}")]
pub worker_local_schedule_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{overflow}")]
pub worker_overflow_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{fd}")]
pub io_driver_fd_registered_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{fd}")]
pub io_driver_fd_deregistered_count: ObserveCounter<u64>,
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
#[metric(unit = "{ready}")]
pub io_driver_ready_count: ObserveCounter<u64>,
}
pub(crate) struct PipelineMetricsMonitor {
start_time: Instant,
#[cfg(all(not(windows), feature = "jemalloc"))]
jemalloc_supported: bool,
#[cfg(all(not(windows), feature = "jemalloc"))]
allocated: Option<ThreadLocal<u64>>,
#[cfg(all(not(windows), feature = "jemalloc"))]
deallocated: Option<ThreadLocal<u64>>,
#[cfg(all(not(windows), feature = "jemalloc"))]
last_allocated: u64,
#[cfg(all(not(windows), feature = "jemalloc"))]
last_deallocated: u64,
rusage_thread_supported: bool,
wall_start: Instant,
cpu_start: ThreadTime,
metrics: MetricSet<PipelineMetrics>,
tokio_rt: Option<tokio::runtime::RuntimeMetrics>,
tokio_metrics: MetricSet<TokioRuntimeMetrics>,
registry: TelemetryRegistryHandle,
}
impl PipelineMetricsMonitor {
pub(crate) fn new(pipeline_ctx: PipelineContext) -> Self {
let now = Instant::now();
#[cfg(all(not(windows), feature = "jemalloc"))]
let (jemalloc_supported, allocated, deallocated, last_allocated, last_deallocated) = {
let jemalloc_init = (|| {
let alloc_mib = thread::allocatedp::mib().ok()?;
let dealloc_mib = thread::deallocatedp::mib().ok()?;
let allocated = alloc_mib.read().ok()?;
let deallocated = dealloc_mib.read().ok()?;
let last_allocated = allocated.get();
let last_deallocated = deallocated.get();
Some((allocated, deallocated, last_allocated, last_deallocated))
})();
if let Some((allocated, deallocated, last_allocated, last_deallocated)) = jemalloc_init
{
(
true,
Some(allocated),
Some(deallocated),
last_allocated,
last_deallocated,
)
} else {
(false, None, None, 0, 0)
}
};
let rusage_thread_supported = Self::init_rusage_baseline();
let tokio_rt = tokio::runtime::Handle::try_current()
.ok()
.map(|handle| handle.metrics());
let entity_key = crate::entity_context::pipeline_entity_key().expect(
"pipeline entity key not set; ensure pipeline entity is registered and instrumented",
);
let metrics = pipeline_ctx.register_metric_set_for_entity::<PipelineMetrics>(entity_key);
let tokio_metrics =
pipeline_ctx.register_metric_set_for_entity::<TokioRuntimeMetrics>(entity_key);
let registry = pipeline_ctx.metrics_registry();
Self {
start_time: now,
#[cfg(all(not(windows), feature = "jemalloc"))]
jemalloc_supported,
#[cfg(all(not(windows), feature = "jemalloc"))]
allocated,
#[cfg(all(not(windows), feature = "jemalloc"))]
deallocated,
#[cfg(all(not(windows), feature = "jemalloc"))]
last_allocated,
#[cfg(all(not(windows), feature = "jemalloc"))]
last_deallocated,
rusage_thread_supported,
wall_start: now,
cpu_start: ThreadTime::now(),
metrics,
tokio_rt,
tokio_metrics,
registry,
}
}
pub const fn metrics_mut(&mut self) -> &mut MetricSet<PipelineMetrics> {
&mut self.metrics
}
pub const fn tokio_metrics_mut(&mut self) -> &mut MetricSet<TokioRuntimeMetrics> {
&mut self.tokio_metrics
}
#[cfg(test)]
#[allow(dead_code)] pub fn update_metrics(&mut self) {
self.update_tokio_metrics();
self.update_pipeline_metrics();
}
pub fn update_pipeline_metrics(&mut self) {
#[cfg(all(not(windows), feature = "jemalloc"))]
if self.jemalloc_supported
&& let (Some(allocated), Some(deallocated)) =
(self.allocated.as_ref(), self.deallocated.as_ref())
{
let cur_alloc = allocated.get();
let cur_dealloc = deallocated.get();
let delta_alloc = cur_alloc.wrapping_sub(self.last_allocated);
let delta_dealloc = cur_dealloc.wrapping_sub(self.last_deallocated);
self.last_allocated = cur_alloc;
self.last_deallocated = cur_dealloc;
self.metrics.memory_allocated.observe(cur_alloc);
self.metrics.memory_freed.observe(cur_dealloc);
self.metrics.memory_allocated_delta.add(delta_alloc);
self.metrics.memory_freed_delta.add(delta_dealloc);
self.metrics
.memory_usage
.observe(cur_alloc.saturating_sub(cur_dealloc));
}
self.update_rusage_metrics();
let now_wall = Instant::now();
let now_cpu = ThreadTime::now();
let uptime = now_wall.duration_since(self.start_time).as_secs_f64();
self.metrics.uptime.set(uptime);
let wall_duration = now_wall.duration_since(self.wall_start);
let cpu_duration = now_cpu.duration_since(self.cpu_start);
let wall_micros = wall_duration.as_micros();
self.metrics.cpu_time.add(cpu_duration.as_secs_f64());
if wall_micros > 0 {
let cpu_micros = cpu_duration.as_micros();
let usage_ratio = cpu_micros as f64 / wall_micros as f64;
let usage_ratio = usage_ratio.clamp(0.0, 1.0);
self.metrics.cpu_utilization.set(usage_ratio);
} else {
self.metrics.cpu_utilization.set(0.0);
}
self.wall_start = now_wall;
self.cpu_start = now_cpu;
}
pub fn update_tokio_metrics(&mut self) {
if self.tokio_rt.is_none() {
self.tokio_rt = tokio::runtime::Handle::try_current()
.ok()
.map(|handle| handle.metrics());
}
let Some(tokio_rt) = self.tokio_rt.as_ref() else {
return;
};
let num_workers_usize = tokio_rt.num_workers();
let num_workers = u64::try_from(num_workers_usize).unwrap_or(u64::MAX);
let num_alive_tasks = u64::try_from(tokio_rt.num_alive_tasks()).unwrap_or(u64::MAX);
let global_queue_depth = u64::try_from(tokio_rt.global_queue_depth()).unwrap_or(u64::MAX);
self.tokio_metrics.worker_count.observe(num_workers);
self.tokio_metrics
.task_active_count
.observe(num_alive_tasks);
self.tokio_metrics
.global_task_queue_size
.observe(global_queue_depth);
#[cfg(target_has_atomic = "64")]
{
let mut total_busy_s = 0.0f64;
let mut total_park_count = 0u64;
let mut total_park_unpark_count = 0u64;
for worker in 0..num_workers_usize {
total_busy_s += tokio_rt.worker_total_busy_duration(worker).as_secs_f64();
total_park_count =
total_park_count.saturating_add(tokio_rt.worker_park_count(worker));
total_park_unpark_count = total_park_unpark_count
.saturating_add(tokio_rt.worker_park_unpark_count(worker));
}
self.tokio_metrics.worker_busy_time.observe(total_busy_s);
self.tokio_metrics
.worker_park_count
.observe(total_park_count);
self.tokio_metrics
.worker_park_unpark_count
.observe(total_park_unpark_count);
}
#[cfg(tokio_unstable)]
{
let blocking_queue_depth =
u64::try_from(tokio_rt.blocking_queue_depth()).unwrap_or(u64::MAX);
let num_blocking_threads =
u64::try_from(tokio_rt.num_blocking_threads()).unwrap_or(u64::MAX);
let num_idle_blocking_threads =
u64::try_from(tokio_rt.num_idle_blocking_threads()).unwrap_or(u64::MAX);
let mut local_queue_depth_sum = 0u64;
for worker in 0..num_workers_usize {
local_queue_depth_sum = local_queue_depth_sum.saturating_add(
u64::try_from(tokio_rt.worker_local_queue_depth(worker)).unwrap_or(u64::MAX),
);
}
self.tokio_metrics
.blocking_task_queue_size
.observe(blocking_queue_depth);
self.tokio_metrics
.blocking_thread_count
.observe(num_blocking_threads);
self.tokio_metrics
.blocking_thread_idle_count
.observe(num_idle_blocking_threads);
self.tokio_metrics
.worker_local_queue_size
.observe(local_queue_depth_sum);
}
#[cfg(all(tokio_unstable, target_has_atomic = "64"))]
{
self.tokio_metrics
.spawned_tasks_count
.observe(tokio_rt.spawned_tasks_count());
self.tokio_metrics
.remote_schedule_count
.observe(tokio_rt.remote_schedule_count());
self.tokio_metrics
.budget_forced_yield_count
.observe(tokio_rt.budget_forced_yield_count());
self.tokio_metrics
.io_driver_fd_registered_count
.observe(tokio_rt.io_driver_fd_registered_count());
self.tokio_metrics
.io_driver_fd_deregistered_count
.observe(tokio_rt.io_driver_fd_deregistered_count());
self.tokio_metrics
.io_driver_ready_count
.observe(tokio_rt.io_driver_ready_count());
let mut worker_noop_sum = 0u64;
let mut worker_steal_sum = 0u64;
let mut worker_steal_ops_sum = 0u64;
let mut worker_poll_sum = 0u64;
let mut worker_local_schedule_sum = 0u64;
let mut worker_overflow_sum = 0u64;
for worker in 0..num_workers_usize {
worker_noop_sum =
worker_noop_sum.saturating_add(tokio_rt.worker_noop_count(worker));
worker_steal_sum =
worker_steal_sum.saturating_add(tokio_rt.worker_steal_count(worker));
worker_steal_ops_sum =
worker_steal_ops_sum.saturating_add(tokio_rt.worker_steal_operations(worker));
worker_poll_sum =
worker_poll_sum.saturating_add(tokio_rt.worker_poll_count(worker));
worker_local_schedule_sum = worker_local_schedule_sum
.saturating_add(tokio_rt.worker_local_schedule_count(worker));
worker_overflow_sum =
worker_overflow_sum.saturating_add(tokio_rt.worker_overflow_count(worker));
}
self.tokio_metrics
.worker_noop_count
.observe(worker_noop_sum);
self.tokio_metrics
.worker_steal_success_count
.observe(worker_steal_sum);
self.tokio_metrics
.worker_steal_attempt_count
.observe(worker_steal_ops_sum);
self.tokio_metrics
.worker_poll_count
.observe(worker_poll_sum);
self.tokio_metrics
.worker_local_schedule_count
.observe(worker_local_schedule_sum);
self.tokio_metrics
.worker_overflow_count
.observe(worker_overflow_sum);
}
}
fn init_rusage_baseline() -> bool {
#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "openbsd"))]
{
if let Ok(_usage) = getrusage(UsageWho::RUSAGE_THREAD) {
return true;
}
}
false
}
#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "openbsd"))]
fn update_rusage_metrics(&mut self) {
if !self.rusage_thread_supported {
return;
}
match getrusage(UsageWho::RUSAGE_THREAD) {
Ok(usage) => {
let voluntary =
u64::try_from(usage.voluntary_context_switches()).unwrap_or_default();
let involuntary =
u64::try_from(usage.involuntary_context_switches()).unwrap_or_default();
let minor_faults = u64::try_from(usage.minor_page_faults()).unwrap_or_default();
let major_faults = u64::try_from(usage.major_page_faults()).unwrap_or_default();
self.metrics.context_switches_voluntary.observe(voluntary);
self.metrics
.context_switches_involuntary
.observe(involuntary);
self.metrics.page_faults_minor.observe(minor_faults);
self.metrics.page_faults_major.observe(major_faults);
}
Err(_) => {
self.rusage_thread_supported = false;
}
}
}
#[cfg(not(any(target_os = "linux", target_os = "freebsd", target_os = "openbsd")))]
fn update_rusage_metrics(&mut self) {
let _ = self.rusage_thread_supported;
}
}
impl Drop for PipelineMetricsMonitor {
fn drop(&mut self) {
let _ = self
.registry
.unregister_metric_set(self.metrics.metric_set_key());
let _ = self
.registry
.unregister_metric_set(self.tokio_metrics.metric_set_key());
}
}
#[cfg(all(test, not(windows), feature = "jemalloc-testing"))]
mod jemalloc_tests {
use super::*;
use crate::context::ControllerContext;
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
use std::hint::black_box;
use std::time::{Duration, Instant};
#[global_allocator]
static GLOBAL: tikv_jemallocator::Jemalloc = tikv_jemallocator::Jemalloc;
#[test]
fn pipeline_metrics_monitor_black_box_updates_jemalloc() {
let telemetry_registry = TelemetryRegistryHandle::new();
let controller = ControllerContext::new(telemetry_registry);
let pipeline_ctx = controller.pipeline_context_with("grp".into(), "pipe".into(), 0, 1, 0);
let pipeline_entity_key = pipeline_ctx.register_pipeline_entity();
let _pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_ctx.metrics_registry(),
pipeline_entity_key,
);
let mut monitor = PipelineMetricsMonitor::new(pipeline_ctx);
monitor.update_metrics();
let cpu0 = monitor.metrics.cpu_time.get();
let cs_vol0 = monitor.metrics.context_switches_voluntary.get();
let cs_invol0 = monitor.metrics.context_switches_involuntary.get();
let pf_min0 = monitor.metrics.page_faults_minor.get();
let pf_maj0 = monitor.metrics.page_faults_major.get();
let mem0 = monitor.metrics.memory_allocated.get();
let mut v = Vec::with_capacity(10_000);
for i in 0..10_000u64 {
v.push(i);
}
let _ = black_box(&v);
let start = Instant::now();
while start.elapsed() < Duration::from_millis(20) {
let _ = black_box(1u64.wrapping_mul(2));
}
std::thread::sleep(Duration::from_millis(5));
monitor.update_metrics();
assert!(monitor.metrics.cpu_time.get() >= cpu0);
assert!(monitor.metrics.cpu_utilization.get() <= 1.0);
assert!(monitor.metrics.memory_allocated.get() >= mem0);
assert!(monitor.metrics.memory_allocated_delta.get() > 0);
if monitor.rusage_thread_supported {
assert!(monitor.metrics.context_switches_voluntary.get() >= cs_vol0);
assert!(monitor.metrics.context_switches_involuntary.get() >= cs_invol0);
assert!(monitor.metrics.page_faults_minor.get() >= pf_min0);
assert!(monitor.metrics.page_faults_major.get() >= pf_maj0);
} else {
assert_eq!(monitor.metrics.context_switches_voluntary.get(), cs_vol0);
assert_eq!(
monitor.metrics.context_switches_involuntary.get(),
cs_invol0
);
assert_eq!(monitor.metrics.page_faults_minor.get(), pf_min0);
assert_eq!(monitor.metrics.page_faults_major.get(), pf_maj0);
}
}
}
#[cfg(all(test, any(windows, not(feature = "jemalloc"))))]
mod non_jemalloc_tests {
use super::*;
use crate::context::ControllerContext;
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
use std::hint::black_box;
use std::time::{Duration, Instant};
#[test]
fn pipeline_metrics_monitor_does_not_update_memory_without_jemalloc() {
let telemetry_registry = TelemetryRegistryHandle::new();
let controller = ControllerContext::new(telemetry_registry);
let pipeline_ctx = controller.pipeline_context_with("grp".into(), "pipe".into(), 0, 1, 0);
let pipeline_entity_key = pipeline_ctx.register_pipeline_entity();
let _pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_ctx.metrics_registry(),
pipeline_entity_key,
);
let mut monitor = PipelineMetricsMonitor::new(pipeline_ctx);
monitor.update_metrics();
let cpu0 = monitor.metrics.cpu_time.get();
let cs_vol0 = monitor.metrics.context_switches_voluntary.get();
let cs_invol0 = monitor.metrics.context_switches_involuntary.get();
let pf_min0 = monitor.metrics.page_faults_minor.get();
let pf_maj0 = monitor.metrics.page_faults_major.get();
let mut v = Vec::with_capacity(10_000);
for i in 0..10_000u64 {
v.push(i);
}
let _ = black_box(&v);
let start = Instant::now();
while start.elapsed() < Duration::from_millis(20) {
let _ = black_box(1u64.wrapping_mul(2));
}
std::thread::sleep(Duration::from_millis(5));
monitor.update_metrics();
assert!(monitor.metrics.cpu_time.get() >= cpu0);
assert!(monitor.metrics.cpu_utilization.get() <= 1.0);
if monitor.rusage_thread_supported {
assert!(monitor.metrics.context_switches_voluntary.get() >= cs_vol0);
assert!(monitor.metrics.context_switches_involuntary.get() >= cs_invol0);
assert!(monitor.metrics.page_faults_minor.get() >= pf_min0);
assert!(monitor.metrics.page_faults_major.get() >= pf_maj0);
} else {
assert_eq!(monitor.metrics.context_switches_voluntary.get(), cs_vol0);
assert_eq!(
monitor.metrics.context_switches_involuntary.get(),
cs_invol0
);
assert_eq!(monitor.metrics.page_faults_minor.get(), pf_min0);
assert_eq!(monitor.metrics.page_faults_major.get(), pf_maj0);
}
assert_eq!(monitor.metrics.memory_allocated.get(), 0);
assert_eq!(monitor.metrics.memory_freed.get(), 0);
assert_eq!(monitor.metrics.memory_usage.get(), 0);
assert_eq!(monitor.metrics.memory_allocated_delta.get(), 0);
assert_eq!(monitor.metrics.memory_freed_delta.get(), 0);
}
}