Skip to main content

moirai_executor/metrics/
mod.rs

1//! Performance metrics and monitoring for the executor system.
2//!
3//! This module provides comprehensive metrics collection and reporting
4//! for monitoring executor performance and identifying bottlenecks.
5
6use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
7use std::time::{Duration, Instant};
8
9/// Percentage conversion factor to maintain precision across metrics
10const PERCENTAGE_PRECISION_FACTOR: f64 = 100.0;
11
12/// Maximum success rate when no tasks have failed
13const MAX_SUCCESS_RATE: f64 = 100.0;
14
15/// Default utilization when no workers are available
16const DEFAULT_UTILIZATION: f64 = 0.0;
17
18/// Executor performance metrics.
19///
20/// Every field is written by a real production path; untracked quantities
21/// (process memory, CPU utilization) deliberately have no fields here.
22#[derive(Debug)]
23pub struct ExecutorMetrics {
24    // Task counters
25    /// Tasks accepted by the executor.
26    pub tasks_spawned: AtomicU64,
27    /// Tasks whose body ran to completion.
28    pub tasks_completed: AtomicU64,
29    /// Tasks whose body panicked or reported failure.
30    pub tasks_failed: AtomicU64,
31    /// Tasks whose cancel request was honored before the body ran.
32    pub tasks_cancelled: AtomicU64,
33
34    // Timing metrics
35    /// Accumulated task body execution time in nanoseconds.
36    pub total_execution_time: AtomicU64,
37
38    // Thread pool metrics
39    /// Workers currently running a task body.
40    pub active_workers: AtomicUsize,
41    /// Workers parked waiting for work.
42    pub idle_workers: AtomicUsize,
43    /// Workers owned by the pool.
44    pub total_workers: AtomicUsize,
45
46    // Queue metrics
47    /// Tasks queued and not yet started.
48    pub pending_tasks: AtomicUsize,
49    /// High-water mark of the pending queue.
50    pub max_queue_depth: AtomicUsize,
51
52    // System metrics
53    /// Executor construction instant; the metrics epoch.
54    pub started_at: Instant,
55    /// Nanoseconds after `started_at` of the last metrics update.
56    pub last_updated_after_ns: AtomicU64,
57}
58
59impl ExecutorMetrics {
60    /// Create new metrics tracker
61    pub fn new() -> Self {
62        let now = Instant::now();
63        Self {
64            tasks_spawned: AtomicU64::new(0),
65            tasks_completed: AtomicU64::new(0),
66            tasks_failed: AtomicU64::new(0),
67            tasks_cancelled: AtomicU64::new(0),
68            total_execution_time: AtomicU64::new(0),
69            active_workers: AtomicUsize::new(0),
70            idle_workers: AtomicUsize::new(0),
71            total_workers: AtomicUsize::new(0),
72            pending_tasks: AtomicUsize::new(0),
73            max_queue_depth: AtomicUsize::new(0),
74            started_at: now,
75            last_updated_after_ns: AtomicU64::new(0),
76        }
77    }
78
79    /// Record a new task being spawned
80    pub fn record_task_spawned(&self) {
81        self.tasks_spawned.fetch_add(1, Ordering::Relaxed);
82        self.update_timestamp();
83    }
84
85    /// Record a task completion with execution time
86    pub fn record_task_completed(&self, execution_time: Duration) {
87        let exec_nanos = execution_time.as_nanos() as u64;
88        self.tasks_completed.fetch_add(1, Ordering::Relaxed);
89        self.total_execution_time
90            .fetch_add(exec_nanos, Ordering::Relaxed);
91
92        self.update_timestamp();
93    }
94
95    /// Record a task failure
96    pub fn record_task_failed(&self) {
97        self.tasks_failed.fetch_add(1, Ordering::Relaxed);
98        self.update_timestamp();
99    }
100
101    /// Record a task whose cancel request was honored before its body ran.
102    pub fn record_task_cancelled(&self) {
103        self.tasks_cancelled.fetch_add(1, Ordering::Relaxed);
104        self.update_timestamp();
105    }
106
107    /// Update worker count metrics
108    pub fn update_worker_counts(&self, active: usize, idle: usize, total: usize) {
109        self.active_workers.store(active, Ordering::Relaxed);
110        self.idle_workers.store(idle, Ordering::Relaxed);
111        self.total_workers.store(total, Ordering::Relaxed);
112        self.update_timestamp();
113    }
114
115    /// Update queue metrics
116    pub fn update_queue_metrics(&self, pending: usize) {
117        self.pending_tasks.store(pending, Ordering::Relaxed);
118
119        // Update max queue depth if this is a new maximum
120        let current_max = self.max_queue_depth.load(Ordering::Relaxed);
121        if pending > current_max {
122            self.max_queue_depth.store(pending, Ordering::Relaxed);
123        }
124
125        self.update_timestamp();
126    }
127
128    /// Get current throughput (tasks per second)
129    pub fn throughput(&self) -> f64 {
130        let elapsed = self.started_at.elapsed().as_secs_f64();
131        if elapsed > 0.0 {
132            self.tasks_completed.load(Ordering::Relaxed) as f64 / elapsed
133        } else {
134            0.0
135        }
136    }
137
138    /// Get success rate as percentage
139    pub fn success_rate(&self) -> f64 {
140        let completed = self.tasks_completed.load(Ordering::Relaxed);
141        let failed = self.tasks_failed.load(Ordering::Relaxed);
142        let total = completed + failed;
143
144        if total > 0 {
145            (completed as f64 / total as f64) * PERCENTAGE_PRECISION_FACTOR
146        } else {
147            MAX_SUCCESS_RATE
148        }
149    }
150
151    /// Get average task duration
152    pub fn average_task_duration(&self) -> Duration {
153        let completed = self.tasks_completed.load(Ordering::Relaxed);
154        let average = self
155            .total_execution_time
156            .load(Ordering::Relaxed)
157            .checked_div(completed)
158            .unwrap_or(0);
159
160        Duration::from_nanos(average)
161    }
162
163    /// Get worker utilization percentage
164    pub fn worker_utilization(&self) -> f64 {
165        let total = self.total_workers.load(Ordering::Relaxed);
166        if total > 0 {
167            let active = self.active_workers.load(Ordering::Relaxed);
168            (active as f64 / total as f64) * PERCENTAGE_PRECISION_FACTOR
169        } else {
170            DEFAULT_UTILIZATION
171        }
172    }
173
174    /// Get uptime
175    pub fn uptime(&self) -> Duration {
176        self.started_at.elapsed()
177    }
178
179    /// Get the last metrics update timestamp.
180    pub fn last_updated(&self) -> Instant {
181        self.started_at
182            .checked_add(Duration::from_nanos(
183                self.last_updated_after_ns.load(Ordering::Relaxed),
184            ))
185            .unwrap_or(self.started_at)
186    }
187
188    /// Update the last updated timestamp
189    fn update_timestamp(&self) {
190        self.last_updated_after_ns
191            .store(elapsed_nanos_since(self.started_at), Ordering::Relaxed);
192    }
193
194    /// Reset all metrics (useful for testing)
195    pub fn reset(&self) {
196        self.tasks_spawned.store(0, Ordering::Relaxed);
197        self.tasks_completed.store(0, Ordering::Relaxed);
198        self.tasks_failed.store(0, Ordering::Relaxed);
199        self.tasks_cancelled.store(0, Ordering::Relaxed);
200        self.total_execution_time.store(0, Ordering::Relaxed);
201        self.active_workers.store(0, Ordering::Relaxed);
202        self.idle_workers.store(0, Ordering::Relaxed);
203        self.total_workers.store(0, Ordering::Relaxed);
204        self.pending_tasks.store(0, Ordering::Relaxed);
205        self.max_queue_depth.store(0, Ordering::Relaxed);
206        self.update_timestamp();
207    }
208}
209
210fn elapsed_nanos_since(origin: Instant) -> u64 {
211    origin.elapsed().as_nanos().min(u128::from(u64::MAX)) as u64
212}
213
214impl Default for ExecutorMetrics {
215    fn default() -> Self {
216        Self::new()
217    }
218}