1use crate::utils::file_utils::write_file_with_context;
4use anyhow::Result;
5use hashbrown::HashMap;
6use serde::{Deserialize, Serialize};
7use std::path::Path;
8use std::sync::{Arc, Mutex, mpsc};
9use std::time::{Duration, Instant};
10use tokio::sync::RwLock;
11
12#[derive(Debug, Clone, Serialize, Deserialize)]
14pub struct BenchmarkResults {
15 pub test_name: String,
17 pub iterations: u64,
19 pub total_duration: Duration,
21 pub avg_duration_ns: u64,
23 pub min_duration_ns: u64,
25 pub max_duration_ns: u64,
27 pub percentile_95_ns: u64,
29 pub percentile_99_ns: u64,
31 pub throughput_ops_per_sec: f64,
33 pub memory_usage_mb: Option<f64>,
35 pub cpu_usage_percent: Option<f64>,
37}
38
39#[derive(Debug, Clone, Default)]
41pub struct ResourceMetrics {
42 pub memory_used_mb: f64,
44 pub cpu_percent: f64,
46 pub network_bytes_sent: u64,
48 pub network_bytes_received: u64,
50 pub disk_reads: u64,
52 pub disk_writes: u64,
54}
55
56pub struct PerformanceProfiler {
58 sessions: Arc<RwLock<HashMap<String, BenchmarkSession>>>,
60
61 resource_monitor: Arc<ResourceMonitor>,
63
64 history: Arc<RwLock<Vec<BenchmarkResults>>>,
66}
67
68#[derive(Debug)]
70pub struct BenchmarkSession {
71 pub name: String,
73 pub start_time: Instant,
75 pub iterations: u64,
77 pub durations: Vec<Duration>,
79 pub resource_snapshots: Vec<ResourceMetrics>,
81}
82
83pub struct ResourceMonitor {
85 current_metrics: Arc<RwLock<ResourceMetrics>>,
87
88 monitor_interval: Duration,
90
91 monitor_task: Mutex<Option<MonitorTask>>,
95}
96
97struct MonitorTask {
98 stop_sender: mpsc::Sender<()>,
99 handle: std::thread::JoinHandle<()>,
100}
101
102impl MonitorTask {
103 fn stop_and_join(self) {
104 let _ = self.stop_sender.send(());
105 if self.handle.join().is_err() {
106 tracing::warn!("vtcode performance monitor thread panicked during shutdown");
107 }
108 }
109}
110
111impl PerformanceProfiler {
112 pub fn new() -> Self {
114 Self {
115 sessions: Arc::new(RwLock::new(HashMap::new())),
116 resource_monitor: Arc::new(ResourceMonitor::new(Duration::from_millis(100))),
117 history: Arc::new(RwLock::new(Vec::new())),
118 }
119 }
120
121 pub async fn start_benchmark(&self, name: &str) -> Result<()> {
123 let session = BenchmarkSession {
124 name: name.to_string(),
125 start_time: Instant::now(),
126 iterations: 0,
127 durations: Vec::new(),
128 resource_snapshots: Vec::new(),
129 };
130
131 self.sessions.write().await.insert(name.to_string(), session);
132 self.resource_monitor.start_monitoring().await?;
133
134 Ok(())
135 }
136
137 pub async fn record_operation(&self, session_name: &str, duration: Duration) -> Result<()> {
139 let mut sessions = self.sessions.write().await;
140 if let Some(session) = sessions.get_mut(session_name) {
141 session.iterations += 1;
142 session.durations.push(duration);
143
144 if session.iterations % 100 == 0 {
146 let metrics = self.resource_monitor.get_current_metrics().await;
147 session.resource_snapshots.push(metrics);
148 }
149 }
150
151 Ok(())
152 }
153
154 pub async fn end_benchmark(&self, session_name: &str) -> Result<BenchmarkResults> {
156 let session = {
157 let mut sessions = self.sessions.write().await;
158 sessions
159 .remove(session_name)
160 .ok_or_else(|| anyhow::anyhow!("Benchmark session '{session_name}' not found"))?
161 };
162
163 self.resource_monitor.stop_monitoring().await?;
164
165 let results = self.calculate_results(session).await;
166
167 self.history.write().await.push(results.clone());
169
170 Ok(results)
171 }
172
173 async fn calculate_results(&self, session: BenchmarkSession) -> BenchmarkResults {
175 let total_duration = session.start_time.elapsed();
176 let mut durations_ns: Vec<u64> = session.durations.iter().map(|d| d.as_nanos() as u64).collect();
177
178 durations_ns.sort_unstable();
179
180 let avg_duration_ns = if !durations_ns.is_empty() {
181 durations_ns.iter().sum::<u64>() / durations_ns.len() as u64
182 } else {
183 0
184 };
185
186 let min_duration_ns = durations_ns.first().copied().unwrap_or(0);
187 let max_duration_ns = durations_ns.last().copied().unwrap_or(0);
188
189 let percentile_95_ns = if !durations_ns.is_empty() {
190 #[allow(
191 clippy::cast_sign_loss,
192 reason = "Intentional compatibility, platform, or test-only suppression."
193 )]
194 let index = (durations_ns.len() as f64 * 0.95) as usize;
195 durations_ns.get(index.min(durations_ns.len() - 1)).copied().unwrap_or(0)
196 } else {
197 0
198 };
199
200 let percentile_99_ns = if !durations_ns.is_empty() {
201 #[allow(
202 clippy::cast_sign_loss,
203 reason = "Intentional compatibility, platform, or test-only suppression."
204 )]
205 let index = (durations_ns.len() as f64 * 0.99) as usize;
206 durations_ns.get(index.min(durations_ns.len() - 1)).copied().unwrap_or(0)
207 } else {
208 0
209 };
210
211 let throughput_ops_per_sec = if total_duration.as_secs_f64() > 0.0 {
212 session.iterations as f64 / total_duration.as_secs_f64()
213 } else {
214 0.0
215 };
216
217 let avg_memory_mb = if !session.resource_snapshots.is_empty() {
219 Some(
220 session.resource_snapshots.iter().map(|m| m.memory_used_mb).sum::<f64>()
221 / session.resource_snapshots.len() as f64,
222 )
223 } else {
224 None
225 };
226
227 let avg_cpu_percent = if !session.resource_snapshots.is_empty() {
228 Some(
229 session.resource_snapshots.iter().map(|m| m.cpu_percent).sum::<f64>()
230 / session.resource_snapshots.len() as f64,
231 )
232 } else {
233 None
234 };
235
236 BenchmarkResults {
237 test_name: session.name,
238 iterations: session.iterations,
239 total_duration,
240 avg_duration_ns,
241 min_duration_ns,
242 max_duration_ns,
243 percentile_95_ns,
244 percentile_99_ns,
245 throughput_ops_per_sec,
246 memory_usage_mb: avg_memory_mb,
247 cpu_usage_percent: avg_cpu_percent,
248 }
249 }
250
251 pub async fn get_history(&self) -> Vec<BenchmarkResults> {
253 self.history.read().await.clone()
254 }
255
256 pub fn compare_results(&self, baseline: &BenchmarkResults, current: &BenchmarkResults) -> ComparisonReport {
258 let throughput_change = if baseline.throughput_ops_per_sec > 0.0 {
259 ((current.throughput_ops_per_sec - baseline.throughput_ops_per_sec) / baseline.throughput_ops_per_sec)
260 * 100.0
261 } else {
262 0.0
263 };
264
265 let avg_latency_change = if baseline.avg_duration_ns > 0 {
266 ((current.avg_duration_ns as f64 - baseline.avg_duration_ns as f64) / baseline.avg_duration_ns as f64)
267 * 100.0
268 } else {
269 0.0
270 };
271
272 let memory_change = match (baseline.memory_usage_mb, current.memory_usage_mb) {
273 (Some(baseline_mem), Some(current_mem)) => Some(((current_mem - baseline_mem) / baseline_mem) * 100.0),
274 _ => None,
275 };
276
277 ComparisonReport {
278 baseline_name: baseline.test_name.clone(),
279 current_name: current.test_name.clone(),
280 throughput_change_percent: throughput_change,
281 avg_latency_change_percent: avg_latency_change,
282 memory_change_percent: memory_change,
283 is_improvement: throughput_change > 0.0 && avg_latency_change < 0.0,
284 }
285 }
286
287 pub async fn export_results(&self, file_path: &str) -> Result<()> {
289 let history = self.get_history().await;
290 let json = serde_json::to_string_pretty(&history)?;
291 write_file_with_context(Path::new(file_path), &json, "benchmark results").await?;
292 Ok(())
293 }
294}
295
296#[derive(Debug, Clone, Serialize, Deserialize)]
298pub struct ComparisonReport {
299 pub baseline_name: String,
301 pub current_name: String,
303 pub throughput_change_percent: f64,
305 pub avg_latency_change_percent: f64,
307 pub memory_change_percent: Option<f64>,
309 pub is_improvement: bool,
311}
312
313impl Drop for ResourceMonitor {
314 fn drop(&mut self) {
315 let task = match self.monitor_task.lock() {
316 Ok(mut task_slot) => task_slot.take(),
317 Err(poisoned) => poisoned.into_inner().take(),
318 };
319 if let Some(task) = task {
320 task.stop_and_join();
321 }
322 }
323}
324
325impl ResourceMonitor {
326 pub fn new(monitor_interval: Duration) -> Self {
328 Self {
329 current_metrics: Arc::new(RwLock::new(ResourceMetrics::default())),
330 monitor_interval,
331 monitor_task: Mutex::new(None),
332 }
333 }
334
335 pub async fn start_monitoring(&self) -> Result<()> {
337 let mut task_slot = self.monitor_task.lock().unwrap_or_else(|e| e.into_inner());
338 if task_slot.is_some() {
339 return Ok(()); }
341
342 let current_metrics = Arc::clone(&self.current_metrics);
343 let interval = self.monitor_interval;
344 let (stop_sender, stop_receiver) = mpsc::channel();
345
346 let monitor_thread = std::thread::Builder::new()
347 .name("vtcode-perf-monitor".to_string())
348 .spawn(move || {
349 loop {
350 match stop_receiver.recv_timeout(interval) {
351 Ok(()) | Err(mpsc::RecvTimeoutError::Disconnected) => break,
352 Err(mpsc::RecvTimeoutError::Timeout) => {
353 let sample = Self::collect_system_metrics_sync();
354 *current_metrics.blocking_write() = sample;
355 }
356 }
357 }
358 })
359 .map_err(|error| anyhow::Error::new(error).context("failed to spawn resource monitor thread"))?;
360
361 *task_slot = Some(MonitorTask { stop_sender, handle: monitor_thread });
362 Ok(())
363 }
364
365 pub async fn stop_monitoring(&self) -> Result<()> {
367 let mut task_slot = self.monitor_task.lock().unwrap_or_else(|e| e.into_inner());
368 if let Some(task) = task_slot.take() {
369 task.stop_and_join();
372 }
373 Ok(())
374 }
375
376 pub async fn get_current_metrics(&self) -> ResourceMetrics {
378 self.current_metrics.read().await.clone()
379 }
380
381 fn collect_system_metrics_sync() -> ResourceMetrics {
383 let memory_used_mb = Self::get_memory_usage_mb();
384 ResourceMetrics {
385 memory_used_mb,
386 cpu_percent: Self::get_cpu_usage_percent(),
387 network_bytes_sent: 0,
388 network_bytes_received: 0,
389 disk_reads: 0,
390 disk_writes: 0,
391 }
392 }
393
394 fn get_memory_usage_mb() -> f64 {
396 vtcode_commons::memory::sample_rss_mb()
397 }
398
399 fn get_cpu_usage_percent() -> f64 {
401 0.0
404 }
405}
406
407#[macro_export]
409macro_rules! benchmark {
410 ($profiler:expr, $name:expr, $code:block) => {{
411 let start = std::time::Instant::now();
412 let result = $code;
413 let duration = start.elapsed();
414 $profiler.record_operation($name, duration).await?;
415 result
416 }};
417}
418
419pub struct BenchmarkUtils;
421
422impl BenchmarkUtils {
423 pub async fn benchmark_function<F, R>(
425 profiler: &PerformanceProfiler,
426 name: &str,
427 iterations: u64,
428 mut func: F,
429 ) -> Result<BenchmarkResults>
430 where
431 F: FnMut() -> R,
432 {
433 profiler.start_benchmark(name).await?;
434
435 for _ in 0..iterations {
436 let start = Instant::now();
437 let _ = func();
438 let duration = start.elapsed();
439 profiler.record_operation(name, duration).await?;
440 }
441
442 profiler.end_benchmark(name).await
443 }
444
445 pub async fn benchmark_async_function<F, Fut, R>(
447 profiler: &PerformanceProfiler,
448 name: &str,
449 iterations: u64,
450 mut func: F,
451 ) -> Result<BenchmarkResults>
452 where
453 F: FnMut() -> Fut,
454 Fut: Future<Output = R>,
455 {
456 profiler.start_benchmark(name).await?;
457
458 for _ in 0..iterations {
459 let start = Instant::now();
460 let _ = func().await;
461 let duration = start.elapsed();
462 profiler.record_operation(name, duration).await?;
463 }
464
465 profiler.end_benchmark(name).await
466 }
467
468 pub async fn regression_test(
470 profiler: &PerformanceProfiler,
471 baseline_name: &str,
472 current_name: &str,
473 max_regression_percent: f64,
474 ) -> Result<bool> {
475 let history = profiler.get_history().await;
476
477 let baseline = history
478 .iter()
479 .find(|r| r.test_name == baseline_name)
480 .ok_or_else(|| anyhow::anyhow!("Baseline '{baseline_name}' not found"))?;
481
482 let current = history
483 .iter()
484 .find(|r| r.test_name == current_name)
485 .ok_or_else(|| anyhow::anyhow!("Current '{current_name}' not found"))?;
486
487 let comparison = profiler.compare_results(baseline, current);
488
489 let regression = comparison.avg_latency_change_percent > max_regression_percent
491 || comparison.throughput_change_percent < -max_regression_percent;
492
493 if regression {
494 tracing::warn!(
495 latency_change_percent = comparison.avg_latency_change_percent,
496 throughput_change_percent = comparison.throughput_change_percent,
497 "Performance regression detected"
498 );
499 }
500
501 Ok(!regression)
502 }
503}
504
505impl Default for PerformanceProfiler {
506 fn default() -> Self {
507 Self::new()
508 }
509}
510
511#[cfg(test)]
512mod tests {
513 use super::*;
514 use tokio::time::sleep;
515
516 #[tokio::test]
517 async fn test_benchmark_session() -> Result<()> {
518 let profiler = PerformanceProfiler::new();
519
520 profiler.start_benchmark("test_session").await?;
521
522 for i in 0..10 {
524 let duration = Duration::from_millis(10 + i);
525 profiler.record_operation("test_session", duration).await?;
526 }
527
528 let results = profiler.end_benchmark("test_session").await?;
529
530 assert_eq!(results.test_name, "test_session");
531 assert_eq!(results.iterations, 10);
532 assert!(results.avg_duration_ns > 0);
533
534 Ok(())
535 }
536
537 #[tokio::test]
538 async fn monitor_stop_interrupts_worker_and_clears_handle() {
539 let monitor = ResourceMonitor::new(Duration::from_secs(5));
540 monitor.start_monitoring().await.expect("start monitoring");
541
542 let started = Instant::now();
543 monitor.stop_monitoring().await.expect("stop monitoring");
544 assert!(started.elapsed() < Duration::from_secs(1));
545 assert!(monitor.monitor_task.lock().expect("monitor state lock").is_none());
546
547 monitor.start_monitoring().await.expect("restart monitoring");
548 let restarted = Instant::now();
549 monitor.stop_monitoring().await.expect("stop restarted monitoring");
550 assert!(restarted.elapsed() < Duration::from_secs(1));
551 }
552
553 #[tokio::test]
554 async fn test_benchmark_utils() -> Result<()> {
555 let profiler = PerformanceProfiler::new();
556
557 let results = BenchmarkUtils::benchmark_function(&profiler, "test_function", 100, || {
558 std::thread::sleep(Duration::from_micros(100));
560 42
561 })
562 .await?;
563
564 assert_eq!(results.iterations, 100);
565 assert!(results.throughput_ops_per_sec > 0.0);
566
567 Ok(())
568 }
569
570 #[tokio::test]
571 async fn test_async_benchmark() -> Result<()> {
572 let profiler = PerformanceProfiler::new();
573
574 let results = BenchmarkUtils::benchmark_async_function(&profiler, "test_async_function", 50, || async {
575 sleep(Duration::from_micros(200)).await;
576 "result"
577 })
578 .await?;
579
580 assert_eq!(results.iterations, 50);
581 assert!(results.avg_duration_ns > 0);
582
583 Ok(())
584 }
585}