moirai_executor/metrics/
mod.rs1use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
7use std::time::{Duration, Instant};
8
9const PERCENTAGE_PRECISION_FACTOR: f64 = 100.0;
11
12const MAX_SUCCESS_RATE: f64 = 100.0;
14
15const DEFAULT_UTILIZATION: f64 = 0.0;
17
18#[derive(Debug)]
23pub struct ExecutorMetrics {
24 pub tasks_spawned: AtomicU64,
27 pub tasks_completed: AtomicU64,
29 pub tasks_failed: AtomicU64,
31 pub tasks_cancelled: AtomicU64,
33
34 pub total_execution_time: AtomicU64,
37
38 pub active_workers: AtomicUsize,
41 pub idle_workers: AtomicUsize,
43 pub total_workers: AtomicUsize,
45
46 pub pending_tasks: AtomicUsize,
49 pub max_queue_depth: AtomicUsize,
51
52 pub started_at: Instant,
55 pub last_updated_after_ns: AtomicU64,
57}
58
59impl ExecutorMetrics {
60 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 pub fn record_task_spawned(&self) {
81 self.tasks_spawned.fetch_add(1, Ordering::Relaxed);
82 self.update_timestamp();
83 }
84
85 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 pub fn record_task_failed(&self) {
97 self.tasks_failed.fetch_add(1, Ordering::Relaxed);
98 self.update_timestamp();
99 }
100
101 pub fn record_task_cancelled(&self) {
103 self.tasks_cancelled.fetch_add(1, Ordering::Relaxed);
104 self.update_timestamp();
105 }
106
107 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 pub fn update_queue_metrics(&self, pending: usize) {
117 self.pending_tasks.store(pending, Ordering::Relaxed);
118
119 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 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 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 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 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 pub fn uptime(&self) -> Duration {
176 self.started_at.elapsed()
177 }
178
179 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 fn update_timestamp(&self) {
190 self.last_updated_after_ns
191 .store(elapsed_nanos_since(self.started_at), Ordering::Relaxed);
192 }
193
194 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}