use crate::traits::{EvaluationResult, QualityScore};
use crate::EvaluationError;
use std::collections::{HashMap, VecDeque};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use tokio::sync::mpsc;
use tokio::task;
use voirs_sdk::AudioBuffer;
#[derive(Debug, Clone)]
pub struct StreamingConfig {
pub chunk_size: usize,
pub overlap_size: usize,
pub max_buffer_chunks: usize,
pub target_latency_ms: u64,
pub enable_quality_monitoring: bool,
pub enable_adaptive_processing: bool,
pub enable_quality_prediction: bool,
pub enable_network_adaptation: bool,
pub enable_parallel_processing: bool,
pub enable_anomaly_detection: bool,
pub prediction_horizon: usize,
pub network_monitor_interval_ms: u64,
pub anomaly_sensitivity: f32,
pub enable_predictive_buffering: bool,
}
impl Default for StreamingConfig {
fn default() -> Self {
Self {
chunk_size: 1024, overlap_size: 256, max_buffer_chunks: 10, target_latency_ms: 100, enable_quality_monitoring: true,
enable_adaptive_processing: true,
enable_quality_prediction: true,
enable_network_adaptation: false, enable_parallel_processing: true,
enable_anomaly_detection: true,
prediction_horizon: 5, network_monitor_interval_ms: 1000, anomaly_sensitivity: 0.7, enable_predictive_buffering: true,
}
}
}
#[derive(Debug, Clone)]
pub struct AudioChunk {
pub samples: Vec<f32>,
pub sample_rate: u32,
pub timestamp: Instant,
pub sequence: u64,
}
impl AudioChunk {
pub fn new(samples: Vec<f32>, sample_rate: u32, sequence: u64) -> Self {
Self {
samples,
sample_rate,
timestamp: Instant::now(),
sequence,
}
}
pub fn duration(&self) -> f64 {
self.samples.len() as f64 / self.sample_rate as f64
}
pub fn to_audio_buffer(&self) -> AudioBuffer {
AudioBuffer::mono(self.samples.clone(), self.sample_rate)
}
}
#[derive(Debug, Clone)]
pub struct StreamingQualityMetrics {
pub snr_estimate: f32,
pub dynamic_range: f32,
pub spectral_flatness: f32,
pub energy_level: f32,
pub clipping_detected: bool,
pub silence_detected: bool,
pub processing_latency_ms: f64,
pub timestamp: Instant,
}
impl Default for StreamingQualityMetrics {
fn default() -> Self {
Self {
snr_estimate: 0.0,
dynamic_range: 0.0,
spectral_flatness: 0.0,
energy_level: 0.0,
clipping_detected: false,
silence_detected: false,
processing_latency_ms: 0.0,
timestamp: Instant::now(),
}
}
}
#[derive(Debug, Clone)]
pub struct QualityPrediction {
pub predicted_score: f32,
pub confidence: f32,
pub trend: f32,
pub risk_level: RiskLevel,
pub recommendations: Vec<String>,
}
#[derive(Debug, Clone, PartialEq)]
pub enum RiskLevel {
Low,
Medium,
High,
Critical,
}
#[derive(Debug, Clone)]
pub struct NetworkCondition {
pub bandwidth_estimate: u64,
pub rtt_ms: f64,
pub packet_loss_rate: f32,
pub jitter_ms: f64,
pub quality_score: f32,
pub timestamp: Instant,
}
impl Default for NetworkCondition {
fn default() -> Self {
Self {
bandwidth_estimate: 1_000_000, rtt_ms: 50.0,
packet_loss_rate: 0.0,
jitter_ms: 5.0,
quality_score: 1.0,
timestamp: Instant::now(),
}
}
}
#[derive(Debug, Clone)]
pub struct AnomalyDetection {
pub anomaly_detected: bool,
pub anomaly_score: f32,
pub anomaly_type: AnomalyType,
pub description: String,
pub severity: AnomalySeverity,
pub recommended_actions: Vec<String>,
}
#[derive(Debug, Clone, PartialEq)]
pub enum AnomalyType {
None,
SuddenQualityDrop,
UnexpectedSilence,
ExcessiveClipping,
FrequencyImbalance,
ProcessingDelay,
NetworkIssue,
Unknown,
}
#[derive(Debug, Clone, PartialEq)]
pub enum AnomalySeverity {
Info,
Warning,
Error,
Critical,
}
#[derive(Debug, Clone)]
pub struct AdvancedProcessingStats {
pub basic_stats: ProcessingStats,
pub prediction_accuracy: f32,
pub network_adaptation_efficiency: f32,
pub anomaly_detection_rate: f32,
pub parallel_efficiency: f32,
pub buffer_utilization: BufferUtilization,
}
#[derive(Debug, Clone)]
pub struct BufferUtilization {
pub average_fill_percentage: f32,
pub peak_usage: f32,
pub underruns: u64,
pub overruns: u64,
pub adaptive_adjustments: u64,
}
pub struct StreamingEvaluator {
config: StreamingConfig,
buffer: Arc<Mutex<VecDeque<AudioChunk>>>,
sequence: Arc<Mutex<u64>>,
metrics_history: Arc<Mutex<VecDeque<StreamingQualityMetrics>>>,
quality_sender: Option<mpsc::UnboundedSender<StreamingQualityMetrics>>,
processing_stats: Arc<Mutex<ProcessingStats>>,
quality_predictor: Arc<Mutex<QualityPredictor>>,
network_monitor: Arc<Mutex<NetworkMonitor>>,
anomaly_detector: Arc<Mutex<AnomalyDetector>>,
advanced_stats: Arc<Mutex<AdvancedProcessingStats>>,
parallel_pool: Option<Arc<tokio::task::JoinSet<()>>>,
}
#[derive(Debug, Default, Clone)]
pub struct ProcessingStats {
pub total_chunks_processed: u64,
pub total_processing_time_ms: f64,
pub average_latency_ms: f64,
pub peak_latency_ms: f64,
pub dropped_chunks: u64,
}
#[derive(Debug)]
struct QualityPredictor {
quality_history: VecDeque<f32>,
prediction_models: HashMap<String, PredictionModel>,
prediction_accuracy: f32,
}
#[derive(Debug)]
struct PredictionModel {
model_type: String,
parameters: Vec<f32>,
accuracy_history: VecDeque<f32>,
}
#[derive(Debug)]
struct NetworkMonitor {
current_condition: NetworkCondition,
condition_history: VecDeque<NetworkCondition>,
last_check: Instant,
}
#[derive(Debug)]
struct AnomalyDetector {
baseline_metrics: StreamingQualityMetrics,
anomaly_history: VecDeque<AnomalyDetection>,
thresholds: AnomalyThresholds,
}
#[derive(Debug)]
struct AnomalyThresholds {
quality_drop_threshold: f32,
silence_duration_threshold: Duration,
clipping_rate_threshold: f32,
latency_threshold: f64,
energy_deviation_threshold: f32,
}
impl StreamingEvaluator {
pub fn new(config: StreamingConfig) -> Self {
let quality_predictor = QualityPredictor {
quality_history: VecDeque::new(),
prediction_models: {
let mut models = HashMap::new();
models.insert(
"moving_average".to_string(),
PredictionModel {
model_type: "moving_average".to_string(),
parameters: vec![0.2, 0.8], accuracy_history: VecDeque::new(),
},
);
models.insert(
"trend_analysis".to_string(),
PredictionModel {
model_type: "trend_analysis".to_string(),
parameters: vec![0.1, 0.3, 0.6], accuracy_history: VecDeque::new(),
},
);
models
},
prediction_accuracy: 0.5,
};
let network_monitor = NetworkMonitor {
current_condition: NetworkCondition::default(),
condition_history: VecDeque::new(),
last_check: Instant::now(),
};
let anomaly_detector = AnomalyDetector {
baseline_metrics: StreamingQualityMetrics::default(),
anomaly_history: VecDeque::new(),
thresholds: AnomalyThresholds {
quality_drop_threshold: 0.3,
silence_duration_threshold: Duration::from_secs(2),
clipping_rate_threshold: 0.1,
latency_threshold: 200.0, energy_deviation_threshold: 0.5,
},
};
let advanced_stats = AdvancedProcessingStats {
basic_stats: ProcessingStats::default(),
prediction_accuracy: 0.5,
network_adaptation_efficiency: 1.0,
anomaly_detection_rate: 0.0,
parallel_efficiency: 1.0,
buffer_utilization: BufferUtilization {
average_fill_percentage: 0.0,
peak_usage: 0.0,
underruns: 0,
overruns: 0,
adaptive_adjustments: 0,
},
};
Self {
config,
buffer: Arc::new(Mutex::new(VecDeque::new())),
sequence: Arc::new(Mutex::new(0)),
metrics_history: Arc::new(Mutex::new(VecDeque::new())),
quality_sender: None,
processing_stats: Arc::new(Mutex::new(ProcessingStats::default())),
quality_predictor: Arc::new(Mutex::new(quality_predictor)),
network_monitor: Arc::new(Mutex::new(network_monitor)),
anomaly_detector: Arc::new(Mutex::new(anomaly_detector)),
advanced_stats: Arc::new(Mutex::new(advanced_stats)),
parallel_pool: None,
}
}
pub fn setup_quality_monitoring(&mut self) -> mpsc::UnboundedReceiver<StreamingQualityMetrics> {
let (sender, receiver) = mpsc::unbounded_channel();
self.quality_sender = Some(sender);
receiver
}
pub async fn process_chunk(&mut self, chunk: AudioChunk) -> EvaluationResult<()> {
let start_time = Instant::now();
{
let mut buffer = self
.buffer
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock buffer".to_string(),
source: None,
})?;
buffer.push_back(chunk.clone());
while buffer.len() > self.config.max_buffer_chunks {
buffer.pop_front();
let mut stats =
self.processing_stats
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock processing stats".to_string(),
source: None,
})?;
stats.dropped_chunks += 1;
}
}
if self.config.enable_quality_monitoring {
let metrics = self.calculate_chunk_quality_metrics(&chunk).await?;
let processing_time = start_time.elapsed().as_secs_f64() * 1000.0;
{
let mut stats =
self.processing_stats
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock processing stats".to_string(),
source: None,
})?;
stats.total_chunks_processed += 1;
stats.total_processing_time_ms += processing_time;
stats.average_latency_ms =
stats.total_processing_time_ms / stats.total_chunks_processed as f64;
if processing_time > stats.peak_latency_ms {
stats.peak_latency_ms = processing_time;
}
}
{
let mut history =
self.metrics_history
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock metrics history".to_string(),
source: None,
})?;
history.push_back(metrics.clone());
while history.len() > 100 {
history.pop_front();
}
}
if let Some(ref sender) = self.quality_sender {
let mut updated_metrics = metrics;
updated_metrics.processing_latency_ms = processing_time;
let _ = sender.send(updated_metrics);
}
}
if self.config.enable_adaptive_processing {
self.adapt_processing_parameters().await?;
}
Ok(())
}
async fn calculate_chunk_quality_metrics(
&self,
chunk: &AudioChunk,
) -> EvaluationResult<StreamingQualityMetrics> {
let samples = &chunk.samples;
if samples.is_empty() {
return Ok(StreamingQualityMetrics::default());
}
let energy_level = samples.iter().map(|&x| x * x).sum::<f32>() / samples.len() as f32;
let energy_level = energy_level.sqrt();
let max_val = samples.iter().map(|&x| x.abs()).fold(0.0f32, f32::max);
let rms = energy_level;
let dynamic_range = if rms > 0.0 {
20.0 * (max_val / rms).log10()
} else {
0.0
};
let clipping_threshold = 0.95;
let clipping_detected = samples.iter().any(|&x| x.abs() > clipping_threshold);
let silence_threshold = 0.001;
let silence_detected = energy_level < silence_threshold;
let noise_floor = 0.01;
let snr_estimate = if energy_level > noise_floor {
20.0 * (energy_level / noise_floor).log10()
} else {
0.0
};
let spectral_flatness = self.calculate_spectral_flatness(samples);
Ok(StreamingQualityMetrics {
snr_estimate,
dynamic_range,
spectral_flatness,
energy_level,
clipping_detected,
silence_detected,
processing_latency_ms: 0.0, timestamp: Instant::now(),
})
}
fn calculate_spectral_flatness(&self, samples: &[f32]) -> f32 {
if samples.len() < 32 {
return 0.5; }
let window_size = 32;
let num_windows = samples.len() / window_size;
if num_windows < 2 {
return 0.5;
}
let mut window_energies = Vec::new();
for i in 0..num_windows {
let start = i * window_size;
let end = (start + window_size).min(samples.len());
let energy: f32 = samples[start..end].iter().map(|&x| x * x).sum();
window_energies.push(energy.max(1e-10)); }
let geometric_mean =
window_energies.iter().map(|&x| x.ln()).sum::<f32>() / window_energies.len() as f32;
let geometric_mean = geometric_mean.exp();
let arithmetic_mean = window_energies.iter().sum::<f32>() / window_energies.len() as f32;
if arithmetic_mean > 0.0 {
(geometric_mean / arithmetic_mean).min(1.0)
} else {
0.0
}
}
async fn adapt_processing_parameters(&mut self) -> EvaluationResult<()> {
let stats = {
let stats_guard =
self.processing_stats
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock processing stats".to_string(),
source: None,
})?;
ProcessingStats {
total_chunks_processed: stats_guard.total_chunks_processed,
total_processing_time_ms: stats_guard.total_processing_time_ms,
average_latency_ms: stats_guard.average_latency_ms,
peak_latency_ms: stats_guard.peak_latency_ms,
dropped_chunks: stats_guard.dropped_chunks,
}
};
if stats.average_latency_ms > self.config.target_latency_ms as f64 * 1.5 {
if self.config.chunk_size < 4096 {
self.config.chunk_size = (self.config.chunk_size as f32 * 1.2) as usize;
self.config.overlap_size = self.config.chunk_size / 4; }
} else if stats.average_latency_ms < self.config.target_latency_ms as f64 * 0.5 {
if self.config.chunk_size > 256 {
self.config.chunk_size = (self.config.chunk_size as f32 * 0.8) as usize;
self.config.overlap_size = self.config.chunk_size / 4;
}
}
Ok(())
}
pub fn get_current_metrics(&self) -> EvaluationResult<Option<StreamingQualityMetrics>> {
let history =
self.metrics_history
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock metrics history".to_string(),
source: None,
})?;
Ok(history.back().cloned())
}
pub fn get_processing_stats(&self) -> EvaluationResult<ProcessingStats> {
let stats = self
.processing_stats
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock processing stats".to_string(),
source: None,
})?;
Ok(ProcessingStats {
total_chunks_processed: stats.total_chunks_processed,
total_processing_time_ms: stats.total_processing_time_ms,
average_latency_ms: stats.average_latency_ms,
peak_latency_ms: stats.peak_latency_ms,
dropped_chunks: stats.dropped_chunks,
})
}
pub fn get_buffered_audio(&self) -> EvaluationResult<Option<AudioBuffer>> {
let buffer = self
.buffer
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock buffer".to_string(),
source: None,
})?;
if buffer.is_empty() {
return Ok(None);
}
let first_chunk = &buffer[0];
let sample_rate = first_chunk.sample_rate;
let mut combined_samples = Vec::new();
for chunk in buffer.iter() {
combined_samples.extend_from_slice(&chunk.samples);
}
Ok(Some(AudioBuffer::mono(combined_samples, sample_rate)))
}
pub fn reset(&mut self) -> EvaluationResult<()> {
{
let mut buffer = self
.buffer
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock buffer".to_string(),
source: None,
})?;
buffer.clear();
}
{
let mut sequence =
self.sequence
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock sequence".to_string(),
source: None,
})?;
*sequence = 0;
}
{
let mut history =
self.metrics_history
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock metrics history".to_string(),
source: None,
})?;
history.clear();
}
{
let mut stats =
self.processing_stats
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock processing stats".to_string(),
source: None,
})?;
*stats = ProcessingStats::default();
}
Ok(())
}
pub async fn predict_quality(&mut self) -> EvaluationResult<QualityPrediction> {
if !self.config.enable_quality_prediction {
return Ok(QualityPrediction {
predicted_score: 0.5,
confidence: 0.0,
trend: 0.0,
risk_level: RiskLevel::Low,
recommendations: vec!["Quality prediction disabled".to_string()],
});
}
let mut predictor =
self.quality_predictor
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock quality predictor".to_string(),
source: None,
})?;
let metrics_history =
self.metrics_history
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock metrics history".to_string(),
source: None,
})?;
if metrics_history.len() < 3 {
return Ok(QualityPrediction {
predicted_score: 0.5,
confidence: 0.3,
trend: 0.0,
risk_level: RiskLevel::Medium,
recommendations: vec!["Insufficient data for prediction".to_string()],
});
}
let recent_scores: Vec<f32> = metrics_history
.iter()
.rev()
.take(self.config.prediction_horizon)
.map(|m| (m.snr_estimate + m.dynamic_range + m.energy_level * 10.0) / 3.0)
.collect();
for &score in &recent_scores {
predictor.quality_history.push_back(score);
if predictor.quality_history.len() > 20 {
predictor.quality_history.pop_front();
}
}
let mut predictions = Vec::new();
let mut confidences = Vec::new();
if let Some(ma_model) = predictor.prediction_models.get("moving_average") {
let short_ma = recent_scores.iter().take(3).sum::<f32>() / 3.0;
let long_ma = recent_scores.iter().sum::<f32>() / recent_scores.len() as f32;
let ma_prediction =
ma_model.parameters[0] * short_ma + ma_model.parameters[1] * long_ma;
predictions.push(ma_prediction);
confidences.push(0.7);
}
if let Some(trend_model) = predictor.prediction_models.get("trend_analysis") {
let trend = if recent_scores.len() >= 2 {
(recent_scores[0] - recent_scores[recent_scores.len() - 1])
/ recent_scores.len() as f32
} else {
0.0
};
let trend_prediction = recent_scores[0] + trend * trend_model.parameters[1];
predictions.push(trend_prediction);
confidences.push(0.6);
}
let predicted_score = if !predictions.is_empty() {
predictions.iter().sum::<f32>() / predictions.len() as f32
} else {
0.5
};
let confidence = if !confidences.is_empty() {
confidences.iter().sum::<f32>() / confidences.len() as f32
} else {
0.3
};
let trend = if recent_scores.len() >= 2 {
recent_scores[0] - recent_scores[recent_scores.len() - 1]
} else {
0.0
};
let risk_level = match predicted_score {
s if s < 0.2 => RiskLevel::Critical,
s if s < 0.4 => RiskLevel::High,
s if s < 0.6 => RiskLevel::Medium,
_ => RiskLevel::Low,
};
let mut recommendations = Vec::new();
if predicted_score < 0.5 {
recommendations
.push("Consider reducing chunk size for better responsiveness".to_string());
}
if trend < -0.1 {
recommendations.push("Quality trending downward - monitor closely".to_string());
}
if confidence < 0.5 {
recommendations.push("Low prediction confidence - collect more data".to_string());
}
Ok(QualityPrediction {
predicted_score,
confidence,
trend,
risk_level,
recommendations,
})
}
pub async fn monitor_network_conditions(&mut self) -> EvaluationResult<NetworkCondition> {
if !self.config.enable_network_adaptation {
return Ok(NetworkCondition::default());
}
let mut monitor =
self.network_monitor
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock network monitor".to_string(),
source: None,
})?;
let now = Instant::now();
if now.duration_since(monitor.last_check).as_millis()
< self.config.network_monitor_interval_ms as u128
{
return Ok(monitor.current_condition.clone());
}
let mut condition = NetworkCondition::default();
let random_factor = (now.elapsed().as_secs() % 100) as f64 / 100.0;
condition.bandwidth_estimate = (1_000_000.0 * (0.5 + random_factor)) as u64;
condition.rtt_ms = 20.0 + random_factor * 100.0;
condition.packet_loss_rate = (random_factor * 0.05) as f32;
condition.jitter_ms = random_factor * 20.0;
let bandwidth_score = (condition.bandwidth_estimate as f32 / 2_000_000.0).min(1.0);
let latency_score = 1.0 - (condition.rtt_ms / 200.0).min(1.0) as f32;
let loss_score = 1.0 - condition.packet_loss_rate * 20.0;
let jitter_score = 1.0 - (condition.jitter_ms / 50.0).min(1.0) as f32;
condition.quality_score =
(bandwidth_score + latency_score + loss_score + jitter_score) / 4.0;
condition.timestamp = now;
monitor.condition_history.push_back(condition.clone());
if monitor.condition_history.len() > 50 {
monitor.condition_history.pop_front();
}
monitor.current_condition = condition.clone();
monitor.last_check = now;
drop(monitor);
self.adapt_to_network_conditions(&condition).await?;
Ok(condition)
}
pub async fn detect_anomalies(
&mut self,
current_metrics: &StreamingQualityMetrics,
) -> EvaluationResult<AnomalyDetection> {
if !self.config.enable_anomaly_detection {
return Ok(AnomalyDetection {
anomaly_detected: false,
anomaly_score: 0.0,
anomaly_type: AnomalyType::None,
description: "Anomaly detection disabled".to_string(),
severity: AnomalySeverity::Info,
recommended_actions: Vec::new(),
});
}
let mut detector =
self.anomaly_detector
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock anomaly detector".to_string(),
source: None,
})?;
if detector.baseline_metrics.energy_level == 0.0 {
detector.baseline_metrics = current_metrics.clone();
return Ok(AnomalyDetection {
anomaly_detected: false,
anomaly_score: 0.0,
anomaly_type: AnomalyType::None,
description: "Establishing baseline".to_string(),
severity: AnomalySeverity::Info,
recommended_actions: Vec::new(),
});
}
let mut anomaly_score = 0.0;
let mut detected_anomalies = Vec::new();
let quality_diff = detector.baseline_metrics.snr_estimate - current_metrics.snr_estimate;
if quality_diff > detector.thresholds.quality_drop_threshold {
anomaly_score += quality_diff * 2.0;
detected_anomalies.push(AnomalyType::SuddenQualityDrop);
}
if current_metrics.clipping_detected {
anomaly_score += 0.5;
detected_anomalies.push(AnomalyType::ExcessiveClipping);
}
if current_metrics.silence_detected && detector.baseline_metrics.energy_level > 0.01 {
anomaly_score += 0.3;
detected_anomalies.push(AnomalyType::UnexpectedSilence);
}
if current_metrics.processing_latency_ms > detector.thresholds.latency_threshold {
anomaly_score += 0.4;
detected_anomalies.push(AnomalyType::ProcessingDelay);
}
let energy_diff =
(current_metrics.energy_level - detector.baseline_metrics.energy_level).abs();
if energy_diff > detector.thresholds.energy_deviation_threshold {
anomaly_score += energy_diff;
detected_anomalies.push(AnomalyType::FrequencyImbalance);
}
anomaly_score *= self.config.anomaly_sensitivity;
let (anomaly_type, severity) = if detected_anomalies.is_empty() {
(AnomalyType::None, AnomalySeverity::Info)
} else {
let primary_anomaly = detected_anomalies[0].clone();
let severity = match anomaly_score {
s if s > 1.0 => AnomalySeverity::Critical,
s if s > 0.7 => AnomalySeverity::Error,
s if s > 0.4 => AnomalySeverity::Warning,
_ => AnomalySeverity::Info,
};
(primary_anomaly, severity)
};
let anomaly_detected = anomaly_score > 0.3;
let description = match &anomaly_type {
AnomalyType::None => "No anomalies detected".to_string(),
AnomalyType::SuddenQualityDrop => format!("Quality dropped by {:.2}", quality_diff),
AnomalyType::ExcessiveClipping => "Audio clipping detected".to_string(),
AnomalyType::UnexpectedSilence => "Unexpected silence period".to_string(),
AnomalyType::ProcessingDelay => format!(
"Processing latency: {:.1}ms",
current_metrics.processing_latency_ms
),
AnomalyType::FrequencyImbalance => "Energy level deviation detected".to_string(),
_ => "Unknown anomaly detected".to_string(),
};
let mut recommended_actions = Vec::new();
if anomaly_detected {
match anomaly_type {
AnomalyType::SuddenQualityDrop => {
recommended_actions.push("Check audio source quality".to_string());
recommended_actions.push("Verify network conditions".to_string());
}
AnomalyType::ExcessiveClipping => {
recommended_actions.push("Reduce input gain".to_string());
recommended_actions.push("Check for audio saturation".to_string());
}
AnomalyType::ProcessingDelay => {
recommended_actions.push("Increase chunk size".to_string());
recommended_actions.push("Check system resources".to_string());
}
_ => {
recommended_actions.push("Monitor closely".to_string());
}
}
}
let detection = AnomalyDetection {
anomaly_detected,
anomaly_score,
anomaly_type,
description,
severity,
recommended_actions,
};
detector.anomaly_history.push_back(detection.clone());
if detector.anomaly_history.len() > 100 {
detector.anomaly_history.pop_front();
}
Ok(detection)
}
async fn adapt_to_network_conditions(
&mut self,
condition: &NetworkCondition,
) -> EvaluationResult<()> {
if condition.bandwidth_estimate < 500_000 {
if self.config.chunk_size > 512 {
self.config.chunk_size = 512;
self.config.overlap_size = 128;
}
} else if condition.bandwidth_estimate > 2_000_000 {
if self.config.chunk_size < 2048 {
self.config.chunk_size = 2048;
self.config.overlap_size = 512;
}
}
if condition.rtt_ms > 100.0 {
self.config.target_latency_ms = (condition.rtt_ms * 2.0) as u64;
}
if condition.jitter_ms > 20.0 {
self.config.max_buffer_chunks = (15.0 * condition.jitter_ms / 20.0) as usize;
}
Ok(())
}
pub fn get_advanced_stats(&self) -> EvaluationResult<AdvancedProcessingStats> {
let stats = self
.advanced_stats
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock advanced stats".to_string(),
source: None,
})?;
Ok(stats.clone())
}
pub fn get_recent_anomalies(&self, limit: usize) -> EvaluationResult<Vec<AnomalyDetection>> {
let detector =
self.anomaly_detector
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock anomaly detector".to_string(),
source: None,
})?;
Ok(detector
.anomaly_history
.iter()
.rev()
.take(limit)
.cloned()
.collect())
}
pub fn get_network_history(&self, limit: usize) -> EvaluationResult<Vec<NetworkCondition>> {
let monitor =
self.network_monitor
.lock()
.map_err(|_| EvaluationError::ProcessingError {
message: "Failed to lock network monitor".to_string(),
source: None,
})?;
Ok(monitor
.condition_history
.iter()
.rev()
.take(limit)
.cloned()
.collect())
}
}
pub fn chunk_audio_buffer(
audio: &AudioBuffer,
chunk_size: usize,
overlap_size: usize,
) -> Vec<AudioChunk> {
let samples = audio.samples();
let sample_rate = audio.sample_rate();
let mut chunks = Vec::new();
let step_size = chunk_size - overlap_size;
let mut sequence = 0;
let mut start = 0;
while start < samples.len() {
let end = (start + chunk_size).min(samples.len());
let chunk_samples = samples[start..end].to_vec();
if chunk_samples.len() >= chunk_size / 2 {
chunks.push(AudioChunk::new(chunk_samples, sample_rate, sequence));
sequence += 1;
}
start += step_size;
if start >= samples.len() {
break;
}
}
chunks
}
#[cfg(test)]
mod tests {
use super::*;
use voirs_sdk::AudioBuffer;
#[tokio::test]
async fn test_streaming_evaluator_creation() {
let config = StreamingConfig::default();
let evaluator = StreamingEvaluator::new(config);
let stats = evaluator.get_processing_stats().unwrap();
assert_eq!(stats.total_chunks_processed, 0);
}
#[tokio::test]
async fn test_chunk_processing() {
let config = StreamingConfig::default();
let mut evaluator = StreamingEvaluator::new(config);
let samples = vec![0.1, 0.2, -0.1, -0.2, 0.3, -0.3]; let chunk = AudioChunk::new(samples, 16000, 0);
let result = evaluator.process_chunk(chunk).await;
assert!(result.is_ok());
let metrics = evaluator.get_current_metrics().unwrap();
assert!(metrics.is_some());
let metrics = metrics.unwrap();
assert!(metrics.energy_level > 0.0);
assert!(!metrics.silence_detected);
}
#[tokio::test]
async fn test_audio_chunking() {
let samples = vec![0.1; 1000]; let audio = AudioBuffer::mono(samples, 16000);
let chunks = chunk_audio_buffer(&audio, 256, 64);
assert!(chunks.len() > 1);
assert_eq!(chunks[0].samples.len(), 256);
assert_eq!(chunks[0].sample_rate, 16000);
assert_eq!(chunks[0].sequence, 0);
}
#[tokio::test]
async fn test_quality_monitoring() {
let config = StreamingConfig {
enable_quality_monitoring: true,
..Default::default()
};
let mut evaluator = StreamingEvaluator::new(config);
let _receiver = evaluator.setup_quality_monitoring();
let samples = vec![0.5; 512]; let chunk = AudioChunk::new(samples, 16000, 0);
let result = evaluator.process_chunk(chunk).await;
assert!(result.is_ok());
let metrics = evaluator.get_current_metrics().unwrap().unwrap();
assert!(metrics.energy_level > 0.4);
assert!(!metrics.silence_detected);
assert!(!metrics.clipping_detected);
}
#[tokio::test]
async fn test_silence_detection() {
let config = StreamingConfig::default();
let mut evaluator = StreamingEvaluator::new(config);
let samples = vec![0.0001; 512]; let chunk = AudioChunk::new(samples, 16000, 0);
let result = evaluator.process_chunk(chunk).await;
assert!(result.is_ok());
let metrics = evaluator.get_current_metrics().unwrap().unwrap();
assert!(metrics.silence_detected);
assert!(metrics.energy_level < 0.001);
}
#[tokio::test]
async fn test_clipping_detection() {
let config = StreamingConfig::default();
let mut evaluator = StreamingEvaluator::new(config);
let samples = vec![0.98, -0.97, 0.99, -0.98]; let chunk = AudioChunk::new(samples, 16000, 0);
let result = evaluator.process_chunk(chunk).await;
assert!(result.is_ok());
let metrics = evaluator.get_current_metrics().unwrap().unwrap();
assert!(metrics.clipping_detected);
}
#[test]
fn test_streaming_config_defaults() {
let config = StreamingConfig::default();
assert_eq!(config.chunk_size, 1024);
assert_eq!(config.overlap_size, 256);
assert_eq!(config.max_buffer_chunks, 10);
assert_eq!(config.target_latency_ms, 100);
assert!(config.enable_quality_monitoring);
assert!(config.enable_adaptive_processing);
}
}