1use serde::{Deserialize, Serialize};
7use std::collections::HashMap;
8use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
9use std::sync::Arc;
10use std::time::{Duration, SystemTime, UNIX_EPOCH};
11use tokio::sync::RwLock;
12
13#[derive(Debug, Clone)]
15pub struct MetricsCollector {
16 request_counts: Arc<RwLock<HashMap<String, AtomicU64>>>,
18 response_times: Arc<RwLock<HashMap<String, Vec<Duration>>>>,
20 error_counts: Arc<RwLock<HashMap<String, AtomicU64>>>,
22 resource_usage: Arc<ResourceUsageTracker>,
24 connection_metrics: Arc<ConnectionMetrics>,
26 health_status: Arc<RwLock<HealthStatus>>,
28}
29
30#[derive(Debug)]
32pub struct ResourceUsageTracker {
33 memory_usage_mb: AtomicUsize,
35 peak_memory_mb: AtomicUsize,
37 cpu_samples: RwLock<Vec<f64>>,
39 active_requests: AtomicUsize,
41 peak_concurrent_requests: AtomicUsize,
43}
44
45#[derive(Debug)]
47pub struct ConnectionMetrics {
48 total_connections: AtomicU64,
50 active_connections: AtomicUsize,
52 peak_connections: AtomicUsize,
54 connection_errors: AtomicU64,
56 connection_durations: RwLock<Vec<Duration>>,
58}
59
60#[derive(Debug, Clone, Serialize, Deserialize)]
62pub struct HealthStatus {
63 pub health_score: u8,
65 pub components: HashMap<String, ComponentHealth>,
67 pub last_check: SystemTime,
69 pub alerts: Vec<HealthAlert>,
71}
72
73#[derive(Debug, Clone, Serialize, Deserialize)]
75pub struct ComponentHealth {
76 pub status: ComponentStatus,
78 pub message: Option<String>,
80 pub last_update: SystemTime,
82 pub metrics: HashMap<String, f64>,
84}
85
86#[derive(Debug, Clone, Serialize, Deserialize)]
88pub enum ComponentStatus {
89 Healthy,
91 Warning,
93 Degraded,
95 Critical,
97 Unavailable,
99}
100
101#[derive(Debug, Clone, Serialize, Deserialize)]
103pub struct HealthAlert {
104 pub severity: AlertSeverity,
106 pub message: String,
108 pub component: String,
110 pub timestamp: SystemTime,
112 pub id: String,
114}
115
116#[derive(Debug, Clone, Serialize, Deserialize)]
118pub enum AlertSeverity {
119 Info,
121 Warning,
123 Error,
125 Critical,
127}
128
129#[derive(Debug, Clone, Serialize, Deserialize)]
131pub struct MetricsSnapshot {
132 pub timestamp: SystemTime,
134 pub requests: RequestMetrics,
136 pub resources: ResourceMetrics,
138 pub connections: ConnectionMetricsSnapshot,
140 pub health: HealthStatus,
142}
143
144#[derive(Debug, Clone, Serialize, Deserialize)]
146pub struct RequestMetrics {
147 pub total_requests: u64,
149 pub requests_per_second: f64,
151 pub avg_response_time_ms: f64,
153 pub p95_response_time_ms: f64,
155 pub error_rate: f64,
157 pub top_request_types: Vec<(String, u64)>,
159}
160
161#[derive(Debug, Clone, Serialize, Deserialize)]
163pub struct ResourceMetrics {
164 pub current_memory_mb: usize,
166 pub peak_memory_mb: usize,
168 pub avg_cpu_usage: f64,
170 pub active_requests: usize,
172 pub peak_concurrent_requests: usize,
174}
175
176#[derive(Debug, Clone, Serialize, Deserialize)]
178pub struct ConnectionMetricsSnapshot {
179 pub total_connections: u64,
181 pub active_connections: usize,
183 pub peak_connections: usize,
185 pub connection_error_rate: f64,
187 pub avg_connection_duration_ms: f64,
189}
190
191impl MetricsCollector {
192 pub fn new() -> Self {
194 Self {
195 request_counts: Arc::new(RwLock::new(HashMap::new())),
196 response_times: Arc::new(RwLock::new(HashMap::new())),
197 error_counts: Arc::new(RwLock::new(HashMap::new())),
198 resource_usage: Arc::new(ResourceUsageTracker::new()),
199 connection_metrics: Arc::new(ConnectionMetrics::new()),
200 health_status: Arc::new(RwLock::new(HealthStatus::new())),
201 }
202 }
203
204 pub async fn record_request(&self, method: &str, response_time: Duration) {
206 {
208 let mut counts = self.request_counts.write().await;
209 let counter = counts
210 .entry(method.to_string())
211 .or_insert_with(|| AtomicU64::new(0));
212 counter.fetch_add(1, Ordering::Relaxed);
213 }
214
215 {
217 let mut times = self.response_times.write().await;
218 let method_times = times.entry(method.to_string()).or_insert_with(Vec::new);
219 method_times.push(response_time);
220
221 if method_times.len() > 1000 {
223 method_times.drain(..method_times.len() - 1000);
224 }
225 }
226 }
227
228 pub async fn record_error(&self, error_type: &str) {
230 let mut counts = self.error_counts.write().await;
231 let counter = counts
232 .entry(error_type.to_string())
233 .or_insert_with(|| AtomicU64::new(0));
234 counter.fetch_add(1, Ordering::Relaxed);
235 }
236
237 pub fn update_memory_usage(&self, current_mb: usize) {
239 self.resource_usage
240 .memory_usage_mb
241 .store(current_mb, Ordering::Relaxed);
242
243 let current_peak = self.resource_usage.peak_memory_mb.load(Ordering::Relaxed);
245 if current_mb > current_peak {
246 self.resource_usage
247 .peak_memory_mb
248 .store(current_mb, Ordering::Relaxed);
249 }
250 }
251
252 pub async fn record_cpu_usage(&self, cpu_percent: f64) {
254 let mut samples = self.resource_usage.cpu_samples.write().await;
255 samples.push(cpu_percent);
256
257 let len = samples.len();
259 if len > 100 {
260 samples.drain(..len - 100);
261 }
262 }
263
264 pub fn start_request(&self) {
266 let current = self
267 .resource_usage
268 .active_requests
269 .fetch_add(1, Ordering::Relaxed)
270 + 1;
271
272 let current_peak = self
274 .resource_usage
275 .peak_concurrent_requests
276 .load(Ordering::Relaxed);
277 if current > current_peak {
278 self.resource_usage
279 .peak_concurrent_requests
280 .store(current, Ordering::Relaxed);
281 }
282 }
283
284 pub fn complete_request(&self) {
286 self.resource_usage
287 .active_requests
288 .fetch_sub(1, Ordering::Relaxed);
289 }
290
291 pub fn record_connection(&self) {
293 self.connection_metrics
294 .total_connections
295 .fetch_add(1, Ordering::Relaxed);
296 let current = self
297 .connection_metrics
298 .active_connections
299 .fetch_add(1, Ordering::Relaxed)
300 + 1;
301
302 let current_peak = self
304 .connection_metrics
305 .peak_connections
306 .load(Ordering::Relaxed);
307 if current > current_peak {
308 self.connection_metrics
309 .peak_connections
310 .store(current, Ordering::Relaxed);
311 }
312 }
313
314 pub fn record_connection_close(&self, duration: Duration) {
316 self.connection_metrics
317 .active_connections
318 .fetch_sub(1, Ordering::Relaxed);
319
320 let duration_holder = Arc::clone(&self.connection_metrics);
322 tokio::spawn(async move {
323 let mut durations_guard = duration_holder.connection_durations.write().await;
324 durations_guard.push(duration);
325
326 let len = durations_guard.len();
328 if len > 1000 {
329 durations_guard.drain(..len - 1000);
330 }
331 });
332 }
333
334 pub fn record_connection_error(&self) {
336 self.connection_metrics
337 .connection_errors
338 .fetch_add(1, Ordering::Relaxed);
339 }
340
341 pub async fn get_snapshot(&self) -> MetricsSnapshot {
343 let request_metrics = self.calculate_request_metrics().await;
344 let resource_metrics = self.calculate_resource_metrics().await;
345 let connection_metrics = self.calculate_connection_metrics().await;
346 let health_status = self.health_status.read().await.clone();
347
348 MetricsSnapshot {
349 timestamp: SystemTime::now(),
350 requests: request_metrics,
351 resources: resource_metrics,
352 connections: connection_metrics,
353 health: health_status,
354 }
355 }
356
357 pub async fn perform_health_check(&self) -> HealthStatus {
359 let mut health = HealthStatus::new();
360
361 health
363 .components
364 .insert("memory".to_string(), self.check_memory_health().await);
365 health.components.insert(
366 "connections".to_string(),
367 self.check_connection_health().await,
368 );
369 health
370 .components
371 .insert("requests".to_string(), self.check_request_health().await);
372 health
373 .components
374 .insert("errors".to_string(), self.check_error_health().await);
375
376 health.health_score = self.calculate_health_score(&health.components);
378
379 health.alerts = self.generate_health_alerts(&health.components).await;
381
382 *self.health_status.write().await = health.clone();
384
385 health
386 }
387
388 async fn calculate_request_metrics(&self) -> RequestMetrics {
389 let counts = self.request_counts.read().await;
390 let total_requests: u64 = counts.values().map(|c| c.load(Ordering::Relaxed)).sum();
391
392 let times = self.response_times.read().await;
393 let all_times: Vec<Duration> = times.values().flatten().cloned().collect();
394
395 let avg_response_time = if all_times.is_empty() {
396 0.0
397 } else {
398 all_times.iter().map(|d| d.as_millis() as f64).sum::<f64>() / all_times.len() as f64
399 };
400
401 let p95_response_time = if all_times.is_empty() {
402 0.0
403 } else {
404 let mut sorted_times = all_times.clone();
405 sorted_times.sort();
406 let p95_index = (sorted_times.len() as f64 * 0.95) as usize;
407 sorted_times
408 .get(p95_index.min(sorted_times.len() - 1))
409 .unwrap_or(&Duration::from_millis(0))
410 .as_millis() as f64
411 };
412
413 let error_counts = self.error_counts.read().await;
414 let total_errors: u64 = error_counts
415 .values()
416 .map(|c| c.load(Ordering::Relaxed))
417 .sum();
418 let error_rate = if total_requests > 0 {
419 (total_errors as f64 / total_requests as f64) * 100.0
420 } else {
421 0.0
422 };
423
424 let top_request_types: Vec<(String, u64)> = counts
425 .iter()
426 .map(|(method, count)| (method.clone(), count.load(Ordering::Relaxed)))
427 .collect();
428
429 RequestMetrics {
430 total_requests,
431 requests_per_second: 0.0, avg_response_time_ms: avg_response_time,
433 p95_response_time_ms: p95_response_time,
434 error_rate,
435 top_request_types,
436 }
437 }
438
439 async fn calculate_resource_metrics(&self) -> ResourceMetrics {
440 let cpu_samples = self.resource_usage.cpu_samples.read().await;
441 let avg_cpu_usage = if cpu_samples.is_empty() {
442 0.0
443 } else {
444 cpu_samples.iter().sum::<f64>() / cpu_samples.len() as f64
445 };
446
447 ResourceMetrics {
448 current_memory_mb: self.resource_usage.memory_usage_mb.load(Ordering::Relaxed),
449 peak_memory_mb: self.resource_usage.peak_memory_mb.load(Ordering::Relaxed),
450 avg_cpu_usage,
451 active_requests: self.resource_usage.active_requests.load(Ordering::Relaxed),
452 peak_concurrent_requests: self
453 .resource_usage
454 .peak_concurrent_requests
455 .load(Ordering::Relaxed),
456 }
457 }
458
459 async fn calculate_connection_metrics(&self) -> ConnectionMetricsSnapshot {
460 let durations = self.connection_metrics.connection_durations.read().await;
461 let avg_duration = if durations.is_empty() {
462 0.0
463 } else {
464 durations.iter().map(|d| d.as_millis() as f64).sum::<f64>() / durations.len() as f64
465 };
466
467 let total_connections = self
468 .connection_metrics
469 .total_connections
470 .load(Ordering::Relaxed);
471 let connection_errors = self
472 .connection_metrics
473 .connection_errors
474 .load(Ordering::Relaxed);
475 let error_rate = if total_connections > 0 {
476 (connection_errors as f64 / total_connections as f64) * 100.0
477 } else {
478 0.0
479 };
480
481 ConnectionMetricsSnapshot {
482 total_connections,
483 active_connections: self
484 .connection_metrics
485 .active_connections
486 .load(Ordering::Relaxed),
487 peak_connections: self
488 .connection_metrics
489 .peak_connections
490 .load(Ordering::Relaxed),
491 connection_error_rate: error_rate,
492 avg_connection_duration_ms: avg_duration,
493 }
494 }
495
496 async fn check_memory_health(&self) -> ComponentHealth {
497 let current_mb = self.resource_usage.memory_usage_mb.load(Ordering::Relaxed);
498 let peak_mb = self.resource_usage.peak_memory_mb.load(Ordering::Relaxed);
499
500 let status = if current_mb > 1000 {
501 ComponentStatus::Critical
503 } else if current_mb > 500 {
504 ComponentStatus::Warning
506 } else {
507 ComponentStatus::Healthy
508 };
509
510 let message = if current_mb > 1000 {
511 Some("High memory usage detected".to_string())
512 } else {
513 None
514 };
515
516 let mut metrics = HashMap::new();
517 metrics.insert("current_mb".to_string(), current_mb as f64);
518 metrics.insert("peak_mb".to_string(), peak_mb as f64);
519
520 ComponentHealth {
521 status,
522 message,
523 last_update: SystemTime::now(),
524 metrics,
525 }
526 }
527
528 async fn check_connection_health(&self) -> ComponentHealth {
529 let active = self
530 .connection_metrics
531 .active_connections
532 .load(Ordering::Relaxed);
533 let peak = self
534 .connection_metrics
535 .peak_connections
536 .load(Ordering::Relaxed);
537 let errors = self
538 .connection_metrics
539 .connection_errors
540 .load(Ordering::Relaxed);
541 let total = self
542 .connection_metrics
543 .total_connections
544 .load(Ordering::Relaxed);
545
546 let error_rate = if total > 0 {
547 (errors as f64 / total as f64) * 100.0
548 } else {
549 0.0
550 };
551
552 let status = if error_rate > 10.0 {
553 ComponentStatus::Critical
554 } else if error_rate > 5.0 || active > 100 {
555 ComponentStatus::Warning
556 } else {
557 ComponentStatus::Healthy
558 };
559
560 let message = if error_rate > 10.0 {
561 Some(format!("High connection error rate: {error_rate:.1}%"))
562 } else if active > 100 {
563 Some("High number of active connections".to_string())
564 } else {
565 None
566 };
567
568 let mut metrics = HashMap::new();
569 metrics.insert("active_connections".to_string(), active as f64);
570 metrics.insert("peak_connections".to_string(), peak as f64);
571 metrics.insert("error_rate".to_string(), error_rate);
572
573 ComponentHealth {
574 status,
575 message,
576 last_update: SystemTime::now(),
577 metrics,
578 }
579 }
580
581 async fn check_request_health(&self) -> ComponentHealth {
582 let active_requests = self.resource_usage.active_requests.load(Ordering::Relaxed);
583
584 let status = if active_requests > 1000 {
585 ComponentStatus::Critical
586 } else if active_requests > 500 {
587 ComponentStatus::Warning
588 } else {
589 ComponentStatus::Healthy
590 };
591
592 let message = if active_requests > 1000 {
593 Some("Very high request load".to_string())
594 } else if active_requests > 500 {
595 Some("High request load".to_string())
596 } else {
597 None
598 };
599
600 let mut metrics = HashMap::new();
601 metrics.insert("active_requests".to_string(), active_requests as f64);
602
603 ComponentHealth {
604 status,
605 message,
606 last_update: SystemTime::now(),
607 metrics,
608 }
609 }
610
611 async fn check_error_health(&self) -> ComponentHealth {
612 let error_counts = self.error_counts.read().await;
613 let total_errors: u64 = error_counts
614 .values()
615 .map(|c| c.load(Ordering::Relaxed))
616 .sum();
617
618 let request_counts = self.request_counts.read().await;
619 let total_requests: u64 = request_counts
620 .values()
621 .map(|c| c.load(Ordering::Relaxed))
622 .sum();
623
624 let error_rate = if total_requests > 0 {
625 (total_errors as f64 / total_requests as f64) * 100.0
626 } else {
627 0.0
628 };
629
630 let status = if error_rate > 10.0 {
631 ComponentStatus::Critical
632 } else if error_rate > 5.0 {
633 ComponentStatus::Warning
634 } else {
635 ComponentStatus::Healthy
636 };
637
638 let message = if error_rate > 10.0 {
639 Some(format!("High error rate: {error_rate:.1}%"))
640 } else if error_rate > 5.0 {
641 Some(format!("Elevated error rate: {error_rate:.1}%"))
642 } else {
643 None
644 };
645
646 let mut metrics = HashMap::new();
647 metrics.insert("error_rate".to_string(), error_rate);
648 metrics.insert("total_errors".to_string(), total_errors as f64);
649
650 ComponentHealth {
651 status,
652 message,
653 last_update: SystemTime::now(),
654 metrics,
655 }
656 }
657
658 fn calculate_health_score(&self, components: &HashMap<String, ComponentHealth>) -> u8 {
659 let total_components = components.len();
660 if total_components == 0 {
661 return 100;
662 }
663
664 let score_sum: u32 = components
665 .values()
666 .map(|health| match health.status {
667 ComponentStatus::Healthy => 100,
668 ComponentStatus::Warning => 75,
669 ComponentStatus::Degraded => 50,
670 ComponentStatus::Critical => 25,
671 ComponentStatus::Unavailable => 0,
672 })
673 .sum();
674
675 (score_sum / total_components as u32) as u8
676 }
677
678 async fn generate_health_alerts(
679 &self,
680 components: &HashMap<String, ComponentHealth>,
681 ) -> Vec<HealthAlert> {
682 let mut alerts = Vec::new();
683
684 for (component_name, health) in components {
685 let (severity, should_alert) = match health.status {
686 ComponentStatus::Critical => (AlertSeverity::Critical, true),
687 ComponentStatus::Warning => (AlertSeverity::Warning, true),
688 ComponentStatus::Degraded => (AlertSeverity::Warning, true),
689 ComponentStatus::Unavailable => (AlertSeverity::Critical, true),
690 ComponentStatus::Healthy => (AlertSeverity::Info, false),
691 };
692
693 if should_alert {
694 let message = health
695 .message
696 .clone()
697 .unwrap_or_else(|| format!("{component_name} component is not healthy"));
698
699 alerts.push(HealthAlert {
700 severity,
701 message,
702 component: component_name.clone(),
703 timestamp: SystemTime::now(),
704 id: format!(
705 "{}_{}",
706 component_name,
707 SystemTime::now()
708 .duration_since(UNIX_EPOCH)
709 .unwrap_or(Duration::from_secs(0))
710 .as_secs()
711 ),
712 });
713 }
714 }
715
716 alerts
717 }
718}
719
720impl Default for MetricsCollector {
721 fn default() -> Self {
722 Self::new()
723 }
724}
725
726impl ResourceUsageTracker {
727 fn new() -> Self {
728 Self {
729 memory_usage_mb: AtomicUsize::new(0),
730 peak_memory_mb: AtomicUsize::new(0),
731 cpu_samples: RwLock::new(Vec::new()),
732 active_requests: AtomicUsize::new(0),
733 peak_concurrent_requests: AtomicUsize::new(0),
734 }
735 }
736}
737
738impl ConnectionMetrics {
739 fn new() -> Self {
740 Self {
741 total_connections: AtomicU64::new(0),
742 active_connections: AtomicUsize::new(0),
743 peak_connections: AtomicUsize::new(0),
744 connection_errors: AtomicU64::new(0),
745 connection_durations: RwLock::new(Vec::new()),
746 }
747 }
748}
749
750impl HealthStatus {
751 fn new() -> Self {
752 Self {
753 health_score: 100,
754 components: HashMap::new(),
755 last_check: SystemTime::now(),
756 alerts: Vec::new(),
757 }
758 }
759}
760
761#[cfg(test)]
762mod tests {
763 use super::*;
764 use tokio::time::Duration;
765
766 #[tokio::test]
767 async fn test_metrics_collection() {
768 let metrics = MetricsCollector::new();
769
770 metrics
772 .record_request("textDocument/completion", Duration::from_millis(100))
773 .await;
774 metrics
775 .record_request("textDocument/hover", Duration::from_millis(50))
776 .await;
777 metrics.record_error("timeout").await;
778
779 metrics.update_memory_usage(100);
780 metrics.record_cpu_usage(25.5).await;
781
782 metrics.record_connection();
783
784 let snapshot = metrics.get_snapshot().await;
785
786 assert!(snapshot.requests.total_requests > 0);
787 assert!(snapshot.resources.current_memory_mb == 100);
788 assert!(snapshot.connections.total_connections > 0);
789 }
790
791 #[tokio::test]
792 async fn test_health_check() {
793 let metrics = MetricsCollector::new();
794
795 metrics.update_memory_usage(100); metrics.record_cpu_usage(25.0).await;
798
799 let health = metrics.perform_health_check().await;
800
801 assert!(health.health_score > 50);
802 assert!(health.components.contains_key("memory"));
803 assert!(health.components.contains_key("connections"));
804 }
805
806 #[tokio::test]
807 async fn test_resource_limits() {
808 let metrics = MetricsCollector::new();
809
810 metrics.update_memory_usage(1200); metrics
815 .resource_usage
816 .active_requests
817 .store(1500, Ordering::Relaxed);
818
819 let health = metrics.perform_health_check().await;
820 let memory_health = health.components.get("memory").unwrap();
821 let requests_health = health.components.get("requests").unwrap();
822
823 assert!(matches!(memory_health.status, ComponentStatus::Critical));
824 assert!(matches!(requests_health.status, ComponentStatus::Critical));
825
826 assert!(health.health_score < 75);
828 assert!(!health.alerts.is_empty());
829 }
830}