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    pub tasks_spawned: AtomicU64,
26    pub tasks_completed: AtomicU64,
27    pub tasks_failed: AtomicU64,
28    /// Tasks whose cancel request was honored before the body ran.
29    pub tasks_cancelled: AtomicU64,
30
31    // Timing metrics
32    pub total_execution_time: AtomicU64, // in nanoseconds
33
34    // Thread pool metrics
35    pub active_workers: AtomicUsize,
36    pub idle_workers: AtomicUsize,
37    pub total_workers: AtomicUsize,
38
39    // Queue metrics
40    pub pending_tasks: AtomicUsize,
41    pub max_queue_depth: AtomicUsize,
42
43    // System metrics
44    pub started_at: Instant,
45    pub last_updated_after_ns: AtomicU64,
46}
47
48impl ExecutorMetrics {
49    /// Create new metrics tracker
50    pub fn new() -> Self {
51        let now = Instant::now();
52        Self {
53            tasks_spawned: AtomicU64::new(0),
54            tasks_completed: AtomicU64::new(0),
55            tasks_failed: AtomicU64::new(0),
56            tasks_cancelled: AtomicU64::new(0),
57            total_execution_time: AtomicU64::new(0),
58            active_workers: AtomicUsize::new(0),
59            idle_workers: AtomicUsize::new(0),
60            total_workers: AtomicUsize::new(0),
61            pending_tasks: AtomicUsize::new(0),
62            max_queue_depth: AtomicUsize::new(0),
63            started_at: now,
64            last_updated_after_ns: AtomicU64::new(0),
65        }
66    }
67
68    /// Record a new task being spawned
69    pub fn record_task_spawned(&self) {
70        self.tasks_spawned.fetch_add(1, Ordering::Relaxed);
71        self.update_timestamp();
72    }
73
74    /// Record a task completion with execution time
75    pub fn record_task_completed(&self, execution_time: Duration) {
76        let exec_nanos = execution_time.as_nanos() as u64;
77        self.tasks_completed.fetch_add(1, Ordering::Relaxed);
78        self.total_execution_time
79            .fetch_add(exec_nanos, Ordering::Relaxed);
80
81        self.update_timestamp();
82    }
83
84    /// Record a task failure
85    pub fn record_task_failed(&self) {
86        self.tasks_failed.fetch_add(1, Ordering::Relaxed);
87        self.update_timestamp();
88    }
89
90    /// Record a task whose cancel request was honored before its body ran.
91    pub fn record_task_cancelled(&self) {
92        self.tasks_cancelled.fetch_add(1, Ordering::Relaxed);
93        self.update_timestamp();
94    }
95
96    /// Update worker count metrics
97    pub fn update_worker_counts(&self, active: usize, idle: usize, total: usize) {
98        self.active_workers.store(active, Ordering::Relaxed);
99        self.idle_workers.store(idle, Ordering::Relaxed);
100        self.total_workers.store(total, Ordering::Relaxed);
101        self.update_timestamp();
102    }
103
104    /// Update queue metrics
105    pub fn update_queue_metrics(&self, pending: usize) {
106        self.pending_tasks.store(pending, Ordering::Relaxed);
107
108        // Update max queue depth if this is a new maximum
109        let current_max = self.max_queue_depth.load(Ordering::Relaxed);
110        if pending > current_max {
111            self.max_queue_depth.store(pending, Ordering::Relaxed);
112        }
113
114        self.update_timestamp();
115    }
116
117    /// Get current throughput (tasks per second)
118    pub fn throughput(&self) -> f64 {
119        let elapsed = self.started_at.elapsed().as_secs_f64();
120        if elapsed > 0.0 {
121            self.tasks_completed.load(Ordering::Relaxed) as f64 / elapsed
122        } else {
123            0.0
124        }
125    }
126
127    /// Get success rate as percentage
128    pub fn success_rate(&self) -> f64 {
129        let completed = self.tasks_completed.load(Ordering::Relaxed);
130        let failed = self.tasks_failed.load(Ordering::Relaxed);
131        let total = completed + failed;
132
133        if total > 0 {
134            (completed as f64 / total as f64) * PERCENTAGE_PRECISION_FACTOR
135        } else {
136            MAX_SUCCESS_RATE
137        }
138    }
139
140    /// Get average task duration
141    pub fn average_task_duration(&self) -> Duration {
142        let completed = self.tasks_completed.load(Ordering::Relaxed);
143        let average = self
144            .total_execution_time
145            .load(Ordering::Relaxed)
146            .checked_div(completed)
147            .unwrap_or(0);
148
149        Duration::from_nanos(average)
150    }
151
152    /// Get worker utilization percentage
153    pub fn worker_utilization(&self) -> f64 {
154        let total = self.total_workers.load(Ordering::Relaxed);
155        if total > 0 {
156            let active = self.active_workers.load(Ordering::Relaxed);
157            (active as f64 / total as f64) * PERCENTAGE_PRECISION_FACTOR
158        } else {
159            DEFAULT_UTILIZATION
160        }
161    }
162
163    /// Get uptime
164    pub fn uptime(&self) -> Duration {
165        self.started_at.elapsed()
166    }
167
168    /// Get the last metrics update timestamp.
169    pub fn last_updated(&self) -> Instant {
170        self.started_at
171            .checked_add(Duration::from_nanos(
172                self.last_updated_after_ns.load(Ordering::Relaxed),
173            ))
174            .unwrap_or(self.started_at)
175    }
176
177    /// Update the last updated timestamp
178    fn update_timestamp(&self) {
179        self.last_updated_after_ns
180            .store(elapsed_nanos_since(self.started_at), Ordering::Relaxed);
181    }
182
183    /// Reset all metrics (useful for testing)
184    pub fn reset(&self) {
185        self.tasks_spawned.store(0, Ordering::Relaxed);
186        self.tasks_completed.store(0, Ordering::Relaxed);
187        self.tasks_failed.store(0, Ordering::Relaxed);
188        self.tasks_cancelled.store(0, Ordering::Relaxed);
189        self.total_execution_time.store(0, Ordering::Relaxed);
190        self.active_workers.store(0, Ordering::Relaxed);
191        self.idle_workers.store(0, Ordering::Relaxed);
192        self.total_workers.store(0, Ordering::Relaxed);
193        self.pending_tasks.store(0, Ordering::Relaxed);
194        self.max_queue_depth.store(0, Ordering::Relaxed);
195        self.update_timestamp();
196    }
197}
198
199fn elapsed_nanos_since(origin: Instant) -> u64 {
200    origin.elapsed().as_nanos().min(u128::from(u64::MAX)) as u64
201}
202
203impl Default for ExecutorMetrics {
204    fn default() -> Self {
205        Self::new()
206    }
207}