use crate::audio::{RingBuffer, SincResampler, SAMPLE_RATE};
use crate::vad::{VadConfig, VoiceActivityDetector};
#[cfg(not(feature = "std"))]
use alloc::vec::Vec;
pub const DEFAULT_CHUNK_DURATION: f32 = 30.0;
pub const DEFAULT_CHUNK_OVERLAP: f32 = 1.0;
pub const MIN_SPEECH_DURATION_MS: u32 = 500;
pub const LOW_LATENCY_CHUNK_DURATION: f32 = 0.5;
pub const LOW_LATENCY_CHUNK_OVERLAP: f32 = 0.05;
pub const LOW_LATENCY_MIN_SPEECH_MS: u32 = 100;
pub const LOW_LATENCY_PARTIAL_THRESHOLD: f32 = 0.25;
pub const LOW_LATENCY_BUFFER_DURATION: f32 = 5.0;
pub const LOW_LATENCY_FRAME_SIZE_MS: u32 = 10;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum LatencyMode {
#[default]
Standard,
LowLatency,
UltraLow,
Custom,
}
#[derive(Debug, Clone)]
pub struct StreamingConfig {
pub input_sample_rate: u32,
pub output_sample_rate: u32,
pub chunk_duration: f32,
pub chunk_overlap: f32,
pub enable_vad: bool,
pub vad_threshold: f32,
pub min_speech_duration_ms: u32,
pub buffer_duration: f32,
pub latency_mode: LatencyMode,
}
impl Default for StreamingConfig {
fn default() -> Self {
Self {
input_sample_rate: 44100,
output_sample_rate: SAMPLE_RATE,
chunk_duration: DEFAULT_CHUNK_DURATION,
chunk_overlap: DEFAULT_CHUNK_OVERLAP,
enable_vad: true,
vad_threshold: 0.5,
min_speech_duration_ms: MIN_SPEECH_DURATION_MS,
buffer_duration: 120.0, latency_mode: LatencyMode::Standard,
}
}
}
impl StreamingConfig {
#[must_use]
pub fn with_sample_rate(input_sample_rate: u32) -> Self {
Self {
input_sample_rate,
..Default::default()
}
}
#[must_use]
pub fn low_latency() -> Self {
Self {
input_sample_rate: 44100,
output_sample_rate: SAMPLE_RATE,
chunk_duration: LOW_LATENCY_CHUNK_DURATION,
chunk_overlap: LOW_LATENCY_CHUNK_OVERLAP,
enable_vad: true,
vad_threshold: 0.5,
min_speech_duration_ms: LOW_LATENCY_MIN_SPEECH_MS,
buffer_duration: LOW_LATENCY_BUFFER_DURATION,
latency_mode: LatencyMode::LowLatency,
}
}
#[must_use]
pub fn ultra_low_latency() -> Self {
Self {
input_sample_rate: 44100,
output_sample_rate: SAMPLE_RATE,
chunk_duration: 0.25, chunk_overlap: 0.025, enable_vad: true,
vad_threshold: 0.5,
min_speech_duration_ms: 50, buffer_duration: 2.0, latency_mode: LatencyMode::UltraLow,
}
}
#[must_use]
pub fn custom_latency(
chunk_duration: f32,
chunk_overlap: f32,
min_speech_duration_ms: u32,
buffer_duration: f32,
) -> Self {
Self {
input_sample_rate: 44100,
output_sample_rate: SAMPLE_RATE,
chunk_duration,
chunk_overlap,
enable_vad: true,
vad_threshold: 0.5,
min_speech_duration_ms,
buffer_duration,
latency_mode: LatencyMode::Custom,
}
}
#[must_use]
pub fn with_latency_mode(mut self, mode: LatencyMode) -> Self {
self.latency_mode = mode;
self
}
#[must_use]
pub const fn latency_mode(&self) -> LatencyMode {
self.latency_mode
}
#[must_use]
pub fn expected_latency_ms(&self) -> f32 {
self.chunk_duration * 1000.0
}
#[must_use]
pub const fn is_low_latency(&self) -> bool {
matches!(
self.latency_mode,
LatencyMode::LowLatency | LatencyMode::UltraLow
)
}
#[must_use]
pub fn with_vad(mut self) -> Self {
self.enable_vad = true;
self
}
#[must_use]
pub fn without_vad(mut self) -> Self {
self.enable_vad = false;
self
}
#[must_use]
pub fn vad_threshold(mut self, threshold: f32) -> Self {
self.vad_threshold = threshold;
self
}
#[must_use]
pub fn chunk_duration(mut self, duration: f32) -> Self {
self.chunk_duration = duration;
self
}
#[must_use]
pub fn chunk_overlap(mut self, overlap: f32) -> Self {
self.chunk_overlap = overlap;
self
}
#[must_use]
pub fn min_speech_duration_ms(mut self, duration: u32) -> Self {
self.min_speech_duration_ms = duration;
self
}
#[must_use]
pub fn chunk_samples(&self) -> usize {
(self.chunk_duration * self.output_sample_rate as f32) as usize
}
#[must_use]
pub fn overlap_samples(&self) -> usize {
(self.chunk_overlap * self.output_sample_rate as f32) as usize
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProcessorState {
WaitingForSpeech,
AccumulatingSpeech,
PartialResultReady,
ChunkReady,
Processing,
Error,
}
#[derive(Debug, Clone, PartialEq)]
pub enum StreamingEvent {
SpeechStart,
SpeechEnd,
PartialReady {
accumulated_samples: usize,
duration_secs: f32,
},
ChunkReady {
duration_secs: f32,
},
ProcessingStarted,
ProcessingCompleted,
Error(String),
Reset,
}
#[derive(Debug)]
pub struct StreamingProcessor {
config: StreamingConfig,
input_buffer: RingBuffer,
resampler: Option<SincResampler>,
vad: VoiceActivityDetector,
chunk_buffer: Vec<f32>,
overlap_buffer: Vec<f32>,
state: ProcessorState,
prev_state: ProcessorState,
speech_frames: u32,
silence_frames: u32,
samples_processed: u64,
events: Vec<StreamingEvent>,
partial_threshold_samples: usize,
last_partial_position: usize,
}
const DEFAULT_PARTIAL_THRESHOLD_SECS: f32 = 3.0;
impl StreamingProcessor {
#[must_use]
pub fn new(config: StreamingConfig) -> Self {
let input_buffer =
RingBuffer::for_duration(config.buffer_duration, config.input_sample_rate);
let resampler = if config.input_sample_rate == config.output_sample_rate {
None
} else {
SincResampler::new(config.input_sample_rate, config.output_sample_rate).ok()
};
let vad_config = VadConfig {
energy_threshold: config.vad_threshold * 4.0, ..VadConfig::default()
};
let vad = VoiceActivityDetector::new(vad_config);
let chunk_capacity = config.chunk_samples() + config.overlap_samples();
let partial_threshold_samples =
(DEFAULT_PARTIAL_THRESHOLD_SECS * config.output_sample_rate as f32) as usize;
Self {
config,
input_buffer,
resampler,
vad,
chunk_buffer: Vec::with_capacity(chunk_capacity),
overlap_buffer: Vec::new(),
state: ProcessorState::WaitingForSpeech,
prev_state: ProcessorState::WaitingForSpeech,
speech_frames: 0,
silence_frames: 0,
samples_processed: 0,
events: Vec::new(),
partial_threshold_samples,
last_partial_position: 0,
}
}
#[must_use]
pub fn with_sample_rate(sample_rate: u32) -> Self {
Self::new(StreamingConfig::with_sample_rate(sample_rate))
}
#[must_use]
pub const fn state(&self) -> ProcessorState {
self.state
}
#[must_use]
pub const fn samples_processed(&self) -> u64 {
self.samples_processed
}
#[must_use]
pub fn chunk_len(&self) -> usize {
self.chunk_buffer.len()
}
#[must_use]
pub fn chunk_progress(&self) -> f32 {
self.chunk_buffer.len() as f32 / self.config.chunk_samples() as f32
}
#[must_use]
pub fn has_chunk(&self) -> bool {
self.state == ProcessorState::ChunkReady
|| self.chunk_buffer.len() >= self.config.chunk_samples()
}
#[must_use]
pub fn overlap_len(&self) -> usize {
self.overlap_buffer.len()
}
#[must_use]
pub fn overlap_duration(&self) -> f32 {
self.overlap_buffer.len() as f32 / self.config.output_sample_rate as f32
}
#[must_use]
pub fn has_overlap(&self) -> bool {
!self.overlap_buffer.is_empty()
}
#[must_use]
pub fn configured_overlap_samples(&self) -> usize {
self.config.overlap_samples()
}
#[must_use]
pub fn configured_overlap_duration(&self) -> f32 {
self.config.chunk_overlap
}
pub fn clear_overlap(&mut self) {
self.overlap_buffer.clear();
}
#[must_use]
pub fn get_overlap_buffer(&self) -> Vec<f32> {
self.overlap_buffer.clone()
}
pub fn set_overlap_buffer(&mut self, overlap: Vec<f32>) {
self.overlap_buffer = overlap;
}
#[must_use]
pub fn has_events(&self) -> bool {
!self.events.is_empty()
}
#[must_use]
pub fn event_count(&self) -> usize {
self.events.len()
}
pub fn drain_events(&mut self) -> Vec<StreamingEvent> {
core::mem::take(&mut self.events)
}
pub fn pop_event(&mut self) -> Option<StreamingEvent> {
if self.events.is_empty() {
None
} else {
Some(self.events.remove(0))
}
}
#[must_use]
pub fn peek_event(&self) -> Option<&StreamingEvent> {
self.events.first()
}
pub fn clear_events(&mut self) {
self.events.clear();
}
#[must_use]
pub fn has_partial(&self) -> bool {
self.state == ProcessorState::PartialResultReady
|| (self.state == ProcessorState::AccumulatingSpeech
&& self.chunk_buffer.len() >= self.partial_threshold_samples
&& self.chunk_buffer.len() > self.last_partial_position)
}
pub fn get_partial(&mut self) -> Option<Vec<f32>> {
if !self.has_partial() {
return None;
}
self.last_partial_position = self.chunk_buffer.len();
Some(self.chunk_buffer.clone())
}
#[must_use]
pub fn partial_duration(&self) -> f32 {
self.chunk_buffer.len() as f32 / self.config.output_sample_rate as f32
}
pub fn set_partial_threshold(&mut self, seconds: f32) {
self.partial_threshold_samples = (seconds * self.config.output_sample_rate as f32) as usize;
}
#[must_use]
pub fn partial_threshold(&self) -> f32 {
self.partial_threshold_samples as f32 / self.config.output_sample_rate as f32
}
pub fn mark_processing_started(&mut self) {
if self.state == ProcessorState::ChunkReady
|| self.state == ProcessorState::PartialResultReady
{
self.prev_state = self.state;
self.state = ProcessorState::Processing;
self.events.push(StreamingEvent::ProcessingStarted);
}
}
pub fn mark_processing_completed(&mut self) {
if self.state == ProcessorState::Processing {
self.state = ProcessorState::WaitingForSpeech;
self.events.push(StreamingEvent::ProcessingCompleted);
}
}
pub fn mark_error(&mut self, message: &str) {
self.prev_state = self.state;
self.state = ProcessorState::Error;
self.events.push(StreamingEvent::Error(message.to_string()));
}
pub fn recover_from_error(&mut self) {
if self.state == ProcessorState::Error {
self.state = ProcessorState::WaitingForSpeech;
self.chunk_buffer.clear();
self.last_partial_position = 0;
}
}
#[must_use]
pub const fn prev_state(&self) -> ProcessorState {
self.prev_state
}
fn emit_event(&mut self, event: StreamingEvent) {
self.events.push(event);
}
pub fn push_audio(&mut self, samples: &[f32]) {
self.input_buffer.write_overwrite(samples);
self.samples_processed += samples.len() as u64;
}
pub fn process(&mut self) {
let available = self.input_buffer.available_read();
if available == 0 {
return;
}
let frame_size = (0.030 * self.config.input_sample_rate as f32) as usize;
let mut input_frame = vec![0.0; frame_size];
while self.input_buffer.available_read() >= frame_size {
let read = self.input_buffer.read(&mut input_frame);
if read < frame_size {
break;
}
let resampled = if let Some(ref resampler) = self.resampler {
match resampler.resample(&input_frame) {
Ok(samples) => samples,
Err(_) => continue,
}
} else {
input_frame.clone()
};
let is_speech = if self.config.enable_vad {
let event = self.vad.process_frame(&resampled);
matches!(
event,
crate::vad::VadEvent::SpeechStart | crate::vad::VadEvent::Continue
) && self.vad.state() == crate::vad::VadState::Speech
} else {
true
};
self.update_state(is_speech, &resampled);
}
}
fn update_state(&mut self, is_speech: bool, samples: &[f32]) {
let prev_speech_frames = self.speech_frames;
let prev_silence_frames = self.silence_frames;
if is_speech {
self.speech_frames += 1;
self.silence_frames = 0;
} else {
self.silence_frames += 1;
self.speech_frames = 0;
}
self.prev_state = self.state;
match self.state {
ProcessorState::WaitingForSpeech => {
if is_speech && self.speech_frames >= self.min_speech_frames() {
self.state = ProcessorState::AccumulatingSpeech;
self.chunk_buffer.clear();
self.chunk_buffer.extend(&self.overlap_buffer);
self.chunk_buffer.extend_from_slice(samples);
self.last_partial_position = 0;
self.emit_event(StreamingEvent::SpeechStart);
}
}
ProcessorState::AccumulatingSpeech => {
self.chunk_buffer.extend_from_slice(samples);
if self.chunk_buffer.len() >= self.partial_threshold_samples
&& self.chunk_buffer.len() > self.last_partial_position
&& self.last_partial_position == 0
{
self.emit_event(StreamingEvent::PartialReady {
accumulated_samples: self.chunk_buffer.len(),
duration_secs: self.partial_duration(),
});
}
if self.chunk_buffer.len() >= self.config.chunk_samples() {
self.state = ProcessorState::ChunkReady;
let duration =
self.chunk_buffer.len() as f32 / self.config.output_sample_rate as f32;
self.emit_event(StreamingEvent::ChunkReady {
duration_secs: duration,
});
}
else if !is_speech && self.silence_frames >= self.max_silence_frames() {
self.emit_event(StreamingEvent::SpeechEnd);
if self.chunk_buffer.len() >= self.config.overlap_samples() * 2 {
self.state = ProcessorState::ChunkReady;
let duration =
self.chunk_buffer.len() as f32 / self.config.output_sample_rate as f32;
self.emit_event(StreamingEvent::ChunkReady {
duration_secs: duration,
});
} else {
self.state = ProcessorState::WaitingForSpeech;
self.chunk_buffer.clear();
self.last_partial_position = 0;
}
}
}
ProcessorState::PartialResultReady => {
self.chunk_buffer.extend_from_slice(samples);
if self.chunk_buffer.len() >= self.config.chunk_samples() {
self.state = ProcessorState::ChunkReady;
let duration =
self.chunk_buffer.len() as f32 / self.config.output_sample_rate as f32;
self.emit_event(StreamingEvent::ChunkReady {
duration_secs: duration,
});
}
}
ProcessorState::ChunkReady | ProcessorState::Processing | ProcessorState::Error => {
}
}
let _ = prev_speech_frames;
let _ = prev_silence_frames;
}
fn min_speech_frames(&self) -> u32 {
let frame_duration_ms = 30;
self.config.min_speech_duration_ms / frame_duration_ms
}
fn max_silence_frames(&self) -> u32 {
let _ = self; let frame_duration_ms = 30;
1000 / frame_duration_ms
}
pub fn get_chunk(&mut self) -> Option<Vec<f32>> {
if !self.has_chunk() {
return None;
}
let overlap_size = self.config.overlap_samples();
if self.chunk_buffer.len() > overlap_size {
let start = self.chunk_buffer.len() - overlap_size;
self.overlap_buffer = self.chunk_buffer[start..].to_vec();
}
let target_size = self.config.chunk_samples();
if self.chunk_buffer.len() < target_size {
self.chunk_buffer.resize(target_size, 0.0);
}
let chunk = core::mem::take(&mut self.chunk_buffer);
self.prev_state = self.state;
self.state = ProcessorState::WaitingForSpeech;
self.speech_frames = 0;
self.silence_frames = 0;
self.last_partial_position = 0;
Some(chunk)
}
pub fn flush(&mut self) -> Option<Vec<f32>> {
if self.chunk_buffer.is_empty() {
return None;
}
self.process();
let min_size = self.config.overlap_samples() * 2;
if self.chunk_buffer.len() < min_size {
return None;
}
let target_size = self.config.chunk_samples();
self.chunk_buffer.resize(target_size, 0.0);
let chunk = core::mem::take(&mut self.chunk_buffer);
self.prev_state = self.state;
self.state = ProcessorState::WaitingForSpeech;
self.overlap_buffer.clear();
self.last_partial_position = 0;
Some(chunk)
}
pub fn reset(&mut self) {
self.input_buffer.clear();
self.chunk_buffer.clear();
self.overlap_buffer.clear();
self.prev_state = self.state;
self.state = ProcessorState::WaitingForSpeech;
self.speech_frames = 0;
self.silence_frames = 0;
self.last_partial_position = 0;
self.emit_event(StreamingEvent::Reset);
}
#[must_use]
pub fn stats(&self) -> ProcessorStats {
ProcessorStats {
samples_processed: self.samples_processed,
buffer_available: self.input_buffer.available_read(),
buffer_capacity: self.input_buffer.capacity(),
chunk_progress: self.chunk_progress(),
state: self.state,
}
}
}
#[derive(Debug, Clone)]
pub struct ProcessorStats {
pub samples_processed: u64,
pub buffer_available: usize,
pub buffer_capacity: usize,
pub chunk_progress: f32,
pub state: ProcessorState,
}
impl ProcessorStats {
#[must_use]
pub fn buffer_fill(&self) -> f32 {
if self.buffer_capacity == 0 {
0.0
} else {
self.buffer_available as f32 / self.buffer_capacity as f32 * 100.0
}
}
#[must_use]
pub fn duration_processed(&self, sample_rate: u32) -> f32 {
self.samples_processed as f32 / sample_rate as f32
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_config_default() {
let config = StreamingConfig::default();
assert_eq!(config.input_sample_rate, 44100);
assert_eq!(config.output_sample_rate, 16000);
assert!((config.chunk_duration - 30.0).abs() < f32::EPSILON);
assert!(config.enable_vad);
}
#[test]
fn test_config_with_sample_rate() {
let config = StreamingConfig::with_sample_rate(48000);
assert_eq!(config.input_sample_rate, 48000);
}
#[test]
fn test_config_without_vad() {
let config = StreamingConfig::default().without_vad();
assert!(!config.enable_vad);
}
#[test]
fn test_config_chunk_samples() {
let config = StreamingConfig::default();
assert_eq!(config.chunk_samples(), 480000);
}
#[test]
fn test_config_overlap_samples() {
let config = StreamingConfig::default();
assert_eq!(config.overlap_samples(), 16000);
}
#[test]
fn test_processor_new() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert_eq!(processor.samples_processed(), 0);
assert_eq!(processor.chunk_len(), 0);
}
#[test]
fn test_processor_with_sample_rate() {
let processor = StreamingProcessor::with_sample_rate(48000);
assert_eq!(processor.config.input_sample_rate, 48000);
}
#[test]
fn test_processor_same_sample_rate() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
..Default::default()
};
let processor = StreamingProcessor::new(config);
assert!(processor.resampler.is_none());
}
#[test]
fn test_processor_different_sample_rate() {
let config = StreamingConfig::default(); let processor = StreamingProcessor::new(config);
assert!(processor.resampler.is_some());
}
#[test]
fn test_push_audio() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
let samples = vec![0.0; 1000];
processor.push_audio(&samples);
assert_eq!(processor.samples_processed(), 1000);
}
#[test]
fn test_push_audio_multiple() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.push_audio(&vec![0.0; 500]);
processor.push_audio(&vec![0.0; 500]);
assert_eq!(processor.samples_processed(), 1000);
}
#[test]
fn test_initial_state() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert!(!processor.has_chunk());
}
#[test]
fn test_chunk_progress_empty() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert!((processor.chunk_progress() - 0.0).abs() < f32::EPSILON);
}
#[test]
fn test_reset() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.push_audio(&vec![0.1; 10000]);
processor.reset();
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert_eq!(processor.chunk_len(), 0);
}
#[test]
fn test_stats() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.push_audio(&vec![0.0; 1000]);
let stats = processor.stats();
assert_eq!(stats.samples_processed, 1000);
assert_eq!(stats.state, ProcessorState::WaitingForSpeech);
}
#[test]
fn test_stats_buffer_fill() {
let stats = ProcessorStats {
samples_processed: 0,
buffer_available: 500,
buffer_capacity: 1000,
chunk_progress: 0.0,
state: ProcessorState::WaitingForSpeech,
};
assert!((stats.buffer_fill() - 50.0).abs() < f32::EPSILON);
}
#[test]
fn test_stats_duration_processed() {
let stats = ProcessorStats {
samples_processed: 16000,
buffer_available: 0,
buffer_capacity: 1000,
chunk_progress: 0.0,
state: ProcessorState::WaitingForSpeech,
};
assert!((stats.duration_processed(16000) - 1.0).abs() < f32::EPSILON);
}
#[test]
fn test_vad_disabled() {
let config = StreamingConfig::default().without_vad();
let processor = StreamingProcessor::new(config);
assert!(!processor.config.enable_vad);
}
#[test]
fn test_flush_empty() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
assert!(processor.flush().is_none());
}
#[test]
fn test_get_chunk_not_ready() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
assert!(processor.get_chunk().is_none());
}
#[test]
fn test_process_silence() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
enable_vad: true,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
let silence = vec![0.0; 4800]; processor.push_audio(&silence);
processor.process();
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
}
#[test]
fn test_process_with_vad_disabled() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
enable_vad: false,
chunk_duration: 0.5, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
let audio = vec![0.1; 8000]; processor.push_audio(&audio);
processor.process();
assert!(processor.chunk_len() > 0 || processor.has_chunk());
}
#[test]
fn test_process_empty_buffer() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.process();
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert_eq!(processor.chunk_len(), 0);
}
#[test]
fn test_config_vad_threshold() {
let config = StreamingConfig::default().vad_threshold(0.8);
assert!((config.vad_threshold - 0.8).abs() < f32::EPSILON);
}
#[test]
fn test_config_chunk_duration() {
let config = StreamingConfig::default().chunk_duration(10.0);
assert!((config.chunk_duration - 10.0).abs() < f32::EPSILON);
}
#[test]
fn test_stats_buffer_fill_zero_capacity() {
let stats = ProcessorStats {
samples_processed: 0,
buffer_available: 0,
buffer_capacity: 0,
chunk_progress: 0.0,
state: ProcessorState::WaitingForSpeech,
};
assert!((stats.buffer_fill() - 0.0).abs() < f32::EPSILON);
}
#[test]
fn test_update_state_accumulating() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
enable_vad: false,
chunk_duration: 0.1, min_speech_duration_ms: 0, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
let samples = vec![0.1; 320]; processor.state = ProcessorState::AccumulatingSpeech;
processor.update_state(true, &samples);
assert_eq!(processor.state(), ProcessorState::AccumulatingSpeech);
assert_eq!(processor.chunk_buffer.len(), 320);
}
#[test]
fn test_update_state_chunk_ready() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
enable_vad: false,
chunk_duration: 0.02, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.chunk_buffer = vec![0.1; 320];
processor.state = ProcessorState::AccumulatingSpeech;
let samples = vec![0.1; 320];
processor.update_state(true, &samples);
assert_eq!(processor.state(), ProcessorState::ChunkReady);
}
#[test]
fn test_update_state_silence_ends_speech() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 30.0,
chunk_overlap: 0.1, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 5000]; processor.silence_frames = 50;
let samples = vec![0.0; 320];
processor.update_state(false, &samples);
assert_eq!(processor.state(), ProcessorState::ChunkReady);
}
#[test]
fn test_update_state_short_segment_discarded() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 30.0,
chunk_overlap: 1.0, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 100]; processor.silence_frames = 50;
let samples = vec![0.0; 320];
processor.update_state(false, &samples);
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert!(processor.chunk_buffer.is_empty());
}
#[test]
fn test_update_state_chunk_ready_stays() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::ChunkReady;
let samples = vec![0.1; 320];
processor.update_state(true, &samples);
assert_eq!(processor.state(), ProcessorState::ChunkReady);
}
#[test]
fn test_get_chunk_ready() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.02, chunk_overlap: 0.01, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.chunk_buffer = vec![0.1; 400];
processor.state = ProcessorState::ChunkReady;
let chunk = processor.get_chunk();
assert!(chunk.is_some());
let chunk = chunk.expect("chunk should exist");
assert_eq!(chunk.len(), 400);
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
}
#[test]
fn test_get_chunk_saves_overlap() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.02, chunk_overlap: 0.01, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.chunk_buffer = vec![0.1; 400];
processor.state = ProcessorState::ChunkReady;
let _ = processor.get_chunk();
assert_eq!(processor.overlap_buffer.len(), 160);
}
#[test]
fn test_get_chunk_pads_short_chunk() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.02, chunk_overlap: 0.0,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.chunk_buffer = vec![0.1; 200];
processor.state = ProcessorState::ChunkReady;
let chunk = processor.get_chunk().expect("should get chunk");
assert_eq!(chunk.len(), 320);
}
#[test]
fn test_flush_with_accumulated_data() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.1, chunk_overlap: 0.01, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.chunk_buffer = vec![0.1; 500]; processor.state = ProcessorState::AccumulatingSpeech;
let flushed = processor.flush();
assert!(flushed.is_some());
let flushed = flushed.expect("should flush");
assert_eq!(flushed.len(), 1600);
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert!(processor.overlap_buffer.is_empty());
}
#[test]
fn test_flush_too_short() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.1,
chunk_overlap: 0.05, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.chunk_buffer = vec![0.1; 100];
let flushed = processor.flush();
assert!(flushed.is_none());
}
#[test]
fn test_min_speech_frames() {
let config = StreamingConfig {
min_speech_duration_ms: 300,
..Default::default()
};
let processor = StreamingProcessor::new(config);
assert_eq!(processor.min_speech_frames(), 10);
}
#[test]
fn test_max_silence_frames() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert_eq!(processor.max_silence_frames(), 33);
}
#[test]
fn test_waiting_transitions_to_accumulating() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
min_speech_duration_ms: 0, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.overlap_buffer = vec![0.05; 100];
processor.speech_frames = 1;
let samples = vec![0.1; 320];
processor.update_state(true, &samples);
assert_eq!(processor.state(), ProcessorState::AccumulatingSpeech);
assert!(processor.chunk_buffer.len() >= 320);
}
#[test]
fn test_streaming_event_variants() {
let speech_start = StreamingEvent::SpeechStart;
let speech_end = StreamingEvent::SpeechEnd;
let partial_ready = StreamingEvent::PartialReady {
accumulated_samples: 48000,
duration_secs: 3.0,
};
let chunk_ready = StreamingEvent::ChunkReady {
duration_secs: 30.0,
};
let processing_started = StreamingEvent::ProcessingStarted;
let processing_completed = StreamingEvent::ProcessingCompleted;
let error = StreamingEvent::Error("test error".to_string());
let reset = StreamingEvent::Reset;
assert!(format!("{speech_start:?}").contains("SpeechStart"));
assert!(format!("{speech_end:?}").contains("SpeechEnd"));
assert!(format!("{partial_ready:?}").contains("PartialReady"));
assert!(format!("{chunk_ready:?}").contains("ChunkReady"));
assert!(format!("{processing_started:?}").contains("ProcessingStarted"));
assert!(format!("{processing_completed:?}").contains("ProcessingCompleted"));
assert!(format!("{error:?}").contains("Error"));
assert!(format!("{reset:?}").contains("Reset"));
let cloned = speech_start.clone();
assert_eq!(cloned, StreamingEvent::SpeechStart);
}
#[test]
fn test_processor_state_new_variants() {
let partial_ready = ProcessorState::PartialResultReady;
let processing = ProcessorState::Processing;
let error = ProcessorState::Error;
assert!(format!("{partial_ready:?}").contains("PartialResultReady"));
assert!(format!("{processing:?}").contains("Processing"));
assert!(format!("{error:?}").contains("Error"));
assert_ne!(
ProcessorState::PartialResultReady,
ProcessorState::Processing
);
assert_ne!(ProcessorState::Processing, ProcessorState::Error);
}
#[test]
fn test_event_handling_initial() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert!(!processor.has_events());
assert_eq!(processor.event_count(), 0);
assert!(processor.peek_event().is_none());
}
#[test]
fn test_event_pop_and_drain() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.events.push(StreamingEvent::SpeechStart);
processor.events.push(StreamingEvent::SpeechEnd);
assert!(processor.has_events());
assert_eq!(processor.event_count(), 2);
let event = processor.pop_event();
assert!(event.is_some());
assert_eq!(event, Some(StreamingEvent::SpeechStart));
assert_eq!(processor.event_count(), 1);
let remaining = processor.drain_events();
assert_eq!(remaining.len(), 1);
assert_eq!(remaining[0], StreamingEvent::SpeechEnd);
assert!(!processor.has_events());
}
#[test]
fn test_event_peek() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.events.push(StreamingEvent::Reset);
let peeked = processor.peek_event();
assert!(peeked.is_some());
assert_eq!(peeked, Some(&StreamingEvent::Reset));
assert_eq!(processor.event_count(), 1);
}
#[test]
fn test_clear_events() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.events.push(StreamingEvent::SpeechStart);
processor.events.push(StreamingEvent::SpeechEnd);
processor.clear_events();
assert!(!processor.has_events());
assert_eq!(processor.event_count(), 0);
}
#[test]
fn test_partial_result_threshold() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
assert!((processor.partial_threshold() - 3.0).abs() < 0.01);
processor.set_partial_threshold(5.0);
assert!((processor.partial_threshold() - 5.0).abs() < 0.01);
}
#[test]
fn test_partial_duration() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
assert!((processor.partial_duration() - 0.0).abs() < 0.01);
processor.chunk_buffer = vec![0.1; 16000]; assert!((processor.partial_duration() - 1.0).abs() < 0.01);
}
#[test]
fn test_has_partial_not_accumulating() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert!(!processor.has_partial()); }
#[test]
fn test_has_partial_below_threshold() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 16000];
assert!(!processor.has_partial());
}
#[test]
fn test_has_partial_above_threshold() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 64000];
assert!(processor.has_partial());
}
#[test]
fn test_get_partial() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 64000];
let partial = processor.get_partial();
assert!(partial.is_some());
assert_eq!(partial.as_ref().map(|p| p.len()), Some(64000));
assert_eq!(processor.last_partial_position, 64000);
assert!(!processor.has_partial());
}
#[test]
fn test_processing_state_transitions() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::ChunkReady;
processor.mark_processing_started();
assert_eq!(processor.state(), ProcessorState::Processing);
assert!(processor.has_events());
let event = processor.pop_event();
assert_eq!(event, Some(StreamingEvent::ProcessingStarted));
processor.mark_processing_completed();
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
let event = processor.pop_event();
assert_eq!(event, Some(StreamingEvent::ProcessingCompleted));
}
#[test]
fn test_error_state() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.mark_error("Test error message");
assert_eq!(processor.state(), ProcessorState::Error);
let event = processor.pop_event();
assert!(matches!(event, Some(StreamingEvent::Error(msg)) if msg == "Test error message"));
}
#[test]
fn test_error_recovery() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::Error;
processor.chunk_buffer = vec![0.1; 1000];
processor.last_partial_position = 500;
processor.recover_from_error();
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert!(processor.chunk_buffer.is_empty());
assert_eq!(processor.last_partial_position, 0);
}
#[test]
fn test_prev_state_tracking() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
min_speech_duration_ms: 0,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
assert_eq!(processor.prev_state(), ProcessorState::WaitingForSpeech);
processor.speech_frames = 1;
processor.update_state(true, &vec![0.1; 320]);
assert_eq!(processor.prev_state(), ProcessorState::WaitingForSpeech);
assert_eq!(processor.state(), ProcessorState::AccumulatingSpeech);
}
#[test]
fn test_speech_start_event() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
min_speech_duration_ms: 0, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.speech_frames = 1;
processor.update_state(true, &vec![0.1; 320]);
assert!(processor.has_events());
let event = processor.pop_event();
assert_eq!(event, Some(StreamingEvent::SpeechStart));
}
#[test]
fn test_speech_end_and_chunk_ready_events() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 30.0,
chunk_overlap: 0.1,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 5000];
processor.silence_frames = 50;
processor.update_state(false, &vec![0.0; 320]);
let events = processor.drain_events();
assert!(events.iter().any(|e| *e == StreamingEvent::SpeechEnd));
assert!(events
.iter()
.any(|e| matches!(e, StreamingEvent::ChunkReady { .. })));
}
#[test]
fn test_reset_emits_event() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 1000];
processor.reset();
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert!(processor.chunk_buffer.is_empty());
let event = processor.pop_event();
assert_eq!(event, Some(StreamingEvent::Reset));
}
#[test]
fn test_partial_ready_event_on_threshold() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
min_speech_duration_ms: 0,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.set_partial_threshold(0.1);
processor.speech_frames = 1;
processor.update_state(true, &vec![0.1; 320]);
processor.drain_events();
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 1500];
processor.update_state(true, &vec![0.1; 320]);
let events = processor.drain_events();
assert!(events
.iter()
.any(|e| matches!(e, StreamingEvent::PartialReady { .. })));
}
#[test]
fn test_processing_state_ignores_transitions() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::Processing;
processor.chunk_buffer = vec![0.1; 1000];
processor.update_state(true, &vec![0.1; 320]);
assert_eq!(processor.state(), ProcessorState::Processing);
}
#[test]
fn test_error_state_ignores_transitions() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.state = ProcessorState::Error;
processor.update_state(true, &vec![0.1; 320]);
assert_eq!(processor.state(), ProcessorState::Error);
}
#[test]
fn test_partial_result_ready_continues_accumulating() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.1, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.state = ProcessorState::PartialResultReady;
processor.chunk_buffer = vec![0.1; 1000];
processor.update_state(true, &vec![0.1; 320]);
assert_eq!(processor.chunk_buffer.len(), 1320);
}
#[test]
fn test_partial_to_chunk_ready_transition() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.1, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.state = ProcessorState::PartialResultReady;
processor.chunk_buffer = vec![0.1; 1500];
processor.update_state(true, &vec![0.1; 320]);
assert_eq!(processor.state(), ProcessorState::ChunkReady);
}
#[test]
fn test_get_chunk_resets_partial_position() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.02,
chunk_overlap: 0.01,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.state = ProcessorState::ChunkReady;
processor.chunk_buffer = vec![0.1; 400];
processor.last_partial_position = 200;
let _ = processor.get_chunk();
assert_eq!(processor.last_partial_position, 0);
}
#[test]
fn test_flush_resets_partial_position() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.1,
chunk_overlap: 0.01,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 500];
processor.last_partial_position = 100;
let _ = processor.flush();
assert_eq!(processor.last_partial_position, 0);
}
#[test]
fn test_default_partial_threshold_constant() {
assert!((DEFAULT_PARTIAL_THRESHOLD_SECS - 3.0).abs() < f32::EPSILON);
}
#[test]
fn test_config_chunk_overlap_builder() {
let config = StreamingConfig::default().chunk_overlap(0.5);
assert!((config.chunk_overlap - 0.5).abs() < 0.01);
}
#[test]
fn test_config_min_speech_duration_builder() {
let config = StreamingConfig::default().min_speech_duration_ms(500);
assert_eq!(config.min_speech_duration_ms, 500);
}
#[test]
fn test_overlap_len_initial() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert_eq!(processor.overlap_len(), 0);
}
#[test]
fn test_overlap_duration_empty() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert!((processor.overlap_duration() - 0.0).abs() < 0.01);
}
#[test]
fn test_has_overlap_initial() {
let processor = StreamingProcessor::new(StreamingConfig::default());
assert!(!processor.has_overlap());
}
#[test]
fn test_configured_overlap_samples() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_overlap: 0.5, ..Default::default()
};
let processor = StreamingProcessor::new(config);
assert_eq!(processor.configured_overlap_samples(), 8000); }
#[test]
fn test_configured_overlap_duration() {
let config = StreamingConfig {
chunk_overlap: 0.75,
..Default::default()
};
let processor = StreamingProcessor::new(config);
assert!((processor.configured_overlap_duration() - 0.75).abs() < 0.01);
}
#[test]
fn test_clear_overlap() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.overlap_buffer = vec![0.1; 1000];
processor.clear_overlap();
assert!(!processor.has_overlap());
assert_eq!(processor.overlap_len(), 0);
}
#[test]
fn test_get_overlap_buffer() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.overlap_buffer = vec![0.5; 100];
let overlap = processor.get_overlap_buffer();
assert_eq!(overlap.len(), 100);
assert!((overlap[0] - 0.5).abs() < 0.01);
}
#[test]
fn test_set_overlap_buffer() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
let custom_overlap = vec![0.25; 200];
processor.set_overlap_buffer(custom_overlap.clone());
assert!(processor.has_overlap());
assert_eq!(processor.overlap_len(), 200);
let buffer = processor.get_overlap_buffer();
assert!((buffer[0] - 0.25).abs() < 0.01);
}
#[test]
fn test_overlap_preserved_after_get_chunk() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.02, chunk_overlap: 0.01, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.chunk_buffer = (0..400).map(|i| i as f32 * 0.001).collect();
processor.state = ProcessorState::ChunkReady;
let _ = processor.get_chunk();
assert!(processor.has_overlap());
assert_eq!(processor.overlap_len(), 160);
let overlap = processor.get_overlap_buffer();
assert!((overlap[0] - 0.240).abs() < 0.001);
}
#[test]
fn test_overlap_used_when_starting_accumulation() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
min_speech_duration_ms: 0, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.overlap_buffer = vec![0.05; 100];
processor.speech_frames = 1;
let samples = vec![0.1; 320];
processor.update_state(true, &samples);
assert!(processor.chunk_buffer.len() >= 420);
assert!((processor.chunk_buffer[0] - 0.05).abs() < 0.01);
}
#[test]
fn test_overlap_duration_with_samples() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.overlap_buffer = vec![0.1; 8000];
assert!((processor.overlap_duration() - 0.5).abs() < 0.01);
}
#[test]
fn test_overlap_reset_clears_buffer() {
let mut processor = StreamingProcessor::new(StreamingConfig::default());
processor.overlap_buffer = vec![0.1; 1000];
processor.reset();
assert!(!processor.has_overlap());
}
#[test]
fn test_flush_clears_overlap() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.1,
chunk_overlap: 0.01,
..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.state = ProcessorState::AccumulatingSpeech;
processor.chunk_buffer = vec![0.1; 500];
processor.overlap_buffer = vec![0.2; 100];
let _ = processor.flush();
assert!(!processor.has_overlap());
}
#[test]
fn test_multiple_chunks_preserve_overlap_chain() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.02, chunk_overlap: 0.005, min_speech_duration_ms: 0, ..Default::default()
};
let mut processor = StreamingProcessor::new(config);
processor.chunk_buffer = vec![0.1; 400];
processor.state = ProcessorState::ChunkReady;
let _ = processor.get_chunk();
let first_overlap = processor.get_overlap_buffer();
assert_eq!(first_overlap.len(), 80);
processor.speech_frames = 1;
processor.state = ProcessorState::WaitingForSpeech;
processor.update_state(true, &vec![0.2; 320]);
assert!(processor.chunk_buffer.len() >= 320);
assert!((processor.chunk_buffer[0] - 0.1).abs() < 0.01);
}
#[test]
fn test_latency_mode_enum() {
let standard = LatencyMode::Standard;
let low = LatencyMode::LowLatency;
let ultra = LatencyMode::UltraLow;
let custom = LatencyMode::Custom;
assert_eq!(standard, LatencyMode::Standard);
assert_eq!(low, LatencyMode::LowLatency);
assert_eq!(ultra, LatencyMode::UltraLow);
assert_eq!(custom, LatencyMode::Custom);
assert!(format!("{standard:?}").contains("Standard"));
assert!(format!("{low:?}").contains("LowLatency"));
assert!(format!("{ultra:?}").contains("UltraLow"));
assert!(format!("{custom:?}").contains("Custom"));
let copied = standard;
assert_eq!(copied, LatencyMode::Standard);
}
#[test]
fn test_latency_mode_default() {
let mode = LatencyMode::default();
assert_eq!(mode, LatencyMode::Standard);
}
#[test]
fn test_low_latency_constants() {
assert!((LOW_LATENCY_CHUNK_DURATION - 0.5).abs() < f32::EPSILON);
assert!((LOW_LATENCY_CHUNK_OVERLAP - 0.05).abs() < f32::EPSILON);
assert_eq!(LOW_LATENCY_MIN_SPEECH_MS, 100);
assert!((LOW_LATENCY_PARTIAL_THRESHOLD - 0.25).abs() < f32::EPSILON);
assert!((LOW_LATENCY_BUFFER_DURATION - 5.0).abs() < f32::EPSILON);
assert_eq!(LOW_LATENCY_FRAME_SIZE_MS, 10);
}
#[test]
fn test_streaming_config_default_latency_mode() {
let config = StreamingConfig::default();
assert_eq!(config.latency_mode, LatencyMode::Standard);
}
#[test]
fn test_streaming_config_low_latency() {
let config = StreamingConfig::low_latency();
assert_eq!(config.latency_mode, LatencyMode::LowLatency);
assert!((config.chunk_duration - 0.5).abs() < f32::EPSILON);
assert!((config.chunk_overlap - 0.05).abs() < f32::EPSILON);
assert_eq!(config.min_speech_duration_ms, 100);
assert!((config.buffer_duration - 5.0).abs() < f32::EPSILON);
assert!(config.enable_vad);
}
#[test]
fn test_streaming_config_ultra_low_latency() {
let config = StreamingConfig::ultra_low_latency();
assert_eq!(config.latency_mode, LatencyMode::UltraLow);
assert!((config.chunk_duration - 0.25).abs() < f32::EPSILON);
assert!((config.chunk_overlap - 0.025).abs() < f32::EPSILON);
assert_eq!(config.min_speech_duration_ms, 50);
assert!((config.buffer_duration - 2.0).abs() < f32::EPSILON);
}
#[test]
fn test_streaming_config_custom_latency() {
let config = StreamingConfig::custom_latency(0.75, 0.1, 200, 10.0);
assert_eq!(config.latency_mode, LatencyMode::Custom);
assert!((config.chunk_duration - 0.75).abs() < f32::EPSILON);
assert!((config.chunk_overlap - 0.1).abs() < f32::EPSILON);
assert_eq!(config.min_speech_duration_ms, 200);
assert!((config.buffer_duration - 10.0).abs() < f32::EPSILON);
}
#[test]
fn test_streaming_config_with_latency_mode() {
let config = StreamingConfig::default().with_latency_mode(LatencyMode::Custom);
assert_eq!(config.latency_mode, LatencyMode::Custom);
}
#[test]
fn test_streaming_config_latency_mode_getter() {
let config = StreamingConfig::low_latency();
assert_eq!(config.latency_mode(), LatencyMode::LowLatency);
}
#[test]
fn test_expected_latency_ms_standard() {
let config = StreamingConfig::default();
assert!((config.expected_latency_ms() - 30000.0).abs() < 1.0);
}
#[test]
fn test_expected_latency_ms_low_latency() {
let config = StreamingConfig::low_latency();
assert!((config.expected_latency_ms() - 500.0).abs() < 1.0);
}
#[test]
fn test_expected_latency_ms_ultra_low() {
let config = StreamingConfig::ultra_low_latency();
assert!((config.expected_latency_ms() - 250.0).abs() < 1.0);
}
#[test]
fn test_is_low_latency_standard() {
let config = StreamingConfig::default();
assert!(!config.is_low_latency());
}
#[test]
fn test_is_low_latency_low() {
let config = StreamingConfig::low_latency();
assert!(config.is_low_latency());
}
#[test]
fn test_is_low_latency_ultra() {
let config = StreamingConfig::ultra_low_latency();
assert!(config.is_low_latency());
}
#[test]
fn test_is_low_latency_custom() {
let config = StreamingConfig::custom_latency(0.5, 0.05, 100, 5.0);
assert!(!config.is_low_latency()); }
#[test]
fn test_low_latency_chunk_samples() {
let config = StreamingConfig::low_latency();
assert_eq!(config.chunk_samples(), 8000);
}
#[test]
fn test_ultra_low_latency_chunk_samples() {
let config = StreamingConfig::ultra_low_latency();
assert_eq!(config.chunk_samples(), 4000);
}
#[test]
fn test_low_latency_overlap_samples() {
let config = StreamingConfig::low_latency();
assert_eq!(config.overlap_samples(), 800);
}
#[test]
fn test_processor_with_low_latency_config() {
let config = StreamingConfig::low_latency();
let processor = StreamingProcessor::new(config);
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
}
#[test]
fn test_processor_low_latency_min_speech_frames() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
min_speech_duration_ms: 100, ..StreamingConfig::low_latency()
};
let processor = StreamingProcessor::new(config);
assert_eq!(processor.min_speech_frames(), 3);
}
#[test]
fn test_low_latency_end_to_end() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
chunk_duration: 0.5,
chunk_overlap: 0.05,
enable_vad: false, min_speech_duration_ms: 0,
..StreamingConfig::low_latency()
};
let mut processor = StreamingProcessor::new(config);
let audio = vec![0.1; 8000];
processor.push_audio(&audio);
processor.process();
assert!(processor.has_chunk() || processor.chunk_len() > 0);
}
#[test]
fn test_ultra_low_latency_end_to_end() {
let config = StreamingConfig {
input_sample_rate: 16000,
output_sample_rate: 16000,
enable_vad: false,
min_speech_duration_ms: 0,
..StreamingConfig::ultra_low_latency()
};
let mut processor = StreamingProcessor::new(config);
let audio = vec![0.1; 4000];
processor.push_audio(&audio);
processor.process();
assert!(processor.has_chunk() || processor.chunk_len() > 0);
}
#[test]
fn test_low_latency_builder_chain() {
let config = StreamingConfig::low_latency()
.without_vad()
.vad_threshold(0.3);
assert_eq!(config.latency_mode, LatencyMode::LowLatency);
assert!(!config.enable_vad);
assert!((config.vad_threshold - 0.3).abs() < f32::EPSILON);
}
#[test]
fn test_latency_mode_equality() {
assert_eq!(LatencyMode::Standard, LatencyMode::Standard);
assert_ne!(LatencyMode::Standard, LatencyMode::LowLatency);
assert_ne!(LatencyMode::LowLatency, LatencyMode::UltraLow);
assert_ne!(LatencyMode::UltraLow, LatencyMode::Custom);
}
#[test]
fn test_low_latency_config_sample_rates() {
let config = StreamingConfig::low_latency();
assert_eq!(config.input_sample_rate, 44100);
assert_eq!(config.output_sample_rate, SAMPLE_RATE); }
#[test]
fn test_processor_stats_with_low_latency() {
let config = StreamingConfig::low_latency();
let mut processor = StreamingProcessor::new(config);
processor.push_audio(&vec![0.1; 1000]);
let stats = processor.stats();
assert_eq!(stats.samples_processed, 1000);
assert_eq!(stats.state, ProcessorState::WaitingForSpeech);
}
#[test]
fn test_low_latency_reset() {
let config = StreamingConfig::low_latency();
let mut processor = StreamingProcessor::new(config);
processor.push_audio(&vec![0.1; 5000]);
processor.process();
processor.reset();
assert_eq!(processor.state(), ProcessorState::WaitingForSpeech);
assert_eq!(processor.chunk_len(), 0);
assert!(!processor.has_overlap());
}
}