| | |
| | | )) |
| | | } |
| | | |
| | | #[derive(Debug, Clone, Default)] |
| | | pub struct PcmS16leStreamChunkResult { |
| | | pub frames: Vec<PcmFrame>, |
| | | pub dropped_tail_bytes: usize, |
| | | pub buffered_source_bytes: usize, |
| | | pub buffered_source_samples: usize, |
| | | } |
| | | |
| | | #[derive(Debug, Clone, Default)] |
| | | pub struct PcmS16leStreamDebugDump { |
| | | pub debug_source_path: Option<String>, |
| | | pub debug_pcm_wav_path: Option<String>, |
| | | pub debug_pcm_wav_size_bytes: Option<u64>, |
| | | } |
| | | |
| | | #[derive(Debug, Clone)] |
| | | pub struct PcmS16leStreamDecoder { |
| | | source_sample_rate_hz: u32, |
| | | source_num_channels: u16, |
| | | target_sample_rate_hz: u32, |
| | | target_num_channels: u16, |
| | | source_byte_buffer: Vec<u8>, |
| | | source_sample_buffer: Vec<i16>, |
| | | debug_source_samples: Vec<i16>, |
| | | debug_target_samples: Vec<i16>, |
| | | } |
| | | |
| | | impl PcmS16leStreamDecoder { |
| | | pub fn new( |
| | | source_sample_rate_hz: u32, |
| | | source_num_channels: u32, |
| | | target_sample_rate_hz: u32, |
| | | target_num_channels: u16, |
| | | ) -> Result<Self> { |
| | | if source_sample_rate_hz == 0 { |
| | | return Err(anyhow!("pcm_s16le source sample rate is zero")); |
| | | } |
| | | if source_num_channels == 0 { |
| | | return Err(anyhow!("pcm_s16le source channel count is zero")); |
| | | } |
| | | let source_num_channels = u16::try_from(source_num_channels) |
| | | .map_err(|_| anyhow!("unsupported pcm_s16le channel count {source_num_channels}"))?; |
| | | Ok(Self { |
| | | source_sample_rate_hz, |
| | | source_num_channels, |
| | | target_sample_rate_hz, |
| | | target_num_channels, |
| | | source_byte_buffer: Vec::new(), |
| | | source_sample_buffer: Vec::new(), |
| | | debug_source_samples: Vec::new(), |
| | | debug_target_samples: Vec::new(), |
| | | }) |
| | | } |
| | | |
| | | pub fn push_bytes( |
| | | &mut self, |
| | | pcm_bytes: &[u8], |
| | | source_sample_rate_hz: u32, |
| | | source_num_channels: u32, |
| | | last: bool, |
| | | ) -> Result<PcmS16leStreamChunkResult> { |
| | | let source_num_channels = u16::try_from(source_num_channels) |
| | | .map_err(|_| anyhow!("unsupported pcm_s16le channel count {source_num_channels}"))?; |
| | | if source_sample_rate_hz != self.source_sample_rate_hz |
| | | || source_num_channels != self.source_num_channels |
| | | { |
| | | return Err(anyhow!( |
| | | "pcm_s16le stream format changed from {}Hz/{}ch to {}Hz/{}ch", |
| | | self.source_sample_rate_hz, |
| | | self.source_num_channels, |
| | | source_sample_rate_hz, |
| | | source_num_channels |
| | | )); |
| | | } |
| | | |
| | | self.source_byte_buffer.extend_from_slice(pcm_bytes); |
| | | let source_frame_alignment_bytes = usize::from(self.source_num_channels) * 2; |
| | | let aligned_len = self.source_byte_buffer.len() |
| | | - self.source_byte_buffer.len() % source_frame_alignment_bytes; |
| | | if aligned_len > 0 { |
| | | let aligned_bytes: Vec<u8> = self.source_byte_buffer.drain(..aligned_len).collect(); |
| | | let source_samples: Vec<i16> = aligned_bytes |
| | | .chunks_exact(2) |
| | | .map(|chunk| i16::from_le_bytes([chunk[0], chunk[1]])) |
| | | .collect(); |
| | | self.debug_source_samples.extend_from_slice(&source_samples); |
| | | let remapped = remap_channels( |
| | | source_samples, |
| | | self.source_num_channels, |
| | | self.target_num_channels, |
| | | ); |
| | | self.source_sample_buffer.extend(remapped); |
| | | } |
| | | |
| | | let dropped_tail_bytes = if last && !self.source_byte_buffer.is_empty() { |
| | | let dropped = self.source_byte_buffer.len(); |
| | | self.source_byte_buffer.clear(); |
| | | dropped |
| | | } else { |
| | | 0 |
| | | }; |
| | | |
| | | let mut frames = self.drain_full_frames(); |
| | | if last { |
| | | if let Some(frame) = self.drain_final_partial_frame() { |
| | | frames.push(frame); |
| | | } |
| | | } |
| | | |
| | | Ok(PcmS16leStreamChunkResult { |
| | | frames, |
| | | dropped_tail_bytes, |
| | | buffered_source_bytes: self.source_byte_buffer.len(), |
| | | buffered_source_samples: self.source_sample_buffer.len(), |
| | | }) |
| | | } |
| | | |
| | | pub fn write_debug_dump( |
| | | &self, |
| | | dir: &str, |
| | | call_id: &str, |
| | | debug_label: &str, |
| | | ) -> Result<PcmS16leStreamDebugDump> { |
| | | let base_dir = PathBuf::from(dir); |
| | | fs::create_dir_all(&base_dir) |
| | | .with_context(|| format!("failed to create audio debug dump dir {dir}"))?; |
| | | let safe_call_id = sanitize_file_segment(call_id); |
| | | let safe_label = sanitize_file_segment(debug_label); |
| | | |
| | | let source_path = base_dir.join(format!("{safe_call_id}-{safe_label}-source.wav")); |
| | | write_debug_wav( |
| | | &source_path, |
| | | &self.debug_source_samples, |
| | | self.source_sample_rate_hz, |
| | | self.source_num_channels, |
| | | )?; |
| | | |
| | | let pcm_path = base_dir.join(format!("{safe_call_id}-{safe_label}-target.wav")); |
| | | write_debug_wav( |
| | | &pcm_path, |
| | | &self.debug_target_samples, |
| | | self.target_sample_rate_hz, |
| | | self.target_num_channels, |
| | | )?; |
| | | let pcm_size = fs::metadata(&pcm_path) |
| | | .with_context(|| format!("failed to stat audio debug wav {}", pcm_path.display()))? |
| | | .len(); |
| | | |
| | | Ok(PcmS16leStreamDebugDump { |
| | | debug_source_path: Some(source_path.to_string_lossy().to_string()), |
| | | debug_pcm_wav_path: Some(pcm_path.to_string_lossy().to_string()), |
| | | debug_pcm_wav_size_bytes: Some(pcm_size), |
| | | }) |
| | | } |
| | | |
| | | fn drain_full_frames(&mut self) -> Vec<PcmFrame> { |
| | | let source_samples_per_frame = self.source_samples_per_20ms_frame(); |
| | | let mut frames = Vec::new(); |
| | | while self.source_sample_buffer.len() >= source_samples_per_frame { |
| | | let source_samples: Vec<i16> = self |
| | | .source_sample_buffer |
| | | .drain(..source_samples_per_frame) |
| | | .collect(); |
| | | frames.push(self.convert_source_samples_to_frame(source_samples, true)); |
| | | } |
| | | frames |
| | | } |
| | | |
| | | fn drain_final_partial_frame(&mut self) -> Option<PcmFrame> { |
| | | if self.source_sample_buffer.is_empty() { |
| | | return None; |
| | | } |
| | | let source_samples: Vec<i16> = self.source_sample_buffer.drain(..).collect(); |
| | | Some(self.convert_source_samples_to_frame(source_samples, true)) |
| | | } |
| | | |
| | | fn convert_source_samples_to_frame( |
| | | &mut self, |
| | | source_samples: Vec<i16>, |
| | | pad_to_frame: bool, |
| | | ) -> PcmFrame { |
| | | let mut target_samples = resample_linear( |
| | | source_samples, |
| | | self.source_sample_rate_hz, |
| | | self.target_sample_rate_hz, |
| | | self.target_num_channels, |
| | | ); |
| | | let target_frame_samples = self.target_samples_per_20ms_frame(); |
| | | if target_samples.len() > target_frame_samples { |
| | | target_samples.truncate(target_frame_samples); |
| | | } else if pad_to_frame && target_samples.len() < target_frame_samples { |
| | | target_samples.resize(target_frame_samples, 0); |
| | | } |
| | | self.debug_target_samples.extend_from_slice(&target_samples); |
| | | PcmFrame::new( |
| | | target_samples, |
| | | self.target_sample_rate_hz, |
| | | u32::from(self.target_num_channels), |
| | | self.target_samples_per_channel_per_20ms_frame() as u32, |
| | | ) |
| | | } |
| | | |
| | | fn source_samples_per_20ms_frame(&self) -> usize { |
| | | let source_frames = |
| | | ((f64::from(self.source_sample_rate_hz) / 50.0).round() as usize).max(1); |
| | | source_frames * usize::from(self.target_num_channels) |
| | | } |
| | | |
| | | fn target_samples_per_20ms_frame(&self) -> usize { |
| | | self.target_samples_per_channel_per_20ms_frame() * usize::from(self.target_num_channels) |
| | | } |
| | | |
| | | fn target_samples_per_channel_per_20ms_frame(&self) -> usize { |
| | | ((self.target_sample_rate_hz / 1000) * 20).max(1) as usize |
| | | } |
| | | } |
| | | |
| | | fn decode_audio_frames( |
| | | audio_bytes: &[u8], |
| | | source_kind: &'static str, |
| | |
| | | |
| | | #[cfg(test)] |
| | | mod tests { |
| | | use super::{load_pre_recorded_frames, pcm_s16le_bytes_to_frames}; |
| | | use super::{PcmS16leStreamDecoder, load_pre_recorded_frames, pcm_s16le_bytes_to_frames}; |
| | | |
| | | #[tokio::test] |
| | | async fn load_pre_recorded_frames_reads_local_wav() { |
| | |
| | | assert_eq!(frames[0].num_channels, 1); |
| | | assert_eq!(frames[0].samples_per_channel, 960); |
| | | } |
| | | |
| | | #[test] |
| | | fn pcm_s16le_stream_decoder_does_not_pad_each_network_chunk() { |
| | | let mut decoder = PcmS16leStreamDecoder::new(44_100, 1, 48_000, 1).expect("stream decoder"); |
| | | let mut pcm = Vec::new(); |
| | | for index in 0..4410 { |
| | | let sample = if index % 2 == 0 { 1024_i16 } else { -1024_i16 }; |
| | | pcm.extend_from_slice(&sample.to_le_bytes()); |
| | | } |
| | | |
| | | let mut frames = Vec::new(); |
| | | let chunk_size = 698 * 2; |
| | | for (index, chunk) in pcm.chunks(chunk_size).enumerate() { |
| | | let last = (index + 1) * chunk_size >= pcm.len(); |
| | | let result = decoder |
| | | .push_bytes(chunk, 44_100, 1, last) |
| | | .expect("push pcm chunk"); |
| | | frames.extend(result.frames); |
| | | } |
| | | |
| | | assert_eq!(5, frames.len()); |
| | | assert_eq!( |
| | | 4_800, |
| | | frames.iter().map(|frame| frame.data.len()).sum::<usize>() |
| | | ); |
| | | assert!(frames.iter().all(|frame| frame.sample_rate == 48_000)); |
| | | assert!(frames.iter().all(|frame| frame.num_channels == 1)); |
| | | assert!(frames.iter().all(|frame| frame.samples_per_channel == 960)); |
| | | } |
| | | |
| | | #[test] |
| | | fn pcm_s16le_stream_decoder_waits_until_full_20ms_frame() { |
| | | let mut decoder = PcmS16leStreamDecoder::new(44_100, 1, 48_000, 1).expect("stream decoder"); |
| | | let ten_ms = vec![0_u8; 441 * 2]; |
| | | |
| | | let first = decoder |
| | | .push_bytes(&ten_ms, 44_100, 1, false) |
| | | .expect("first chunk"); |
| | | assert!(first.frames.is_empty()); |
| | | assert_eq!(441, first.buffered_source_samples); |
| | | |
| | | let second = decoder |
| | | .push_bytes(&ten_ms, 44_100, 1, false) |
| | | .expect("second chunk"); |
| | | assert_eq!(1, second.frames.len()); |
| | | assert_eq!(960, second.frames[0].samples_per_channel); |
| | | assert_eq!(0, second.buffered_source_samples); |
| | | } |
| | | |
| | | #[test] |
| | | fn pcm_s16le_stream_decoder_can_keep_16k_native_audio_profile() { |
| | | let mut decoder = PcmS16leStreamDecoder::new(16_000, 1, 16_000, 1).expect("stream decoder"); |
| | | let mut pcm = Vec::new(); |
| | | for index in 0..1_600 { |
| | | let sample = if index % 2 == 0 { 768_i16 } else { -768_i16 }; |
| | | pcm.extend_from_slice(&sample.to_le_bytes()); |
| | | } |
| | | |
| | | let mut frames = Vec::new(); |
| | | let chunk_size = 250 * 2; |
| | | for (index, chunk) in pcm.chunks(chunk_size).enumerate() { |
| | | let last = (index + 1) * chunk_size >= pcm.len(); |
| | | let result = decoder |
| | | .push_bytes(chunk, 16_000, 1, last) |
| | | .expect("push pcm chunk"); |
| | | frames.extend(result.frames); |
| | | } |
| | | |
| | | assert_eq!(5, frames.len()); |
| | | assert!(frames.iter().all(|frame| frame.sample_rate == 16_000)); |
| | | assert!(frames.iter().all(|frame| frame.num_channels == 1)); |
| | | assert!(frames.iter().all(|frame| frame.samples_per_channel == 320)); |
| | | assert_eq!( |
| | | 1_600, |
| | | frames.iter().map(|frame| frame.data.len()).sum::<usize>() |
| | | ); |
| | | } |
| | | } |