use crate::error::{ObservabilityError, Result};
use chrono::{DateTime, Duration, Utc};
use parking_lot::RwLock;
use serde::{Deserialize, Serialize};
use std::collections::VecDeque;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
pub const DEFAULT_AVAILABILITY_TARGET: f64 = 99.9;
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub struct AvailabilitySample {
pub timestamp: DateTime<Utc>,
pub total_requests: u64,
pub successful_requests: u64,
pub failed_requests: u64,
pub availability_pct: f64,
}
impl AvailabilitySample {
pub fn new(
timestamp: DateTime<Utc>,
total_requests: u64,
successful_requests: u64,
) -> Self {
let failed_requests = total_requests.saturating_sub(successful_requests);
let availability_pct = if total_requests > 0 {
(successful_requests as f64 / total_requests as f64) * 100.0
} else {
100.0 };
Self {
timestamp,
total_requests,
successful_requests,
failed_requests,
availability_pct,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AvailabilityStatus {
pub service_name: String,
pub current_availability: f64,
pub target_availability: f64,
pub is_slo_met: bool,
pub total_requests: u64,
pub successful_requests: u64,
pub failed_requests: u64,
pub window_duration: Duration,
pub evaluated_at: DateTime<Utc>,
}
pub struct AvailabilityTracker {
service_name: String,
target_availability: f64,
total_requests: AtomicU64,
successful_requests: AtomicU64,
failed_requests: AtomicU64,
samples: Arc<RwLock<VecDeque<AvailabilitySample>>>,
max_samples: usize,
retention_duration: Option<Duration>,
}
impl AvailabilityTracker {
pub fn new(service_name: impl Into<String>) -> Self {
Self {
service_name: service_name.into(),
target_availability: DEFAULT_AVAILABILITY_TARGET,
total_requests: AtomicU64::new(0),
successful_requests: AtomicU64::new(0),
failed_requests: AtomicU64::new(0),
samples: Arc::new(RwLock::new(VecDeque::with_capacity(1000))),
max_samples: 1000,
retention_duration: Some(Duration::days(30)),
}
}
pub fn with_target(mut self, target_pct: f64) -> Self {
self.target_availability = target_pct;
self
}
pub fn with_max_samples(mut self, max_samples: usize) -> Self {
self.max_samples = max_samples;
self
}
pub fn with_retention(mut self, duration: Duration) -> Self {
self.retention_duration = Some(duration);
self
}
pub fn service_name(&self) -> &str {
&self.service_name
}
pub fn target_availability(&self) -> f64 {
self.target_availability
}
pub fn record_success(&self) {
self.total_requests.fetch_add(1, Ordering::SeqCst);
self.successful_requests.fetch_add(1, Ordering::SeqCst);
}
pub fn record_failure(&self) {
self.total_requests.fetch_add(1, Ordering::SeqCst);
self.failed_requests.fetch_add(1, Ordering::SeqCst);
}
pub fn record_request(&self, is_success: bool) {
if is_success {
self.record_success();
} else {
self.record_failure();
}
}
pub fn record_batch(&self, total: u64, successful: u64) {
self.total_requests.fetch_add(total, Ordering::SeqCst);
self.successful_requests.fetch_add(successful, Ordering::SeqCst);
let failed = total.saturating_sub(successful);
self.failed_requests.fetch_add(failed, Ordering::SeqCst);
}
pub fn take_snapshot(&self) -> AvailabilitySample {
let total = self.total_requests.load(Ordering::SeqCst);
let successful = self.successful_requests.load(Ordering::SeqCst);
let sample = AvailabilitySample::new(Utc::now(), total, successful);
let mut samples = self.samples.write();
if let Some(retention) = self.retention_duration {
let cutoff = Utc::now() - retention;
while samples.front().map_or(false, |s| s.timestamp < cutoff) {
samples.pop_front();
}
}
while samples.len() >= self.max_samples {
samples.pop_front();
}
samples.push_back(sample);
sample
}
pub fn reset_counters(&self) {
self.total_requests.store(0, Ordering::SeqCst);
self.successful_requests.store(0, Ordering::SeqCst);
self.failed_requests.store(0, Ordering::SeqCst);
}
pub fn current_availability(&self) -> f64 {
let total = self.total_requests.load(Ordering::SeqCst);
let successful = self.successful_requests.load(Ordering::SeqCst);
if total == 0 {
return 100.0;
}
(successful as f64 / total as f64) * 100.0
}
pub fn is_slo_met(&self) -> bool {
self.current_availability() >= self.target_availability
}
pub fn get_status(&self) -> AvailabilityStatus {
let total = self.total_requests.load(Ordering::SeqCst);
let successful = self.successful_requests.load(Ordering::SeqCst);
let failed = self.failed_requests.load(Ordering::SeqCst);
let availability = self.current_availability();
AvailabilityStatus {
service_name: self.service_name.clone(),
current_availability: availability,
target_availability: self.target_availability,
is_slo_met: availability >= self.target_availability,
total_requests: total,
successful_requests: successful,
failed_requests: failed,
window_duration: self.retention_duration.unwrap_or_else(|| Duration::days(30)),
evaluated_at: Utc::now(),
}
}
pub fn availability_in_window(&self, window: Duration) -> Result<f64> {
let samples = self.samples.read();
let cutoff = Utc::now() - window;
let (total, successful) = samples
.iter()
.filter(|s| s.timestamp >= cutoff)
.fold((0u64, 0u64), |(t, s), sample| {
(t + sample.total_requests, s + sample.successful_requests)
});
if total == 0 {
return Err(ObservabilityError::SloCalculationError(
"No samples available in the specified window".to_string(),
));
}
Ok((successful as f64 / total as f64) * 100.0)
}
pub fn get_samples(&self) -> Vec<AvailabilitySample> {
self.samples.read().iter().cloned().collect()
}
pub fn samples_in_range(
&self,
start: DateTime<Utc>,
end: DateTime<Utc>,
) -> Vec<AvailabilitySample> {
self.samples
.read()
.iter()
.filter(|s| s.timestamp >= start && s.timestamp <= end)
.cloned()
.collect()
}
pub fn estimated_downtime_minutes(&self, total_window_minutes: f64) -> f64 {
let availability = self.current_availability();
let downtime_fraction = (100.0 - availability) / 100.0;
total_window_minutes * downtime_fraction
}
pub fn error_rate(&self) -> f64 {
100.0 - self.current_availability()
}
pub fn sample_count(&self) -> usize {
self.samples.read().len()
}
}
impl Default for AvailabilityTracker {
fn default() -> Self {
Self::new("default_service")
}
}
pub struct CompositeAvailability {
trackers: Vec<Arc<AvailabilityTracker>>,
}
impl CompositeAvailability {
pub fn new() -> Self {
Self {
trackers: Vec::new(),
}
}
pub fn add_tracker(&mut self, tracker: Arc<AvailabilityTracker>) {
self.trackers.push(tracker);
}
pub fn serial_availability(&self) -> f64 {
if self.trackers.is_empty() {
return 100.0;
}
self.trackers
.iter()
.map(|t| t.current_availability() / 100.0)
.product::<f64>()
* 100.0
}
pub fn parallel_availability(&self) -> f64 {
if self.trackers.is_empty() {
return 100.0;
}
let failure_product: f64 = self
.trackers
.iter()
.map(|t| 1.0 - (t.current_availability() / 100.0))
.product();
(1.0 - failure_product) * 100.0
}
pub fn weighted_average(&self, weights: &[f64]) -> Result<f64> {
if self.trackers.len() != weights.len() {
return Err(ObservabilityError::SloCalculationError(
"Weights count must match tracker count".to_string(),
));
}
if self.trackers.is_empty() {
return Ok(100.0);
}
let total_weight: f64 = weights.iter().sum();
if total_weight <= 0.0 {
return Err(ObservabilityError::SloCalculationError(
"Total weight must be positive".to_string(),
));
}
let weighted_sum: f64 = self
.trackers
.iter()
.zip(weights.iter())
.map(|(t, w)| t.current_availability() * w)
.sum();
Ok(weighted_sum / total_weight)
}
pub fn get_all_statuses(&self) -> Vec<AvailabilityStatus> {
self.trackers.iter().map(|t| t.get_status()).collect()
}
}
impl Default for CompositeAvailability {
fn default() -> Self {
Self::new()
}
}
pub struct UptimeCalculator {
start_time: DateTime<Utc>,
downtime_duration: Arc<RwLock<Duration>>,
current_downtime_start: Arc<RwLock<Option<DateTime<Utc>>>>,
}
impl UptimeCalculator {
pub fn new() -> Self {
Self {
start_time: Utc::now(),
downtime_duration: Arc::new(RwLock::new(Duration::zero())),
current_downtime_start: Arc::new(RwLock::new(None)),
}
}
pub fn mark_down(&self) {
let mut start = self.current_downtime_start.write();
if start.is_none() {
*start = Some(Utc::now());
}
}
pub fn mark_up(&self) {
let mut start = self.current_downtime_start.write();
if let Some(down_start) = start.take() {
let downtime = Utc::now() - down_start;
*self.downtime_duration.write() = *self.downtime_duration.read() + downtime;
}
}
pub fn is_down(&self) -> bool {
self.current_downtime_start.read().is_some()
}
pub fn total_downtime(&self) -> Duration {
let mut total = *self.downtime_duration.read();
if let Some(down_start) = *self.current_downtime_start.read() {
total = total + (Utc::now() - down_start);
}
total
}
pub fn total_uptime(&self) -> Duration {
let total_duration = Utc::now() - self.start_time;
total_duration - self.total_downtime()
}
pub fn uptime_percentage(&self) -> f64 {
let total_duration = (Utc::now() - self.start_time).num_seconds() as f64;
if total_duration <= 0.0 {
return 100.0;
}
let uptime_seconds = self.total_uptime().num_seconds() as f64;
(uptime_seconds / total_duration) * 100.0
}
pub fn reset(&self) {
*self.downtime_duration.write() = Duration::zero();
*self.current_downtime_start.write() = None;
}
}
impl Default for UptimeCalculator {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_availability_sample() {
let sample = AvailabilitySample::new(Utc::now(), 100, 99);
assert_eq!(sample.total_requests, 100);
assert_eq!(sample.successful_requests, 99);
assert_eq!(sample.failed_requests, 1);
assert_eq!(sample.availability_pct, 99.0);
}
#[test]
fn test_availability_tracker() {
let tracker = AvailabilityTracker::new("test_service").with_target(99.0);
for _ in 0..99 {
tracker.record_success();
}
tracker.record_failure();
assert_eq!(tracker.current_availability(), 99.0);
assert!(tracker.is_slo_met());
let status = tracker.get_status();
assert_eq!(status.total_requests, 100);
assert_eq!(status.successful_requests, 99);
assert_eq!(status.failed_requests, 1);
}
#[test]
fn test_composite_availability() {
let tracker1 = Arc::new(AvailabilityTracker::new("service1"));
let tracker2 = Arc::new(AvailabilityTracker::new("service2"));
for _ in 0..99 {
tracker1.record_success();
}
tracker1.record_failure();
for _ in 0..99 {
tracker2.record_success();
}
tracker2.record_failure();
let mut composite = CompositeAvailability::new();
composite.add_tracker(tracker1);
composite.add_tracker(tracker2);
let serial = composite.serial_availability();
assert!((serial - 98.01).abs() < 0.01);
let parallel = composite.parallel_availability();
assert!((parallel - 99.99).abs() < 0.01);
}
#[test]
fn test_uptime_calculator() {
let calculator = UptimeCalculator::new();
assert!(!calculator.is_down());
assert!(calculator.uptime_percentage() > 99.9);
calculator.mark_down();
assert!(calculator.is_down());
calculator.mark_up();
assert!(!calculator.is_down());
}
#[test]
fn test_batch_recording() {
let tracker = AvailabilityTracker::new("batch_test");
tracker.record_batch(1000, 950);
assert_eq!(tracker.current_availability(), 95.0);
assert_eq!(tracker.error_rate(), 5.0);
}
}