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,
26 pub tasks_completed: AtomicU64,
27 pub tasks_failed: AtomicU64,
28 pub tasks_cancelled: AtomicU64,
30
31 pub total_execution_time: AtomicU64, pub active_workers: AtomicUsize,
36 pub idle_workers: AtomicUsize,
37 pub total_workers: AtomicUsize,
38
39 pub pending_tasks: AtomicUsize,
41 pub max_queue_depth: AtomicUsize,
42
43 pub started_at: Instant,
45 pub last_updated_after_ns: AtomicU64,
46}
47
48impl ExecutorMetrics {
49 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 pub fn record_task_spawned(&self) {
70 self.tasks_spawned.fetch_add(1, Ordering::Relaxed);
71 self.update_timestamp();
72 }
73
74 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 pub fn record_task_failed(&self) {
86 self.tasks_failed.fetch_add(1, Ordering::Relaxed);
87 self.update_timestamp();
88 }
89
90 pub fn record_task_cancelled(&self) {
92 self.tasks_cancelled.fetch_add(1, Ordering::Relaxed);
93 self.update_timestamp();
94 }
95
96 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 pub fn update_queue_metrics(&self, pending: usize) {
106 self.pending_tasks.store(pending, Ordering::Relaxed);
107
108 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 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 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 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 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 pub fn uptime(&self) -> Duration {
165 self.started_at.elapsed()
166 }
167
168 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 fn update_timestamp(&self) {
179 self.last_updated_after_ns
180 .store(elapsed_nanos_since(self.started_at), Ordering::Relaxed);
181 }
182
183 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}